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. 概述
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. 最佳实践
- 幂等性: 确保所有任务和行动是幂等的——重复执行产生相同结果
- Gas 预留: 始终为紧急操作(如清算预防)预留足够的 Gas
- 冷却优先: 避免在冷却期内重复触发同一条规则
- 检查点: 对长时间工作流启用链上检查点,支持断点续跑
- 事件去重: 使用事件 ID 去重,防止同一事件触发多次
- 分布式锁: 多个 Agent 实例时使用链上锁防止重复执行
- 渐进重试: 失败的 Job 使用指数退避冷却
- 监控告警: 注册 on_error 回调,将错误发送到监控系统
- 状态快照: 在执行关键操作前保存状态快照
- 版本控制: 工作流定义使用语义版本号,支持向后兼容
本文档为 MSG Chain AI Agent 工作流自动化系统的完整参考指南。
链标识:msg-chain-1| Bech32 前缀:msg
许可证: Apache-2.0
