axon-llm::swarm 多 Agent 编排¶
适用版本:
axon-llmv0.6.0+ 状态:已实现(0.6.0 工作流 B + 0.11.0 投票共识) 上游 plan:docs/superpowers/plans/2026-07-18-axon-quant-0.6.0.md§3,docs/superpowers/specs/2026-08-01-axon-quant-0.11.0-multi-agent-design.md
SwarmOrchestrator 把 4 类 agent (Market / Risk / Execution / Audit) 串成可运行 pipeline, 配合 HarnessBridge 做最终裁决,TradingBackend 真下单到 MockTradingBackend / 交易所。 Python 端可直接构造 + 启动 + 注入信号,完整可观测。
1. 4-Agent 架构¶
┌────────────────────────────┐
│ SwarmOrchestrator │
│ ┌──────────────────────┐ │
│ │ run_loop_arc │ │◀──── inject_market_signal
│ │ dispatch() │ │ inject_vote_response
│ └──────────────────────┘ │
│ ConsensusManager + Stats │
└─────┬──────┬──────┬─────────┘
│ │ │ (mpsc inbox)
┌──────────────┘ │ └────────────┐
│ │ │
┌────▼────┐ ┌─────▼─────┐ ┌────▼────┐
│ Market │ ───▶ │ Risk │ ───▶ │Execution│ ──▶ PlaceOrderTool
│ Agent │ Market │ Agent │ Risk │ Agent │ QueryPortfolioTool
│ │ Signal │ │ Assess │ │ TradingBackend
└─────────┘ └───────────┘ └─────────┘
│ │
│ ▼
│ ┌──────────┐
│ │ Audit │
│ │ Agent │
│ └──────────┘
▼
MarketDataSource (Mock / WS / CSV)
关键设计:
- 每个 agent 持
Box<dyn DeclarativeAgentRunner>,复用 DeclarativeAgent 抽象 - Agent 间通信用
tokio::mpsc,不用共享状态 - Orchestrator 主循环监听 inbox,按
MessageContent路由 HarnessBridge由 SwarmOrchestrator 持有 + 各 agent 共享同一份Arc
2. 消息路由表(run_loop 内)¶
| 收到的消息 | 处理动作 |
|---|---|
MarketAnalysis(signal) | 创建 TradeDecision 投票,广播 VoteRequest 给 Risk + Execution |
RiskAssessment{approved=true} | 转发给 Execution agent 生成 ExecutionRequest |
RiskAssessment{approved=false} | 广播给 Audit 记录拒绝原因 |
ExecutionResult(result) | 广播给 Audit 审计执行结果 |
VoteResult{passed=true} | 调 HarnessBridge.adjudicate() 做最终裁决 |
VoteResponse(response) | 投到 ConsensusManager,达法定人数回 VoteResult |
Shutdown | 设置 shutdown_requested=true,主循环退出 |
Heartbeat / 其他 | 忽略 |
投票共识与 Harness 集成¶
VoteResponse → ConsensusManager
│
▼ SimpleMajority (≥2 votes)
VoteResult{passed=true}
│
▼
HarnessBridge.adjudicate(intent, ctx)
│
┌────────────┼────────────┬────────────┐
▼ ▼ ▼ ▼
Approved Rejected CircuitBreak NeedRevision
(执行) (审计) (shutdown) (回滚重分析)
HarnessBridge 缺省时降级为 Adjudication::Approved(零侵入模式,投票通过即批准)。
3. 关键模块¶
3.1 DeclarativeAgentRunner trait¶
#[async_trait]
pub trait DeclarativeAgentRunner: Send + Sync {
fn id(&self) -> &AgentId;
fn role(&self) -> AgentRole;
fn status(&self) -> AgentStatus;
async fn handle_message(&mut self, msg: AgentMessage) -> Result<RunnerOutput, SwarmError>;
}
- Object safety:trait 满足 object-safety,允许
Arc<dyn DeclarativeAgentRunner>跨 task 持有 - Sync 约束:
Status返回Copy,handle_message&mut self+ 异步,允许并发
3.2 4 个 Agent¶
| Agent | 输入 | 输出 | 配置 |
|---|---|---|---|
MarketAgent | MarketDataSource (tick) | MarketSignal | symbols + price_change_threshold |
RiskAgent | MarketSignal | RiskAssessment | RiskAgentConfig (默认阈值) |
ExecutionAgent | RiskAssessment{approved=true} | ExecutionResult (通过 PlaceOrderTool) | TradingTools { place_order, query_portfolio } |
AuditAgent | ExecutionResult | (审计记录) | AuditAgentConfig (默认阈值) |
3.3 PaperTradingBackend¶
PaperTradingBackend 实现了 Stage K TradingBackend trait,模拟真实交易:
- 滑点:
slippage_bps(基点,买入上浮/卖出下浮) - 手续费:
commission_bps(按 notional 收) - 状态:
cash+positions: HashMap<symbol, (qty, entry_price)>+last_prices - 价格更新:
place_order后last_prices[symbol] = fill_price - 现金流:Buy 扣
notional * (1 + commission_bps),Sell 加notional * (1 - commission_bps)
get_balance() 返回 cash + Σ(qty * last_price)(实时 NAV),get_positions() 返回 (symbol, qty, entry_price, current_price, unrealized_pnl)。
3.4 SwarmOrchestrator::run_loop_arc¶
Arc<TokioMutex<SwarmOrchestrator>> 跨 owner 共享,主循环:
pub async fn run_loop_arc(
orchestrator: Arc<TokioMutex<Self>>,
mut inbox_rx: mpsc::Receiver<AgentMessage>,
) {
let tick = Duration::from_millis(guard.config.loop_tick_ms);
loop {
if guard.shutdown_requested { break; }
let next = timeout(tick, inbox_rx.recv()).await;
match next {
Ok(Some(msg)) => dispatch(msg).await,
Ok(None) => break, // channel 关闭
Err(_) => continue, // timeout,检查 shutdown
}
}
}
退出条件:Shutdown 消息 / request_shutdown() / inbox channel 关闭(所有 agent outbox drop)。
4. Python 接入¶
4.1 典型用法¶
from axon_quant.llm import (
SwarmConfig, SwarmOrchestrator, MarketSignal, SignalType,
TradingTools,
)
from axon_quant.trading import (
MockTradingBackend, PlaceOrderTool, QueryPortfolioTool, RiskLimits,
)
# 1. 构造 orchestrator
config = SwarmConfig(vote_timeout_ms=5000, loop_tick_ms=100)
orch = SwarmOrchestrator(config)
# 2. 注册 4 类 agent(register_*_agent 自动启动 run_loop)
orch.register_market_agent(agent_id="m0", symbols=["BTC-USDT"])
orch.register_risk_agent(agent_id="r0")
# ExecutionAgent 必须传 tools(否则走 mock 模式)
backend = MockTradingBackend()
risk = RiskLimits(allowed_symbols=["BTC-USDT"])
place = PlaceOrderTool(backend=backend, mode="dry_run", risk=risk)
query = QueryPortfolioTool(backend=backend)
tools = TradingTools(place_order=place, query_portfolio=query)
orch.register_execution_agent(agent_id="e0", tools=tools)
orch.register_audit_agent(agent_id="a0")
# 3. 注入 MarketSignal → orchestrator 触发投票 + 转发
orch.inject_market_signal(MarketSignal(
symbol="BTC-USDT",
signal_type=SignalType.Buy,
confidence=0.9,
reasoning="momentum breakout",
))
# 4. 读统计
import time; time.sleep(0.5)
stats = orch.stats()
print(stats["market_signals"], stats["votes_created"])
# 5. 关闭
orch.stop()
4.2 完整 Pipeline demo¶
参见 examples/18_harness/swarm_demo.py。
5. 测试覆盖¶
| 测试文件 | 数量 | 内容 |
|---|---|---|
python/tests/test_swarm_pipeline_e2e.py | 25/25 ✅ | 枚举/数据结构 + 4 类 agent 注册 + lifecycle + inject + stats |
crates/axon-llm/src/swarm/ lib unittests | 全部通过 | DeclarativeAgentRunner / Orchestrator / Vote / 4 agent / market_data |
crates/axon-llm/src/trading/paper_backend.rs lib unittests | 全部通过 | place_order / balance / position 滑点+手续费验证 |
总测试数:322 lib unittests + 74 integration + 3 doctests + 25 Python E2E = 424 全过。
6. 现状与未实现项(基于 0.6.0)¶
已完成的 0.3.x / 0.4.x 路线图¶
PlaceOrderTool三模式:DryRun/TwoPhase/Direct三种SafetyMode全部实现。DryRun为默认安全模式(有意设计,防止 LLM 直发订单)。TwoPhase内部用pending: Mutex<HashMap<token, PendingOrder>>跟踪待确认订单,4 个 e2e 测试覆盖:首调用返回confirm_token/ 二次带 token 真发 / 错误 token 拒绝 / token 单次消费。Direct直接调 backend 无拦截。RiskAgent基础限额: 已实现max_order_notional+quantity > 0检查;合规时approved=true+risk_score=0.1+ 空violations,违规时附带违规列表。- 真实交易所接入:
ExchangeTradingBackend(crates/axon-llm/src/trading/exchange.rs)完整实现,把ExchangeAdapter(Binance / OKX)适配为TradingBackend;依赖trading-exchangefeature,SymbolMap提供 LLM symbol ↔ 交易所 symbol 双向映射。注:同文件测试模块里有 8 处unimplemented!(),是#[cfg(test)]内的MockAdapterstub(测试路径不调用),非生产代码缺口。
未实现(0.6.0+ 路线图)¶
RiskAgent高级风控:RiskAgentConfig已定义max_position/max_drawdown字段但未做检查;risk_score目前只是二元(0.1 / 0.9);波动率 / VaR / 历史回撤窗口 / 仓位集中度等指标未实现。0.6.0 收口部分能力:axon-risk0.6.0 新增跨 leg 风险约束(check_leg_pair(portfolio, &LegPair) -> RiskResult+RiskReason::LegPairNetExposureExceeded+per_leg_var+stress_pair/stress_portfolio),RiskConfig.max_leg_pair_net_exposure默认 0.0(严格 delta 中性)。RiskAgent接入这些 API 替换 happy path 是后续工作。- 4-Agent pipeline 跨进程协调(分布式 swarm): 当前
SwarmOrchestrator是单进程内 mpsc 通道;跨进程协调、共识状态机ConsensusManager持久化、axon-distributed的 Ray Actor 化包装未实现。计划在 0.7.0+ 路线图。 - 单 Agent 跨进程复用:
MarketAgent/RiskAgent/AuditAgent目前以独立tokio::task::spawn运行,跨进程调度 / 共享 LLM client / 全局 prompt cache 仍待设计。
7. 0.11.0 Multi-Agent 投票共识层¶
新增于 v0.11.0 位置:
crates/axon-llm/src/swarm/consensus.rs
0.11.0 在已有 4-Agent pipeline 之上新增了 投票共识决策层,实现多 trader 独立决策 → 加权投票聚合 → risk 一票否决的完整流程。
7.1 决策流程¶
Bar 到达
│
├──→ Trader A (ReAct/规则) ──→ AgentVote{Buy, 0.8}
├──→ Trader B (ReAct/规则) ──→ AgentVote{Buy, 0.6}
└──→ Trader C (ReAct/规则) ──→ AgentVote{Hold, 0.0}
│
▼
┌─────────────────────────┐
│ VotingStrategy │
│ weighted majority: │
│ Buy = 0.8+0.6 = 1.4 │
│ Hold = 0.0 │
│ → 聚合结果: Buy(0.7) │
└────────────┬────────────┘
▼
┌─────────────────────────┐
│ ConsensusRiskAgent │
│ 检查: │
│ - 当前持仓/敞口 │
│ - 连续亏损次数 │
│ - 单 bar 最大仓位 │
│ → approve / veto │
└────────────┬────────────┘
▼
Final Action (Buy 0.7) 或 Hold (被否决)
7.2 核心类型¶
// 投票
pub struct AgentVote {
pub agent_id: String,
pub action: TraderAction, // Buy / Sell / Hold
pub confidence: f64,
pub reasoning: String,
}
// 投票策略 trait
pub trait VotingStrategy: Send + Sync {
fn aggregate(&self, votes: &[AgentVote]) -> ConsensusDecision;
}
// 内置策略
// 1. WeightedMajorityVote — confidence 加权,超阈值(0.5)胜出
// 2. UnanimousVote — 全票一致才通过(保守模式)
// 风控审核(纯规则,不走 LLM)
pub struct ConsensusRiskAgent {
pub max_position: f64,
pub max_consecutive_loss: u32,
pub max_drawdown: f64,
}
// 编排器
pub struct VotingOrchestrator {
traders: Vec<Box<dyn TraderCallback>>,
risk_agent: ConsensusRiskAgent,
voting: Box<dyn VotingStrategy>,
}
7.3 Python 接入¶
from axon_quant._native import PyVotingOrchestrator
orch = PyVotingOrchestrator(
traders=[trader_a, trader_b, trader_c],
risk_config={"max_position": 0.5, "max_consecutive_loss": 3},
voting="weighted_majority",
)
decision = orch.on_bar(bar_dict)
# decision: {"action": ..., "votes": [...], "risk_verdict": ..., "confidence": ...}
7.4 与 0.6.0 SwarmOrchestrator 的关系¶
| 维度 | 0.6.0 SwarmOrchestrator | 0.11.0 VotingOrchestrator |
|---|---|---|
| 定位 | 4-Agent 异步消息 pipeline | bar-by-bar 同步投票决策 |
| 通信 | tokio mpsc 异步 | 同步回调 (TraderCallback) |
| 风控 | RiskAgent (LLM 可介入) | ConsensusRiskAgent (纯规则) |
| 适用 | 生产级异步交易 | 回测/评估/多策略对比 |
两者共存于 swarm/ 模块,不互相依赖。
7.5 Trajectory Schema v0.11.0¶
在 0.10.0 schema 基础上扩展: - 新增 votes 数组:每个 trader 的 action + confidence + reasoning - 新增 aggregation 字段:投票策略名 + 聚合结果 - 新增 risk_verdict 字段:approve/veto + reason - 新增 token_usage 字段:per-bar 累计 input/output tokens