diff --git a/.env.example b/.env.example index cc4cf5b..bf0c365 100644 --- a/.env.example +++ b/.env.example @@ -24,6 +24,14 @@ TUSHARE_TOKEN=your_token_here # FINNHUB_API_KEY=your_key_here # TAVILY_API_KEY=your_key_here +# ── 交易(可选,Futu)── +# 自主交易操作员:trade_execute 一步下单/撤单,不经过人工确认。 +# real 账户解锁只能在 Futu OpenD 客户端手工完成,代码侧零密码,不在此处配置交易密码。 +# ATHENACLAW_TRADE_BROKER=futu +# ATHENACLAW_FUTU_HOST=127.0.0.1 +# ATHENACLAW_FUTU_PORT=11111 +# ATHENACLAW_FUTU_SECURITY_FIRM= + # ── Telegram(可选)── # TELEGRAM_BOT_TOKEN=123456:abcdef # TELEGRAM_ALLOWED_USER_IDS=123456789 diff --git a/docs/agent-design.md b/docs/agent-design.md index 525d40d..07afa26 100644 --- a/docs/agent-design.md +++ b/docs/agent-design.md @@ -132,8 +132,7 @@ Pi/bub/ampcode 的核心洞见:**read/write/edit/bash 是最小完备工具集 | `portfolio(action, ...)` | 领域增强 | 持仓快照 | 维护结构化当前持仓,不记录交易历史 | | `watchlist(action, ...)` | 领域增强 | 自选快照 | 维护结构化当前自选列表,不记录历史观察日志 | | `trade_account(action, ...)` | 交易域 | 远端账户快照 | 读取远端 broker 账户、持仓、未完成订单、订单状态 | -| `trade_plan(operation, ...)` | 交易域 | 执行计划 | 生成交易计划,不直接产生外部副作用 | -| `trade_apply(plan_id)` | 交易域 | 动作执行 | 执行交易计划,执行前必须确认 | +| `trade_execute(operation, ...)` | 交易域 | 自主执行 | 单次调用内完成校验+preview+风控裁决+下单/撤单,不经过人工确认 | | `compute(code)` | 领域增强 | 计算器 | 沙箱化 Python(安全版 bash) | | `market_ohlcv(symbol, interval, mode)` | 领域核心 | 眼睛 | 内核数据原语,支持日线/分钟/history/latest | | `recall(query)` | 领域增强 | 回忆 | 全文搜索 memory + notebook | @@ -146,7 +145,7 @@ memory.write 就是 `write("memory.md", content)`,read 就是 `read("memory.md 但当前持仓和具体自选列表是例外:它们需要被 Agent 稳定读取,也适合被自动化或 UI 直接消费,因此分别用 `portfolio` 维护 `portfolio.json`、用 `watchlist` 维护 `watchlist.json`,而不是继续把明细散写在 memory.md 里。 -远端 broker 交易又是另一类例外:它不是工作区快照维护,而是受确认与审计约束的外部动作。因此新增一层交易边界层,由 `trade_account / trade_plan / trade_apply` 暴露给 LLM,详细规格见 [trading.md](./trading.md)。 +远端 broker 交易又是另一类例外:它不是工作区快照维护,而是受审计约束的外部动作。安全边界是程序化的 `RiskGuard`,不是人类确认——Agent 是自主交易操作员。因此新增一层交易边界层,由 `trade_account / trade_execute` 暴露给 LLM,详细规格见 [trading.md](./trading.md)。 read/write/edit 的覆盖范围: - 写研究报告 → `write("notebook/research/宁德时代/2024-01-15.md", content)` @@ -733,8 +732,8 @@ Agent 输出任务设计稿: | market_ohlcv | ✅ | MarketAdapter + DataStore | | bash | ✅ | shell 执行 + 超时 + 进程树清理 | | trade_account | ✅ | 远端 broker 账户只读访问 | -| trade_plan | ✅ | 交易计划生成 | -| trade_apply | ✅ | 交易计划执行 | +| trade_execute | ✅ | 自主下单/撤单,单步执行,无人工确认 | +| market_snapshot | ✅ | 富途实时行情快照 | ### 预连通管道 diff --git a/docs/automation.md b/docs/automation.md index 6fcb396..933943b 100644 --- a/docs/automation.md +++ b/docs/automation.md @@ -228,8 +228,9 @@ Reaction 不是 workflow,而是“一次执行”。 交易工具的自动化规则: - `trade_account` 允许只读使用 -- `trade_plan` 永久禁止 -- `trade_apply` 永久禁止 +- `trade_execute` 允许调用——Agent 已是自主交易操作员,automation 可自主下单/撤单 + +⚠️ 风控专题首要 TODO:无人值守 automation 对 real 账户裸奔风险最高。`TradeOrchestrator` 传给 `RiskGuard` 的 `RiskContext.automation` 字段就是为它预留的裁决挂钩,但这期 `AllowAllGuard` 恒 ALLOW,当前没有任何拦截。详见 [trading.md](./trading.md) §6。 ## 11. 幂等与恢复 diff --git a/docs/tools.md b/docs/tools.md index fe9a28a..05a7f99 100644 --- a/docs/tools.md +++ b/docs/tools.md @@ -13,8 +13,8 @@ | `portfolio` | 维护结构化当前持仓快照 `portfolio.json` | | `watchlist` | 维护结构化自选列表快照 `watchlist.json` | | `trade_account` | 读取远端 broker 账户、持仓、未完成订单、订单状态 | -| `trade_plan` | 生成交易执行计划,不直接产生外部副作用 | -| `trade_apply` | 执行 `trade_plan` 生成的计划 | +| `trade_execute` | 自主交易执行:单次调用内完成校验+preview+风控裁决+下单/撤单,不经过人工确认 | +| `market_snapshot` | 富途实时行情快照(last/bid/ask),区别于 `market_ohlcv` 的 K 线 | | `compute` | 在沙箱中对已加载 OHLCV 做 Python 分析 | | `read` | 读工作区文件 | | `write` | 写工作区文件 | @@ -22,7 +22,7 @@ | `bash` | 执行 shell 命令(按权限控制) | | `web_search` / `web_fetch` | 可选 Web 搜索与抓取 | -其中最容易用错的是 `portfolio`、`watchlist`、`trade_account`、`trade_plan`、`trade_apply`、`market_ohlcv` 和 `compute`。 +其中最容易用错的是 `portfolio`、`watchlist`、`trade_account`、`trade_execute`、`market_ohlcv` 和 `compute`。 ## portfolio @@ -114,18 +114,19 @@ `memory.md` 适合记录高层关注方向和长期偏好。 具体自选 symbol 清单、观察理由和加入时间需要结构化读取,因此单独维护在 `watchlist.json`。 -## trade_account / trade_plan / trade_apply +## trade_account / trade_execute ### 核心语义 -这三者共同组成远端 broker 交易闭环: +这两者共同组成远端 broker 交易闭环,自主交易操作员范式——中间不经过人工确认: -- `trade_account`:读取远端账户状态 -- `trade_plan`:创建可确认、可审计的执行计划 -- `trade_apply`:执行计划 +- `trade_account`:读取远端账户状态(只读) +- `trade_execute`:唯一执行入口,`operation=submit_limit|cancel`,单次调用内完成校验+preview+`RiskGuard.evaluate`+下单/撤单+状态回读 它们不是 `portfolio` 的替代物,也不会自动维护 `portfolio.json`。 +安全边界是程序化的 `RiskGuard`(这期 `AllowAllGuard` 占位,恒 ALLOW,真实风控是独立专题),不是人类确认——详见 [trading.md](./trading.md) §6。 + ### V1 边界 只支持: @@ -147,17 +148,16 @@ ### 正确用法 1. 先 `trade_account.list_accounts` -2. 再 `trade_plan.submit_limit` 或 `trade_plan.cancel` -3. 最后 `trade_apply` +2. 直接 `trade_execute(operation="submit_limit", ...)` 或 `trade_execute(operation="cancel", ...)` -不能跳过 `trade_plan` 直接执行。 +一次调用即完成下单/撤单,没有二段式确认步骤。 显式参数规则: - `trade_account.get_positions/get_summary/get_open_orders` 必须显式传 `account_ref` - `trade_account.get_order_status` 必须显式传 `order_ref` -- `trade_plan.submit_limit` 必须显式传 `account_ref` -- `trade_plan.cancel` 必须显式传 `order_ref` +- `trade_execute(operation="submit_limit")` 必须显式传 `account_ref` +- `trade_execute(operation="cancel")` 必须显式传 `order_ref` - 不会自动承接最近账户、最近订单,也不会返回 suggestion ### 账户发现 @@ -184,19 +184,16 @@ `trade_account.get_positions` 成功后,会把该账户快照写入 `Kernel.data["account"]`。 后续 `compute` 读取到的 `account/cash/equity/positions` 就来自这个当前活动账户快照。 -### 计划与执行返回 - -`trade_plan.submit_limit` 返回的 plan 除了 `plan_id` 外,还会返回: - -- `normalized_intent` +### 执行返回 -其中 `normalized_intent.limit_price` 是后续 `trade_apply` 唯一允许执行的价格。 -如果 provider 对输入价格做了规范化,plan 的 `warnings` 里会明确写出原始价格和规范化后的价格。 - -`trade_apply` 返回除了 `order_status` 外,还会返回: +`trade_execute(operation="submit_limit")` 一次调用内完成校验+preview+风控裁决+下单+状态回读,返回除了 `order_status` 外,还会返回: +- `normalized_intent`:如果 provider 对输入价格做了规范化,`warnings` 里会明确写出原始价格和规范化后的价格 - `finalized` - `warnings` +- `plan_id`:语义已从"待确认计划 id"改为一次性执行 id,纯审计用途 + +`trade_execute(operation="cancel")` 返回同样包含 `order_status`/`finalized`/`warnings`。 语义是: diff --git a/docs/trading.md b/docs/trading.md index e1d6ace..5aecfcf 100644 --- a/docs/trading.md +++ b/docs/trading.md @@ -2,9 +2,9 @@ ## 1. 目标 -交易 V1 的目标不是接完某家券商的全部能力,而是在 AthenaClaw 中建立一套稳定、可扩展、可审计的交易骨架。 +交易系统的目标不是接完某家券商的全部能力,而是在 AthenaClaw 中建立一套稳定、可扩展、可审计的交易骨架。**自主交易是核心目标**:Agent 就是要替代人类完成完全自主的下单/撤单决策与执行,中间不经过人工确认。 -V1 只覆盖最小人工下单闭环: +V1 覆盖的执行闭环: - 列远端 broker 可用账户 - 读取当前持仓 @@ -12,22 +12,25 @@ V1 只覆盖最小人工下单闭环: - 提交股票/ETF 限价单 - 查询订单当前状态 - 撤销未完成订单 +- 读取同 broker 的实时快照(`market_snapshot`),供下单前判断盘口 - 成交后刷新当前活动账户快照,供 `compute` 使用 非目标: -- 不做自动交易 - 不做独立交易服务进程 - 不自动把远端 broker 状态写回 `portfolio.json` - 不做 `MARKET`、止损单、条件单、TWAP/VWAP、期权/期货/融资融券 +- 不在代码里持有交易密码、不调用 `unlock_trade`(见 §7) +- 这期不做真实风控拦截(见 §6) ## 2. 分层架构 ```mermaid flowchart LR - LLM["LLM / Session"] --> Tools["trade_account / trade_plan / trade_apply"] + LLM["LLM / Session"] --> Tools["trade_account / trade_execute"] Tools --> Orch["TradeOrchestrator"] - Orch --> Store["TradePlanStore / TradeAuditLog"] + Orch --> Guard["RiskGuard (seam, AllowAllGuard)"] + Orch --> Audit["TradeAuditLog"] Orch --> Adapter["TradeBrokerAdapter"] Adapter --> Broker["Futu OpenD / future brokers"] Tools --> Kernel["Kernel.data['account']"] @@ -35,36 +38,36 @@ flowchart LR 三层职责: -- `tool` 层:LLM 接口、参数 schema、确认交互、结果格式化、`Kernel.data` 注入 -- `TradeOrchestrator`:交易域规则、plan 生命周期、canonical 状态/错误、apply 后回读 +- `tool` 层:LLM 接口、参数 schema、结果格式化、`Kernel.data` 注入 +- `TradeOrchestrator`:交易域规则、canonical 状态/错误、风控裁决入口、执行后回读 - `TradeBrokerAdapter`:broker SDK/API 接入和字段翻译 这里的“交易边界层”是一个**进程内模块**,不是独立微服务。代码主对象统一叫 `TradeOrchestrator`。 ## 3. `TradeOrchestrator` 的定位 -它的唯一职责是:把 Agent 产生的交易意图,转成**确定性、可确认、可执行、可回读**的 canonical 流程。 +它的唯一职责是:把 Agent 产生的交易意图,转成**确定性、经风控裁决、可执行、可回读**的 canonical 流程。 属于它的职责: - 校验 V1 公共约束:股票/ETF、`LIMIT`、`BUY/SELL/CANCEL` - 统一 canonical 输入输出 -- 生成 `TradePlan` -- 生成 `plan_summary` 与 `confirm_text` -- 管理 `plan_id`、TTL、单次消费、重复 apply 幂等 +- 价格规范化 + preview 前置拦截 +- 单次调用内完成:校验 → preview → `RiskGuard.evaluate` → 下单/撤单 → 状态回读 - 统一 canonical 订单状态与错误码 -- apply 后补查订单状态,必要时补查持仓 +- 执行后补查订单状态,必要时补查持仓 - 生成标准化账户快照,交给 tool 层写入 `Kernel.data["account"]` +- 把每次风控裁决与执行结果写入 `TradeAuditLog` 不属于它的职责: - 自然语言理解 - prompt 拼装 -- CLI/TUI/IM 确认 UI +- 人工确认 UI(交易链路已彻底不需要) - OpenD 连接、SDK session、provider 原始枚举 - `portfolio.json` / `watchlist.json` 维护 -- 自动交易、策略执行、回调 worker -- 行情拉取 +- 真实风控规则本身(见 §6,这期只留 seam) +- 行情拉取(`market_snapshot` 走独立的 `SnapshotAdapter`,不经过 `TradeOrchestrator`) 边界判定原则: @@ -101,7 +104,7 @@ flowchart LR 这些字段的分工如下: - 公共字段负责跨 broker 的基础选户语义,例如“这个账户是否 active”“是否支持 US 市场”“是否更像股票户还是期权户” -- `extra` 负责承载 provider 专有信息;V1 只在账户发现落地,不默认扩散到 orders/positions/plan/receipt +- `extra` 负责承载 provider 专有信息;V1 只在账户发现落地,不默认扩散到 orders/positions/receipt Futu 账户发现至少要把这些专有信息放进 `extra`: @@ -118,9 +121,8 @@ canonical 标识: - `account_ref` - `order_ref` -- `plan_id` -这三类引用都是系统生成的 opaque string。LLM 只能传递,不能猜测、拼接或修改。 +这两类引用都是系统生成的 opaque string。LLM 只能传递,不能猜测、拼接或修改。`TradeApplyResult.plan_id` 字段依旧存在,但语义已从“待确认的计划 id”改为“这次一次性执行的 id”,纯粹用于审计追踪,不再有生命周期或消费语义。 显式参数原则: @@ -143,55 +145,67 @@ canonical 状态只保留: 补充结果语义: - `TradePreview` 可以返回 `normalized_limit_price` 与 `normalization_reason` -- `TradePlan` 可以返回 `normalized_intent` - `TradeApplyResult` 返回 `finalized` 与 `warnings` - `status=ok` 只表示工具执行成功;业务上是否已进入终态,要看 `finalized + order_status` ## 5. 工具协议 -对外只保留 3 个工具: +对外保留 3 个工具: -- `trade_account` +- `trade_account`(只读) - `list_accounts` - `get_positions` - `get_open_orders` - `get_order_status` - `get_summary` -- `trade_plan` - - `submit_limit` - - `cancel` -- `trade_apply` - - 只接 `plan_id` +- `trade_execute`(唯一执行入口) + - `submit_limit`:一次调用直接提交限价单,内部完成校验 + preview + 风控裁决 + 下单 + 状态回读 + - `cancel`:一次调用直接撤单,内部完成终态检查 + 风控裁决 + 撤单 + 有界轮询确认 +- `market_snapshot`(若对应 broker 已接入行情):读取实时 last/bid/ask 快照,与 `market_ohlcv` 的历史序列语义不同,见 §7 关键约束: -- 所有执行动作必须先 `trade_plan`,后 `trade_apply` -- `trade_apply` 永远只消费 `plan_id` -- tool 不得直连 adapter 的 mutating 方法 +- `trade_execute` 调用即视为 Agent 已完成决策判断,不经过人工确认;唯一的安全边界是 `TradeOrchestrator` 内的 `RiskGuard.evaluate`(见 §6) +- tool 不得直连 adapter 的 mutating 方法,必须经过 `TradeOrchestrator` - `Kernel.data["account"]` 在 V1 里只表示“当前活动账户快照” - `trade_account.list_accounts` 必须返回足够的账户能力信息,让用户和 Agent 在不打开券商客户端的前提下也能判断哪个账户支持目标市场 - 缺少 `account_ref` / `order_ref` 时返回 `missing_*` 错误,不做隐式补参 +- 被 `RiskGuard` 拒绝时返回 `error_code=permission_denied` -## 6. Prompt 与安全 +## 6. 安全模型:风控 vs 确认 -系统会在交易工具注册时条件注入 `TRADE_GUIDE`。其规则是: +**确认(human-in-the-loop)与风控(risk control)是两件不同的事,交易链路只保留后者:** -- 先 `trade_plan`,后 `trade_apply` -- 不得伪造 `account_ref/order_ref/plan_id` +- 确认:`request_confirm`、plan→apply 两阶段、TTL 令牌、`confirm_text` —— 这些都是为“中间有人类点头”设计的仪式。自主操作员场景下这个人类不存在,仪式退化为纯开销,交易链路已**彻底删除**对 `request_confirm` 的依赖。 +- 风控:程序化硬边界,代码裁决、LLM 不可绕过。这是自主操作员真正需要的纪律,落地为 `TradeOrchestrator` 在下单/撤单前调用的 `RiskGuard.evaluate(RiskContext) -> RiskDecision`。 + +**这期的状态**:`RiskGuard` 只是一个 seam,唯一实现 `AllowAllGuard` 恒返回 `ALLOW`,不拦截任何单。真实风控(notional 上限、限价偏离市价保护、日内下单频率、单标的集中度、禁买清单……)是独立的大设计,留给后续专题重新做。这里显式标注,不假装安全。 + +`RiskContext` 携带的字段是为未来风控专题预留的契约点: + +- `operation`:`submit_limit` / `cancel` +- `env`:`simulate` / `real`,来自 `account_ref`/`order_ref` 解码,单点决定,不在别处重复判断 +- `automation`:本次调用是否发生在无人值守的 automation reaction 中,来自 `AutomationToolPolicy.automation` 标记 +- `intent`:原始下单/撤单意图 + +`RiskAction` 预留三态:`ALLOW` / `DENY` / `ESCALATE`(real 大额场景的逃生口,未来落地)。 + +`TRADE_GUIDE` 在交易工具注册时条件注入,规则要点: + +- Agent 是自主交易操作员:`trade_execute` 一步下单/撤单,直接对远端 broker 账户生效,调用即视为已完成决策判断 +- 执行前会经过风控裁决;被拒绝时返回 `error_code=permission_denied`,不要重试,向用户说明原因 +- 不得伪造 `account_ref/order_ref` - V1 只支持限价单 - `portfolio` 不是远端 broker 账户 - 下单前先看 `trade_account.list_accounts` 里的 `supported_markets`、`account_status`、`account_kind`;`extra` 可用于解释 provider 特有限制 - 查询具体账户或订单时必须显式携带 `account_ref` / `order_ref` -- 在 `trade_apply` 成功前,不能宣称“已提交/已成交” -- 没有同 broker 的新鲜行情或明确 market-state 证据时,不得推断“更容易成交”“当前处于常规交易时段” +- 在 `trade_execute` 返回 `status=ok` 且订单进入终态前,不能宣称“已提交/已成交” +- 没有同 broker 的新鲜行情或明确 market-state 证据时,不得推断“更容易成交”“当前处于常规交易时段”;若已注册 `market_snapshot`,下单前优先用它拿一次同 broker 的实时快照 - broker 交易不会自动改写 `portfolio.json` -这些规则在 prompt 中做引导,但最终安全边界在代码里执行,不依赖 LLM 自觉遵守。 - -确认交互统一改成 message-based confirm。 -文件系统工具继续传路径消息;交易工具传动作消息,例如: +这些规则在 prompt 中做引导,但最终安全边界在代码里执行(`RiskGuard.evaluate`),不依赖 LLM 自觉遵守。 -- `确认提交 BUY 100 AAPL 限价 180.00 吗?` +**automation carve-out**:`AutomationToolPolicy._ALWAYS_DENIED` 不再包含交易工具,automation reaction 可以直接执行 `trade_execute`。⚠️ 无人值守 automation 对 real 账户裸奔风险最高,`RiskContext.automation` 字段就是为它预留的挂钩,但 `AllowAllGuard` 当前不做任何拦截——这是风控专题的首要 TODO。 ## 7. Futu 适配规则 @@ -207,24 +221,23 @@ Futu 作为首个 provider,只接证券账户交易。 - `modify_order(ModifyOrderOp.CANCEL)` -> `cancel_order` - `accinfo_query` -> `get_account_summary` - `acctradinginfo_query` -> `preview_limit_order` -- `get_market_snapshot(price_spread)` -> provider 内部价格合法化辅助,不暴露为公共工具 +- `get_market_snapshot` -> 行情侧 `FutuAdapter.snapshot()`,支撑 `market_snapshot` 工具;不需要订阅,一次调用拿当前盘口 last/bid/ask Futu 账户发现不做 provider 侧隐式过滤;仍返回全部账户。选户正确性由两层保障: - `list_accounts` 直接暴露账户能力与 `extra` -- `trade_plan.submit_limit` 在 plan 阶段先做价格规范化,再通过 `preview_limit_order` 前置拦截“不支持该市场/账户已失效/账户类型不适合/最大可买卖不足”等硬失败 -- `trade_apply(cancel)` 会在内部做短时、有界的状态确认;若短时内未进入终态,返回 `finalized=false`,而不是假装已经撤单成功 +- `trade_execute(submit_limit)` 内部先做价格规范化,再通过 `preview_limit_order` 前置拦截“不支持该市场/账户已失效/账户类型不适合/最大可买卖不足”等硬失败 +- `trade_execute(cancel)` 会在内部做短时、有界的状态确认;若短时内未进入终态,返回 `finalized=false`,而不是假装已经撤单成功 + +**真实账户解锁:手工 OpenD 解锁是官方 real 路径,代码零密码**。`FutuTradeConfig` 不包含交易密码字段,代码侧不持有、不传递、不调用 `unlock_trade`。当账户处于未解锁状态时,下单/撤单返回 `TradeErrorCode.TRADE_LOCKED`,错误信息会明确指引“请在 Futu OpenD 客户端手工解锁交易”。这不是缺失的功能,是刻意的设计选择:安全边界必须在代码里体现为“不做什么”,不留 LLM 或自动化流程自由裁量的空间。 V1 不做: -- `unlock_trade` - `market order` - `history deals` - `fees` - callback worker -真实账户解锁仍由用户在 OpenD 侧手工完成。 - ## 8. 扩展原则 新增 broker: @@ -247,13 +260,19 @@ V1 不做: - 要么作为 adapter 内部增强 - 要么独立为 provider-specific tool +真实风控落地时的扩展点:只需要新增一个 `RiskGuard` 实现(替换 `AllowAllGuard`)并在 `runtime/bundle.py` 注入,`TradeOrchestrator` 与工具层不需要改动。 + ## 9. 测试与验收 单测: -- `TradeOrchestrator` 的 plan/apply/TTL/幂等/状态刷新 +- `TradeOrchestrator.execute_limit` / `execute_cancel` 的单步执行、状态刷新 +- `RiskGuard` DENY 时阻断下单/撤单,且审计日志仍记录裁决 +- `RiskContext.automation` 从 `AutomationToolPolicy` 正确透传 - canonical 错误码和状态映射 - 当前活动账户快照语义 +- Futu `TRADE_LOCKED` 的手工解锁提示文案 +- Futu `snapshot()` 字段规范化 合约测试: @@ -261,13 +280,14 @@ V1 不做: 工具测试: -- `trade_plan -> trade_apply -> get_order_status` -- confirm 拒绝后不执行 +- `trade_execute(submit_limit) -> get_order_status` 单步闭环 +- `trade_execute` 被 `RiskGuard` 拒绝时返回 `permission_denied` 验收标准: -- 模拟账户可完成限价单下单、查状态、撤单 +- 模拟账户可完成限价单下单、查状态、撤单,全程无人工确认弹窗 - 成交后 `compute` 能读取最新 `account` -- 自动化任务不能执行 `trade_plan` / `trade_apply` +- automation 任务可以执行 `trade_execute`(carve-out 已放开);`bash`/`task_*`/`create_subagent` 仍被拒绝 - 缺少 ref 时返回 `missing_*`,而不是 `invalid_*` - 不在陈旧行情或弱提示上推断成交概率或当前市场状态 +- real 账户未解锁时,错误信息清楚指向 OpenD 手工解锁,而不是暗示系统可以自动解锁 diff --git a/pyproject.toml b/pyproject.toml index 0a915ca..65f54bf 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -35,6 +35,9 @@ finnhub = [ anthropic = [ "anthropic>=0.40", ] +futu = [ + "futu-api>=9.0", +] [project.scripts] athenaclaw = "athenaclaw.interfaces.cli:main" diff --git a/src/athenaclaw/CLAUDE.md b/src/athenaclaw/CLAUDE.md index c5446ef..adf122f 100644 --- a/src/athenaclaw/CLAUDE.md +++ b/src/athenaclaw/CLAUDE.md @@ -9,9 +9,9 @@ AthenaClaw 主产品包。 - `kernel/`:Kernel、Session、Permission、prompt 组装 - `runtime/`:AgentConfig、KernelBundle、trace wiring、session store - `llm/`:消息模型、provider、context 压缩 -- `tools/`:filesystem、compute、market、portfolio、watchlist、shell、web +- `tools/`:filesystem、compute、market(含 market_snapshot)、trade(trade_account/trade_execute)、portfolio、watchlist、shell、web - `automation/`:任务定义、执行、worker -- `trading/`:交易域 canonical 类型、TradeOrchestrator、plan/audit store、账户快照映射 +- `trading/`:交易域 canonical 类型、TradeOrchestrator(execute_limit/execute_cancel 单步执行)、RiskGuard seam、audit store、账户快照映射 - `skills/`:skills 发现与展开 - `subagents/`:子代理定义、运行和系统集成 - `interfaces/`:CLI、Telegram、Discord、TUI、IM diff --git a/src/athenaclaw/automation/policy.py b/src/athenaclaw/automation/policy.py index 69da66c..cb484a1 100644 --- a/src/athenaclaw/automation/policy.py +++ b/src/athenaclaw/automation/policy.py @@ -13,18 +13,24 @@ from athenaclaw.tools.filesystem.path import resolve_path +# trade_execute 不在硬禁清单里:Agent 已是自主交易操作员,automation 可执行交易。 +# ⚠️ 风控专题首要 TODO:无人值守 automation 对 real 账户裸奔风险最高—— +# 当前 orchestrator 传给 RiskGuard 的 RiskContext.automation 字段就是为它预留的挂钩, +# 但 AllowAllGuard 恒 ALLOW,此刻没有任何拦截。方案在此显式标注,不假装安全。 _ALWAYS_DENIED = { "bash", "task_plan", "task_apply", "task_control", "create_subagent", - "trade_plan", - "trade_apply", } class AutomationToolPolicy(ToolAccessPolicy): + # automation 标记:交易工具据此把 RiskContext.automation 置 True。 + # 单点标记,不散落到别处判断“这是不是自动化触发”。 + automation = True + def __init__( self, *, diff --git a/src/athenaclaw/integrations/futu/config.py b/src/athenaclaw/integrations/futu/config.py index 7d50750..a9b16b7 100644 --- a/src/athenaclaw/integrations/futu/config.py +++ b/src/athenaclaw/integrations/futu/config.py @@ -19,4 +19,12 @@ class FutuConfig: @dataclass(frozen=True) class FutuTradeConfig(FutuConfig): - """交易侧沿用的兼容别名。""" + """交易侧沿用的兼容别名。 + + 刻意不包含交易密码字段。SIMULATE/REAL 两个环境的解锁都只能在本机 + Futu OpenD 客户端手工完成(点一次,直到 OpenD 进程重启前一直有效)。 + 代码侧不持有密码、不调用 unlock_trade、不做自动解锁 —— 这是设计选择, + 不是遗漏。账户处于未解锁状态时,下单/撤单会失败并返回 + TradeErrorCode.TRADE_LOCKED,错误信息里会指引去 OpenD 手工解锁。 + """ + diff --git a/src/athenaclaw/integrations/futu/trade_adapter.py b/src/athenaclaw/integrations/futu/trade_adapter.py index 607f9bd..825dcc3 100644 --- a/src/athenaclaw/integrations/futu/trade_adapter.py +++ b/src/athenaclaw/integrations/futu/trade_adapter.py @@ -429,7 +429,7 @@ def _ensure_frame(ret: int, data: Any, op: str, *, allow_empty: bool = False) -> if ret != futu.RET_OK: text = str(data) code = _translate_error_code(text, op=op) - raise TradeError(code, f"{op} 失败: {text}") + raise TradeError(code, _error_message(code, op=op, raw=text)) if isinstance(data, pd.DataFrame): return data if allow_empty: @@ -617,6 +617,17 @@ def _translate_error_code(text: str, *, op: str) -> TradeErrorCode: return TradeErrorCode.PROVIDER_ERROR +def _error_message(code: TradeErrorCode, *, op: str, raw: str) -> str: + if code is TradeErrorCode.TRADE_LOCKED: + # 本系统零密码:REAL 账户解锁只能在 Futu OpenD 客户端手工完成,代码侧不持有、 + # 不传递、不代为调用 unlock_trade。这里只负责把"去哪解锁"说清楚。 + return ( + f"{op} 失败: 交易账户未解锁({raw})。" + "请在 Futu OpenD 客户端手工解锁交易,本系统不持有交易密码、不会自动解锁。" + ) + return f"{op} 失败: {raw}" + + def _snap_to_increment(value: Decimal, increment: Decimal) -> Decimal: if increment <= 0: return value diff --git a/src/athenaclaw/integrations/market/futu.py b/src/athenaclaw/integrations/market/futu.py index 5a25155..a441731 100644 --- a/src/athenaclaw/integrations/market/futu.py +++ b/src/athenaclaw/integrations/market/futu.py @@ -1,7 +1,8 @@ """ [INPUT]: pandas, datetime, athenaclaw.integrations.futu.*, athenaclaw.tools.market.schema -[OUTPUT]: FutuAdapter — Futu OpenD Quote OHLCV 适配器 -[POS]: MarketAdapter 实现,通过 Futu request_history_kline/get_cur_kline 获取 CN/HK/US OHLCV +[OUTPUT]: FutuAdapter — Futu OpenD Quote OHLCV + 实时快照适配器 +[POS]: MarketAdapter 实现,通过 Futu request_history_kline/get_cur_kline 获取 CN/HK/US OHLCV; + 同时实现 SnapshotAdapter.snapshot,经 get_market_snapshot 提供下单前的实时快照/报价 [PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md """ @@ -128,6 +129,19 @@ def _fetch_latest(self, quote_ctx, futu, query: MarketQuery) -> pd.DataFrame: df = _ensure_frame(ret, data, futu=futu, op="get_cur_kline") return _normalize_kline_frame(df).tail(1).reset_index(drop=True) + def snapshot(self, symbols: list[str]) -> list[dict[str, Any]]: + """一批 symbol 的实时快照 — 不需要订阅,直接一次调用拿当前盘口。""" + try: + futu = _load_futu() + quote_ctx = self._manager.quote_context() + codes = [to_futu_code(symbol) for symbol in symbols] + ret, data = quote_ctx.get_market_snapshot(codes) + if ret != futu.RET_OK: + raise ValueError(f"get_market_snapshot 失败: {data}") + return _normalize_snapshot_frame(data) + except Exception as exc: + raise ValueError(_friendly_error_message(exc)) from exc + def _default_window(query: MarketQuery) -> tuple[str | None, str | None]: if query.start_dt is not None or query.end_dt is not None: @@ -170,6 +184,42 @@ def _empty_ohlcv_frame() -> pd.DataFrame: return pd.DataFrame(columns=["date", "open", "high", "low", "close", "volume"]) +def _normalize_snapshot_frame(frame: pd.DataFrame) -> list[dict[str, Any]]: + if frame is None or frame.empty: + return [] + records: list[dict[str, Any]] = [] + for _, row in frame.iterrows(): + records.append({ + "symbol": str(row.get("code")), + "name": row.get("name"), + "last_price": _to_float(row.get("last_price")), + "open_price": _to_float(row.get("open_price")), + "high_price": _to_float(row.get("high_price")), + "low_price": _to_float(row.get("low_price")), + "prev_close_price": _to_float(row.get("prev_close_price")), + "volume": _to_float(row.get("volume")), + "turnover": _to_float(row.get("turnover")), + "turnover_rate": _to_float(row.get("turnover_rate")), + "bid_price": _to_float(row.get("bid_price")), + "ask_price": _to_float(row.get("ask_price")), + "bid_vol": _to_float(row.get("bid_vol")), + "ask_vol": _to_float(row.get("ask_vol")), + "price_spread": _to_float(row.get("price_spread")), + "suspension": bool(row.get("suspension")) if pd.notna(row.get("suspension")) else None, + "update_time": row.get("update_time") if pd.notna(row.get("update_time")) else None, + }) + return records + + +def _to_float(value: Any) -> float | None: + if value is None or (isinstance(value, float) and pd.isna(value)): + return None + try: + return float(value) + except (TypeError, ValueError): + return None + + def _ensure_frame(ret: int, data: Any, *, futu, op: str) -> pd.DataFrame: if ret != futu.RET_OK: raise ValueError(f"{op} 失败: {data}") diff --git a/src/athenaclaw/kernel/service.py b/src/athenaclaw/kernel/service.py index 85e0806..fc84aa3 100644 --- a/src/athenaclaw/kernel/service.py +++ b/src/athenaclaw/kernel/service.py @@ -68,7 +68,7 @@ portfolio.json — 用户的结构化当前持仓快照。 ✅ 写入:账户、币种现金、当前持仓、更新时间 ✅ 触发:用户发完整持仓截图、直接给出当前持仓,或明确说某笔已执行交易后当前仓位变了 - ✅ 例外:如果这次变化来自 trade_account/trade_plan/trade_apply 等远端 broker 工具,不自动改写 portfolio.json;只有用户明确要求同步工作区快照时才更新 + ✅ 例外:如果这次变化来自 trade_account/trade_execute 等远端 broker 工具,不自动改写 portfolio.json;只有用户明确要求同步工作区快照时才更新 ❌ 不写:想买/想卖的计划、watchlist、模糊推测、不完整截图 组织:单文件 JSON,按账户维护当前状态,不记录交易历史。 读取:当任务涉及仓位分析、风险分析、集中度、相关性、个性化建议时,优先调用 portfolio 工具读取。 @@ -171,14 +171,15 @@ 1. 远端 broker 账户不等于 portfolio.json。portfolio.json 是工作区快照,不会被交易工具自动改写。 2. 读取远端账户状态时,用 trade_account,不要猜账户,也不要把远端账户内容写成 portfolio.json。 -3. 所有执行动作必须先 trade_plan,再 trade_apply。不要直接调用 trade_apply,也不要把原始下单参数传给 trade_apply。 +3. 你是自主交易操作员:trade_execute 一步下单/撤单,直接对远端 broker 账户生效,不经过人工确认,调用即视为你已完成决策判断。执行前会经过风控裁决;被拒绝时返回 error_code=permission_denied,此时不要重试,向用户说明原因。 4. V1 只支持股票/ETF 的 LIMIT 限价单,以及撤销未完成订单。不支持 market/stop/条件单/衍生品。 -5. account_ref/order_ref/plan_id 都是系统返回的 opaque 引用。只能传递,不能猜测、拼接或修改。 -6. 在 trade_apply 返回成功之前,不能告诉用户“已提交”“已成交”。 +5. account_ref/order_ref 都是系统返回的 opaque 引用。只能传递,不能猜测、拼接或修改。 +6. 在 trade_execute 返回 status=ok 且订单进入终态(filled/partially_filled/cancelled 等)前,不能告诉用户“已提交”“已成交”。 7. 下单前先看 trade_account.list_accounts 返回的 supported_markets、account_status、account_kind;优先选择 account_status=active 且支持目标市场、类型适合股票/ETF 的账户。extra 是 provider 信息层,可用于解释限制,但不是通用契约。 8. 除 list_accounts 外,读取具体账户或订单时必须显式携带 account_ref 或 order_ref。不要省略,也不要假设系统会自动承接最近账户或最近订单。 9. 没有同 broker 的新鲜行情或明确 market-state 证据时,不得推断“更容易成交”“当前处于常规交易时段”。类似 session=RTH 只能理解为订单会话语义,不能理解为当前市场状态。 10. trade_account.get_positions 成功后,会把当前活动账户快照写入 Kernel.data['account'],供后续 compute 使用;这个 account 只代表当前活动账户,不是多账户容器。 +11. 如果这次会话里注册了 market_snapshot,下限价单前优先用它拿一次同 broker 的实时快照,用 update_time 判断新鲜度、用 last_price/bid_price/ask_price 判断限价是否偏离盘口过远;market_snapshot 只是当前一个时间点,不能替代 market_ohlcv 的历史序列。 """ @@ -387,7 +388,7 @@ def _assemble_system_prompt(self) -> None: from athenaclaw.kernel.seed import SEED_PROMPT identity = SEED_PROMPT parts = [identity, WORKSPACE_GUIDE, AUTOMATION_GUIDE] - if {"trade_account", "trade_plan", "trade_apply"} & set(self._tools): + if {"trade_account", "trade_execute"} & set(self._tools): parts.append(TRADE_GUIDE) runtime_paths_xml = self._runtime_paths_prompt() if runtime_paths_xml: diff --git a/src/athenaclaw/runtime/bundle.py b/src/athenaclaw/runtime/bundle.py index bb41a39..a0715f8 100644 --- a/src/athenaclaw/runtime/bundle.py +++ b/src/athenaclaw/runtime/bundle.py @@ -20,7 +20,7 @@ from athenaclaw.llm.providers import AnthropicProvider, LLMProvider, OpenAIChatProvider from athenaclaw.runtime.session_store import JsonSessionStore, SessionStore from athenaclaw.tools import bash, compute, edit, market, portfolio, read, trade, watchlist, web, write -from athenaclaw.trading import TradeAuditLog, TradeOrchestrator, TradePlanStore +from athenaclaw.trading import AllowAllGuard, TradeAuditLog, TradeOrchestrator from athenaclaw.subagents import SubAgentDef @@ -340,10 +340,14 @@ def build_kernel_bundle( if trade_adapter is not None: orchestrator = TradeOrchestrator( adapter=trade_adapter, - plan_store=TradePlanStore(state), + guard=AllowAllGuard(), audit_log=TradeAuditLog(state), ) trade.register(kernel, orchestrator) + if config.trade_broker == "futu": + # 快照走同一个 broker 的行情线路,和 market_cn/hk/us 选了谁无关 — + # 下单前要看的是"我即将成交的这条线路"当前价格,不是研究用的 OHLCV 数据源。 + market.register_snapshot(kernel, _make_adapter("futu", config)) read.register(kernel, workspace, cwd=cwd) write.register(kernel, workspace, cwd=cwd) edit.register(kernel, workspace, cwd=cwd) diff --git a/src/athenaclaw/tools/market/__init__.py b/src/athenaclaw/tools/market/__init__.py index 576967d..c260834 100644 --- a/src/athenaclaw/tools/market/__init__.py +++ b/src/athenaclaw/tools/market/__init__.py @@ -1,3 +1,4 @@ +from athenaclaw.tools.market.snapshot import register as register_snapshot from athenaclaw.tools.market.tool import register -__all__ = ["register"] +__all__ = ["register", "register_snapshot"] diff --git a/src/athenaclaw/tools/market/snapshot.py b/src/athenaclaw/tools/market/snapshot.py new file mode 100644 index 0000000..6db553d --- /dev/null +++ b/src/athenaclaw/tools/market/snapshot.py @@ -0,0 +1,64 @@ +""" +[INPUT]: typing.Protocol, athenaclaw.kernel (Kernel) +[OUTPUT]: SnapshotAdapter Protocol + register() +[POS]: 领域核心工具,与 market_ohlcv 同级;提供实时快照/报价,供下单前判断行情新鲜度,不落 DataStore 历史 +[PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md +""" + +from __future__ import annotations + +from typing import Any, Protocol + + +# ───────────────────────────────────────────────────────────────────────────── +# SnapshotAdapter Protocol +# ───────────────────────────────────────────────────────────────────────────── + +class SnapshotAdapter(Protocol): + """实时快照适配器接口 — 一次调用返回一批 symbol 的当前盘口,不需要订阅、不返回时间序列""" + + name: str + + def snapshot(self, symbols: list[str]) -> list[dict[str, Any]]: ... + + +# ───────────────────────────────────────────────────────────────────────────── +# 注册 +# ───────────────────────────────────────────────────────────────────────────── + +def register(kernel: object, adapter: SnapshotAdapter) -> None: + """向 Kernel 注册 market_snapshot 工具""" + + def market_snapshot(args: dict) -> dict: + symbols = args.get("symbols") or [] + if not isinstance(symbols, list) or not symbols: + return {"error": "缺少参数: symbols(非空数组)"} + try: + quotes = adapter.snapshot([str(item).strip() for item in symbols]) + except Exception as exc: + return {"error": str(exc)} + return {"status": "ok", "quotes": quotes} + + kernel.tool( + name="market_snapshot", + description=( + "获取一批 symbol 的实时快照:最新价、开高低、昨收、成交量额、买卖价等。" + "直接来自下单 broker 的行情线路,不需要订阅,不写入 market_ohlcv 使用的 DataStore 历史," + "只返回当前这一个时间点,不是时间序列。" + "何时使用: 下单前确认当前价格新鲜度、判断限价是否偏离盘口过远、撤单前重新核对最新成交价。" + "何时不要用: 需要历史 K 线或分钟级序列时用 market_ohlcv。" + "返回的 update_time 是 broker 侧行情时间戳,用它判断数据新鲜度,不得凭经验推断当前是否在正常交易时段。" + ), + parameters={ + "type": "object", + "properties": { + "symbols": { + "type": "array", + "items": {"type": "string"}, + "description": "标的代码数组,如 [\"AAPL\", \"00700.HK\", \"600519.SH\"]", + }, + }, + "required": ["symbols"], + }, + handler=market_snapshot, + ) diff --git a/src/athenaclaw/tools/trade/tool.py b/src/athenaclaw/tools/trade/tool.py index afd96ab..084b2dc 100644 --- a/src/athenaclaw/tools/trade/tool.py +++ b/src/athenaclaw/tools/trade/tool.py @@ -1,7 +1,7 @@ """ [INPUT]: athenaclaw.kernel, athenaclaw.trading [OUTPUT]: register() -[POS]: 交易工具层;暴露 trade_account/trade_plan/trade_apply 给 LLM +[POS]: 交易工具层;暴露 trade_account(只读)/trade_execute(下单执行)给 LLM [PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md """ @@ -64,46 +64,33 @@ def trade_account_handler(args: dict[str, Any]) -> dict[str, Any]: return error_payload(exc) return {"error": "未知 action,可用值: list_accounts/get_positions/get_open_orders/get_order_status/get_summary"} - def trade_plan_handler(args: dict[str, Any]) -> dict[str, Any]: + def trade_execute_handler(args: dict[str, Any]) -> dict[str, Any]: operation = str(args.get("operation") or "").strip().lower() + # automation 单点判定:本次调用是否来自无人值守 reaction,交给 RiskContext.automation, + # 供未来风控专题识别“automation + real”这个最高风险组合。 + automation = bool(getattr(getattr(kernel, "_tool_policy", None), "automation", False)) try: if operation == "submit_limit": account_ref = str(args.get("account_ref") or "").strip() if not account_ref: return _missing_account_ref() - plan = orchestrator.plan_submit_limit( + result = orchestrator.execute_limit( account_ref=account_ref, symbol=str(args.get("symbol") or "").strip(), side=str(args.get("side") or "").strip(), quantity=float(args.get("quantity")), limit_price=float(args.get("limit_price")), + automation=automation, ) - return {"status": "ok", **plan.to_dict()} - if operation == "cancel": + elif operation == "cancel": order_ref = str(args.get("order_ref") or "").strip() if not order_ref: return _missing_order_ref() - plan = orchestrator.plan_cancel(order_ref=order_ref) - return {"status": "ok", **plan.to_dict()} + result = orchestrator.execute_cancel(order_ref=order_ref, automation=automation) + else: + return {"error": "未知 operation,可用值: submit_limit/cancel"} except (TypeError, ValueError): - return {"error": "trade_plan 参数不合法"} - except TradeError as exc: - return error_payload(exc) - return {"error": "未知 operation,可用值: submit_limit/cancel"} - - def trade_apply_handler(args: dict[str, Any]) -> dict[str, Any]: - plan_id = str(args.get("plan_id") or "").strip() - if not plan_id: - return {"error": "缺少参数: plan_id"} - record = orchestrator.get_plan(plan_id) - if record is None: - return {"error": f"未找到 plan: {plan_id}", "error_code": "plan_not_found"} - plan = record["plan"] - confirm_text = str(plan.get("confirm_text") or f"确认执行 {plan_id} 吗?") - if not kernel.request_confirm(confirm_text): - return {"status": "cancelled", "plan_id": plan_id, "message": "用户取消确认"} - try: - result = orchestrator.apply(plan_id) + return {"error": "trade_execute 参数不合法"} except TradeError as exc: return error_payload(exc) @@ -144,12 +131,15 @@ def trade_apply_handler(args: dict[str, Any]) -> dict[str, Any]: ) kernel.tool( - name="trade_plan", + name="trade_execute", description=( - "创建交易执行计划,不直接产生外部副作用。只支持股票/ETF 的 LIMIT 限价单与撤单。" - "先调用 trade_plan,再调用 trade_apply。不要跳过 plan。" - "plan 返回的 plan_id 是一次性短期令牌;不要伪造或猜测。" - "submit_limit 必须显式提供 account_ref,cancel 必须显式提供 order_ref。" + "直接对远端 broker 账户执行交易动作:提交股票/ETF 限价单,或撤销未完成订单。" + "你是自主交易操作员:本工具一步下单/撤单,不经过人工确认,调用即视为你已完成决策判断。" + "执行前会经过风控裁决;被拒绝时返回 error_code=permission_denied。" + "submit_limit 必须显式提供 account_ref/symbol/side/quantity/limit_price;" + "cancel 必须显式提供 order_ref。account_ref/order_ref 是系统返回的 opaque 引用,只能传递,不能猜测或拼接。" + "执行成功后会自动刷新订单当前状态;若订单已部分或全部成交,会刷新当前活动账户快照。" + "status=ok 只表示工具执行成功,不表示订单动作已进入终态;请结合 finalized 与 order_status 判断。" ), parameters={ "type": "object", @@ -164,23 +154,5 @@ def trade_apply_handler(args: dict[str, Any]) -> dict[str, Any]: }, "required": ["operation"], }, - handler=trade_plan_handler, - ) - - kernel.tool( - name="trade_apply", - description=( - "执行 trade_plan 生成的计划。只接受 plan_id,执行前会请求用户确认。" - "不要把原始下单参数直接传给 trade_apply。" - "执行成功后会自动刷新订单当前状态;若订单已部分或全部成交,会刷新当前活动账户快照。" - "status=ok 只表示工具执行成功,不表示订单动作已进入终态;请结合 finalized 与 order_status 判断。" - ), - parameters={ - "type": "object", - "properties": { - "plan_id": {"type": "string", "description": "trade_plan 返回的一次性 plan_id"}, - }, - "required": ["plan_id"], - }, - handler=trade_apply_handler, + handler=trade_execute_handler, ) diff --git a/src/athenaclaw/trading/CLAUDE.md b/src/athenaclaw/trading/CLAUDE.md new file mode 100644 index 0000000..b2b550c --- /dev/null +++ b/src/athenaclaw/trading/CLAUDE.md @@ -0,0 +1,22 @@ +# trading/ +> L2 | 父级: ../CLAUDE.md + +交易边界层:把 Agent 的下单/撤单意图转成确定性、经风控裁决、可执行、可回读的 canonical 流程。自主交易操作员范式——`orchestrator` 由旧的 plan/apply 二段式改为单步 `execute_limit`/`execute_cancel`,中间不经过人工确认。 + +## 成员清单 + +types.py: canonical dataclass(`TradeAccountDescriptor`/`TradeAccountSnapshot`/`TradeAccountSummary`/`TradeApplyResult`/`TradeCapabilities`/`TradeOpenOrder`/`TradeOrderSnapshot`/`TradePosition`/`TradePreview`/`TradeReceipt`/`SubmitLimitOrderIntent`)与 `account_ref`/`order_ref` 编解码。`TradeApplyResult.plan_id` 语义已从"待确认计划 id"改为"一次性执行 id",纯审计用途 +errors.py: `TradeErrorCode`/`TradeError`/`error_payload` —— 交易域统一错误语义 +protocol.py: `TradeBrokerAdapter` Protocol —— broker adapter 必须实现的公共能力集合 +policy.py: `RiskGuard` seam —— `RiskAction`(ALLOW/DENY/ESCALATE)/`RiskContext`(operation/env/automation/intent)/`RiskDecision`/`AllowAllGuard`(这期唯一实现,恒 ALLOW,占位 + TODO,真实风控留后续专题) +orchestrator.py: `TradeOrchestrator` —— 核心编排对象。`execute_limit`/`execute_cancel` 一次调用内完成校验+preview+`guard.evaluate`+下单/撤单+状态回读,guard 裁决与执行结果均写入 `TradeAuditLog` +store.py: `TradeAuditLog` —— 交易审计的轻量持久层(append-only jsonl) +snapshots.py: `build_kernel_account` —— `TradeAccountSnapshot` 到 `Kernel.data["account"]` 的映射 + +## 安全模型 + +确认(human-in-the-loop)与风控(risk control)是两件事:交易链路已彻底删除 `request_confirm` 依赖;`RiskGuard.evaluate` 是唯一的程序化安全边界,`orchestrator.execute_*` 在下单/撤单前必调用。这期 `AllowAllGuard` 不拦任何单——真实风控(notional/价格偏离/频率/集中度/禁买清单)是独立专题,`RiskContext.automation` 字段已为 automation×real 场景预留裁决挂钩。 + +详见 `docs/trading.md`。 + +[PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md diff --git a/src/athenaclaw/trading/__init__.py b/src/athenaclaw/trading/__init__.py index 96dfca2..1f997f5 100644 --- a/src/athenaclaw/trading/__init__.py +++ b/src/athenaclaw/trading/__init__.py @@ -1,8 +1,9 @@ from athenaclaw.trading.errors import TradeError, TradeErrorCode, error_payload from athenaclaw.trading.orchestrator import TradeOrchestrator +from athenaclaw.trading.policy import AllowAllGuard, RiskAction, RiskContext, RiskDecision, RiskGuard from athenaclaw.trading.protocol import TradeBrokerAdapter from athenaclaw.trading.snapshots import build_kernel_account -from athenaclaw.trading.store import TradeAuditLog, TradePlanStore +from athenaclaw.trading.store import TradeAuditLog from athenaclaw.trading.types import ( TradeAccountDescriptor, TradeAccountSnapshot, @@ -11,7 +12,6 @@ TradeCapabilities, TradeOpenOrder, TradeOrderSnapshot, - TradePlan, TradePosition, TradePreview, TradeReceipt, @@ -23,7 +23,12 @@ ) __all__ = [ + "AllowAllGuard", "SubmitLimitOrderIntent", + "RiskAction", + "RiskContext", + "RiskDecision", + "RiskGuard", "TradeAccountDescriptor", "TradeAccountSnapshot", "TradeAccountSummary", @@ -36,8 +41,6 @@ "TradeOpenOrder", "TradeOrchestrator", "TradeOrderSnapshot", - "TradePlan", - "TradePlanStore", "TradePosition", "TradePreview", "TradeReceipt", diff --git a/src/athenaclaw/trading/errors.py b/src/athenaclaw/trading/errors.py index 056e2e0..f6e5d6c 100644 --- a/src/athenaclaw/trading/errors.py +++ b/src/athenaclaw/trading/errors.py @@ -21,9 +21,6 @@ class TradeErrorCode(str, Enum): MISSING_ORDER_REF = "missing_order_ref" ORDER_NOT_FOUND = "order_not_found" ORDER_NOT_CANCELLABLE = "order_not_cancellable" - PLAN_NOT_FOUND = "plan_not_found" - PLAN_EXPIRED = "plan_expired" - PLAN_ALREADY_APPLIED = "plan_already_applied" INVALID_ACCOUNT_REF = "invalid_account_ref" INVALID_ORDER_REF = "invalid_order_ref" INVALID_SIDE = "invalid_side" diff --git a/src/athenaclaw/trading/orchestrator.py b/src/athenaclaw/trading/orchestrator.py index b567942..c0f5db1 100644 --- a/src/athenaclaw/trading/orchestrator.py +++ b/src/athenaclaw/trading/orchestrator.py @@ -1,21 +1,22 @@ """ [INPUT]: uuid, datetime, athenaclaw.trading.* [OUTPUT]: TradeOrchestrator -[POS]: 交易边界层的核心编排对象;负责 plan/apply/状态回读 +[POS]: 交易边界层的核心编排对象;负责账户/订单读取与单步下单/撤单执行 [PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md """ from __future__ import annotations from decimal import Decimal, InvalidOperation -from datetime import datetime, timedelta, timezone +from datetime import datetime, timezone from time import sleep from uuid import uuid4 from athenaclaw.tools.market.schema import normalize_symbol from athenaclaw.trading.errors import TradeError, TradeErrorCode +from athenaclaw.trading.policy import RiskAction, RiskContext, RiskGuard from athenaclaw.trading.protocol import TradeBrokerAdapter -from athenaclaw.trading.store import TradeAuditLog, TradePlanStore +from athenaclaw.trading.store import TradeAuditLog from athenaclaw.trading.types import ( SubmitLimitOrderIntent, TradeAccountSnapshot, @@ -23,9 +24,9 @@ TradeApplyResult, TradeOpenOrder, TradeOrderSnapshot, - TradePlan, - TradePreview, TradeReceipt, + decode_account_ref, + decode_order_ref, ) @@ -37,15 +38,13 @@ def __init__( self, *, adapter: TradeBrokerAdapter, - plan_store: TradePlanStore, + guard: RiskGuard, audit_log: TradeAuditLog, - plan_ttl_sec: int = 120, cancel_confirm_delays: tuple[float, ...] = (0.2, 0.5, 1.0), ) -> None: self._adapter = adapter - self._plan_store = plan_store + self._guard = guard self._audit_log = audit_log - self._plan_ttl_sec = plan_ttl_sec self._cancel_confirm_delays = cancel_confirm_delays def list_accounts(self): @@ -81,7 +80,7 @@ def get_order_status(self, order_ref: str) -> TradeOrderSnapshot: self._require_order_ref(order_ref) return self._adapter.get_order_status(order_ref) - def plan_submit_limit( + def execute_limit( self, *, account_ref: str, @@ -89,7 +88,13 @@ def plan_submit_limit( side: str, quantity: float, limit_price: float, - ) -> TradePlan: + automation: bool = False, + ) -> TradeApplyResult: + """校验 + 预检规整 + 风控裁决 + 下单,一步到底。 + + 自主操作员不再有「plan 给人看,apply 再执行」的中间站——这里就是 + 唯一的执行入口,guard.evaluate 是唯一的安全边界。 + """ self._require_account_ref(account_ref) requested_limit_price = self._validate_price(limit_price) intent = SubmitLimitOrderIntent( @@ -120,30 +125,33 @@ def plan_submit_limit( 0, f"limit_price normalized from {_format_price(intent.limit_price)} to {_format_price(normalized_intent.limit_price)}{suffix}", ) - created_at = _utc_now_iso() - expires_at = _utc_expiry(self._plan_ttl_sec) - plan = TradePlan( - plan_id=f"plan-{uuid4().hex}", + + env = self._account_parts(account_ref)["env"] + ctx = RiskContext(operation="submit_limit", env=env, automation=automation, intent=normalized_intent.to_dict()) + decision = self._guard.evaluate(ctx) + self._audit_log.append({"event": "trade.guard.decision", "context": ctx.to_dict(), "decision": decision.to_dict()}) + if decision.action is not RiskAction.ALLOW: + raise TradeError(TradeErrorCode.PERMISSION_DENIED, decision.reason or "交易未授权") + + receipt = self._adapter.submit_limit_order(normalized_intent) + order_status = self._adapter.get_order_status(receipt.order_ref) + account_snapshot = None + if order_status.status in {"partially_filled", "filled"}: + account_snapshot = self.get_positions(normalized_intent.account_ref) + result = TradeApplyResult( + plan_id=f"exec-{uuid4().hex}", operation="submit_limit", - plan_summary=_submit_plan_summary(normalized_intent), - confirm_text=_submit_confirm_text(normalized_intent), + result_summary=_submit_result_summary(normalized_intent, order_status), + receipt=receipt, + order_status=order_status, + account_snapshot=account_snapshot, + finalized=order_status.status in _TERMINAL_STATUSES, warnings=tuple(warnings), - created_at=created_at, - expires_at=expires_at, - normalized_intent=normalized_intent.to_dict(), ) - self._plan_store.save( - plan, - payload={ - "intent": normalized_intent.to_dict(), - "requested_limit_price": intent.limit_price, - "preview": _preview_payload(preview), - }, - ) - self._audit_log.append({"event": "trade.plan.created", "plan": plan.to_dict()}) - return plan + self._audit_log.append({"event": "trade.executed", "result": result.to_dict(), "guard": decision.to_dict()}) + return result - def plan_cancel(self, *, order_ref: str) -> TradePlan: + def execute_cancel(self, *, order_ref: str, automation: bool = False) -> TradeApplyResult: self._require_order_ref(order_ref) order = self.get_order_status(order_ref) if order.status in _TERMINAL_STATUSES: @@ -152,82 +160,31 @@ def plan_cancel(self, *, order_ref: str) -> TradePlan: f"订单当前状态不可撤销: {order.status}", details={"order_ref": order_ref, "status": order.status}, ) - created_at = _utc_now_iso() - plan = TradePlan( - plan_id=f"plan-{uuid4().hex}", + + env = decode_order_ref(order_ref)["env"] + ctx = RiskContext(operation="cancel", env=env, automation=automation, intent={"order_ref": order_ref}) + decision = self._guard.evaluate(ctx) + self._audit_log.append({"event": "trade.guard.decision", "context": ctx.to_dict(), "decision": decision.to_dict()}) + if decision.action is not RiskAction.ALLOW: + raise TradeError(TradeErrorCode.PERMISSION_DENIED, decision.reason or "交易未授权") + + receipt = self._adapter.cancel_order(order_ref) + order_status = self._confirm_order_status(order_ref) + finalized = order_status.status in _TERMINAL_STATUSES + warnings = ("cancel_requested_not_finalized",) if not finalized else () + result = TradeApplyResult( + plan_id=f"exec-{uuid4().hex}", operation="cancel", - plan_summary=_cancel_plan_summary(order), - confirm_text=_cancel_confirm_text(order), - warnings=(), - created_at=created_at, - expires_at=_utc_expiry(self._plan_ttl_sec), + result_summary=_cancel_result_summary(order_status, finalized=finalized), + receipt=receipt, + order_status=order_status, + account_snapshot=None, + finalized=finalized, + warnings=warnings, ) - self._plan_store.save(plan, payload={"order_ref": order_ref}) - self._audit_log.append({"event": "trade.plan.created", "plan": plan.to_dict()}) - return plan - - def apply(self, plan_id: str) -> TradeApplyResult: - record = self._plan_store.load(plan_id) - if record is None: - raise TradeError(TradeErrorCode.PLAN_NOT_FOUND, f"未找到 plan: {plan_id}") - status = str(record.get("status") or "") - if status == "applied": - raise TradeError(TradeErrorCode.PLAN_ALREADY_APPLIED, "该 plan 已执行") - plan = record["plan"] - if _is_expired(str(plan["expires_at"])): - self._plan_store.mark_expired(plan_id) - raise TradeError(TradeErrorCode.PLAN_EXPIRED, "该 plan 已过期,请重新 plan") - - operation = str(plan["operation"]) - if operation == "submit_limit": - payload = record["payload"]["intent"] - intent = SubmitLimitOrderIntent( - account_ref=str(payload["account_ref"]), - symbol=str(payload["symbol"]), - side=str(payload["side"]), - quantity=float(payload["quantity"]), - limit_price=float(payload["limit_price"]), - ) - receipt = self._adapter.submit_limit_order(intent) - order_status = self._adapter.get_order_status(receipt.order_ref) - account_snapshot = None - if order_status.status in {"partially_filled", "filled"}: - account_snapshot = self.get_positions(intent.account_ref) - result = TradeApplyResult( - plan_id=plan_id, - operation=operation, - result_summary=_submit_result_summary(intent, order_status), - receipt=receipt, - order_status=order_status, - account_snapshot=account_snapshot, - finalized=order_status.status in _TERMINAL_STATUSES, - ) - elif operation == "cancel": - order_ref = str(record["payload"]["order_ref"]) - receipt = self._adapter.cancel_order(order_ref) - order_status = self._confirm_order_status(order_ref) - finalized = order_status.status in _TERMINAL_STATUSES - warnings = ("cancel_requested_not_finalized",) if not finalized else () - result = TradeApplyResult( - plan_id=plan_id, - operation=operation, - result_summary=_cancel_result_summary(order_status, finalized=finalized), - receipt=receipt, - order_status=order_status, - account_snapshot=None, - finalized=finalized, - warnings=warnings, - ) - else: - raise TradeError(TradeErrorCode.UNSUPPORTED_OPERATION, f"不支持的 plan 操作: {operation}") - - self._plan_store.mark_applied(plan_id, result=result.to_dict(), applied_at=_utc_now_iso()) - self._audit_log.append({"event": "trade.plan.applied", "plan_id": plan_id, "result": result.to_dict()}) + self._audit_log.append({"event": "trade.executed", "result": result.to_dict(), "guard": decision.to_dict()}) return result - def get_plan(self, plan_id: str) -> dict[str, object] | None: - return self._plan_store.load(plan_id) - def _adapter_capabilities(self): return self._adapter.capabilities() @@ -269,8 +226,6 @@ def _confirm_order_status(self, order_ref: str) -> TradeOrderSnapshot: @staticmethod def _require_account_ref(account_ref: str) -> None: - from athenaclaw.trading.types import decode_account_ref - try: decode_account_ref(account_ref) except Exception as exc: # pragma: no cover - thin validation wrapper @@ -278,8 +233,6 @@ def _require_account_ref(account_ref: str) -> None: @staticmethod def _require_order_ref(order_ref: str) -> None: - from athenaclaw.trading.types import decode_order_ref - try: decode_order_ref(order_ref) except Exception as exc: # pragma: no cover - thin validation wrapper @@ -287,8 +240,6 @@ def _require_order_ref(order_ref: str) -> None: @staticmethod def _account_parts(account_ref: str) -> dict[str, str]: - from athenaclaw.trading.types import decode_account_ref - try: return decode_account_ref(account_ref) except Exception as exc: # pragma: no cover - validated by _require_account_ref @@ -299,42 +250,6 @@ def _utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat() -def _utc_expiry(ttl_sec: int) -> str: - return (datetime.now(timezone.utc) + timedelta(seconds=ttl_sec)).isoformat() - - -def _is_expired(ts: str) -> bool: - return datetime.fromisoformat(ts) < datetime.now(timezone.utc) - - -def _preview_payload(preview: TradePreview | None) -> dict | None: - if preview is None: - return None - return preview.to_dict() - - -def _submit_plan_summary(intent: SubmitLimitOrderIntent) -> str: - return ( - f"计划提交 {intent.side.upper()} {intent.quantity:g} {intent.symbol} " - f"限价 {_format_price(intent.limit_price)}" - ) - - -def _submit_confirm_text(intent: SubmitLimitOrderIntent) -> str: - return ( - f"确认提交 {intent.side.upper()} {intent.quantity:g} {intent.symbol} " - f"限价 {_format_price(intent.limit_price)} 吗?" - ) - - -def _cancel_plan_summary(order: TradeOrderSnapshot) -> str: - return f"计划撤销 {order.symbol} {order.side.upper()} {order.quantity:g} 的未完成订单" - - -def _cancel_confirm_text(order: TradeOrderSnapshot) -> str: - return f"确认撤销 {order.symbol} {order.side.upper()} {order.quantity:g} 的订单吗?" - - def _submit_result_summary(intent: SubmitLimitOrderIntent, order: TradeOrderSnapshot) -> str: return ( f"已提交 {intent.side.upper()} {intent.quantity:g} {intent.symbol} " diff --git a/src/athenaclaw/trading/policy.py b/src/athenaclaw/trading/policy.py new file mode 100644 index 0000000..62bcc61 --- /dev/null +++ b/src/athenaclaw/trading/policy.py @@ -0,0 +1,69 @@ +""" +[INPUT]: dataclasses, enum, typing.Protocol +[OUTPUT]: RiskAction, RiskContext, RiskDecision, RiskGuard, AllowAllGuard +[POS]: 交易执行前的风控裁决单点 seam;orchestrator.execute_* 在下单前必须过 guard.evaluate +[PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md +""" + +from __future__ import annotations + +from dataclasses import dataclass +from enum import Enum +from typing import Any, Protocol + + +class RiskAction(str, Enum): + ALLOW = "allow" + DENY = "deny" + ESCALATE = "escalate" # 契约预留:real 大额逃生口,未来风控专题落地 + + +@dataclass(frozen=True) +class RiskContext: + """裁决输入。intent 携带原始下单意图(symbol/side/quantity/limit_price 等), + automation 标记该次执行是否发生在无人值守的 automation reaction 中—— + 这是未来给 automation×real 做 DENY/ESCALATE 的契约点,这期只透传不使用。""" + + operation: str + env: str + automation: bool = False + intent: dict[str, Any] | None = None + + def to_dict(self) -> dict[str, Any]: + return { + "operation": self.operation, + "env": self.env, + "automation": self.automation, + "intent": self.intent, + } + + +@dataclass(frozen=True) +class RiskDecision: + action: RiskAction + reason: str | None = None + + def to_dict(self) -> dict[str, Any]: + return {"action": self.action.value, "reason": self.reason} + + +class RiskGuard(Protocol): + def evaluate(self, ctx: RiskContext) -> RiskDecision: ... + + +@dataclass +class AllowAllGuard: + """[TODO 风控专题] 占位护栏,恒 ALLOW,不拦任何单。 + + 这期用户明确要求不做风控拦截——真实风控是独立的大设计(notional 上限、 + 限价偏离市价保护、日内下单频率、单标的集中度、禁买清单……任何单一维度 + 都无法满足真实需求),留给后续专题重新设计。 + + 这个类存在的意义是把 orchestrator.execute_* 里的裁决点固定下来:风控 + 专题落地时只需替换这个实现(或注入另一个 RiskGuard),不改动 orchestrator + 与工具层任何一行。契约方向:对 simulate 与 real 都生效;real 大额场景走 + RiskAction.ESCALATE;automation=True(无人值守)场景应优先收紧 real。 + """ + + def evaluate(self, ctx: RiskContext) -> RiskDecision: + return RiskDecision(RiskAction.ALLOW) diff --git a/src/athenaclaw/trading/store.py b/src/athenaclaw/trading/store.py index edbd698..c036465 100644 --- a/src/athenaclaw/trading/store.py +++ b/src/athenaclaw/trading/store.py @@ -1,7 +1,7 @@ """ [INPUT]: json, pathlib, datetime -[OUTPUT]: TradePlanStore, TradeAuditLog -[POS]: 交易 plan 与审计的轻量持久层 +[OUTPUT]: TradeAuditLog +[POS]: 交易审计的轻量持久层 [PROTOCOL]: 变更时更新此头部,然后检查 CLAUDE.md """ @@ -12,58 +12,6 @@ from pathlib import Path from typing import Any -from athenaclaw.trading.types import TradePlan - - -class TradePlanStore: - def __init__(self, state_dir: Path) -> None: - self._plans_dir = state_dir / "trade" / "plans" - - def save(self, plan: TradePlan, *, payload: dict[str, Any]) -> None: - self._plans_dir.mkdir(parents=True, exist_ok=True) - record = { - "status": "pending", - "plan": plan.to_dict(), - "payload": payload, - "applied_at": None, - "result": None, - } - self._path(plan.plan_id).write_text( - json.dumps(record, ensure_ascii=False, indent=2) + "\n", - encoding="utf-8", - ) - - def load(self, plan_id: str) -> dict[str, Any] | None: - path = self._path(plan_id) - if not path.exists(): - return None - return json.loads(path.read_text(encoding="utf-8")) - - def mark_applied(self, plan_id: str, *, result: dict[str, Any], applied_at: str | None = None) -> None: - record = self.load(plan_id) - if record is None: - return - record["status"] = "applied" - record["applied_at"] = applied_at or datetime.now(timezone.utc).isoformat() - record["result"] = result - self._path(plan_id).write_text( - json.dumps(record, ensure_ascii=False, indent=2) + "\n", - encoding="utf-8", - ) - - def mark_expired(self, plan_id: str) -> None: - record = self.load(plan_id) - if record is None: - return - record["status"] = "expired" - self._path(plan_id).write_text( - json.dumps(record, ensure_ascii=False, indent=2) + "\n", - encoding="utf-8", - ) - - def _path(self, plan_id: str) -> Path: - return self._plans_dir / f"{plan_id}.json" - class TradeAuditLog: def __init__(self, state_dir: Path) -> None: diff --git a/src/athenaclaw/trading/types.py b/src/athenaclaw/trading/types.py index 2385ae1..91d2dca 100644 --- a/src/athenaclaw/trading/types.py +++ b/src/athenaclaw/trading/types.py @@ -225,31 +225,6 @@ def to_dict(self) -> dict[str, Any]: } -@dataclass(frozen=True) -class TradePlan: - plan_id: str - operation: str - plan_summary: str - confirm_text: str - warnings: tuple[str, ...] - created_at: str - expires_at: str - normalized_intent: dict[str, Any] | None = None - - def to_dict(self) -> dict[str, Any]: - return { - "plan_id": self.plan_id, - "operation": self.operation, - "plan_summary": self.plan_summary, - "confirm_text": self.confirm_text, - "warnings": list(self.warnings), - "created_at": self.created_at, - "expires_at": self.expires_at, - "normalized_intent": self.normalized_intent, - "requires_confirmation": True, - } - - @dataclass(frozen=True) class TradeReceipt: order_ref: str diff --git a/tests/integration/test_runtime_permissions.py b/tests/integration/test_runtime_permissions.py index af585be..071662c 100644 --- a/tests/integration/test_runtime_permissions.py +++ b/tests/integration/test_runtime_permissions.py @@ -157,6 +157,6 @@ def preview_limit_order(self, intent): ) assert "trade_account" in bundle.kernel._tools - assert "trade_plan" in bundle.kernel._tools - assert "trade_apply" in bundle.kernel._tools + assert "trade_execute" in bundle.kernel._tools + assert "market_snapshot" in bundle.kernel._tools assert "" in (bundle.kernel._system_prompt or "") diff --git a/tests/unit/test_futu_market_adapter.py b/tests/unit/test_futu_market_adapter.py index 763b7c1..e131224 100644 --- a/tests/unit/test_futu_market_adapter.py +++ b/tests/unit/test_futu_market_adapter.py @@ -55,6 +55,8 @@ def __init__(self) -> None: self.history_calls: list[dict] = [] self.latest_calls: list[dict] = [] self.subscribe_calls: list[dict] = [] + self.snapshot_response: tuple[int, object] = (_FakeFutu.RET_OK, pd.DataFrame()) + self.snapshot_calls: list[list[str]] = [] def request_history_kline(self, **kwargs): self.history_calls.append(kwargs) @@ -74,6 +76,10 @@ def get_cur_kline(self, **kwargs): self.latest_calls.append(kwargs) return self.latest_response + def get_market_snapshot(self, code_list): + self.snapshot_calls.append(list(code_list)) + return self.snapshot_response + def _make_adapter(monkeypatch, quote_ctx: _FakeQuoteContext) -> FutuAdapter: monkeypatch.setattr("athenaclaw.integrations.market.futu._load_futu", lambda: _FakeFutu) @@ -222,3 +228,65 @@ def test_futu_history_permission_error_is_translated(monkeypatch): def test_market_query_uses_hk_timezone(): query = build_market_query(symbol="00700.HK", interval="1m", mode="history") assert query.timezone == "Asia/Hong_Kong" + + +def test_futu_snapshot_normalizes_fields_without_subscription(monkeypatch): + quote_ctx = _FakeQuoteContext() + quote_ctx.snapshot_response = ( + _FakeFutu.RET_OK, + pd.DataFrame( + [ + { + "code": "US.AAPL", + "name": "Apple Inc", + "last_price": 180.12, + "open_price": 178.0, + "high_price": 181.0, + "low_price": 177.5, + "prev_close_price": 179.0, + "volume": 1_200_000, + "turnover": 216_000_000.0, + "turnover_rate": 0.012, + "bid_price": 180.10, + "ask_price": 180.14, + "bid_vol": 100, + "ask_vol": 200, + "price_spread": 0.04, + "suspension": False, + "update_time": "2026-07-11 21:59:59", + } + ] + ), + ) + adapter = _make_adapter(monkeypatch, quote_ctx) + + result = adapter.snapshot(["AAPL"]) + + assert quote_ctx.snapshot_calls == [["US.AAPL"]] + assert len(result) == 1 + row = result[0] + assert row["symbol"] == "US.AAPL" + assert row["last_price"] == 180.12 + assert row["bid_price"] == 180.10 + assert row["ask_price"] == 180.14 + assert row["suspension"] is False + assert row["update_time"] == "2026-07-11 21:59:59" + + +def test_futu_snapshot_empty_frame_returns_empty_list(monkeypatch): + quote_ctx = _FakeQuoteContext() + quote_ctx.snapshot_response = (_FakeFutu.RET_OK, pd.DataFrame()) + adapter = _make_adapter(monkeypatch, quote_ctx) + + result = adapter.snapshot(["AAPL"]) + + assert result == [] + + +def test_futu_snapshot_provider_error_is_translated(monkeypatch): + quote_ctx = _FakeQuoteContext() + quote_ctx.snapshot_response = (1, "没有权限或者额度不足") + adapter = _make_adapter(monkeypatch, quote_ctx) + + with pytest.raises(ValueError): + adapter.snapshot(["AAPL"]) diff --git a/tests/unit/test_futu_trade_adapter.py b/tests/unit/test_futu_trade_adapter.py index 2e39841..c0c6064 100644 --- a/tests/unit/test_futu_trade_adapter.py +++ b/tests/unit/test_futu_trade_adapter.py @@ -120,6 +120,13 @@ def acctradinginfo_query(self, **_kwargs): ) +class _UnlockRejectedContext: + """place_order 返回未解锁错误 —— 模拟 REAL 账户在 OpenD 侧尚未手工解锁的场景。""" + + def place_order(self, **_kwargs): + return (1, "unlock_trade required before placing order") + + def _make_adapter(context) -> FutuTradeAdapter: adapter = FutuTradeAdapter(config=FutuTradeConfig()) adapter._manager = _FakeManager(context) @@ -189,3 +196,23 @@ def test_futu_preview_returns_soft_warning_and_limits(monkeypatch): assert preview.normalized_limit_price == 174.03 assert preview.normalization_reason == "fallback_us_default" assert context.preview_calls == 1 + + +def test_futu_submit_limit_order_trade_locked_reports_manual_unlock_hint(monkeypatch): + """零密码设计的落地验证:REAL 未解锁时报错必须指向 OpenD 手工解锁,不能暗示代码能自动解锁。""" + monkeypatch.setattr("athenaclaw.integrations.futu.trade_adapter._load_futu", lambda: _FakeFutu) + adapter = _make_adapter(_UnlockRejectedContext()) + intent = SubmitLimitOrderIntent( + account_ref=encode_account_ref(broker="futu", env="real", account_id="1001"), + symbol="AAPL", + side="buy", + quantity=1, + limit_price=180.0, + ) + + with pytest.raises(TradeError) as exc: + adapter.submit_limit_order(intent) + + assert exc.value.code == TradeErrorCode.TRADE_LOCKED + assert "OpenD 客户端手工解锁" in exc.value.message + assert "不持有交易密码" in exc.value.message diff --git a/tests/unit/test_trading.py b/tests/unit/test_trading.py index f6d8b2c..664daec 100644 --- a/tests/unit/test_trading.py +++ b/tests/unit/test_trading.py @@ -7,16 +7,18 @@ from athenaclaw.automation.policy import AutomationToolPolicy from athenaclaw.kernel import Kernel from athenaclaw.trading import ( + AllowAllGuard, + RiskAction, + RiskContext, + RiskDecision, SubmitLimitOrderIntent, TradeAccountDescriptor, - TradeAccountSnapshot, TradeAccountSummary, TradeAuditLog, TradeCapabilities, TradeOpenOrder, TradeOrchestrator, TradeOrderSnapshot, - TradePlanStore, TradePosition, TradePreview, TradeReceipt, @@ -163,23 +165,35 @@ def _current_order_status(self) -> str: return "filled" if self._submitted else "submitted" -def _make_orchestrator(tmp_path: Path) -> tuple[TradeOrchestrator, _FakeTradeAdapter]: +class _DenyGuard: + """总是拒绝的 RiskGuard 测试替身 —— 验证 guard-gate 与审计日志的接线。""" + + def __init__(self, reason: str = "risk denied") -> None: + self.reason = reason + self.seen: list[RiskContext] = [] + + def evaluate(self, ctx: RiskContext) -> RiskDecision: + self.seen.append(ctx) + return RiskDecision(RiskAction.DENY, self.reason) + + +def _make_orchestrator(tmp_path: Path, *, guard=None) -> tuple[TradeOrchestrator, _FakeTradeAdapter]: adapter = _FakeTradeAdapter() orchestrator = TradeOrchestrator( adapter=adapter, - plan_store=TradePlanStore(tmp_path / "state"), + guard=guard or AllowAllGuard(), audit_log=TradeAuditLog(tmp_path / "state"), cancel_confirm_delays=(0.0, 0.0, 0.0), ) return orchestrator, adapter -def test_trade_orchestrator_plan_and_apply_submit_limit(tmp_path): +def test_trade_orchestrator_execute_limit_submits_and_finalizes(tmp_path): orchestrator, adapter = _make_orchestrator(tmp_path) adapter.preview_normalized_limit_price = 174.03 adapter.preview_normalization_reason = "fallback_us_default" - plan = orchestrator.plan_submit_limit( + result = orchestrator.execute_limit( account_ref=adapter.account_ref, symbol="AAPL", side="buy", @@ -188,48 +202,17 @@ def test_trade_orchestrator_plan_and_apply_submit_limit(tmp_path): ) assert adapter.preview_calls == 1 - assert plan.normalized_intent is not None - assert plan.normalized_intent["limit_price"] == 174.03 - assert "normalized from 174.034 to 174.03" in plan.warnings[0] - assert tuple(plan.warnings[1:]) == ("order_session=RTH",) - - result = orchestrator.apply(plan.plan_id) - assert result.operation == "submit_limit" assert result.finalized is True assert result.order_status is not None assert result.order_status.status == "filled" assert result.account_snapshot is not None assert result.account_snapshot.positions[0].quantity == 10 - stored = orchestrator.get_plan(plan.plan_id) - assert stored is not None - assert stored["status"] == "applied" + assert "normalized from 174.034 to 174.03" in result.warnings[0] + assert tuple(result.warnings[1:]) == ("order_session=RTH",) -def test_trade_orchestrator_rejects_expired_plan(tmp_path): - adapter = _FakeTradeAdapter() - orchestrator = TradeOrchestrator( - adapter=adapter, - plan_store=TradePlanStore(tmp_path / "state"), - audit_log=TradeAuditLog(tmp_path / "state"), - plan_ttl_sec=-1, - ) - - plan = orchestrator.plan_submit_limit( - account_ref=adapter.account_ref, - symbol="AAPL", - side="buy", - quantity=10, - limit_price=180, - ) - - with pytest.raises(TradeError) as exc: - orchestrator.apply(plan.plan_id) - - assert exc.value.code == TradeErrorCode.PLAN_EXPIRED - - -def test_trade_orchestrator_rejects_preview_failure_during_plan(tmp_path): +def test_trade_orchestrator_rejects_preview_failure(tmp_path): orchestrator, adapter = _make_orchestrator(tmp_path) adapter.preview_error = TradeError( TradeErrorCode.ACCOUNT_MARKET_UNSUPPORTED, @@ -237,7 +220,7 @@ def test_trade_orchestrator_rejects_preview_failure_during_plan(tmp_path): ) with pytest.raises(TradeError) as exc: - orchestrator.plan_submit_limit( + orchestrator.execute_limit( account_ref=adapter.account_ref, symbol="AAPL", side="buy", @@ -249,11 +232,78 @@ def test_trade_orchestrator_rejects_preview_failure_during_plan(tmp_path): assert exc.value.code == TradeErrorCode.ACCOUNT_MARKET_UNSUPPORTED +def test_trade_orchestrator_guard_deny_blocks_submit_and_is_audited(tmp_path): + guard = _DenyGuard("real 大额禁止裸奔") + orchestrator, adapter = _make_orchestrator(tmp_path, guard=guard) + + with pytest.raises(TradeError) as exc: + orchestrator.execute_limit( + account_ref=adapter.account_ref, + symbol="AAPL", + side="buy", + quantity=10, + limit_price=180, + ) + + assert exc.value.code == TradeErrorCode.PERMISSION_DENIED + assert exc.value.message == "real 大额禁止裸奔" + # DENY 时订单不应真的提交到 broker + assert adapter._submitted is False + # guard 收到的 RiskContext 携带了完整意图,供未来风控专题消费 + assert guard.seen[-1].operation == "submit_limit" + assert guard.seen[-1].env == "simulate" + assert guard.seen[-1].automation is False + + +def test_trade_orchestrator_guard_deny_blocks_cancel(tmp_path): + guard = _DenyGuard() + orchestrator, adapter = _make_orchestrator(tmp_path, guard=guard) + adapter._submitted = True + adapter.order_status_sequence = ["submitted"] + + with pytest.raises(TradeError) as exc: + orchestrator.execute_cancel(order_ref=adapter.order_ref) + + assert exc.value.code == TradeErrorCode.PERMISSION_DENIED + assert guard.seen[-1].operation == "cancel" + + +def test_trade_orchestrator_execute_limit_passes_automation_flag_through(tmp_path): + guard = _DenyGuard() + orchestrator, adapter = _make_orchestrator(tmp_path, guard=guard) + + with pytest.raises(TradeError): + orchestrator.execute_limit( + account_ref=adapter.account_ref, + symbol="AAPL", + side="buy", + quantity=10, + limit_price=180, + automation=True, + ) + + assert guard.seen[-1].automation is True + + +def test_trade_orchestrator_execute_cancel_non_finalized_returns_warning(tmp_path): + orchestrator, adapter = _make_orchestrator(tmp_path) + adapter._submitted = True + adapter.order_status_sequence = ["submitted", "submitted", "submitted", "submitted"] + + result = orchestrator.execute_cancel(order_ref=adapter.order_ref) + + assert result.operation == "cancel" + assert result.finalized is False + assert tuple(result.warnings) == ("cancel_requested_not_finalized",) + assert result.order_status is not None + assert result.order_status.status == "submitted" + assert "撤单请求已发送" in result.result_summary + + def test_trade_tools_inject_active_account_snapshot(tmp_path): orchestrator, adapter = _make_orchestrator(tmp_path) kernel = Kernel(api_key="test") register_trade_tools(kernel, orchestrator) - kernel.on_confirm(lambda message: True) positions = kernel._tools["trade_account"].handler({"action": "get_positions", "account_ref": adapter.account_ref}) assert positions["status"] == "ok" @@ -262,7 +312,7 @@ def test_trade_tools_inject_active_account_snapshot(tmp_path): assert account["cash"] == 1000.0 assert account["positions"]["AAPL"]["quantity"] == 0 - plan = kernel._tools["trade_plan"].handler( + executed = kernel._tools["trade_execute"].handler( { "operation": "submit_limit", "account_ref": adapter.account_ref, @@ -272,11 +322,10 @@ def test_trade_tools_inject_active_account_snapshot(tmp_path): "limit_price": 180, } ) - applied = kernel._tools["trade_apply"].handler({"plan_id": plan["plan_id"]}) - assert applied["status"] == "ok" + assert executed["status"] == "ok" account = kernel.data.get("account") assert account["positions"]["AAPL"]["quantity"] == 10 - assert kernel.data.get("trade:last_result")["plan_id"] == plan["plan_id"] + assert kernel.data.get("trade:last_result")["plan_id"] == executed["plan_id"] def test_trade_tools_missing_refs_return_missing_error_codes(tmp_path): @@ -290,24 +339,27 @@ def test_trade_tools_missing_refs_return_missing_error_codes(tmp_path): order_status = kernel._tools["trade_account"].handler({"action": "get_order_status"}) assert order_status["error_code"] == "missing_order_ref" - cancel_plan = kernel._tools["trade_plan"].handler({"operation": "cancel"}) - assert cancel_plan["error_code"] == "missing_order_ref" + cancel = kernel._tools["trade_execute"].handler({"operation": "cancel"}) + assert cancel["error_code"] == "missing_order_ref" -def test_trade_orchestrator_cancel_non_finalized_returns_warning(tmp_path): - orchestrator, adapter = _make_orchestrator(tmp_path) - adapter._submitted = True - adapter.order_status_sequence = ["submitted", "submitted", "submitted", "submitted"] - - plan = orchestrator.plan_cancel(order_ref=adapter.order_ref) - result = orchestrator.apply(plan.plan_id) +def test_trade_tools_execute_reports_permission_denied(tmp_path): + guard = _DenyGuard("simulate 也不放行") + orchestrator, adapter = _make_orchestrator(tmp_path, guard=guard) + kernel = Kernel(api_key="test") + register_trade_tools(kernel, orchestrator) - assert result.operation == "cancel" - assert result.finalized is False - assert tuple(result.warnings) == ("cancel_requested_not_finalized",) - assert result.order_status is not None - assert result.order_status.status == "submitted" - assert "撤单请求已发送" in result.result_summary + result = kernel._tools["trade_execute"].handler( + { + "operation": "submit_limit", + "account_ref": adapter.account_ref, + "symbol": "AAPL", + "side": "buy", + "quantity": 10, + "limit_price": 180, + } + ) + assert result["error_code"] == "permission_denied" def test_trade_account_list_accounts_exposes_account_capabilities_and_extra(tmp_path): @@ -334,13 +386,14 @@ def test_trade_guide_injected_when_trade_tools_registered(tmp_path): assert kernel._system_prompt is not None assert "" in kernel._system_prompt - assert "trade_apply" in kernel._system_prompt + assert "trade_execute" in kernel._system_prompt + assert "自主交易操作员" in kernel._system_prompt assert "portfolio.json" in kernel._system_prompt assert "必须显式携带 account_ref 或 order_ref" in kernel._system_prompt assert "不能理解为当前市场状态" in kernel._system_prompt -def test_automation_policy_denies_trade_mutations(tmp_path): +def test_automation_policy_allows_trade_execute_but_denies_hard_blocked_tools(tmp_path): policy = AutomationToolPolicy( workspace=tmp_path / "workspace", task_id="task-1", @@ -348,5 +401,30 @@ def test_automation_policy_denies_trade_mutations(tmp_path): ) assert policy.authorize("trade_account", {"action": "list_accounts"}) is None - assert "禁止调用工具" in str(policy.authorize("trade_plan", {})) - assert "禁止调用工具" in str(policy.authorize("trade_apply", {})) + # 彻底自主:automation reaction 现在可以直接执行交易 + assert policy.authorize("trade_execute", {"operation": "submit_limit"}) is None + assert "禁止调用工具" in str(policy.authorize("bash", {})) + assert "禁止调用工具" in str(policy.authorize("create_subagent", {})) + + +def test_trade_execute_marks_automation_context_from_automation_tool_policy(tmp_path): + guard = _DenyGuard() + orchestrator, adapter = _make_orchestrator(tmp_path, guard=guard) + kernel = Kernel(api_key="test") + register_trade_tools(kernel, orchestrator) + kernel.set_tool_policy( + AutomationToolPolicy(workspace=tmp_path / "workspace", task_id="task-1", profile="analysis") + ) + + kernel._tools["trade_execute"].handler( + { + "operation": "submit_limit", + "account_ref": adapter.account_ref, + "symbol": "AAPL", + "side": "buy", + "quantity": 10, + "limit_price": 180, + } + ) + + assert guard.seen[-1].automation is True