dApp Docs/AI Agent 自动化工作流与定时任务指南
Development reference. Not independently verified for production.

MSG Chain AI Agent 自动化工作流与定时任务指南

链标识: msg-chain-1 | Bech32 前缀: msg
适用版本: MSG Chain SDK v0.47+ | CosmWasm 1.5+

⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/


目录

  1. 概述
  2. 定时任务调度器
  3. 链上事件触发器
  4. 工作流编排引擎
  5. 条件与决策引擎
  6. 链上工作流记录
  7. 自动化策略模板
  8. 完整示例
  9. 附录

1. 概述

1.1 为什么需要 Agent 工作流自动化

在 MSG Chain 生态中,AI Agent 不仅仅是被动响应工具——它们是自主运行的经济实体。Agent 需要持续执行链上操作:领取奖励、再平衡仓位、监控清算风险、响应治理提案。

手动触发每个操作违背了 Agent 的自主性设计目标。工作流自动化赋予 Agent 以下能力:

能力 说明
无值守运行 7×24 持续执行策略,无需人工干预
精准时序 毫秒级响应链上事件,不错过窗口
资源优化 批量聚合操作,降低 Gas 开销
状态持久化 每一步执行写入链上,可审计可回放
容错恢复 失败自动重试,断点续跑

1.2 触发类型

Agent 工作流支持四种触发模式,可任意组合:

触发类型
├── 时间触发 (Time)        ← Cron 表达式、固定间隔、指定时间
├── 事件触发 (Event)       ← 链上交易、合约事件、新块、治理提案
├── 条件触发 (Condition)   ← 余额阈值、持仓偏离、APY 变化
└── 手动触发 (Manual)      ← 通过 CLI / Dashboard / 消息调用

1.3 典型用例

用例 触发方式 执行频率
流动性挖矿奖励自动复投 时间 + 条件 每 6 小时 / 奖励 > 阈值
AMM 池再平衡 条件 + 事件 偏离 > 5% / 价格大波动
社交账户自动回复 事件 被提及 / DM
链上数据日报 时间 每天 00:00 UTC
清算监控 + 抢先交易 事件 每个新区块
治理提案自动分析 + 投票 事件 新提案上链
跨链资产转移 条件 + 时间 余额 > 阈值 / 每周
收益率排名与资金迁移 时间 + 条件 每 12 小时 / 排名变化

2. 定时任务调度器

2.1 架构设计

定时任务调度器是 Agent 自动化的核心组件。它负责解析 Cron 表达式、维护任务队列、分发执行信号并记录结果。

┌─────────────────────────────────────────────────────┐
│                  AgentScheduler                      │
│                                                     │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐          │
│  │  Job #1   │  │  Job #2   │  │  Job #3   │  ...  │
│  │ cron: *   │  │ cron: 0  │  │ cron: */  │        │
│  │ /15 * *  │  │ */6 * *  │  │ 30 * * *  │        │
│  │   * *     │  │   * *    │  │   * *     │        │
│  └────┬─────┘  └────┬─────┘  └────┬─────┘          │
│       │              │              │                │
│  ┌────▼──────────────▼──────────────▼─────┐         │
│  │         Scheduler Loop (60s tick)       │         │
│  │   → parse cron → should_run? → execute │         │
│  └────────────────┬───────────────────────┘         │
│                   │                                  │
│  ┌────────────────▼───────────────────────┐         │
│  │      On-chain Execution Logger         │         │
│  │   msg_execute{agent, job_id, status}   │         │
│  └────────────────────────────────────────┘         │
└─────────────────────────────────────────────────────┘

2.2 核心调度器实现

import asyncio
import time
import json
from datetime import datetime
from typing import Dict, List, Callable, Optional, Any, Set
from dataclasses import dataclass, field
from enum import Enum
import hashlib

from croniter import croniter
from msg_sdk import MsgClient, Wallet, Coins
from msg_sdk.agent import AgentExecutor


class JobStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    SUCCESS = "success"
    FAILED = "failed"
    SKIPPED = "skipped"
    PAUSED = "paused"


@dataclass
class ScheduledJob:
    id: str
    cron_expr: str
    task: Callable
    name: str = ""
    description: str = ""
    max_retries: int = 3
    retry_delay: int = 60
    timeout: int = 300
    tags: List[str] = field(default_factory=list)
    created_at: float = 0.0
    last_run: Optional[float] = None
    next_run: Optional[float] = None
    status: JobStatus = JobStatus.PENDING
    run_count: int = 0
    success_count: int = 0
    fail_count: int = 0
    consecutive_failures: int = 0
    cooldown_until: Optional[float] = None
    depends_on: List[str] = field(default_factory=list)


@dataclass
class ExecutionRecord:
    execution_id: str
    job_id: str
    agent_id: str
    started_at: float
    completed_at: Optional[float] = None
    status: JobStatus = JobStatus.PENDING
    error: Optional[str] = None
    gas_used: int = 0
    tx_hash: Optional[str] = None
    result: Optional[dict] = None
    retry_count: int = 0


class AgentScheduler:
    def __init__(
        self,
        agent_id: str,
        client: MsgClient,
        wallet: Wallet,
        executor: Optional[AgentExecutor] = None,
        check_interval: int = 60,
        max_concurrent: int = 10,
        record_on_chain: bool = True,
    ):
        self.agent_id = agent_id
        self.client = client
        self.wallet = wallet
        self.executor = executor or AgentExecutor(agent_id, client, wallet)
        self.check_interval = check_interval
        self.max_concurrent = max_concurrent
        self.record_on_chain = record_on_chain
        self.jobs: Dict[str, ScheduledJob] = {}
        self.execution_history: List[ExecutionRecord] = []
        self.pending_executions: Set[str] = set()
        self._running = False
        self._semaphore: Optional[asyncio.Semaphore] = None
        self._on_error_callbacks: List[Callable] = []
        self._on_success_callbacks: List[Callable] = []

    async def add_cron_job(
        self, job_id: str, cron_expr: str, task: Callable,
        name: str = "", description: str = "",
        max_retries: int = 3, timeout: int = 300,
        tags: Optional[List[str]] = None,
        depends_on: Optional[List[str]] = None,
    ) -> ScheduledJob:
        if job_id in self.jobs:
            raise ValueError(f"Job '{job_id}' already exists")
        if not croniter.is_valid(cron_expr):
            raise ValueError(f"Invalid cron expression: {cron_expr}")
        cron = croniter(cron_expr, datetime.now())
        next_run = cron.get_next(float)
        job = ScheduledJob(
            id=job_id, cron_expr=cron_expr, task=task,
            name=name or job_id, description=description,
            max_retries=max_retries, timeout=timeout,
            tags=tags or [], created_at=time.time(),
            next_run=next_run, depends_on=depends_on or [],
        )
        self.jobs[job_id] = job
        if self.record_on_chain:
            await self._record_job_creation(job)
        return job

    async def add_interval_job(
        self, job_id: str, interval_seconds: int,
        task: Callable, name: str = "", **kwargs,
    ) -> ScheduledJob:
        cron_expr = self._interval_to_cron(interval_seconds)
        return await self.add_cron_job(job_id, cron_expr, task, name=name, **kwargs)

    async def add_one_time_job(
        self, job_id: str, run_at: datetime,
        task: Callable, name: str = "", **kwargs,
    ) -> ScheduledJob:
        cron_expr = f"{run_at.minute} {run_at.hour} {run_at.day} {run_at.month} *"
        job = await self.add_cron_job(job_id, cron_expr, task, name=name, **kwargs)
        job.tags.append("one_time")
        return job

    async def remove_job(self, job_id: str) -> bool:
        return self.jobs.pop(job_id, None) is not None

    async def pause_job(self, job_id: str):
        if job_id in self.jobs:
            self.jobs[job_id].status = JobStatus.PAUSED

    async def resume_job(self, job_id: str):
        if job_id in self.jobs:
            self.jobs[job_id].status = JobStatus.PENDING
            cron = croniter(self.jobs[job_id].cron_expr, datetime.now())
            self.jobs[job_id].next_run = cron.get_next(float)

    async def start(self):
        if self._running:
            raise RuntimeError("Scheduler already running")
        self._running = True
        self._semaphore = asyncio.Semaphore(self.max_concurrent)
        self._scheduler_task = asyncio.create_task(self._run_loop())
        if self.record_on_chain:
            await self._record_scheduler_start()

    async def stop(self, graceful: bool = True):
        self._running = False
        if graceful:
            while self.pending_executions:
                await asyncio.sleep(1)
        if hasattr(self, '_scheduler_task'):
            self._scheduler_task.cancel()
        if self.record_on_chain:
            await self._record_scheduler_stop()

    async def _run_loop(self):
        while self._running:
            try:
                now = datetime.now()
                for job in list(self.jobs.values()):
                    if job.status == JobStatus.PAUSED:
                        continue
                    if job.cooldown_until and time.time() < job.cooldown_until:
                        continue
                    if self.should_run(job, now):
                        if job.id not in self.pending_executions:
                            await self._dispatch_job(job)
                await asyncio.sleep(self.check_interval)
            except asyncio.CancelledError:
                break
            except Exception as e:
                self._emit_error(f"Scheduler loop error: {e}")
                await asyncio.sleep(self.check_interval)

    def should_run(self, job: ScheduledJob, now: datetime) -> bool:
        if not croniter.match(job.cron_expr, now):
            return False
        if job.last_run and (now.timestamp() - job.last_run) < 60:
            return False
        if job.depends_on:
            for dep_id in job.depends_on:
                dep = self.jobs.get(dep_id)
                if dep and dep.status != JobStatus.SUCCESS:
                    return False
        return True

    async def _dispatch_job(self, job: ScheduledJob):
        self.pending_executions.add(job.id)
        job.last_run = time.time()
        job.status = JobStatus.RUNNING
        cron = croniter(job.cron_expr, datetime.now())
        job.next_run = cron.get_next(float)
        asyncio.create_task(self._execute_with_semaphore(job))

    async def _execute_with_semaphore(self, job: ScheduledJob):
        async with self._semaphore:
            await self.execute_job(job)

    async def execute_job(self, job: ScheduledJob):
        execution = ExecutionRecord(
            execution_id=self._gen_exec_id(job.id),
            job_id=job.id, agent_id=self.agent_id,
            started_at=time.time(),
        )
        self.execution_history.append(execution)
        last_error = None
        for attempt in range(job.max_retries + 1):
            if attempt > 0:
                self._emit_log(f"Retry {attempt}/{job.max_retries} for job '{job.id}'")
                await asyncio.sleep(job.retry_delay)
            try:
                result = await asyncio.wait_for(job.task(), timeout=job.timeout)
                execution.status = JobStatus.SUCCESS
                execution.completed_at = time.time()
                execution.result = result
                execution.retry_count = attempt
                job.status = JobStatus.SUCCESS
                job.run_count += 1
                job.success_count += 1
                job.consecutive_failures = 0
                self._emit_success(job, execution)
                await self._record_execution(job, execution)
                self.pending_executions.discard(job.id)
                return
            except asyncio.TimeoutError:
                last_error = f"Timeout after {job.timeout}s"
            except Exception as e:
                last_error = str(e)
            cron = croniter(job.cron_expr, datetime.now())
            job.next_run = cron.get_next(float)
        execution.status = JobStatus.FAILED
        execution.completed_at = time.time()
        execution.error = last_error
        job.status = JobStatus.FAILED
        job.run_count += 1
        job.fail_count += 1
        job.consecutive_failures += 1
        cooldown = min(3600, 60 * (2 ** job.consecutive_failures))
        job.cooldown_until = time.time() + cooldown
        self._emit_error(f"Job '{job.id}' failed after {job.max_retries + 1} attempts: {last_error}")
        await self._record_execution(job, execution)
        self.pending_executions.discard(job.id)

    async def _record_job_creation(self, job: ScheduledJob):
        try:
            msg = {
                "@type": "/msg.agent.v1.MsgRegisterCronJob",
                "agent": self.agent_id,
                "sender": str(self.wallet.address),
                "job_id": job.id,
                "cron_expr": job.cron_expr,
                "name": job.name, "description": job.description,
                "max_retries": str(job.max_retries), "timeout": str(job.timeout),
                "tags": json.dumps(job.tags),
                "depends_on": json.dumps(job.depends_on),
            }
            await self.client.execute_contract(
                self.executor.contract_address, msg, sender=self.wallet,
            )
        except Exception as e:
            self._emit_log(f"Failed to record job creation on chain: {e}")

    async def _record_execution(self, job: ScheduledJob, execution: ExecutionRecord):
        if not self.record_on_chain:
            return
        try:
            msg = {
                "@type": "/msg.agent.v1.MsgRecordExecution",
                "agent": self.agent_id,
                "sender": str(self.wallet.address),
                "execution_id": execution.execution_id,
                "job_id": job.id,
                "status": execution.status.value,
                "started_at": str(int(execution.started_at)),
                "completed_at": str(int(execution.completed_at or 0)),
                "error": execution.error or "",
                "gas_used": str(execution.gas_used),
                "result": json.dumps(execution.result or {}),
                "retry_count": str(execution.retry_count),
            }
            await self.client.execute_contract(
                self.executor.contract_address, msg, sender=self.wallet,
            )
        except Exception as e:
            self._emit_log(f"Failed to record execution on chain: {e}")

    async def _record_scheduler_start(self):
        try:
            msg = {
                "@type": "/msg.agent.v1.MsgSchedulerEvent",
                "agent": self.agent_id,
                "sender": str(self.wallet.address),
                "event_type": "start",
                "timestamp": str(int(time.time())),
                "metadata": json.dumps({"check_interval": self.check_interval}),
            }
            await self.client.execute_contract(
                self.executor.contract_address, msg, sender=self.wallet,
            )
        except Exception:
            pass

    async def _record_scheduler_stop(self):
        try:
            msg = {
                "@type": "/msg.agent.v1.MsgSchedulerEvent",
                "agent": self.agent_id,
                "sender": str(self.wallet.address),
                "event_type": "stop",
                "timestamp": str(int(time.time())),
            }
            await self.client.execute_contract(
                self.executor.contract_address, msg, sender=self.wallet,
            )
        except Exception:
            pass

    async def recover_missed_jobs(self, lookback_hours: int = 24):
        now = datetime.now()
        lookback_start = now.timestamp() - (lookback_hours * 3600)
        recovered = []
        for job in self.jobs.values():
            cron = croniter(job.cron_expr, datetime.fromtimestamp(lookback_start))
            expected_runs = []
            next_time = cron.get_next(float)
            while next_time < now.timestamp():
                expected_runs.append(next_time)
                next_time = cron.get_next(float)
            if not expected_runs:
                continue
            executed_times = {
                e.started_at for e in self.execution_history
                if e.job_id == job.id and e.status == JobStatus.SUCCESS
            }
            missed = [t for t in expected_runs if not self._is_near_any(t, executed_times)]
            if missed:
                self._emit_log(f"Job '{job.id}' missed {len(missed)} executions, catching up...")
                await self.execute_job(job)
                recovered.append({"job_id": job.id, "missed_count": len(missed)})
        return recovered

    def _is_near_any(self, timestamp: float, timestamps: Set[float], tolerance: int = 120) -> bool:
        return any(abs(timestamp - t) <= tolerance for t in timestamps)

    def on_error(self, callback: Callable):
        self._on_error_callbacks.append(callback)

    def on_success(self, callback: Callable):
        self._on_success_callbacks.append(callback)

    def _emit_success(self, job: ScheduledJob, execution: ExecutionRecord):
        for cb in self._on_success_callbacks:
            try:
                cb(job, execution)
            except Exception:
                pass

    def _emit_error(self, message: str):
        for cb in self._on_error_callbacks:
            try:
                cb(message)
            except Exception:
                pass

    def _interval_to_cron(self, seconds: int) -> str:
        if seconds < 60:
            return f"*/{max(1, seconds)} * * * *"
        minutes = seconds // 60
        if minutes < 60:
            return f"*/{minutes} * * * *"
        hours = minutes // 60
        if hours < 24:
            return f"0 */{hours} * * *"
        days = hours // 24
        return f"0 0 */{days} * *"

    def _gen_exec_id(self, job_id: str) -> str:
        raw = f"{job_id}:{time.time()}:{hash(job_id)}"
        return hashlib.sha256(raw.encode()).hexdigest()[:16]

    def _emit_log(self, msg: str):
        print(f"[AgentScheduler:{self.agent_id}] {msg}")

    def get_stats(self) -> dict:
        total = len(self.jobs)
        active = sum(1 for j in self.jobs.values() if j.status == JobStatus.PENDING)
        paused = sum(1 for j in self.jobs.values() if j.status == JobStatus.PAUSED)
        failed = sum(1 for j in self.jobs.values() if j.status == JobStatus.FAILED)
        return {
            "agent_id": self.agent_id,
            "total_jobs": total, "active_jobs": active,
            "paused_jobs": paused, "failed_jobs": failed,
            "total_executions": len(self.execution_history),
            "successful_executions": sum(
                1 for e in self.execution_history if e.status == JobStatus.SUCCESS
            ),
            "failed_executions": sum(
                1 for e in self.execution_history if e.status == JobStatus.FAILED
            ),
            "uptime_seconds": time.time() - getattr(self, '_start_time', 0),
        }

2.3 Cron 表达式高级用法

CRON_EVERY_MINUTE = "* * * * *"
CRON_EVERY_5_MINUTES = "*/5 * * * *"
CRON_EVERY_15_MINUTES = "*/15 * * * *"
CRON_EVERY_30_MINUTES = "*/30 * * * *"
CRON_EVERY_HOUR = "0 * * * *"
CRON_EVERY_2_HOURS = "0 */2 * * *"
CRON_EVERY_6_HOURS = "0 */6 * * *"
CRON_EVERY_12_HOURS = "0 */12 * * *"
CRON_DAILY_MIDNIGHT = "0 0 * * *"
CRON_DAILY_NOON = "0 12 * * *"
CRON_WEEKLY = "0 0 * * 0"
CRON_MONTHLY = "0 0 1 * *"

CRON_PATTERNS = {
    "workday_business_hours": "*/30 9-18 * * 1-5",
    "workday_off_hours": "0 0,2,4,6,8,18,20,22 * * 1-5",
    "weekend": "0 */4 * * 0,6",
    "monthly_1st_15th": "0 0 1,15 * *",
    "quarterly_end": "0 0 31 3,6,9,12 *",
    "beijing_8am": "0 0 * * *",
    "every_2h_30min": "30 */2 * * *",
    "staggered_15min": "1,16,31,46 * * * *",
    "high_activity": "*/10 * * * *",
    "third_wednesday": "0 0 * * 3#3",
}

2.4 分布式调度与去重

class DistributedScheduler(AgentScheduler):
    async def acquire_lock(self, job_id: str, ttl: int = 300) -> bool:
        msg = {
            "@type": "/msg.agent.v1.MsgAcquireLock",
            "agent": self.agent_id,
            "sender": str(self.wallet.address),
            "lock_key": f"scheduler:{job_id}",
            "ttl": str(ttl),
        }
        result = await self.client.execute_contract(
            self.executor.contract_address, msg, sender=self.wallet,
        )
        return result.data.get("acquired", False)

    async def release_lock(self, job_id: str):
        msg = {
            "@type": "/msg.agent.v1.MsgReleaseLock",
            "agent": self.agent_id,
            "sender": str(self.wallet.address),
            "lock_key": f"scheduler:{job_id}",
        }
        await self.client.execute_contract(
            self.executor.contract_address, msg, sender=self.wallet,
        )

    async def _dispatch_job(self, job: ScheduledJob):
        acquired = await self.acquire_lock(job.id)
        if not acquired:
            self._emit_log(f"Job '{job.id}' locked by another instance, skipping")
            return
        try:
            await super()._dispatch_job(job)
        finally:
            await self.release_lock(job.id)

2.5 调度器使用示例

async def demo_scheduler():
    client = MsgClient(rpc_endpoint="https://rpc.msg-chain-1.zone")
    wallet = Wallet.from_mnemonic("your mnemonic here...")
    agent_id = "msg1agentdemo..."
    scheduler = AgentScheduler(
        agent_id=agent_id, client=client, wallet=wallet,
        check_interval=60, max_concurrent=5, record_on_chain=True,
    )
    scheduler.on_error(lambda msg: print(f"ERROR: {msg}"))

    async def auto_compound():
        return {"compounded": True, "amount": "1000000umsg"}
    await scheduler.add_cron_job(
        "auto_compound", CRON_EVERY_6_HOURS, auto_compound,
        name="自动复投", max_retries=5, timeout=600,
    )

    async def generate_report():
        return {"report_id": "daily_通用维护记录"}
    await scheduler.add_cron_job(
        "daily_report", CRON_DAILY_MIDNIGHT, generate_report,
        name="日报生成", tags=["report", "daily"],
    )

    async def health_check():
        print("Health check OK")
    await scheduler.add_interval_job("health_check", 30, health_check)
    await scheduler.start()
    await asyncio.sleep(3600)
    await scheduler.stop()
    print(json.dumps(scheduler.get_stats(), indent=2))

3. 链上事件触发器

3.1 触发器架构

事件触发器监听链上事件流,根据预定义规则触发工作流执行。

区块链
  │
  │  WebSocket / gRPC 订阅
  ▼
┌────────────────────────────────────────────┐
│           Event Stream Listener             │
│  ┌──────────┐ ┌──────────┐ ┌──────────┐   │
│  │ Tx Events│ │ Block    │ │ Contract │   │
│  │ Listener │ │ Listener │ │ Events   │   │
│  └────┬─────┘ └────┬─────┘ └────┬─────┘   │
│       │             │             │          │
│  ┌────▼─────────────▼─────────────▼──────┐  │
│  │         Event Router & Filter          │  │
│  │  → 规则匹配 → 条件评估 → 去重过滤     │  │
│  └────────────────┬──────────────────────┘  │
│                   │                          │
│  ┌────────────────▼──────────────────────┐  │
│  │         Workflow Dispatcher            │  │
│  │  → 触发工作流 → 传递上下文 → 记录日志  │  │
│  └────────────────────────────────────────┘  │
└────────────────────────────────────────────┘

3.2 触发器规则引擎

import asyncio, time, json, fnmatch, re
from typing import Dict, List, Optional, Any, Callable, Set, Tuple
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from msg_sdk import MsgClient, Wallet

class TriggerType(str, Enum):
    BALANCE_THRESHOLD = "balance_threshold"
    NEW_BLOCK = "new_block"
    CONTRACT_EVENT = "contract_event"
    GOVERNANCE_PROPOSAL = "governance_proposal"
    TRANSACTION_PATTERN = "transaction_pattern"
    PRICE_CHANGE = "price_change"
    CUSTOM_QUERY = "custom_query"

@dataclass
class TriggerRule:
    rule_id: str
    trigger_type: TriggerType
    name: str = ""
    description: str = ""
    config: dict = field(default_factory=dict)
    filters: dict = field(default_factory=dict)
    cooldown_seconds: int = 0
    priority: int = 100
    enabled: bool = True
    workflow_id: str = ""
    max_firings_per_window: int = 10
    window_seconds: int = 3600
    created_at: float = 0.0
    last_fired: Optional[float] = None
    firing_count: int = 0

@dataclass
class ChainEvent:
    event_id: str
    block_height: int
    block_time: datetime
    event_type: str
    raw_data: dict
    parsed_data: dict = field(default_factory=dict)
    tx_hash: Optional[str] = None
    contract_address: Optional[str] = None
    sender: Optional[str] = None

class EventTrigger:
    TRIGGER_TEMPLATES = {
        TriggerType.BALANCE_THRESHOLD: {
            "description": "代币余额超过/低于阈值时触发",
            "config": {"token": "umsg", "address": "", "threshold": "1000000", "direction": "above"},
        },
        TriggerType.NEW_BLOCK: {
            "description": "特定区块高度或间隔触发",
            "config": {"interval": 100, "start_height": 0, "end_height": 0, "confirmations": 1},
        },
        TriggerType.CONTRACT_EVENT: {
            "description": "智能合约特定事件触发",
            "config": {"contract": "", "event_type": "wasm", "event_name": "", "filter": {}},
        },
        TriggerType.GOVERNANCE_PROPOSAL: {
            "description": "新治理提案触发",
            "config": {"status_filter": ["PROPOSAL_STATUS_VOTING_PERIOD"], "min_deposit": "0", "proposal_types": []},
        },
        TriggerType.TRANSACTION_PATTERN: {
            "description": "匹配特定交易模式触发",
            "config": {"messages": [], "from_addresses": [], "to_addresses": [], "min_value": "0", "memo_pattern": ""},
        },
        TriggerType.PRICE_CHANGE: {
            "description": "价格变化超过阈值触发",
            "config": {"pair": "MSG/USDC", "oracle": "msg_oracle", "change_percent": 5.0, "direction": "both"},
        },
        TriggerType.CUSTOM_QUERY: {
            "description": "自定义链上查询条件满足时触发",
            "config": {"query_type": "wasm_smart", "contract": "", "query_msg": {}, "condition": "result.value > 1000", "poll_interval": 60},
        },
    }

    def __init__(self, agent_id: str, client: MsgClient, wallet: Wallet, scheduler=None):
        self.agent_id = agent_id
        self.client = client
        self.wallet = wallet
        self.scheduler = scheduler
        self.rules: Dict[str, TriggerRule] = {}
        self._running = False
        self._subscriptions: List[asyncio.Task] = []
        self._rate_limiters: Dict[str, List[float]] = {}
        self._on_trigger_callbacks: List[Callable] = []
        self._trigger_stats: Dict[str, dict] = {}
        self._prev_balances: Dict[str, int] = {}

    async def add_rule(self, rule_id: str, trigger_type: TriggerType, config: dict,
                       workflow_id: str = "", name: str = "", cooldown_seconds: int = 0,
                       priority: int = 100, filters: Optional[dict] = None) -> TriggerRule:
        if rule_id in self.rules:
            raise ValueError(f"Rule '{rule_id}' already exists")
        template = self.TRIGGER_TEMPLATES.get(trigger_type)
        if not template:
            raise ValueError(f"Unknown trigger type: {trigger_type}")
        merged = {**template["config"], **config}
        rule = TriggerRule(rule_id=rule_id, trigger_type=trigger_type, name=name or rule_id,
                          config=merged, filters=filters or {}, workflow_id=workflow_id,
                          cooldown_seconds=cooldown_seconds, priority=priority, created_at=time.time())
        self.rules[rule_id] = rule
        return rule

    async def remove_rule(self, rule_id: str) -> bool:
        return self.rules.pop(rule_id, None) is not None

    async def enable_rule(self, rule_id: str):
        if rule_id in self.rules:
            self.rules[rule_id].enabled = True

    async def disable_rule(self, rule_id: str):
        if rule_id in self.rules:
            self.rules[rule_id].enabled = False

    async def start(self):
        if self._running:
            return
        self._running = True
        groups: Dict[str, List[TriggerRule]] = {}
        for rule in self.rules.values():
            if rule.enabled:
                groups.setdefault(rule.trigger_type.value, []).append(rule)
        for ttype, rules in groups.items():
            self._subscriptions.append(asyncio.create_task(self._start_listener(ttype, rules)))

    async def stop(self):
        self._running = False
        for t in self._subscriptions:
            t.cancel()
        await asyncio.gather(*self._subscriptions, return_exceptions=True)
        self._subscriptions.clear()

    async def _start_listener(self, trigger_type: str, rules: List[TriggerRule]):
        try:
            if trigger_type == "new_block":
                await self._listen_blocks(rules)
            elif trigger_type == "contract_event":
                await self._listen_contract_events(rules)
            elif trigger_type == "transaction_pattern":
                await self._listen_transactions(rules)
            elif trigger_type == "governance_proposal":
                await self._listen_governance(rules)
            elif trigger_type == "balance_threshold":
                await self._poll_balance_thresholds(rules)
            elif trigger_type == "custom_query":
                await self._poll_custom_queries(rules)
            elif trigger_type == "price_change":
                await self._poll_price_changes(rules)
        except asyncio.CancelledError:
            pass
        except Exception as e:
            print(f"[EventTrigger] Listener error for {trigger_type}: {e}")

    async def _listen_blocks(self, rules: List[TriggerRule]):
        ws_url = self.client.rpc_endpoint.replace("http", "ws") + "/websocket"
        last_height = 0
        async with self.client.ws_connect(ws_url) as ws:
            await ws.send(json.dumps({"jsonrpc": "2.0", "method": "subscribe", "params": {"query": "tm.event='NewBlock'"}, "id": 1}))
            while self._running:
                msg = await ws.recv()
                data = json.loads(msg)
                if "result" not in data or "data" not in data["result"]:
                    continue
                block_data = data["result"]["data"]["value"]
                bh = int(block_data["block"]["header"]["height"])
                if bh <= last_height:
                    continue
                last_height = bh
                event = ChainEvent(event_id=f"block_{bh}", block_height=bh, block_time=datetime.utcnow(),
                                  event_type="new_block", raw_data=block_data)
                await self._evaluate_rules(rules, event)

    async def _listen_contract_events(self, rules: List[TriggerRule]):
        ws_url = self.client.rpc_endpoint.replace("http", "ws") + "/websocket"
        async with self.client.ws_connect(ws_url) as ws:
            await ws.send(json.dumps({"jsonrpc": "2.0", "method": "subscribe", "params": {"query": "wasm._contract_address EXISTS"}, "id": 1}))
            while self._running:
                msg = await ws.recv()
                data = json.loads(msg)
                if "result" not in data:
                    continue
                for ev in data["result"]["data"]["value"]["TxResult"]["result"]["events"]:
                    attrs = {a.get("key", ""): a.get("value", "") for a in ev.get("attributes", [])}
                    event = ChainEvent(
                        event_id=f"contract_{attrs.get('_tx_hash', '')}_{ev['type']}",
                        block_height=int(data["result"]["data"]["value"]["TxResult"]["height"]),
                        block_time=datetime.utcnow(), event_type=f"contract:{ev['type']}",
                        raw_data=ev, parsed_data=attrs, tx_hash=attrs.get("_tx_hash"),
                        contract_address=attrs.get("_contract_address", ""),
                    )
                    await self._evaluate_rules(rules, event)

    async def _listen_transactions(self, rules: List[TriggerRule]):
        ws_url = self.client.rpc_endpoint.replace("http", "ws") + "/websocket"
        async with self.client.ws_connect(ws_url) as ws:
            await ws.send(json.dumps({"jsonrpc": "2.0", "method": "subscribe", "params": {"query": "tm.event='Tx'"}, "id": 1}))
            while self._running:
                msg = await ws.recv()
                data = json.loads(msg)
                if "result" not in data:
                    continue
                tx_data = data["result"]["data"]["value"]
                tx_hash = tx_data["TxHash"]
                tx_height = int(tx_data["Height"])
                event = ChainEvent(event_id=f"tx_{tx_hash}", block_height=tx_height,
                                  block_time=datetime.utcnow(), event_type="transaction",
                                  raw_data=tx_data, parsed_data={"tx_hash": tx_hash}, tx_hash=tx_hash)
                await self._evaluate_rules(rules, event)

    async def _listen_governance(self, rules: List[TriggerRule]):
        last_id = 0
        while self._running:
            try:
                proposals = await self.client.get_proposals(status="voting_period")
                for p in proposals:
                    pid = int(p.proposal_id)
                    if pid <= last_id:
                        continue
                    last_id = pid
                    event = ChainEvent(event_id=f"proposal_{pid}", block_height=0, block_time=datetime.utcnow(),
                                      event_type="governance:new_proposal", raw_data=p.to_dict(),
                                      parsed_data={"proposal_id": pid, "title": p.title, "type": p.proposal_type})
                    await self._evaluate_rules(rules, event)
            except Exception as e:
                print(f"[EventTrigger] Governance poll error: {e}")
            await asyncio.sleep(30)

    async def _poll_balance_thresholds(self, rules: List[TriggerRule]):
        while self._running:
            for rule in rules:
                if not rule.enabled:
                    continue
                try:
                    addr = rule.config.get("address", self.wallet.address)
                    token = rule.config.get("token", "umsg")
                    threshold = int(rule.config.get("threshold", "0"))
                    direction = rule.config.get("direction", "above")
                    balance = await self.client.get_balance(addr, token)
                    cur = int(balance.amount)
                    triggered = False
                    if direction == "above" and cur >= threshold:
                        triggered = True
                    elif direction == "below" and cur <= threshold:
                        triggered = True
                    elif direction == "cross":
                        prev = self._prev_balances.get(rule.rule_id, cur)
                        if (prev < threshold <= cur) or (prev > threshold >= cur):
                            triggered = True
                    if triggered:
                        event = ChainEvent(event_id=f"balance_{addr}_{token}_{time.time()}", block_height=0,
                                          block_time=datetime.utcnow(), event_type="balance_threshold",
                                          raw_data={"addr": addr, "balance": cur},
                                          parsed_data={"address": addr, "token": token, "balance": cur, "threshold": threshold})
                        await self._fire_rule(rule, event)
                    self._prev_balances[rule.rule_id] = cur
                except Exception as e:
                    print(f"[EventTrigger] Balance check error for rule '{rule.rule_id}': {e}")
            await asyncio.sleep(60)

    async def _poll_custom_queries(self, rules: List[TriggerRule]):
        while self._running:
            for rule in rules:
                if not rule.enabled:
                    continue
                try:
                    contract = rule.config.get("contract", "")
                    query_msg = rule.config.get("query_msg", {})
                    condition = rule.config.get("condition", "")
                    result = await self.client.query_contract_smart(contract, query_msg)
                    if condition and self._eval_condition(condition, result):
                        event = ChainEvent(event_id=f"custom_{rule.rule_id}_{time.time()}", block_height=0,
                                          block_time=datetime.utcnow(), event_type="custom_query",
                                          raw_data=result, parsed_data={"result": result})
                        await self._fire_rule(rule, event)
                except Exception:
                    pass
            await asyncio.sleep(rules[0].config.get("poll_interval", 60) if rules else 60)

    async def _poll_price_changes(self, rules: List[TriggerRule]):
        last_prices: Dict[str, float] = {}
        while self._running:
            for rule in rules:
                if not rule.enabled:
                    continue
                try:
                    pair = rule.config.get("pair", "")
                    oracle = rule.config.get("oracle", "msg_oracle")
                    change_pct = float(rule.config.get("change_percent", 5.0))
                    direction = rule.config.get("direction", "both")
                    result = await self.client.query_contract_smart(oracle, {"price": {"pair": pair}})
                    price = float(result.get("price", 0))
                    last = last_prices.get(rule.rule_id)
                    if last is not None:
                        change = (price - last) / last * 100
                        triggered = (direction in ("up", "both") and change >= change_pct) or \
                                    (direction in ("down", "both") and change <= -change_pct)
                        if triggered:
                            event = ChainEvent(event_id=f"price_{pair}_{time.time()}", block_height=0,
                                              block_time=datetime.utcnow(), event_type="price_change",
                                              raw_data={"pair": pair, "price": price, "change": change},
                                              parsed_data={"pair": pair, "price": price, "change_percent": round(change, 2)})
                            await self._fire_rule(rule, event)
                    last_prices[rule.rule_id] = price
                except Exception:
                    pass
            await asyncio.sleep(60)

    async def _evaluate_rules(self, rules: List[TriggerRule], event: ChainEvent):
        for rule in sorted(rules, key=lambda r: r.priority):
            if not rule.enabled:
                continue
            if await self._matches_rule(rule, event):
                await self._fire_rule(rule, event)

    async def _matches_rule(self, rule: TriggerRule, event: ChainEvent) -> bool:
        if rule.filters:
            for key, pattern in rule.filters.items():
                value = event.parsed_data.get(key, event.raw_data.get(key))
                if value is None:
                    return False
                if isinstance(pattern, str) and "*" in pattern:
                    if not fnmatch.fnmatch(str(value), pattern):
                        return False
                elif isinstance(pattern, (list, tuple)):
                    if value not in pattern:
                        return False
                elif value != pattern:
                    return False
        if rule.trigger_type == TriggerType.CONTRACT_EVENT:
            ec = rule.config.get("contract", "")
            en = rule.config.get("event_name", "")
            if ec and event.contract_address != ec:
                return False
            if en and en not in event.event_type:
                return False
        return True

    async def _fire_rule(self, rule: TriggerRule, event: ChainEvent):
        if not self._check_rate_limit(rule):
            return
        if rule.last_fired and rule.cooldown_seconds > 0 and (time.time() - rule.last_fired) < rule.cooldown_seconds:
            return
        rule.last_fired = time.time()
        rule.firing_count += 1
        self._trigger_stats.setdefault(rule.rule_id, {"total_firings": 0, "last_event": None})
        self._trigger_stats[rule.rule_id]["total_firings"] += 1
        self._trigger_stats[rule.rule_id]["last_event"] = {"event_id": event.event_id, "time": time.time()}
        for cb in self._on_trigger_callbacks:
            try:
                cb(rule, event)
            except Exception:
                pass
        if rule.workflow_id and self.scheduler:
            asyncio.create_task(self._execute_workflow(rule, event))
        await self._record_trigger(rule, event)

    async def _execute_workflow(self, rule: TriggerRule, event: ChainEvent):
        ctx = {"trigger_rule": rule.rule_id, "event": {k: v for k, v in event.__dict__.items() if k != 'raw_data'}}
        if hasattr(self.scheduler, 'executor') and self.scheduler.executor:
            try:
                await self.scheduler.executor.execute_workflow(workflow_id=rule.workflow_id, context=ctx)
            except Exception as e:
                self._emit_log(f"Workflow '{rule.workflow_id}' failed: {e}")

    def _check_rate_limit(self, rule: TriggerRule) -> bool:
        now = time.time()
        ws = now - rule.window_seconds
        self._rate_limiters.setdefault(rule.rule_id, [])
        self._rate_limiters[rule.rule_id] = [t for t in self._rate_limiters[rule.rule_id] if t > ws]
        if len(self._rate_limiters[rule.rule_id]) >= rule.max_firings_per_window:
            return False
        self._rate_limiters[rule.rule_id].append(now)
        return True

    async def _record_trigger(self, rule: TriggerRule, event: ChainEvent):
        try:
            msg = {"@type": "/msg.agent.v1.MsgTriggerEvent", "agent": self.agent_id,
                   "sender": str(self.wallet.address), "rule_id": rule.rule_id,
                   "trigger_type": rule.trigger_type.value, "event_id": event.event_id,
                   "block_height": str(event.block_height),
                   "event_data": json.dumps(event.parsed_data), "timestamp": str(int(time.time()))}
            await self.client.execute_contract(self.scheduler.executor.contract_address, msg, sender=self.wallet)
        except Exception:
            pass

    def on_trigger(self, callback: Callable):
        self._on_trigger_callbacks.append(callback)

    def _eval_condition(self, condition: str, result: dict) -> bool:
        try:
            return bool(eval(condition, {"__builtins__": {}}, {"result": result}))
        except Exception:
            return False

    def _emit_log(self, msg: str):
        print(f"[EventTrigger:{self.agent_id}] {msg}")

    def get_stats(self) -> dict:
        return {
            "agent_id": self.agent_id, "total_rules": len(self.rules),
            "enabled_rules": sum(1 for r in self.rules.values() if r.enabled),
            "running": self._running, "subscriptions": len(self._subscriptions),
            "rule_stats": self._trigger_stats,
        }

3.3 触发器使用示例

async def demo_event_trigger():
    client = MsgClient(rpc_endpoint="https://rpc.msg-chain-1.zone")
    wallet = Wallet.from_mnemonic("your mnemonic here...")
    trigger = EventTrigger("msg1agentdemo...", client, wallet)

    await trigger.add_rule("balance_monitor", TriggerType.BALANCE_THRESHOLD,
        {"token": "umsg", "address": str(wallet.address), "threshold": "5000000", "direction": "below"},
        workflow_id="top_up_agent", name="余额不足警告", cooldown_seconds=3600)

    await trigger.add_rule("every_100_blocks", TriggerType.NEW_BLOCK,
        {"interval": 100, "confirmations": 1}, workflow_id="health_check")

    await trigger.add_rule("swap_event", TriggerType.CONTRACT_EVENT,
        {"contract": "msg1swapcontract...", "event_name": "wasm-swap"},
        workflow_id="analyze_swap", filters={"sender": str(wallet.address)})

    await trigger.add_rule("new_proposal", TriggerType.GOVERNANCE_PROPOSAL,
        {"status_filter": ["PROPOSAL_STATUS_VOTING_PERIOD"]},
        workflow_id="auto_vote", priority=50)

    trigger.on_trigger(lambda r, e: print(f"Triggered: {r.name} | {e.event_type}"))
    await trigger.start()
    await asyncio.sleep(86400)
    await trigger.stop()

4. 工作流编排引擎

4.1 工作流抽象

工作流是一组有向无环图(DAG)结构化的步骤,支持条件分支、并行执行和错误处理。

                    ┌──────────┐
                    │  Start   │
                    └────┬─────┘
                         │
                    ┌────▼─────┐
                    │ Step A   │
                    │ 获取数据  │
                    └────┬─────┘
                         │
              ┌──────────┼──────────┐
              │          │          │
         ┌────▼────┐ ┌──▼────┐  ┌──▼──────┐
         │ Step B  │ │Step C │  │ Step D  │
         │ 条件判断 │ │ 计算  │  │ 并行任务 │
         └────┬────┘ └──┬────┘  └──┬──────┘
              │          │          │
              └──────────┼──────────┘
                         │
                    ┌────▼─────┐
                    │ Step E   │
                    │ 聚合结果 │
                    └────┬─────┘
                         │
                    ┌────▼─────┐
                    │  End     │
                    └──────────┘

4.2 步骤类型

步骤类型 说明
TASK 普通任务步骤,执行任意异步函数
CONDITION 条件分支,根据表达式决定是否执行后续步骤
PARALLEL 并行执行多个子步骤
WAIT 等待指定时间后继续
SUB_WORKFLOW 嵌套执行另一个工作流
TRANSACTION 发送链上交易
QUERY 执行链上查询
NOTIFICATION 发送通知

4.3 工作流引擎实现

import asyncio, time, json, uuid
from typing import Dict, List, Optional, Any, Callable, Set
from dataclasses import dataclass, field
from enum import Enum
from collections import deque
from msg_sdk import MsgClient, Wallet

class StepStatus(Enum):
    PENDING = "pending"; RUNNING = "running"; SUCCESS = "success"
    FAILED = "failed"; SKIPPED = "skipped"; TIMEOUT = "timeout"

class WorkflowStatus(Enum):
    PENDING = "pending"; RUNNING = "running"; SUCCESS = "success"
    FAILED = "failed"; PAUSED = "paused"; CANCELLED = "cancelled"

class StepType(Enum):
    TASK = "task"; CONDITION = "condition"; PARALLEL = "parallel"
    WAIT = "wait"; SUB_WORKFLOW = "sub_workflow"; TRANSACTION = "transaction"
    QUERY = "query"; NOTIFICATION = "notification"

@dataclass
class StepResult:
    step_id: str; status: StepStatus; output: Any = None
    error: Optional[str] = None; started_at: float = 0.0
    completed_at: Optional[float] = None; duration_ms: float = 0.0; retry_count: int = 0

@dataclass
class WorkflowContext:
    workflow_id: str; execution_id: str; input: dict
    results: Dict[str, StepResult] = field(default_factory=dict)
    variables: Dict[str, Any] = field(default_factory=dict)
    metadata: Dict[str, Any] = field(default_factory=dict)
    started_at: float = 0.0; completed_at: Optional[float] = None

@dataclass
class WorkflowStep:
    id: str; type: StepType; name: str = ""; description: str = ""
    task: Optional[Callable] = None; depends_on: List[str] = field(default_factory=list)
    timeout: int = 300; max_retries: int = 0; retry_delay: int = 5
    condition: Optional[str] = None; metadata: dict = field(default_factory=dict)

@dataclass
class Workflow:
    id: str; name: str = ""; description: str = ""; version: str = "1.0"
    steps: Dict[str, WorkflowStep] = field(default_factory=dict)
    start_step: str = ""; end_steps: List[str] = field(default_factory=list)
    tags: List[str] = field(default_factory=list); timeout: int = 3600
    max_concurrency: int = 10; error_handling: str = "stop"
    created_at: float = 0.0; updated_at: float = 0.0

class WorkflowEngine:
    def __init__(self, agent_id: str, client: MsgClient, wallet: Wallet, record_on_chain: bool = True):
        self.agent_id = agent_id; self.client = client; self.wallet = wallet
        self.record_on_chain = record_on_chain
        self.workflows: Dict[str, Workflow] = {}
        self.active_executions: Dict[str, WorkflowContext] = {}
        self.execution_history: List[dict] = []

    async def define_workflow(self, workflow_id: str, steps: List[WorkflowStep],
                              name: str = "", description: str = "", version: str = "1.0",
                              timeout: int = 3600, error_handling: str = "stop",
                              tags: Optional[List[str]] = None) -> Workflow:
        if workflow_id in self.workflows:
            raise ValueError(f"Workflow '{workflow_id}' already exists")
        step_dict = {s.id: s for s in steps}
        self._validate_dag(step_dict)
        start = self._find_start_step(step_dict)
        end = self._find_end_steps(step_dict)
        wf = Workflow(id=workflow_id, name=name or workflow_id, description=description,
                     version=version, steps=step_dict, start_step=start, end_steps=end,
                     tags=tags or [], timeout=timeout, error_handling=error_handling,
                     created_at=time.time(), updated_at=time.time())
        self.workflows[workflow_id] = wf
        return wf

    async def remove_workflow(self, workflow_id: str) -> bool:
        return self.workflows.pop(workflow_id, None) is not None

    def get_workflow(self, workflow_id: str) -> Optional[Workflow]:
        return self.workflows.get(workflow_id)

    def _validate_dag(self, steps: Dict[str, WorkflowStep]):
        visited, rec = set(), set()
        def dfs(sid: str):
            visited.add(sid); rec.add(sid)
            for dep in steps.get(sid, WorkflowStep(id='', type=StepType.TASK)).depends_on:
                if dep not in steps:
                    raise ValueError(f"Step '{sid}' depends on unknown step '{dep}'")
                if dep not in visited:
                    dfs(dep)
                elif dep in rec:
                    raise ValueError(f"Circular dependency: {sid} -> {dep}")
            rec.discard(sid)
        for sid in steps:
            if sid not in visited:
                dfs(sid)

    def _find_start_step(self, steps: Dict[str, WorkflowStep]) -> str:
        all_deps = set(d for s in steps.values() for d in s.depends_on)
        for sid in steps:
            if sid not in all_deps:
                return sid
        return next(iter(steps.keys()))

    def _find_end_steps(self, steps: Dict[str, WorkflowStep]) -> List[str]:
        depended = set(d for s in steps.values() for d in s.depends_on)
        return [sid for sid in steps if sid not in depended]

    async def execute_workflow(self, workflow_id: str, input_data: Optional[dict] = None,
                                execution_id: Optional[str] = None, parent_context=None) -> dict:
        wf = self.workflows.get(workflow_id)
        if not wf:
            raise ValueError(f"Workflow '{workflow_id}' not found")
        eid = execution_id or f"{workflow_id}_{uuid.uuid4().hex[:8]}"
        ctx = WorkflowContext(workflow_id=workflow_id, execution_id=eid,
                             input=input_data or {}, started_at=time.time())
        self.active_executions[eid] = ctx
        if self.record_on_chain:
            await self._record_execution_start(wf, ctx)
        try:
            await self._execute_dag(wf, ctx)
            ctx.completed_at = time.time()
            if self.record_on_chain:
                await self._record_execution_complete(wf, ctx)
            self.execution_history.append({"execution_id": eid, "workflow_id": workflow_id,
                                         "status": "success", "started_at": ctx.started_at,
                                         "completed_at": ctx.completed_at})
            return {"execution_id": eid, "status": WorkflowStatus.SUCCESS.value,
                    "results": {sid: {"status": sr.status.value, "output": sr.output}
                               for sid, sr in ctx.results.items()},
                    "duration_ms": int((ctx.completed_at - ctx.started_at) * 1000)}
        except Exception as e:
            ctx.completed_at = time.time()
            if self.record_on_chain:
                await self._record_execution_failed(wf, ctx, str(e))
            return {"execution_id": eid, "status": WorkflowStatus.FAILED.value, "error": str(e),
                    "results": {sid: {"status": sr.status.value, "output": sr.output, "error": sr.error}
                               for sid, sr in ctx.results.items()},
                    "duration_ms": int((ctx.completed_at - ctx.started_at) * 1000)}
        finally:
            self.active_executions.pop(eid, None)

    async def _execute_dag(self, workflow: Workflow, context: WorkflowContext):
        steps = workflow.steps
        sorted_steps = self._topological_sort(steps)
        pending, running, completed = set(sorted_steps), set(), set()
        sem = asyncio.Semaphore(workflow.max_concurrency)

        async def exec_step(sid: str):
            async with sem:
                r = await self._run_step(steps[sid], context)
                context.results[sid] = r
                completed.add(sid)

        ready = {workflow.start_step}
        while ready:
            tasks = []
            for sid in ready:
                if sid not in running:
                    running.add(sid)
                    tasks.append(exec_step(sid))
            if not tasks:
                break
            await asyncio.gather(*tasks, return_exceptions=True)
            new_ready = set()
            for sid in list(running):
                if sid in completed:
                    running.discard(sid); pending.discard(sid)
                    for cid in pending:
                        c = steps[cid]
                        if all(d in completed for d in c.depends_on):
                            if c.condition:
                                if not self._eval_step_condition(c.condition, context):
                                    context.results[cid] = StepResult(step_id=cid, status=StepStatus.SKIPPED)
                                    completed.add(cid)
                                    continue
                            new_ready.add(cid)
            ready = new_ready
            if time.time() - context.started_at > workflow.timeout:
                raise TimeoutError(f"Workflow '{workflow.id}' timed out")

    def _topological_sort(self, steps: Dict[str, WorkflowStep]) -> List[str]:
        indeg = {sid: 0 for sid in steps}
        adj = {sid: [] for sid in steps}
        for sid, s in steps.items():
            for d in s.depends_on:
                if d in adj:
                    adj[d].append(sid)
                indeg[sid] = indeg.get(sid, 0) + 1
        q = deque(sid for sid, d in indeg.items() if d == 0)
        res = []
        while q:
            n = q.popleft(); res.append(n)
            for nb in adj.get(n, []):
                indeg[nb] -= 1
                if indeg[nb] == 0:
                    q.append(nb)
        if len(res) != len(steps):
            raise ValueError("Graph has cycles")
        return res

    async def _run_step(self, step: WorkflowStep, context: WorkflowContext) -> StepResult:
        r = StepResult(step_id=step.id, status=StepStatus.PENDING, started_at=time.time())
        if step.type == StepType.WAIT:
            await asyncio.sleep(step.metadata.get("duration", 0))
            r.status = StepStatus.SUCCESS; r.completed_at = time.time(); r.duration_ms = 0
            return r
        if not step.task:
            r.status = StepStatus.SUCCESS; r.completed_at = time.time(); r.duration_ms = 0
            return r
        for attempt in range(step.max_retries + 1):
            if attempt > 0:
                await asyncio.sleep(step.retry_delay)
            try:
                output = await asyncio.wait_for(step.task(context.input, context.results), timeout=step.timeout)
                r.status = StepStatus.SUCCESS; r.output = output; r.retry_count = attempt
                r.completed_at = time.time(); r.duration_ms = (r.completed_at - r.started_at) * 1000
                return r
            except asyncio.TimeoutError:
                last_error = f"Timeout after {step.timeout}s"
            except Exception as e:
                last_error = str(e)
        r.status = StepStatus.FAILED; r.error = last_error
        r.completed_at = time.time(); r.duration_ms = (r.completed_at - r.started_at) * 1000
        return r

    def _eval_step_condition(self, condition: str, ctx: WorkflowContext) -> bool:
        try:
            results = {sid: sr.output for sid, sr in ctx.results.items() if sr.status == StepStatus.SUCCESS}
            return bool(eval(condition, {"__builtins__": {}}, {"results": results, "vars": ctx.variables, "input": ctx.input}))
        except Exception:
            return False

    async def _record_execution_start(self, wf: Workflow, ctx: WorkflowContext):
        try:
            msg = {"@type": "/msg.agent.v1.MsgWorkflowStart", "agent": self.agent_id,
                   "sender": str(self.wallet.address), "workflow_id": wf.id,
                   "execution_id": ctx.execution_id, "version": wf.version,
                   "input": json.dumps(ctx.input), "timestamp": str(int(ctx.started_at))}
            await self.client.execute_contract(self.agent_id, msg, sender=self.wallet)
        except Exception:
            pass

    async def _record_execution_complete(self, wf: Workflow, ctx: WorkflowContext):
        try:
            msg = {"@type": "/msg.agent.v1.MsgWorkflowComplete", "agent": self.agent_id,
                   "sender": str(self.wallet.address), "workflow_id": wf.id,
                   "execution_id": ctx.execution_id, "status": "success",
                   "duration_ms": str(int((ctx.completed_at - ctx.started_at) * 1000))}
            await self.client.execute_contract(self.agent_id, msg, sender=self.wallet)
        except Exception:
            pass

    async def _record_execution_failed(self, wf: Workflow, ctx: WorkflowContext, error: str):
        try:
            msg = {"@type": "/msg.agent.v1.MsgWorkflowFailed", "agent": self.agent_id,
                   "sender": str(self.wallet.address), "workflow_id": wf.id,
                   "execution_id": ctx.execution_id, "error": error}
            await self.client.execute_contract(self.agent_id, msg, sender=self.wallet)
        except Exception:
            pass

    def get_stats(self) -> dict:
        return {"agent_id": self.agent_id, "total_workflows": len(self.workflows),
                "active_executions": len(self.active_executions),
                "total_executions": len(self.execution_history)}

4.4 工作流使用示例

async def demo_workflow():
    engine = WorkflowEngine("msg1agent...", MsgClient("https://rpc.msg-chain-1.zone"), Wallet.from_mnemonic("..."))
    steps = [
        WorkflowStep(id="fetch", type=StepType.QUERY, name="获取池子", task=lambda ctx, res: {"pools": ["pool1", "pool2"]}),
        WorkflowStep(id="analyze", type=StepType.TASK, name="分析", depends_on=["fetch"],
                    task=lambda ctx, res: {"best_apy": 25, "best_pool": "pool1"}),
        WorkflowStep(id="migrate", type=StepType.TRANSACTION, name="迁移", depends_on=["analyze"],
                    task=lambda ctx, res: {"migrated": True}, max_retries=3, timeout=120),
    ]
    await engine.define_workflow("yield_optimizer", steps, name="收益优化器")
    result = await engine.execute_workflow("yield_optimizer", {"min_apy": 15})
    print(json.dumps(result, indent=2, default=str))

5. 条件与决策引擎

5.1 决策引擎实现

import time, json
from typing import Dict, List, Optional, Any, Callable, Tuple
from dataclasses import dataclass, field
from enum import Enum

class ActionPriority(Enum):
    CRITICAL = 0; HIGH = 1; MEDIUM = 2; LOW = 3; BACKGROUND = 4

@dataclass
class RuleAction:
    action_id: str; name: str; description: str = ""
    priority: ActionPriority = ActionPriority.MEDIUM
    cooldown_seconds: int = 3600; max_consecutive: int = 10
    required_balance: Optional[str] = None; params: dict = field(default_factory=dict)

@dataclass
class DecisionRule:
    rule_id: str; name: str; description: str = ""; condition: str = ""
    actions: List[RuleAction] = field(default_factory=list); priority: int = 100
    cooldown_seconds: int = 0; enabled: bool = True
    created_at: float = 0.0; last_evaluated: Optional[float] = None
    eval_count: int = 0; trigger_count: int = 0

@dataclass
class AgentState:
    balances: Dict[str, str] = field(default_factory=dict)
    positions: Dict[str, dict] = field(default_factory=dict)
    pending_rewards: Dict[str, str] = field(default_factory=dict)
    social_mentions: int = 0; unread_dms: int = 0
    open_proposals: int = 0; voted_proposals: List[int] = field(default_factory=list)
    last_action_time: Optional[float] = None

class DecisionEngine:
    BUILTIN_RULES = {
        "auto_compound": {
            "name": "自动复投", "condition": "int(state.pending_rewards.get('umsg', '0')) > constants['MIN_COMPOUND']",
            "priority": 80, "cooldown": 21600,
            "actions": [{"action_id": "compound_rewards", "name": "复投奖励", "priority": "MEDIUM", "cooldown_seconds": 21600}],
        },
        "prevent_liquidation": {
            "name": "清算预防", "condition": "True", "priority": 10, "cooldown": 60,
            "actions": [{"action_id": "add_collateral", "name": "补充抵押品", "priority": "CRITICAL", "cooldown_seconds": 300}],
        },
        "vote_on_proposals": {
            "name": "治理投票", "condition": "state.open_proposals > len(state.voted_proposals)",
            "priority": 50, "cooldown": 60,
            "actions": [{"action_id": "auto_vote", "name": "自动投票", "priority": "HIGH", "cooldown_seconds": 60}],
        },
        "respond_to_mentions": {
            "name": "回复提及", "condition": "state.social_mentions > 0 or state.unread_dms > 0",
            "priority": 40, "cooldown": 300,
            "actions": [{"action_id": "process_mentions", "name": "处理社交提及", "priority": "MEDIUM", "cooldown_seconds": 300}],
        },
    }

    def __init__(self, agent_id: str, state_provider=None, action_executor=None):
        self.agent_id = agent_id
        self.state_provider = state_provider or (lambda: AgentState())
        self.action_executor = action_executor or (lambda a, p: None)
        self.rules: Dict[str, DecisionRule] = {}
        self.cooldowns: Dict[str, float] = {}
        self.action_history: List[dict] = []
        self.execution_counters: Dict[str, int] = {}
        self.constants = {"MIN_COMPOUND": 1_000_000, "REBALANCE_THRESHOLD": 0.05, "LIQUIDATION_THRESHOLD": 1.5}

    def add_rule(self, rule: DecisionRule):
        compile(rule.condition, "<string>", "eval")
        self.rules[rule.rule_id] = rule

    def add_builtin_rules(self, rule_ids: Optional[List[str]] = None):
        for rid in (rule_ids or list(self.BUILTIN_RULES.keys())):
            if rid in self.BUILTIN_RULES:
                t = self.BUILTIN_RULES[rid]
                actions = [RuleAction(**a) for a in t["actions"]]
                self.rules[rid] = DecisionRule(rule_id=rid, name=t["name"], condition=t["condition"],
                                               actions=actions, priority=t["priority"], cooldown_seconds=t["cooldown"],
                                               created_at=time.time())

    def disable_rule(self, rule_id: str):
        if rule_id in self.rules:
            self.rules[rule_id].enabled = False

    def enable_rule(self, rule_id: str):
        if rule_id in self.rules:
            self.rules[rule_id].enabled = True

    async def decide(self, state: Optional[AgentState] = None) -> List[RuleAction]:
        state = state or self.state_provider()
        triggered: List[Tuple[RuleAction, DecisionRule]] = []
        for rule in sorted(self.rules.values(), key=lambda r: r.priority):
            if not rule.enabled:
                continue
            rule.eval_count += 1; rule.last_evaluated = time.time()
            if self._in_cooldown(rule.rule_id, rule.cooldown_seconds):
                continue
            if not self._eval(rule.condition, state):
                continue
            rule.trigger_count += 1
            for action in rule.actions:
                if self._in_cooldown(action.action_id, action.cooldown_seconds):
                    continue
                if action.required_balance and int(state.balances.get("umsg", "0")) < int(action.required_balance):
                    continue
                if self.execution_counters.get(action.action_id, 0) >= action.max_consecutive:
                    continue
                triggered.append((action, rule))
        return self._prioritize(triggered)

    def _prioritize(self, actions: List[Tuple[RuleAction, DecisionRule]]) -> List[RuleAction]:
        order = {ActionPriority.CRITICAL: 0, ActionPriority.HIGH: 1, ActionPriority.MEDIUM: 2, ActionPriority.LOW: 3, ActionPriority.BACKGROUND: 4}
        return [a[0] for a in sorted(actions, key=lambda x: (order.get(x[0].priority, 99), x[1].priority))]

    def _in_cooldown(self, key: str, cd: int) -> bool:
        if cd <= 0:
            return False
        last = self.cooldowns.get(key, 0)
        if time.time() - last < cd:
            return True
        self.cooldowns[key] = time.time()
        return False

    def _eval(self, condition: str, state: AgentState) -> bool:
        try:
            return bool(eval(condition, {"__builtins__": {}},
                           {"state": state, "constants": self.constants, "int": int, "len": len, "min": min, "max": max, "abs": abs}))
        except Exception:
            return False

    async def execute(self, actions: List[RuleAction]) -> List[dict]:
        results = []
        for action in actions:
            try:
                r = await self.action_executor(action.action_id, action.params) if self.action_executor else None
                results.append({"action_id": action.action_id, "status": "success", "result": r})
            except Exception as e:
                results.append({"action_id": action.action_id, "status": "failed", "error": str(e)})
            self.execution_counters[action.action_id] = self.execution_counters.get(action.action_id, 0) + 1
            self.action_history.append({"action_id": action.action_id, "executed_at": time.time()})
        return results

    def get_stats(self) -> dict:
        return {"agent_id": self.agent_id, "total_rules": len(self.rules),
                "total_evaluations": sum(r.eval_count for r in self.rules.values()),
                "total_triggers": sum(r.trigger_count for r in self.rules.values()),
                "total_actions": len(self.action_history), "constants": self.constants}

6. 链上工作流记录

6.1 链上数据结构 (Rust)

use cosmwasm_std::Addr;
use serde::{Deserialize, Serialize};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum WorkflowStatus {
    #[serde(rename = "pending")] Pending,
    #[serde(rename = "running")] Running,
    #[serde(rename = "success")] Success,
    #[serde(rename = "failed")] Failed,
    #[serde(rename = "paused")] Paused,
    #[serde(rename = "cancelled")] Cancelled,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum StepStatus {
    #[serde(rename = "pending")] Pending,
    #[serde(rename = "running")] Running,
    #[serde(rename = "success")] Success,
    #[serde(rename = "failed")] Failed,
    #[serde(rename = "skipped")] Skipped,
    #[serde(rename = "timeout")] Timeout,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct WorkflowRecord {
    pub workflow_id: String,
    pub execution_id: String,
    pub agent: Addr,
    pub status: WorkflowStatus,
    pub steps_completed: Vec<String>,
    pub current_step: Option<String>,
    pub input: String,
    pub output: Option<String>,
    pub started_at: u64,
    pub completed_at: Option<u64>,
    pub error: Option<String>,
    pub version: String,
    pub gas_used: u64,
    pub duration_ms: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct StepRecord {
    pub execution_id: String,
    pub step_id: String,
    pub status: StepStatus,
    pub input: Option<String>,
    pub output: Option<String>,
    pub error: Option<String>,
    pub started_at: u64,
    pub completed_at: u64,
    pub retry_count: u32,
    pub duration_ms: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CronJobRecord {
    pub job_id: String,
    pub agent: Addr,
    pub cron_expr: String,
    pub name: String,
    pub enabled: bool,
    pub max_retries: u32,
    pub timeout: u32,
    pub tags: Vec<String>,
    pub created_at: u64,
    pub last_run: Option<u64>,
    pub next_run: Option<u64>,
    pub run_count: u64,
    pub success_count: u64,
    pub fail_count: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct ExecutionLog {
    pub execution_id: String,
    pub ref_id: String,
    pub exec_type: String,
    pub agent: Addr,
    pub status: String,
    pub started_at: u64,
    pub completed_at: Option<u64>,
    pub gas_used: u64,
    pub tx_hash: Option<String>,
    pub error: Option<String>,
    pub result: Option<String>,
    pub retry_count: u32,
}

6.2 链上查询 (Go)

package keeper

import (
    "time"
    storetypes "cosmossdk.io/store/types"
    "github.com/cosmos/cosmos-sdk/codec"
    sdk "github.com/cosmos/cosmos-sdk/types"
    "msg-chain/x/agent/types"
)

type WorkflowKeeper struct {
    storeKey storetypes.StoreKey
    cdc      codec.BinaryCodec
}

func (k WorkflowKeeper) SaveWorkflowRecord(ctx sdk.Context, record types.WorkflowRecord) {
    store := ctx.KVStore(k.storeKey)
    key := []byte("workflow/" + record.agent.String() + "/" + record.execution_id)
    store.Set(key, k.cdc.MustMarshal(&record))
}

func (k WorkflowKeeper) GetWorkflowRecord(ctx sdk.Context, agent sdk.AccAddress, executionID string) (*types.WorkflowRecord, error) {
    store := ctx.KVStore(k.storeKey)
    bz := store.Get([]byte("workflow/" + agent.String() + "/" + executionID))
    if bz == nil {
        return nil, types.ErrWorkflowNotFound
    }
    var record types.WorkflowRecord
    k.cdc.MustUnmarshal(bz, &record)
    return &record, nil
}

func (k WorkflowKeeper) ListWorkflowRecords(ctx sdk.Context, agent sdk.AccAddress, statusFilter string, limit int) []types.WorkflowRecord {
    store := ctx.KVStore(k.storeKey)
    prefix := []byte("workflow/" + agent.String() + "/")
    iter := store.Iterator(prefix, storetypes.PrefixEndBytes(prefix))
    defer iter.Close()
    var records []types.WorkflowRecord
    for ; iter.Valid(); iter.Next() {
        var r types.WorkflowRecord
        k.cdc.MustUnmarshal(iter.Value(), &r)
        if statusFilter == "" || r.Status.String() == statusFilter {
            records = append(records, r)
            if limit > 0 && len(records) >= limit {
                break
            }
        }
    }
    return records
}

func (k WorkflowKeeper) GetAgentWorkflowStats(ctx sdk.Context, agent sdk.AccAddress) (*types.AgentWorkflowStats, error) {
    records := k.ListWorkflowRecords(ctx, agent, "", 1000)
    stats := &types.AgentWorkflowStats{TotalExecutions: uint64(len(records))}
    cutoff := time.Now().Add(-24 * time.Hour)
    for _, r := range records {
        if r.Status == types.WorkflowStatus_SUCCESS {
            stats.Successful++
        } else if r.Status == types.WorkflowStatus_FAILED {
            stats.Failed++
        }
        stats.TotalGasUsed += r.GasUsed
        if time.Unix(int64(r.StartedAt), 0).After(cutoff) {
            stats.Last24hCount++
        }
    }
    if stats.TotalExecutions > 0 {
        stats.AvgDurationMs = stats.TotalGasUsed / stats.TotalExecutions
    }
    return stats, nil
}

6.3 查询客户端 (Python)

class WorkflowRecordClient:
    def __init__(self, client: MsgClient, agent_address: str):
        self.client = client; self.agent_address = agent_address

    async def get_workflow_record(self, execution_id: str) -> Optional[dict]:
        r = await self.client.query_contract_smart(self.agent_address, {"workflow_record": {"execution_id": execution_id}})
        return r.get("record")

    async def list_records(self, status: str = "", limit: int = 20) -> List[dict]:
        r = await self.client.query_contract_smart(self.agent_address, {"workflow_records": {"status_filter": status, "limit": limit}})
        return r.get("records", [])

    async def get_job_record(self, job_id: str) -> Optional[dict]:
        r = await self.client.query_contract_smart(self.agent_address, {"cron_job": {"job_id": job_id}})
        return r.get("record")

    async def get_stats(self) -> dict:
        r = await self.client.query_contract_smart(self.agent_address, {"workflow_stats": {}})
        return r.get("stats", {})

7. 自动化策略模板

7.1 预置策略

STRATEGIES = {
    "yield_optimizer": {
        "name": "收益优化器",
        "description": "自动管理流动性挖矿奖励的复投和资金迁移",
        "schedule": "0 */6 * * *",
        "triggers": ["reward_claim", "rebalance"],
        "workflow": ["harvest_rewards", "check_pools", "rebalance_if_needed"],
        "decision_rules": ["auto_compound"],
        "params": {"min_compound": 1_000_000, "rebalance_threshold": 0.05, "max_pools": 5},
    },
    "social_manager": {
        "name": "社交管理器",
        "description": "自动回复社交提及和私信",
        "schedule": "*/15 * * * *",
        "triggers": ["social_mention", "direct_message"],
        "workflow": ["fetch_mentions", "analyze_intent", "draft_response", "post_response"],
        "decision_rules": ["respond_to_mentions"],
        "params": {"batch_size": 20, "reply_style": "helpful", "max_daily_replies": 200},
    },
    "liquidation_guard": {
        "name": "清算监控器",
        "description": "监控借贷仓位健康因子,在清算前自动补充抵押品",
        "schedule": "*/1 * * * *",
        "triggers": ["new_block", "balance_threshold"],
        "workflow": ["check_health_factors", "calculate_required_collateral", "add_collateral_if_needed"],
        "decision_rules": ["prevent_liquidation"],
        "params": {"health_threshold": 1.5, "target_health": 2.5, "max_repay_percent": 0.3},
    },
    "data_collector": {
        "name": "链上数据收集器",
        "description": "定期收集链上数据并存储供分析使用",
        "schedule": "0 0 * * *",
        "triggers": ["new_block"],
        "workflow": ["fetch_onchain_data", "process_data", "store_results"],
        "decision_rules": [],
        "params": {"data_types": ["pools", "prices", "volume", "tvl"], "retention_days": 90},
    },
    "governance_agent": {
        "name": "治理投票代理",
        "description": "自动分析治理提案并按策略投票",
        "schedule": "0 */2 * * *",
        "triggers": ["governance_proposal"],
        "workflow": ["fetch_new_proposals", "analyze_proposals", "execute_votes"],
        "decision_rules": ["vote_on_proposals"],
        "params": {"voting_strategy": "follow_largest_validator", "auto_vote": True},
    },
    "airdrop_hunter": {
        "name": "空投猎手",
        "description": "自动检查并领取可用的空投",
        "schedule": "0 */12 * * *",
        "triggers": ["balance_threshold"],
        "workflow": ["check_airdrop_eligibility", "claim_airdrops", "log_results"],
        "decision_rules": [],
        "params": {"check_interval_hours": 12, "auto_claim": True, "min_gas_reserve": "500000"},
    },
}

7.2 策略部署器

class StrategyDeployer:
    def __init__(self, agent_id: str, client: MsgClient, wallet: Wallet):
        self.agent_id = agent_id; self.client = client; self.wallet = wallet
        self.scheduler: Optional[AgentScheduler] = None
        self.trigger: Optional[EventTrigger] = None
        self.engine: Optional[WorkflowEngine] = None
        self.decider: Optional[DecisionEngine] = None

    async def deploy(self, strategy_id: str, custom_params: Optional[dict] = None):
        s = STRATEGIES.get(strategy_id)
        if not s:
            raise ValueError(f"Unknown strategy: {strategy_id}")
        params = {**s["params"], **(custom_params or {})}
        self.scheduler = AgentScheduler(self.agent_id, self.client, self.wallet)
        self.trigger = EventTrigger(self.agent_id, self.client, self.wallet, scheduler=self.scheduler)
        self.engine = WorkflowEngine(self.agent_id, self.client, self.wallet)
        self.decider = DecisionEngine(self.agent_id)

        await self.scheduler.add_cron_job(f"{strategy_id}_main", s["schedule"],
            self._make_task(s, params), name=s["name"])

        trigger_map = {
            "reward_claim": (TriggerType.CONTRACT_EVENT, {"event_name": "wasm-reward"}),
            "rebalance": (TriggerType.BALANCE_THRESHOLD, {"threshold": "0", "direction": "above"}),
            "social_mention": (TriggerType.CONTRACT_EVENT, {"event_name": "wasm-mention"}),
            "new_block": (TriggerType.NEW_BLOCK, {"interval": 1}),
            "balance_threshold": (TriggerType.BALANCE_THRESHOLD, {"direction": "below"}),
            "governance_proposal": (TriggerType.GOVERNANCE_PROPOSAL, {}),
        }
        for t in s["triggers"]:
            if t in trigger_map:
                tt, cfg = trigger_map[t]
                await self.trigger.add_rule(f"{strategy_id}_{t}", tt, cfg)
        for rid in s["decision_rules"]:
            self.decider.add_builtin_rules([rid])
        print(f"Deployed: {strategy_id}")

    def _make_task(self, strategy: dict, params: dict):
        async def task():
            actions = await self.decider.decide()
            if actions:
                await self.decider.execute(actions)
        return task

    async def start(self):
        if self.scheduler: await self.scheduler.start()
        if self.trigger: await self.trigger.start()

    async def stop(self):
        if self.trigger: await self.trigger.stop()
        if self.scheduler: await self.scheduler.stop()

8. 完整示例

8.1 跨 DeFi 协议收益优化 Agent

#!/usr/bin/env python3
import asyncio, json, time, os
from msg_sdk import MsgClient, Wallet

AGENT_ID = os.getenv("AGENT_ADDRESS", "msg1agent...")
MNEMONIC = os.getenv("AGENT_MNEMONIC", "")
RPC = os.getenv("RPC_ENDPOINT", "https://rpc.msg-chain-1.zone")

CONFIG = {
    "compound": {"enabled": True, "interval_hours": 6, "min_reward": 1_000_000},
    "pools": {"enabled": True, "check_interval_hours": 12, "apy_deviation_threshold": 0.05,
              "tracked": [{"id": "msg_usdc", "addr": "msg1poolusdc..."}]},
    "liquidation_guard": {"enabled": True, "health_threshold": 1.5, "target_health": 2.5},
}

class YieldOptimizerAgent:
    def __init__(self):
        self.client = MsgClient(rpc_endpoint=RPC)
        self.wallet = Wallet.from_mnemonic(MNEMONIC)
        self.scheduler: Optional[AgentScheduler] = None
        self.trigger: Optional[EventTrigger] = None
        self.engine: Optional[WorkflowEngine] = None
        self.decider: Optional[DecisionEngine] = None
        self.stats = {"compounds": 0, "migrations": 0, "errors": 0}

    async def initialize(self):
        print(f"[Agent] Initializing {AGENT_ID}")
        self.scheduler = AgentScheduler(AGENT_ID, self.client, self.wallet)
        self.trigger = EventTrigger(AGENT_ID, self.client, self.wallet, scheduler=self.scheduler)
        self.engine = WorkflowEngine(AGENT_ID, self.client, self.wallet)
        self.decider = DecisionEngine(AGENT_ID, state_provider=self._get_state, action_executor=self._exec_action)
        self.decider.add_builtin_rules(["auto_compound", "prevent_liquidation", "vote_on_proposals"])

        if CONFIG["compound"]["enabled"]:
            await self.scheduler.add_cron_job("compound", "0 */6 * * *", self._compound,
                name="自动复投", max_retries=3, timeout=300)
        if CONFIG["pools"]["enabled"]:
            await self.scheduler.add_cron_job("pool_analysis", "0 */12 * * *", self._analyze_pools,
                name="池子分析", max_retries=2, timeout=600)
        if CONFIG["liquidation_guard"]["enabled"]:
            await self.scheduler.add_cron_job("health_check", "*/1 * * * *", self._health_check,
                name="健康检查", max_retries=1, timeout=30)
            await self.trigger.add_rule("balance_drop", TriggerType.BALANCE_THRESHOLD,
                {"direction": "below", "threshold": "5000000", "token": "umsg"},
                name="余额阈值", cooldown_seconds=3600)

        await self.engine.define_workflow("migrate_pool", [
            WorkflowStep(id="query", type=StepType.QUERY, task=lambda ctx, r: {"best": "pool1"}, name="查询"),
            WorkflowStep(id="check", type=StepType.CONDITION, depends_on=["query"],
                        condition="results['query'].get('best', '') != ''"),
            WorkflowStep(id="execute", type=StepType.TRANSACTION, depends_on=["check"],
                        task=lambda ctx, r: self._do_migrate(r["query"].output["best"]),
                        max_retries=3, name="执行迁移"),
        ], name="池子迁移")
        print("[Agent] Initialized")

    async def _compound(self):
        rewards = await self.client.get_balance(str(self.wallet.address), "umsg")
        amt = int(rewards.amount)
        if amt < CONFIG["compound"]["min_reward"]:
            return {"skipped": True, "reason": "below threshold"}
        tx = await self.client.execute_contract(self.wallet.address,
            {"@type": "/msg.agent.v1.MsgCompound", "sender": str(self.wallet.address)})
        self.stats["compounds"] += 1
        return {"compounded": True, "amount": amt, "tx": tx.tx_hash}

    async def _analyze_pools(self):
        results = []
        for pool in CONFIG["pools"]["tracked"]:
            try:
                r = await self.client.query_contract_smart(pool["addr"], {"apy": {}})
                results.append({"id": pool["id"], "apy": float(r.get("apy", 0))})
            except Exception as e:
                results.append({"id": pool["id"], "error": str(e)})
        return {"pools": results}

    async def _health_check(self):
        return {"healthy": True, "timestamp": time.time()}

    async def _do_migrate(self, pool_id: str):
        tx = await self.client.execute_contract(self.wallet.address,
            {"@type": "/msg.agent.v1.MsgMigrate", "pool_id": pool_id, "sender": str(self.wallet.address)})
        self.stats["migrations"] += 1
        return {"migrated": True, "pool": pool_id, "tx": tx.tx_hash}

    def _get_state(self):
        return AgentState()

    async def _exec_action(self, action_id: str, params: dict):
        if action_id == "compound_rewards":
            return await self._compound()
        return {"action": action_id, "status": "simulated"}

    async def run(self):
        await self.initialize()
        await self.scheduler.start()
        await self.trigger.start()
        print("[Agent] Running...")
        try:
            await asyncio.sleep(86400)
        finally:
            await self.scheduler.stop()
            await self.trigger.stop()
            print(f"[Agent] Stats: {json.dumps(self.stats)}")

async def main():
    agent = YieldOptimizerAgent()
    await agent.run()

if __name__ == "__main__":
    asyncio.run(main())

9. 附录

A. Cron 表达式速查表

表达式 说明
* * * * * 每分钟
*/5 * * * * 每 5 分钟
*/15 * * * * 每 15 分钟
*/30 * * * * 每 30 分钟
0 * * * * 每小时
0 */2 * * * 每 2 小时
0 */6 * * * 每 6 小时
0 */12 * * * 每 12 小时
0 0 * * * 每天 00:00 UTC
0 12 * * * 每天 12:00 UTC
0 0 * * 0 每周日
0 0 1 * * 每月 1 日
*/30 9-18 * * 1-5 工作日 9-18 点每 30 分钟
0 0 1,15 * * 每月 1 日和 15 日

B. 消息类型索引

消息类型 用途
/msg.agent.v1.MsgRegisterCronJob 注册定时任务
/msg.agent.v1.MsgRecordExecution 记录执行结果
/msg.agent.v1.MsgSchedulerEvent 调度器事件
/msg.agent.v1.MsgAcquireLock 获取分布式锁
/msg.agent.v1.MsgReleaseLock 释放分布式锁
/msg.agent.v1.MsgWorkflowStart 工作流开始
/msg.agent.v1.MsgWorkflowComplete 工作流完成
/msg.agent.v1.MsgWorkflowFailed 工作流失败
/msg.agent.v1.MsgWorkflowCheckpoint 工作流检查点
/msg.agent.v1.MsgTriggerEvent 事件触发记录

C. 配置文件参考

{
  "agent": {
    "id": "msg1agent...",
    "rpc": "https://rpc.msg-chain-1.zone",
    "denom": "umsg",
    "prefix": "msg"
  },
  "scheduler": {
    "check_interval": 60,
    "max_concurrent": 10,
    "record_on_chain": true,
    "missed_job_recovery_hours": 24
  },
  "triggers": {
    "rate_limit_window_seconds": 3600,
    "max_firings_per_window": 10,
    "default_cooldown_seconds": 60
  },
  "workflows": {
    "max_timeout_seconds": 3600,
    "max_retries": 3,
    "error_handling": "stop"
  },
  "decision_engine": {
    "min_compound": 1000000,
    "rebalance_threshold": 0.05,
    "liquidation_threshold": 1.5
  }
}

D. 最佳实践

  1. 幂等性: 确保所有任务和行动是幂等的——重复执行产生相同结果
  2. Gas 预留: 始终为紧急操作(如清算预防)预留足够的 Gas
  3. 冷却优先: 避免在冷却期内重复触发同一条规则
  4. 检查点: 对长时间工作流启用链上检查点,支持断点续跑
  5. 事件去重: 使用事件 ID 去重,防止同一事件触发多次
  6. 分布式锁: 多个 Agent 实例时使用链上锁防止重复执行
  7. 渐进重试: 失败的 Job 使用指数退避冷却
  8. 监控告警: 注册 on_error 回调,将错误发送到监控系统
  9. 状态快照: 在执行关键操作前保存状态快照
  10. 版本控制: 工作流定义使用语义版本号,支持向后兼容

本文档为 MSG Chain AI Agent 工作流自动化系统的完整参考指南。
链标识: msg-chain-1 | Bech32 前缀: msg
许可证: Apache-2.0