diff --git a/BIONIC_PMS_INTERFACE.md b/BIONIC_PMS_INTERFACE.md new file mode 100644 index 0000000..9dd5303 --- /dev/null +++ b/BIONIC_PMS_INTERFACE.md @@ -0,0 +1,301 @@ +# PMS ↔ 决策系统 (bionic_trader) 接口契约 V1.1 + +> 待办 #9(择时实现 A)与 #10(研判闸)的实现依据。2026-08-03 定稿并双侧落码;同日按用户 +> 审核意见把择时部分整个重做过一版(见 §7 变更记录),当前文本以重做后为准。 +> PMS 侧代码:`app/services/judge.py`(原有)+ `app/services/exec_advisor.py`(新增)。 +> bionic 侧代码:`app/api/main.py` 两个接口 + `workers/tasks_intraday.py` 新增 +> `PMS_JUDGE` 分支 + `app/services/pms_advisor.py` + `config/settings.py` P 段。 +> 机器口径:PMS 跑在 factorevaluation;bionic 跑在 4090 机(API 宿主端口 **38000**)。 + +--- + +## 1. 分工与原则 + +设计 `POSITION_MGMT_DESIGN.md` §1/§7/§8 的落点:持仓系统管「做什么、多少」,决策系统管 +「该不该、何时」。2026-08-03 用户补充的三条原则,两个接口都必须守: + +1. **决策系统还在运行、还有其他职责,改造只做加法。** 新功能全部走新增分支、新增接口、 + 带默认值的新配置;原有六个盘中裁决分支、原有 API、每晚判分、建仓仲裁的输入输出一律 + 不动。凡是「顺手写进共享表」这类会影响既有功能的副作用,都要在本文写明并给出处置。 +2. **提前计算为主、盘中监控为辅,不另做盘中判断。** 大模型推理放在凌晨(每晚的分析本来 + 就在产支撑位、压力位、定性结论),盘中接口只读凌晨结论和当天已有的监控产出,做比对、 + 不做判断。可以接受买不上——现价跑出执行区间就等待,不追。 +3. **拿不到结论时,各按各的纪律降级。** 对 PMS 而言「通过」「驳回」「拿不到」是三件事: + 研判拿不到 → 降级人工确认(绝不当通过,也绝不当驳回——驳回会真的杀掉提议并记一笔 + 驳回账);择时拿不到 → 退回 PMS 内置的实现 B(绝不停出手)。所以 bionic 侧任何给不出 + 结论的情况(解析失败、越界裁决、结论缺失、超时、总开关关)都回 + `{"verdict": "UNAVAILABLE", "reason": ...}`,不许硬编一个看起来像结论的默认值。 + +另外三件事划清归属:配额、14:45 的强制完成、停牌一字板、不追高(当日涨幅与均线偏离)、 +仓位上限、一手、T+1 可卖——这些检查全部留在 PMS 本地(择时的在 `exec_timing.hard_gate`, +其余在规则闸),请求根本不会送到 bionic。两个接口对 PMS 永远回 HTTP 200 + verdict 字段, +PMS 不需要区分状态码(网络层异常自然走 PMS 的超时分支,效果相同)。 + +--- + +## 2. 研判闸 `POST /api/intraday/pms_judge`(#10,走大模型,慢是正常) + +PMS 调用方:`judge.py`(自主提议的二级关口,只审 `PMS_JUDGE_ACTIONS` 里的动作,默认 +FILL/ADD/DCA/SWITCH;命令驱动的动作不过研判闸)。超时 `PMS_JUDGE_TIMEOUT`(默认 90 秒)。 + +### 2.1 请求(PMS → bionic,均为 PMS `judge.request()` 现有产出,未改) + +```json +{ + "direction": "PMS_JUDGE", + "action": "DCA", // FILL / ADD / DCA / SWITCH + "ts_code": "600000.SH", // 点式 + "qty": 200, + "reason": "浮亏 -8.5% 触发第一档补仓评估", // PMS 动作引擎给出的提议理由 + "hard_numbers": {"cushion_pct": -0.085, "stage": 1, "...": "PMS 已算好的硬数字"}, + "context": { + "position": {"ts_code": "...", "avg_cost": 9.1, "total_qty": 1000, + "cushion_pct": -0.085, "support_ref": 8.9, "...": "账本快照"}, + "recent_ledger": [{"at": "...", "action": "...", "arbiter": "...", + "verdict": "...", "price": 9.0, "reason": "..."}] + }, + "must_answer": ["下跌是杀逻辑还是杀情绪"] // DCA 必带 (设计 §7) +} +``` + +### 2.2 bionic 侧处理 + +`process_intraday_audit` 新增 `PMS_JUDGE` 分支,复用既有仲裁哲学与代码路径(同一个大模型 +呼叫、同一条审计轨迹),按动作类型给定性判据: + +| action | 核心问题 | 驳回判据(示例) | +|---|---|---| +| FILL 回踩补足 | 回踩是健康洗盘还是结构转坏 | 放量破位 / 主力持续净流出 / 弱市无个股独立证据 | +| ADD 盈利加仓 | 趋势延续还是冲顶兑现 | 放量滞涨 / 已到压力位 / 主力借涨派发 | +| DCA 补仓 | **必答**:杀逻辑还是杀情绪 | 杀逻辑一律驳回;说不清 → 驳回 | +| SWITCH 调仓 | 换入侧证据是否成立 | — | + +bionic 自动补三样 PMS 没有的上下文:实时资金分布(既有的上游接口)、板块动量与区制快照、 +自己昨夜的 `strategy_daily_results` 结论。裁决收口只认 PASS/REJECT,**越界与解析失败一律 +回 UNAVAILABLE**(§1 原则 3);「证据不足以支持」按提示词导向驳回(宁可错杀)。 + +**同股同动作半小时内不重复研判**(`PMS_JUDGE_RETRY_TTL`,默认 1800 秒,TTL 内直接回上次 +结论,应答多一个 `"cached": true`)。原因:PMS 侧对研判驳回**有意**不做当日去重(其注释 +「研判结论会变」),被驳回的候选可能每分钟重来——节流责任放在 bionic 侧,与建仓仲裁的 +`ENTRY_GATE_RETRY_TTL` 是同一个思路。UNAVAILABLE 不缓存(瞬时故障要重试)。 + +**留痕只写 `strategy_audit_log`**(verdict 记 `PMS_PASS` / `PMS_REJECT` / +`PMS_UNAVAILABLE`,该列 VARCHAR(20) 装得下,已核对建表语句),**不写 `decision_ledger`**: +2026-08-03 复核 `outcome_scorer.py` 确认每晚判分对决策账本是全表扫描、不按裁决类型过滤, +建仓仲裁的历史注入也读整本账——PMS 的裁决写进去会混进这两个既有功能。PMS 自己的评审记录 +本来就落在它侧的 `pms_action_ledger`(arbiter=judge),不需要在 bionic 重复记账。 + +### 2.3 应答(bionic → PMS) + +```json +{"verdict": "PASS" | "REJECT" | "UNAVAILABLE", + "reason": "50字以内一句话", "confidence": 0-100, + "direction": "PMS_JUDGE", "ts_code": "600000.SH", "action": "DCA", + "cached": false} +``` + +PMS 侧映射(`judge.py` 现有逻辑,未改):PASS/APPROVE 一类 → 放行进档位分流;REJECT/DENY +一类 → 驳回并记账;**其余一律按 UNAVAILABLE → 降级人工确认 + ERROR 告警**。 + +### 2.4 执行路径与超时 + +API 收到请求 → 投递 celery 任务(`intraday.process_audit`,队列 `PMS_JUDGE_QUEUE`)→ +同步等结果(`PMS_JUDGE_TASK_TIMEOUT`,默认 75 秒)→ 回 JSON。三层超时必须保持 +**75(bionic 等任务)< 90(PMS 等 HTTP)**,否则 PMS 先断开、bionic 白算。 +队列默认 `brain_queue`(12 并发,零部署改动);若盘中全量重算把它塞住、PMS 侧频繁超时 +降级,再切独立队列(§5.2),这是预埋的第二档不是必选项。 + +--- + +## 3. 盘中择时 `POST /api/intraday/pms_exec`(#9,读凌晨结论,秒回) + +PMS 调用方:`exec_advisor.py`。只在 `PMS_EXEC_IMPL=A` 且过了本地检查后咨询;应答按 +`valid_min` 缓存在指令 `progress.exec_advice`(上限 `PMS_EXEC_ADVICE_TTL_MIN`,默认 10 +交易分钟),失败进入冷却(`PMS_EXEC_FAIL_COOLDOWN_MIN`,默认 5 分钟)。超时 +`PMS_EXEC_TIMEOUT_SEC` 默认 8 秒——这个接口必须秒回,**大模型不进盘中路径**。 + +### 3.1 请求(PMS → bionic) + +```json +{ + "direction": "PMS_EXEC", "ts_code": "600000.SH", + "side": "buy" | "sell", "action": "OPEN|FILL|ADD|DCA|EXIT|TRIM", + "qty_left": 500, // 今日还该投放的量 (仅参考, 配额归 PMS) + "is_last_day": false, "tdays_left": 3, "now": "10:31", + "day": {"price": 10.0, "vwap": 10.05, "...": "PMS 的当日快照, 可整个留空"}, + "refs": {"support": 9.5, "pressure": 11.0, "stop": 9.2, "source": "bionic"}, + "position": {"total_qty": 0, "avail_qty": 0, "avg_cost": null, "cushion_pct": null} +} +``` + +bionic 只用其中的 `ts_code` / `side` / `day.price`(现价缺失时自己从 db13 分钟线取), +其余字段留作排查与将来扩展;区间一律按 bionic 自己库里的昨夜结论算,不用 PMS 转送的参考位 +(PMS 的参考位在结论停更时会换成它自己兜底算的值,那不能当决策系统的结论用)。 + +### 3.2 bionic 侧处理(`pms_advisor.py`,版本标记 `daily_band_v1`) + +**只有三步,全部是读现成的数、做比对,没有任何盘中判断:** + +第一步,读昨夜结论。`strategy_daily_results` 里当晚算好的支撑位、压力位、定性。结论缺失、 +超过 `PMS_EXEC_YSTRAT_MAX_AGE_DAYS`(默认 5 个自然日)、或支撑压力倒挂 → 回 UNAVAILABLE, +PMS 退实现 B。由支撑压力推出两个执行区间(推导是算术,宽度是参数): + +``` +买入区间 = [ 支撑 × (1 − 外沿0.01), 支撑 + (压力 − 支撑) × 0.3 ] +卖出区间 = [ 压力 − (压力 − 支撑) × 0.3, 压力 × (1 + 外沿0.01) ] +只有单边参考位时, 区间宽度按该位的 3% 取, 另一侧回 UNAVAILABLE。 +``` + +第二步,读当天监控。`strategy_audit_log` 今天对该股的裁决(盘前复审与盘中研判本来就在写 +这张表;PMS_ 开头的行是我们替 PMS 写的,读的时候排除)。今天已有离场类结论 +(REVERSAL_SELL / TAKE_PROFIT / FAIL)→ 买入一律等待、卖出直接放行。监控查询失败不拦 +执行(监控是「为辅」),但应答的 `missing` 里必须写明 `today_audit`。 + +第三步,比对现价与区间: + +| 情形 | 应答 | +|---|---| +| 买入,现价在买入区间内 | FIRE,**限价 = 区间上沿**(委托能一直活到价格离开区间,不再是贴着现价的窄缝——2026-08-03 委托全部到期的病根即在此) | +| 买入,现价高于区间上沿 | WAIT「不追,接受买不上」 | +| 买入,现价低于区间下沿(支撑位下方) | WAIT「区间外不买入」 | +| 买入,今天已有离场结论 | WAIT,理由带上是哪条结论 | +| 卖出,现价到达卖出区间(或更高) | FIRE,限价不给(交回 PMS 按它自己的口径贴现价挂出) | +| 卖出,现价未到卖出区间 | WAIT「等高一点卖,当日收尾的强制完成仍由 PMS 兜底」 | +| 卖出,今天已有离场结论 | FIRE(与 PMS 已订阅的卖出信号同源,不会互相矛盾) | + +### 3.3 应答(bionic → PMS) + +```json +{"verdict": "FIRE" | "WAIT" | "UNAVAILABLE", + "limit_price": 9.95, // 买入=区间上沿; 卖出与等待时为 null (PMS 用本地口径) + "valid_min": 10, "reason": "现价 9.7 在买入区间 [9.4, 9.95] 内 (支撑 9.5/压力 11.0), 限价挂区间上沿", + "confidence": null, + "observed": {"price": 9.7, "support": 9.5, "pressure": 11.0, + "buy_band": [9.4, 9.95], "sell_band": [10.55, 11.11], + "y_signal": "BUY", "today_exit": null}, + "missing": [], "advisor": "daily_band_v1"} +``` + +PMS 侧对建议价有一道保护:偏离现价超过 `PMS_EXEC_LIMIT_BAND`(默认 10%,即 A 股单日 +涨跌幅上限)视为异常数据,改回本地口径并在理由里说明。区间边缘离现价百分之几属正常, +正常联调不会触发这条。 + +### 3.4 下一个迭代(本轮明确不做) + +凌晨管线加一步大模型推理,对每只票明示写出买入区间与卖出区间,盘中改为直接读它——届时 +只换本模块第一步的取数来源和 `advisor` 版本标记,应答结构与 PMS 侧都不动。动凌晨管线 +属于重大迭代,做之前会先出细方案确认(覆盖范围、存放位置、算不出来怎么办)。 + +--- + +## 4. 拿不到结论时两侧各自怎么办 + +| 情形 | bionic 回什么 | PMS 怎么办 | 人看哪里 | +|---|---|---|---| +| 研判超时/队列拥堵 | UNAVAILABLE(原因带超时) | 降级人工确认 + ERROR | PMS 日志 `[研判闸]`;频繁出现 → §5.2 切独立队列 | +| 研判裁决解析失败/越界 | UNAVAILABLE | 同上 | bionic 日志 `[PMS_JUDGE]` | +| 该股没有昨夜结论 / 结论过期 | UNAVAILABLE(原因写明) | 本轮退实现 B + 冷却 5 分钟 | 指令 `progress.exec_advice.error`;`last_decision.source` 显示 `B(实现A不可用: ...)` | +| 收盘后 / 无实时价 | UNAVAILABLE「无实时价」 | 同上 | 同上 | +| 择时接口网络不通/超时 | —(HTTP 层异常) | 同上 | 同上 | +| 当天监控表查询失败 | 照常按区间应答,`missing` 带 `today_audit` | 正常执行 | 应答的 missing 字段;bionic 日志 | +| `PMS_API_ENABLED=False` | UNAVAILABLE(原因写明开关) | 各按上述降级 | bionic `.env` | +| PMS 侧未配置(地址空 / `PMS_EXEC_IMPL=B`) | 不会发请求 | 研判闸=人工确认;择时=实现 B | PMS 参数设置页 | + +--- + +## 5. bionic 侧部署(4090 机) + +### 5.1 常规路径(本次交付走这条) + +bionic 的源码是**挂载卷**(compose 里 `.:/app`),与 PMS 打进镜像正相反——**不需要 +build,`git pull` 后重启进程即可**(backend-api 带自动重载,稳妥起见一起重启)。新增 +配置全部有默认值,`.env` 不用动: + +```bash +# 【4090 机 · bionic_trader 仓库根目录】 +git pull +docker compose restart backend-api worker-brain +# 本机自测 (两接口, 不需要 PMS 参与; 研判走大模型要等 30~90 秒): +docker compose exec backend-api python scripts/pms_smoke.py +``` + +改动面:`config/settings.py`(P 段,全默认值)、`app/api/main.py`(两个接口)、 +`workers/tasks_intraday.py`(PMS_JUDGE 分支,含同股同动作的重复研判缓存)、 +`app/services/pms_advisor.py`(新)、`scripts/pms_smoke.py`(新)。 +既有六个盘中裁决分支、既有 API、每晚判分、建仓仲裁的行为一行未动(PMS_JUDGE 是新增分支 +且提前返回,不经过广播与重算;择时接口不写任何表)。 + +### 5.2 研判独立队列(预埋第二档,brain_queue 拥堵时再启用) + +判据:PMS 侧 `[研判闸] 请求失败 ... Timeout` 成为常态(偶发超时降级人工确认是设计内 +行为,不用切)。启用:bionic `.env` 加 `PMS_JUDGE_QUEUE=pms_queue`,compose 加一个 +专属 worker 后 `docker compose up -d worker-pms && docker compose restart backend-api`: + +```yaml + worker-pms: + build: {context: ., dockerfile: docker/Dockerfile} + container_name: trader_worker_pms + command: celery -A workers.celery_app worker -Q pms_queue -l info -c 2 --prefetch-multiplier=1 + env_file: .env + environment: [TZ=Asia/Shanghai] + volumes: ["/etc/localtime:/etc/localtime:ro", ".:/app"] + depends_on: [redis] + networks: [trader_net] + restart: always +``` + +(astock-kg 公告链 2026-07-23 单队列把交互请求饿死的同款处方:交互请求与长任务分开排队。) + +### 5.3 bionic 侧新增配置一览(均有默认值) + +| 键 | 默认 | 说明 | +|---|---|---| +| `PMS_API_ENABLED` | True | 总开关,False = 两接口一律回 UNAVAILABLE | +| `PMS_JUDGE_QUEUE` | brain_queue | 研判任务队列 | +| `PMS_JUDGE_TASK_TIMEOUT` | 75 | API 等任务上限(秒),必须小于 PMS 侧的 90 | +| `PMS_JUDGE_RETRY_TTL` | 1800 | 同股同动作重复研判的间隔(秒),间隔内回上次结论 | +| `PMS_EXEC_BUY_BAND_RATIO` | 0.3 | 买入区间上沿 = 支撑 + (压力−支撑)×此比例 | +| `PMS_EXEC_SELL_BAND_RATIO` | 0.3 | 卖出区间下沿 = 压力 − (压力−支撑)×此比例 | +| `PMS_EXEC_BAND_EDGE` | 0.01 | 区间外沿的放宽(买入下沿=支撑×(1−此值),卖出上沿对称) | +| `PMS_EXEC_ONE_SIDE_BAND` | 0.03 | 只有单边参考位时区间宽度按该位的此比例取 | +| `PMS_EXEC_VALID_MIN` | 10 | 应答建议有效期(交易分钟) | +| `PMS_EXEC_YSTRAT_MAX_AGE_DAYS` | 5 | 昨夜结论超此自然日龄按过期处理 | + +区间比例是初值,联调后按应答里的 `observed`(区间与现价的实际关系)回看再调。 + +--- + +## 6. 联调步骤与判收 + +前提:factorevaluation 能路由到 4090 机的 38000 端口(同网段应当直通;不通先 +`curl http://<4090内网IP>:38000/health` 排网络)。以下 `` 代指 +`http://<4090内网IP>:38000`。 + +```bash +# 1.【4090 机】部署与本机自测 —— §5.1 的三条命令 +# 2.【factorevaluation · tradingSystem 根目录】跨机探活 (只读, 不落表不产生指令): +docker compose run --rm --no-deps pms-web python scripts/probe_bionic.py --base +# 判读: pms_exec 回 FIRE/WAIT 即通, 看应答里的区间对不对得上该股昨夜的支撑压力; +# 回 UNAVAILABLE 看原因 (无实时价=收盘后正常; 没有昨夜结论=该股不在每晚分析 +# 范围, 换一只探)。pms_judge 回 PASS/REJECT 即通, 30~90 秒正常。 +# 3.【factorevaluation】页面「参数设置」填 PMS_JUDGE_API_BASE= → 研判闸生效。 +# 验收: 下一轮自主提议扫描里, 需研判的动作 (FILL/ADD/DCA) 的评审账本出现 +# arbiter=judge 的 PASS/REJECT 行 (make t-gate 看), 而不再是"降级人工确认"。 +# 4.【factorevaluation】要切择时A: 页面把 PMS_EXEC_IMPL 改 A (随时可改回 B)。 +# 验收: make watch 里在途指令的 last_decision 出现 source=A / A缓存; +# make t-ins 里 progress.exec_advice 有结论、区间理由与有效期。 +# 5. 盘中观察一天: 每次咨询 bionic 侧日志有 [PMS 择时] 行; 对比实现A按区间给出的 +# 出手/等待与实现B的差别 (关注限价挂区间上沿后, 委托是否不再成批到期)。 +``` + +判收纪律照旧:**代码就绪 ≠ 接通,接通 ≠ 判收**。#10 判收 = 一笔真实提议走完 +「规则闸 → 研判通过/驳回 → 档位分流」并在 PMS 评审账本留下 arbiter=judge 的行; +#9 判收 = 盘中一条指令按区间出手或等待、且对端故障时当轮自动退实现 B 有留痕 +(`source=B(实现A不可用)`)。 + +## 7. 变更记录 + +| 版本 | 日期 | 内容 | +|---|---|---| +| V1.0 | 2026-08-03 | 首版定稿并双侧落码。择时部分是一套盘中判定规则(资金阈值、动量追买等) | +| V1.1 | 2026-08-03 | 按用户审核意见重做择时:删掉 V1.0 的盘中判定规则(违反「提前计算为主、盘中监控为辅、不另做盘中判断、可以接受买不上」),改为由昨夜支撑压力推出执行区间、盘中只做区间比对与当日监控核对;研判留痕不再写 `decision_ledger`(每晚判分全表扫描,已核实);文档与代码注释清理生造词 | diff --git a/Makefile b/Makefile index a789665..dbe6307 100644 --- a/Makefile +++ b/Makefile @@ -121,10 +121,11 @@ rebuild-accept: ## 只看判收: 安全垫分布 / 批次账 / 行业集中度 # INCREASE_EXPOSURE 升仓 pct% (全局: 垫厚票补到目标 + 候选池新票建仓) # OPEN_TARGET 建仓某股至 x% (个股: 分批 50/25/25, 含一手合并) # -# **这里的「择时」是 PMS 内置的实现 B, 不经决策系统。** 设计里择时有两个实现: -# 实现 A = 委托决策系统盘中择时 —— 待办 #9, 等 bionic 侧接口, **还没接** -# 实现 B = PMS 自己算 (分日配额 / 分笔 / VWAP / 回踩 / 不追高 / 14:45 兜底) —— 已实现 -# 所以跑 A 段不需要动决策系统的任何逻辑, 它这一段本来就不参与。 +# **这里的「择时」默认是 PMS 内置的实现 B, 不经决策系统。** 设计里择时有两个实现: +# 实现 A = 委托决策系统盘中择时 —— 待办 #9, 代码已接完 (BIONIC_PMS_INTERFACE.md), +# **默认关**: 页面把 PMS_EXEC_IMPL 改 A 才生效, 不可用自动退 B +# 实现 B = PMS 自己算 (分日配额 / 分笔 / VWAP / 回踩 / 不追高 / 14:45 兜底) —— 默认档 +# 所以 PMS_EXEC_IMPL=B (默认) 时跑 A 段, 决策系统这一段完全不参与。 A ?= http://127.0.0.1:38100 # 中文必须原样出来: `json.tool` 默认把中文转义成 \uXXXX, 而这些接口的 reason 恰恰是最该 # 看的东西。**不能用 `json.tool --no-ensure-ascii`** —— 这条管道跑在**宿主机**上不是容器里, @@ -203,7 +204,7 @@ t-plans: ## A4 看**方案**("买哪只、多少股、分几批", 还没成指 t-mat: ## A5 方案 → 指令 (先记账后动作: 指令先落表) @curl -s -X POST '$(A)/api/ops/materialize' | $(J) -t-dry: ## A6 择时试算 (实现B, 不经决策系统; dry_run 只算不发) +t-dry: ## A6 择时试算 (按 PMS_EXEC_IMPL 档位走, 默认B; dry_run 只算不发不占研判缓存) @curl -s -X POST '$(A)/api/ops/exec-tick?dry_run=true' | $(J) t-tick: ## A7 真出手 (shadow 下 = 置 DISPATCHED 但不写下游, 你照着在 QMT 手工下) diff --git a/README.md b/README.md index 0dcd813..c1ca30c 100644 --- a/README.md +++ b/README.md @@ -14,6 +14,7 @@ | `QMT_WS_PROTOCOL.md` | **PMS ↔ QMT WebSocket 指令与回报协议 V1.0(定稿)**:传输与重连、Ed25519 签名与幂等、消息集、状态机、断线补发与对账兜底、部署前检查清单。**这是下发通道的唯一实现依据** | | `QMT_INTERFACE_REQUIREMENTS.md` | 与 QMT 侧的数据与接口需求清单 **V2.0**:A 部分(只读数据)与 C 部分(切换约定)有效;**B 部分的表通道已废止**,改由上面的 ws 协议承担 | | `UPSTREAM_PLAN_API.md` | **上游选股计划接口 (`/plan`) 的接入记录**:应答结构、PMS 侧五条口径、参数表,**行业源换 gp_hybk 的定案(§8)**,以及**榜单变化改由 PMS 自算的口径与取舍(§9)**| +| `BIONIC_PMS_INTERFACE.md` | **PMS ↔ 决策系统 (bionic) 接口契约 V1.1**:研判闸 `pms_judge`(大模型仲裁,走 brain_queue)与盘中择时 `pms_exec`(按凌晨结论给执行区间,秒回)两个接口的请求/应答、各种拿不到结论时两侧各自怎么降级、bionic 侧部署与联调步骤。**待办 #9/#10 的实现依据** | | `ddl_pms_v1.sql` | PMS 全部自有表建表语句(153 代理侧,**15 张**:设计 §11 的 10 张 + ws 通道 3 张 + 现金流水 1 张 + 计划名册快照 1 张) | | `Makefile` | 常用操作一句话入口(`make help` 看全部)。改了代码用 `make deploy` / `make up`,**别用 `docker compose restart`**——源码打进镜像,restart 跑的还是旧代码且一声不吭 | | `config/settings.py` | 配置(基础设施键名对齐 bionic;业务参数为初值,页面调参持久化到 `pms_runtime_param` 后优先) | @@ -47,6 +48,8 @@ app/ dispatcher.py 下发通道两适配器: shadow(默认) / ws(落 pms_qmt_order 出口表) proposal_service.py 自主提议: 扫描→规则闸→研判闸→按自主档位分流 (执行/入队) judge.py 研判闸客户端 (决策系统未接通时自动降级为人工确认) + exec_advisor.py 择时实现A客户端: 本地检查先行 → 问决策系统「现价在不在执行区间」 + (应答有缓存与失败冷却) → 不可用退实现B (默认 PMS_EXEC_IMPL=B 整个短路) signal_service.py 盘中信号订阅 (db2 广播 + db3 风控卖出) → 卖出指令或提议 ledger_service.py 成交回放 / 对账 / 除权 / 盘前 / 日终结算 / 运营日报 market.py 行情 (Redis db13) 与参考位 (决策系统主口径 + 兜底自算) @@ -68,6 +71,7 @@ scripts/ test_batch8_units.py 榜单变化: 名册指纹/三种语义/尾部闸/落库往返 60 例 test_batch9_units.py 成本价体检 / 对账按日推进 / 行业闸 / 取整记账 49 例 test_batch10_units.py 静默失败专项: 关键路径不许丢返回值 + 八条实例 48 例 + test_batch11_units.py 择时实现A: 本地检查等价/应答折算/缓存冷却/退B 24 例 test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 58 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练; 含 DDL 体检) check_db.py 实机连通性与表结构自检 (需真实 .env) @@ -76,6 +80,8 @@ scripts/ 高水位、刹车解除日); 影子运行期重来一次, 绝不碰下游表 probe_plan_api.py 上游计划实机探活: 通不通 / 字段口径 / 候选筛选结果 / 有没有价 / 榜单变化 (只读; --snapshot 才落库, 那是它唯一的写操作) + probe_bionic.py 决策系统两接口实机探活 (择时研判 pms_exec + 研判闸 pms_judge, + 只读不落表; 待办 #9/#10 的接通自证) rebuild_ledger.py 账本重建: 预检 → 执行 → 判收 (默认只预检; --yes 才改账) ws_smoke.py ws 联调工具: status/watch/place/cancel/inbox (绕开 dispatch_mode) code_fingerprint.py 当前 .py 源码的短哈希 —— `make test` 拿它比对容器与工作树, @@ -235,7 +241,7 @@ make rebuild-accept # 只看判收:安全垫分布 / 批次账 / 行业集 **A3 之前必须先下一条「任务命令」。** 参数命令(`SET_SCALE` / `SET_PORTFOLIO_CAP` / `SET_STOCK_CAP` / `SET_MAX_NAMES` 这些)只是设约束,**本来就不产生方案**——页面上它们显示「0 条方案」是对的,不是坏了。要走通链路得下 B/C 组的任务命令:`INCREASE_EXPOSURE`(升仓 pct%,全局:垫厚票补到目标 + 候选池新票建仓)或 `OPEN_TARGET`(建仓某股至 x%,个股:分批 50/25/25 含一手合并)。 -**这里的「择时」是 PMS 内置的实现 B,不经决策系统。** 设计里择时有两个实现:实现 A 是委托决策系统盘中择时(待办 #9,等 bionic 侧接口,**还没接**);实现 B 是 PMS 自己算(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底),已实现。所以跑 A 段不需要动决策系统的任何逻辑——它这一段本来就不参与。 +**这里的「择时」默认是 PMS 内置的实现 B,不经决策系统。** 设计里择时有两个实现:实现 A 是委托决策系统盘中择时(待办 #9,**代码已接完、默认关**——页面把 `PMS_EXEC_IMPL` 改成 A 才生效,且不可用自动退 B);实现 B 是 PMS 自己算(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底)。`PMS_EXEC_IMPL=B`(默认)时决策系统完全不参与,跑 A 段行为与接通前一字不差。接口契约与切换步骤见 `BIONIC_PMS_INTERFACE.md`。 **别开 beat**——那会让几个调度位同时动,出了岔子分不清是谁干的。`Makefile` 里有一组 `t-*` 目标,就是把那几个调度位改成手动逐跳触发,顺序与生产一致: @@ -304,7 +310,7 @@ make t-gate # 随时: 规则闸/研判闸拒了什么、为什么 ## 已实现 / 待开发 -**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**榜单变化提示**(名册快照 + 新进/掉榜/档位升降/覆盖翻转/名次跳变,持仓票单列,榜尾截断噪音闸);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 401 例**。 +**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**榜单变化提示**(名册快照 + 新进/掉榜/档位升降/覆盖翻转/名次跳变,持仓票单列,榜尾截断噪音闸);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**择时实现 A 客户端**(`exec_advisor.py`:本地检查先行 → 问决策系统「现价在不在执行区间」→ 应答缓存与失败冷却 → 不可用退实现 B,默认 `PMS_EXEC_IMPL=B` 整个短路);**单测 425 例**。 ### 静默失败专项(2026-07-31,八条已修) @@ -374,8 +380,8 @@ git pull → make test → ALL SUITES PASS | 6 | 上游计划的剩余待确认口径 | 等上游:`UPSTREAM_PLAN_API.md` §4 剩 5 条 + §7.3 新增两条(同一 `date` 计划不幂等、档位规模与文档不符) | | 7 | ~~榜单变化做成页面提示(新进传导链 / 掉榜)~~ | ✅ 2026-07-31:上游 `changes` 恒为 `null`,所以**改成 PMS 自己算**——每次拉计划落一份名册快照(第 15 张表),差异由纯逻辑比。顺带补上「同一 `date` 多版本、下游无从分辨手上是哪一版」那个洞。详见 `UPSTREAM_PLAN_API.md` §9 | | 8 | T0 做T(二期) | 可做,设计已有,无外部依赖。**建议排在账本重建之后**:它要拿真实持仓验触发与 14:50 平回,账本空着只能跑单测 | -| 9 | 择时实现 A(委托决策系统盘中择时) | 阻塞:等 bionic 侧接口 | -| 10 | 研判闸接通 | 阻塞:等 bionic 侧 `process_intraday_audit` 新增 PMS direction。客户端已就位,接口好了在页面填 `PMS_JUDGE_API_BASE` 即通 | +| 9 | 择时实现 A(委托决策系统盘中择时) | **代码就绪待联调**(2026-08-03,当日按用户原则重做过一版):原则是**提前计算为主、盘中监控为辅、不另做盘中判断、可以接受买不上**。bionic 侧 `/api/intraday/pms_exec` 用凌晨已算好的支撑/压力推出买入、卖出执行区间,盘中只回答「现价在不在区间内」并核对当天监控有没有离场结论,秒回;区间内出手时**限价挂区间边缘**(今天委托全部到期的病根——限价贴着现价×1.002 太窄——在实现 A 下由此解决)。PMS 侧 `exec_advisor.py`(本地检查先行、应答缓存、失败冷却、不可用退 B)。凌晨用 LLM 明示产出区间是下一个迭代。**默认 `PMS_EXEC_IMPL=B` 不生效**;联调步骤见 `BIONIC_PMS_INTERFACE.md` §6,探活 `probe_bionic.py` | +| 10 | 研判闸接通 | **代码就绪待联调**(2026-08-03):bionic 侧 `process_intraday_audit` 已加 `PMS_JUDGE` direction——通过、驳回、研判不可用三种结果严格分开(解析失败绝不折成驳回),同股同动作 30 分钟内直接回上次结论;留痕只写 `strategy_audit_log`,**不写 `decision_ledger`**(每晚判分是全表扫描,写进去会混入判分与建仓仲裁的历史)。API `/api/intraday/pms_judge` 同步应答(走 brain_queue,可配独立队列)。PMS 客户端原本就位——**bionic 部署后在页面填 `PMS_JUDGE_API_BASE` 即通** | | 11 | `trading_buy_plan` 退场:PMS 已不读它 | 待上游确认无其他消费方(`UPSTREAM_PLAN_API.md` Q10) | | 12 | ~~部署便利性:`make deploy` 一句话完成 build + up + 各 profile~~ | ✅ 2026-07-31:`Makefile`(`make help` 看全部)。`deploy` 包 `scripts/deploy.sh`;另有 `test` / `initdb` / `check` / `probe` / `changes` / `industry` / `ws-status`。一次性命令统一带 `--no-deps`,免得跑个单测把 beat/worker 也拉起来 | | 13 | ~~静默失败专项:八条全修 + 关键路径禁止丢弃返回值~~ | ✅ 2026-07-31:见上一节。新增 `test_batch10_units.py` 27 例,其中 [A] 组是静态扫描守卫 | @@ -405,6 +411,7 @@ make watch WIDE=1 # 连方案分布与评审账本一起显示 - **`PMS_TOTAL_SCALE` 与账户实际资金要对得上。** 两个数本来就是不同的东西(scale 是「我打算投多少」的命令参数,账户总资产是「现在有多少钱」),允许不等,但差太远会产出一堆执行不了的方案:规划器只认 scale,按 200 万 × 总仓上限排出 140 万的买入,而账户只有 98 万,多出来的会在**规则闸**被一条条判 `INSUFFICIENT_CASH`——方案排得出来、下不出去,每一步看着都正常。2026-07-31 实测就是这组数。盘前准备(`make t-pre`)现在会把两个数摆在一起,差太多就在 `scale_check` 里说破。 - 静默失败专项修完之后,**几处「以前默默过去」的地方现在会明着报失败**,这是设计如此:开关类命令(`HALT_BUY` / `RESUME_ALL` 等)写不进参数就置 `CANCELLED` 并回 `ok=False`,不再显示「已完成」;`daily_settle` 在对账被拦(爆炸半径 / 两源无应答 / 成本价体检不过)或连续不一致升到 ERROR 时 `ok=False`;`premarket` 在「该踩刹车却没踩上」时把它记进 `errors` 而不是 warnings。**看到这些报错先别急着改代码——多半是数据库或下游真的出问题了,以前只是没人告诉你。** - 新增第 15 张表 `pms_plan_snapshot`(名册快照),**部署后要跑一次 `make initdb`**,否则榜单变化那段会一直报「读名册快照失败:表不存在」——候选池不受影响。第一份快照落下之前不报任何变化(`kind=first`),这是设计如此不是坏了。 +- (2026-08-03 补)**决策系统两接口的接入代码已进库,全部默认不生效**:研判闸要页面填 `PMS_JUDGE_API_BASE` 才接通(空=沿用「降级人工确认」现状);择时实现 A 要页面把 `PMS_EXEC_IMPL` 改 A 才启用(默认 B,行为与接通前一字不差)。两个都不需要重启。接通前先跑 `probe_bionic.py` 探活,判读口径见 `BIONIC_PMS_INTERFACE.md` §6。 ### ws 通道实现清单 diff --git a/app/core/exec_timing.py b/app/core/exec_timing.py index 605b946..dd44d89 100644 --- a/app/core/exec_timing.py +++ b/app/core/exec_timing.py @@ -1,9 +1,19 @@ # -*- coding: utf-8 -*- """ -择时执行器 · 实现B「内置保守择时」(纯逻辑, 无外部依赖, 可单测) -================================================================ -设计 POSITION_MGMT_DESIGN.md §8。实现A (委托决策系统盘中择时研判) 走同一个 -`decide()` 接口, 研判接通后在 services 层换实现即可, 本模块是兜底也是一期主力。 +择时执行器 · 实现B「内置保守择时」+ 实现A的纯逻辑部件 (零外部依赖, 可单测) +======================================================================== +设计 POSITION_MGMT_DESIGN.md §8。两个实现的分工 (2026-08-03 接实现A时定下): + + hard_gate() 事实性检查 —— 无价/停牌/配额尽/非时段/一字板/不追高(当日涨幅)/ + 14:45 兜底。这些检查不委托给任何人, 两个择时实现都必须先过它; + 命中返回决策, 未命中返回 None。 + decide() 实现B 全量判定 = hard_gate() + 内置保守规则 (避开开盘/均价/回踩)。 + 行为与拆分前完全一致, test_batch3 锁着。 + apply_advice() 把决策系统的应答 (FIRE/WAIT + 建议价) 折算成与 decide() 同构的决策。 + 建议价通常是执行区间的边缘; 偏离现价超出保护幅度时按本地口径重定并留痕。 + +实现A的取数与降级在 services/exec_advisor.py: 决策系统不可用 → 整轮退 decide() +(设计 §13「择时退实现B」), 绝不因为对端故障停出手。 规则原文与落点: 每日配额 = 剩余量 ÷ 剩余窗口天数, 向上取整到一手 → daily_quota() @@ -90,17 +100,19 @@ def slice_qty(quota: int, slices: int = 1, lot: int = LOT) -> list: return [q for q in out if q > 0] -def decide(*, side: str, now, day: dict, params: dict, is_last_day: bool, - fired_today: int = 0, quota: int = 0) -> dict: - """单条指令在「此刻」该不该出手。 +def hard_gate(*, side: str, now, day: dict, params: dict, is_last_day: bool, + fired_today: int = 0, quota: int = 0): + """事实性检查 (实现A/B 共用的前置)。命中返回决策 dict, 未命中返回 None。 - day: {price, vwap, halted, limit_up, limit_down, day_chg_from_open, support} - params: {sell_avoid_open_min, buy_halt_dayup, eod_force_time, eod_force_discount} - 返回 {"action", "qty_hint", "limit_price", "reason", "forced"} + 包含: 无价/停牌/配额尽/非时段/一字板/买入不追高(当日涨幅)/14:45 兜底与「兜底后 + 不新开买单」。**不含**避开开盘 N 分钟与均价/回踩 —— 那些是各实现自己的规则。 + + 与拆分前 decide() 的唯一语义差别: 卖出的 14:45 兜底现在排在「避开开盘 30 分钟」 + 之前判。两者只在 eod_force_time 被改到 10:00 之前这种病态配置下才会同时成立, + 且真到那时也该是兜底赢 —— 强制完成这件事永远归 PMS 自己管。 """ now_min = hm_to_min(now) price = float(day.get("price") or 0) - vwap = float(day.get("vwap") or 0) eod_min = hm_to_min(params.get("eod_force_time") or "14:45") disc = float(params.get("eod_force_discount") or 0.998) left = max(0, int(quota) - int(fired_today)) @@ -121,16 +133,10 @@ def decide(*, side: str, now, day: dict, params: dict, is_last_day: bool, if side == "sell": if day.get("limit_down") and not is_last_day: return out(ACT_SKIP, "跌停一字板, 当日跳过顺延") - avoid = int(params.get("sell_avoid_open_min") or 30) - if now_min < OPEN_MIN + avoid: - return out(ACT_WAIT, f"避开开盘 {avoid} 分钟 (至 {_fmt(OPEN_MIN + avoid)})") if now_min >= eod_min: return out(ACT_FIRE, f"{_fmt(eod_min)} 兜底: 限价 = 现价×{disc}", limit=round(price * disc, 2), forced=True) - if vwap > 0 and price >= vwap: - return out(ACT_FIRE, f"现价 {price} ≥ 当日均价 {vwap}, 分笔卖出配额", - limit=round(price * disc, 2)) - return out(ACT_WAIT, f"现价 {price} < 当日均价 {vwap or '—'}, 等更好的价") + return None if side == "buy": if day.get("limit_up"): @@ -145,18 +151,89 @@ def decide(*, side: str, now, day: dict, params: dict, is_last_day: bool, return out(ACT_FIRE, f"窗口末日 {_fmt(eod_min)} 强制完成: 限价 = 现价×{premium}", limit=round(price * premium, 2), forced=True) return out(ACT_WAIT, f"{_fmt(eod_min)} 后不新开买单, 顺延次日") - support = float(day.get("support") or 0) - if vwap > 0 and price <= vwap: - return out(ACT_FIRE, f"现价 {price} ≤ 当日均价 {vwap}, 买入配额", - limit=round(price * premium, 2)) - if support > 0 and price <= support * 1.01: - return out(ACT_FIRE, f"现价 {price} 进入回踩带 (支撑 {support}), 买入配额", - limit=round(price * premium, 2)) - return out(ACT_WAIT, f"现价 {price} > 当日均价 {vwap or '—'} 且未回踩, 等回调") + return None return out(ACT_SKIP, f"方向 {side!r} 非法") +def decide(*, side: str, now, day: dict, params: dict, is_last_day: bool, + fired_today: int = 0, quota: int = 0) -> dict: + """实现B: 单条指令在「此刻」该不该出手 (= 硬闸 + 内置保守看法)。 + + day: {price, vwap, halted, limit_up, limit_down, day_chg_from_open, support} + params: {sell_avoid_open_min, buy_halt_dayup, eod_force_time, eod_force_discount} + 返回 {"action", "qty_hint", "limit_price", "reason", "forced"} + """ + h = hard_gate(side=side, now=now, day=day, params=params, is_last_day=is_last_day, + fired_today=fired_today, quota=quota) + if h is not None: + return h + + now_min = hm_to_min(now) + price = float(day.get("price") or 0) + vwap = float(day.get("vwap") or 0) + disc = float(params.get("eod_force_discount") or 0.998) + left = max(0, int(quota) - int(fired_today)) + + def out(action, reason, limit=None, forced=False): + return {"action": action, "qty_hint": left, "limit_price": limit, + "reason": reason, "forced": forced} + + if side == "sell": + avoid = int(params.get("sell_avoid_open_min") or 30) + if now_min < OPEN_MIN + avoid: + return out(ACT_WAIT, f"避开开盘 {avoid} 分钟 (至 {_fmt(OPEN_MIN + avoid)})") + if vwap > 0 and price >= vwap: + return out(ACT_FIRE, f"现价 {price} ≥ 当日均价 {vwap}, 分笔卖出配额", + limit=round(price * disc, 2)) + return out(ACT_WAIT, f"现价 {price} < 当日均价 {vwap or '—'}, 等更好的价") + + premium = round(2 - disc, 4) + support = float(day.get("support") or 0) + if vwap > 0 and price <= vwap: + return out(ACT_FIRE, f"现价 {price} ≤ 当日均价 {vwap}, 买入配额", + limit=round(price * premium, 2)) + if support > 0 and price <= support * 1.01: + return out(ACT_FIRE, f"现价 {price} 进入回踩带 (支撑 {support}), 买入配额", + limit=round(price * premium, 2)) + return out(ACT_WAIT, f"现价 {price} > 当日均价 {vwap or '—'} 且未回踩, 等回调") + + +def apply_advice(*, side: str, day: dict, params: dict, advice: dict, left: int): + """把决策系统的应答折算成与 decide() 同构的决策 (实现A的落地一跳)。 + + advice: {"verdict": FIRE|WAIT, "limit_price"?, "reason"?} —— 已过 hard_gate 才会走到这。 + 建议价通常是执行区间的边缘, 离现价百分之几属正常; 缺失/非法/偏离现价超过 + advice_limit_band (默认 10%) 才视为异常数据, 按本地口径重定 (买 现价×premium / + 卖 现价×disc) 并在 reason 里留痕 —— 限价最终要过出口表的参数校验, 离谱的建议价 + 与其被拒不如就地纠偏。verdict 无法识别返回 None, 由调用方退实现B。 + """ + price = float(day.get("price") or 0) + disc = float(params.get("eod_force_discount") or 0.998) + premium = round(2 - disc, 4) + band = float(params.get("advice_limit_band") or 0.03) + v = str(advice.get("verdict") or "").strip().upper() + reason = f"[实现A] {advice.get('reason') or '决策系统未给理由'}" + + if v == "WAIT": + return {"action": ACT_WAIT, "qty_hint": left, "limit_price": None, + "reason": reason, "forced": False} + if v == "FIRE": + fallback = round(price * (premium if side == "buy" else disc), 2) + try: + limit = float(advice.get("limit_price") or 0) + except (TypeError, ValueError): + limit = 0.0 + if limit <= 0: + limit = fallback + elif price > 0 and abs(limit / price - 1) > band: + reason += f" (建议价 {limit} 偏离现价超 {band:.0%}, 按本地口径 {fallback})" + limit = fallback + return {"action": ACT_FIRE, "qty_hint": left, "limit_price": round(limit, 2), + "reason": reason, "forced": False} + return None + + def _fmt(m: int) -> str: return f"{m // 60:02d}:{m % 60:02d}" diff --git a/app/services/exec_advisor.py b/app/services/exec_advisor.py new file mode 100644 index 0000000..da741b3 --- /dev/null +++ b/app/services/exec_advisor.py @@ -0,0 +1,188 @@ +# -*- coding: utf-8 -*- +""" +择时实现A客户端 · 委托决策系统 (设计 §8, 待办 #9) +================================================== +接口契约见 BIONIC_PMS_INTERFACE.md。决策系统按它凌晨算好的支撑/压力推出买入区间与 +卖出区间, 盘中只回答「现价在不在区间内 + 限价挂哪里」, 不做任何新的盘中判断 +(2026-08-03 用户定的原则: 提前计算为主、盘中监控为辅, 可以接受买不上)。 + +分工与降级 (三条, 都是纪律不是实现细节): + 1. **本地检查先行**: 配额/兜底/停牌/一字板/不追高(当日涨幅) 由 core.exec_timing 的 + hard_gate 先判, 命中就不咨询 —— 这些检查始终留在 PMS 本地。尤其 14:45 兜底: + 到点必须完成, 决策系统说什么都不算。 + 2. **拿不到不等于有答案**: 咨询失败/超时/对端回 UNAVAILABLE/答复无法识别 → 本轮 + 整体退实现B (设计 §13「决策系统择时不可用 → 择时退实现B」), 并进入冷却期 + (冷却内不再咨询, 免得每分钟 tick 都白等一次超时)。绝不把「拿不到」当成 FIRE 或 WAIT。 + 3. **应答带有效期**: 结果缓存在指令 progress.exec_advice 里 (随既有落表持久化, + 页面/t-ins 可见), 有效期内不重复咨询。对端可用 valid_min 缩短有效期, 只缩不放。 + +默认 PMS_EXEC_IMPL=B —— 本模块整个短路, run_tick 行为与接通前一字不差。 +页面把 PMS_EXEC_IMPL 改成 A (并保证 PMS_EXEC_API_BASE 或 PMS_JUDGE_API_BASE 已填) +即切实现A, 随时可改回 B, 不需要重启。 +""" +from __future__ import annotations + +import logging + +from app.core import exec_timing as et +from app.core import tradedays as td +from app.services import judge, param_store + +logger = logging.getLogger("pms.exec_advisor") + +IMPL_A, IMPL_B = "A", "B" +FIRE, WAIT = "FIRE", "WAIT" + + +def impl() -> str: + v = str(param_store.get("PMS_EXEC_IMPL", IMPL_B) or IMPL_B).strip().upper() + return v if v in (IMPL_A, IMPL_B) else IMPL_B + + +def base_url() -> str: + """实现A的接口根地址; PMS_EXEC_API_BASE 留空时沿用研判闸的 PMS_JUDGE_API_BASE + (两者是同一个 bionic 服务, 不逼着用户填两遍)。""" + b = (param_store.get("PMS_EXEC_API_BASE", "") or "").strip().rstrip("/") + return b or judge.base_url() + + +def available() -> bool: + return impl() == IMPL_A and bool(base_url()) + + +def status() -> dict: + """页面/排查用的一句话状态。""" + if impl() != IMPL_A: + return {"impl": IMPL_B, "available": False, + "note": "内置保守择时 (PMS_EXEC_IMPL=B)。切实现A: 页面改 PMS_EXEC_IMPL=A"} + if not base_url(): + return {"impl": IMPL_A, "available": False, + "note": "PMS_EXEC_IMPL=A 但接口地址为空 (PMS_EXEC_API_BASE 与 " + "PMS_JUDGE_API_BASE 都没填) —— 实际全程退实现B"} + return {"impl": IMPL_A, "available": True, "base": base_url(), + "path": param_store.get("PMS_EXEC_PATH", "/api/intraday/pms_exec"), + "ttl_min": param_store.get_int("PMS_EXEC_ADVICE_TTL_MIN", 10)} + + +def _post(url: str, payload: dict, timeout: int) -> dict: + """HTTP 一跳, 单测在这里打桩。""" + import requests + r = requests.post(url, json=payload, timeout=timeout) + r.raise_for_status() + return r.json() or {} + + +def decide(*, side: str, action: str, ts_code: str, now, day: dict, params: dict, + is_last_day: bool, fired_today: int = 0, quota: int = 0, + pos: dict = None, tdays_left=None, prog: dict = None) -> dict: + """择时判定统一入口 (executor.run_tick 的唯一调用点)。 + + 返回结构与 exec_timing.decide 相同, 另带 source 字段标明这条决定是谁做的: + B 实现B (默认档位, 或实现A未配置) + guard 本地事实性检查 (配额/兜底/停牌/一字板/不追高) —— 与实现无关 + A / A缓存 决策系统应答 (新咨询 / 有效期内复用) + B(实现A不可用: ...) 咨询失败退实现B, 括号里是原因 + prog 由调用方传入指令的 progress dict, 咨询结果/失败冷却会写进 prog["exec_advice"], + 随调用方既有的落表动作持久化; 传 None 则本轮结论不缓存 (dry_run 语义)。 + """ + if not available(): + d = et.decide(side=side, now=now, day=day, params=params, is_last_day=is_last_day, + fired_today=fired_today, quota=quota) + d["source"] = IMPL_B + return d + + left = max(0, int(quota) - int(fired_today)) + h = et.hard_gate(side=side, now=now, day=day, params=params, is_last_day=is_last_day, + fired_today=fired_today, quota=quota) + if h is not None: + h["source"] = "guard" # 本地事实性检查 (配额/兜底等), 与实现无关 + return h + + advice, note = _advice(ts_code=ts_code, side=side, action=action, now=now, day=day, + pos=pos or {}, left=left, is_last_day=is_last_day, + tdays_left=tdays_left, prog=prog) + if advice is not None: + d = et.apply_advice(side=side, day=day, params=params, advice=advice, left=left) + if d is not None: + d["source"] = advice.get("_source") or "A" + return d + note = f"研判动作无法识别: {advice.get('verdict')!r}" + + d = et.decide(side=side, now=now, day=day, params=params, is_last_day=is_last_day, + fired_today=fired_today, quota=quota) + d["source"] = f"B(实现A不可用: {note})" + return d + + +def _advice(*, ts_code, side, action, now, day, pos, left, is_last_day, tdays_left, prog): + """取一份有效研判: 缓存命中 → 直接用; 冷却中 → (None, 原因); 否则咨询一次。 + 返回 (advice|None, 不可用原因)。advice 带 _source 标明 A / A缓存。""" + now_min = et.hm_to_min(now) + today = td.ymd() + ttl = max(1, param_store.get_int("PMS_EXEC_ADVICE_TTL_MIN", 10)) + cached = dict((prog or {}).get("exec_advice") or {}) + + if int(cached.get("ymd") or 0) == today: + if (cached.get("verdict") in (FIRE, WAIT) + and now_min <= int(cached.get("valid_until_min") or -1)): + c = dict(cached) + c["_source"] = "A缓存" + return c, "" + if cached.get("fail_until_min") and now_min <= int(cached["fail_until_min"]): + return None, (f"冷却至 {et._fmt(int(cached['fail_until_min']))}: " + f"{cached.get('error') or '上次咨询失败'}") + + payload = { + "direction": "PMS_EXEC", "ts_code": ts_code, "side": side, "action": action, + "qty_left": left, "is_last_day": bool(is_last_day), "tdays_left": tdays_left, + "now": et._fmt(now_min), + "day": {k: day.get(k) for k in ("price", "vwap", "open", "high", "low", + "day_chg_from_open", "bars")}, + "refs": {"support": pos.get("support_ref"), "pressure": pos.get("pressure_ref"), + "stop": pos.get("stop_ref"), "source": pos.get("ref_source")}, + "position": {"total_qty": pos.get("total_qty"), "avail_qty": pos.get("avail_qty"), + "avg_cost": pos.get("avg_cost"), "cushion_pct": pos.get("cushion_pct")}, + } + to = max(1, param_store.get_int("PMS_EXEC_TIMEOUT_SEC", 8)) + url = base_url() + (param_store.get("PMS_EXEC_PATH", "/api/intraday/pms_exec") or "") + + def _cool(err: str): + cool = max(1, param_store.get_int("PMS_EXEC_FAIL_COOLDOWN_MIN", 5)) + if prog is not None: + prog["exec_advice"] = {"ymd": today, "error": err[:200], + "fail_until_min": et.add_trade_minutes(now_min, cool), + "consulted_at": et._fmt(now_min)} + + try: + data = _post(url, payload, to) + except Exception as e: + err = f"{type(e).__name__}: {e}" + logger.error("[择时A] 咨询失败, 本轮退实现B (%s %s): %s", ts_code, side, err) + _cool(err) + return None, err + + verdict = str(data.get("verdict") or "").strip().upper() + if verdict not in (FIRE, WAIT): + # 对端明说给不出结论 (UNAVAILABLE), 或答复不认识 —— 都按拿不到处理, 退实现B + reason = str(data.get("reason") or f"verdict={verdict or '空'}")[:200] + logger.warning("[择时A] 决策系统给不出结论 (%s %s): %s —— 本轮退实现B", + ts_code, side, reason) + _cool(f"UNAVAILABLE: {reason}") + return None, f"UNAVAILABLE: {reason}" + + try: + valid_min = int(data.get("valid_min") or ttl) + except (TypeError, ValueError): + valid_min = ttl + valid_min = max(1, min(valid_min, ttl)) # 对端只能缩短有效期, 不能放长 + adv = {"ymd": today, "verdict": verdict, + "limit_price": data.get("limit_price"), + "reason": str(data.get("reason") or "")[:200], + "confidence": data.get("confidence"), + "valid_until_min": et.add_trade_minutes(now_min, valid_min), + "consulted_at": et._fmt(now_min)} + if prog is not None: + prog["exec_advice"] = adv + a = dict(adv) + a["_source"] = "A" + return a, "" diff --git a/app/services/executor.py b/app/services/executor.py index bd28693..ea6fac1 100644 --- a/app/services/executor.py +++ b/app/services/executor.py @@ -23,8 +23,8 @@ from app.core import exec_timing as et from app.core import rule_gate from app.core import tradedays as td from app.repo import pms_repo, qmt_repo -from app.services import (command_service, dispatcher, industry, market, param_store, - portfolio) +from app.services import (command_service, dispatcher, exec_advisor, industry, market, + param_store, portfolio) logger = logging.getLogger("pms.exec") @@ -132,6 +132,7 @@ def run_tick(*, now=None, dry_run: bool = False) -> dict: "eod_force_time": param_store.get("PMS_EOD_FORCE_TIME", "14:45"), "eod_force_discount": param_store.get_float("PMS_EOD_FORCE_DISCOUNT", 0.998), "no_chase_ma5": param_store.get_float("PMS_NO_CHASE_MA5", 0.06), + "advice_limit_band": param_store.get_float("PMS_EXEC_LIMIT_BAND", 0.10), } slices = param_store.get_int("PMS_EXEC_SLICES", 1) order_ttl = param_store.get_int("PMS_ORDER_TTL_MIN", 10) # 单个分片挂单有效期(交易分钟) @@ -157,8 +158,15 @@ def run_tick(*, now=None, dry_run: bool = False) -> dict: day_ctx = {**day, "support": pos.get("support_ref"), "limit_up": _limit_up(day), "limit_down": _limit_down(day), "halted": not day or not day.get("price")} - d = et.decide(side=side, now=now, day=day_ctx, params=exec_prm, - is_last_day=is_last, fired_today=fired_today, quota=quota) + # 择时判定统一走 exec_advisor: PMS_EXEC_IMPL=B (默认) 时它就是 et.decide 原样; + # =A 时先过硬闸 (配额/兜底等 PMS 自留地), 再委托决策系统, 不可用退实现B。 + # 咨询结论会写进 prog["exec_advice"], 随下面既有的 update_instruction 落表; + # dry_run 传 None —— 只算不落库, 试算不该占用/刷新研判缓存。 + d = exec_advisor.decide(side=side, action=ins.get("action"), ts_code=code, + now=now, day=day_ctx, params=exec_prm, + is_last_day=is_last, fired_today=fired_today, + quota=quota, pos=pos, tdays_left=tdays_left, + prog=(None if dry_run else prog)) if d["action"] != et.ACT_FIRE: prog["last_decision"] = {"at": now.strftime("%H:%M"), **d} @@ -234,7 +242,7 @@ def run_tick(*, now=None, dry_run: bool = False) -> dict: prog["children"] = children prog.setdefault("dispatched_at", now.strftime("%Y-%m-%d %H:%M:%S")) prog["last_decision"] = {"at": now.strftime("%H:%M"), "action": et.ACT_FIRE, - "reason": d["reason"]} + "reason": d["reason"], "source": d.get("source")} # 这里曾经多传了一个 limit_price=None (2026-08-03 实机): 真 repo 的 # update_instruction 没有这个形参 → TypeError → 被本函数外层的 # `except Exception` 吞成 out["errors"] 里的一条。后果是**单子已经发到 QMT diff --git a/app/services/param_store.py b/app/services/param_store.py index 19dc609..5682764 100644 --- a/app/services/param_store.py +++ b/app/services/param_store.py @@ -99,6 +99,13 @@ DESC = { "PMS_SECTOR_MAX_NAMES": "同行业最大持仓只数 (硬拦截)", "PMS_SECTOR_MAX_RATIO": "同行业最大占总仓比例 (硬拦截)", "PMS_EXEC_WINDOW_TDAYS": "任务命令默认执行窗口 (交易日)", + "PMS_EXEC_IMPL": "择时实现: B=内置保守择时 / A=委托决策系统按凌晨结论给执行区间 (不可用自动退B)", + "PMS_EXEC_API_BASE": "择时实现A接口根地址; 留空=沿用 PMS_JUDGE_API_BASE", + "PMS_EXEC_PATH": "择时接口路径 (bionic 侧 /api/intraday/pms_exec)", + "PMS_EXEC_TIMEOUT_SEC": "择时咨询超时 (秒), 超时本轮退实现B", + "PMS_EXEC_ADVICE_TTL_MIN": "择时应答有效期 (交易分钟), 期内不重复咨询", + "PMS_EXEC_FAIL_COOLDOWN_MIN": "择时咨询失败后的冷却 (交易分钟), 冷却内直接走B", + "PMS_EXEC_LIMIT_BAND": "建议价偏离现价超此幅度视为异常数据改用本地口径 (建议价=区间边缘, 偏离几个点属正常)", "PMS_SELL_AVOID_OPEN_MIN": "卖出避开开盘 N 分钟", "PMS_BUY_HALT_DAYUP": "当日涨幅超此停止买入 (不追高)", "PMS_EOD_FORCE_TIME": "当日配额兜底时点", "PMS_EOD_FORCE_DISCOUNT": "兜底限价系数 (卖出)", @@ -288,6 +295,8 @@ _RANGES = { "PMS_DCA_MAX_RATIO": (0, 1), "PMS_SECTOR_MAX_RATIO": (0, 1), "PMS_T0_RATIO_MAX": (0, 0.3334), "PMS_MAX_NAMES": (1, 200), "PMS_TOTAL_SCALE": (0, 10 ** 12), "PMS_EXEC_WINDOW_TDAYS": (1, 20), "PMS_BRAKE_DAYS": (0, 30), + "PMS_EXEC_TIMEOUT_SEC": (1, 60), "PMS_EXEC_ADVICE_TTL_MIN": (1, 120), + "PMS_EXEC_FAIL_COOLDOWN_MIN": (1, 120), "PMS_EXEC_LIMIT_BAND": (0, 0.2), "PMS_PLAN_TOP_N": (1, 1000), "PMS_PLAN_TIMEOUT": (1, 120), "PMS_PLAN_CACHE_SEC": (5, 86400), "PMS_PLAN_STALE_TDAYS": (0, 20), "PMS_PLAN_MIN_SOURCES": (0, 100), "PMS_PLAN_TOP": (0, 5000), @@ -299,6 +308,8 @@ _RANGES = { def _range_check(key, v): if key == "PMS_AUTONOMY" and v not in ("full", "propose_only", "off"): return "PMS_AUTONOMY 只能是 full / propose_only / off" + if key == "PMS_EXEC_IMPL" and str(v).strip().upper() not in ("A", "B"): + return "PMS_EXEC_IMPL 只能是 A (委托决策系统) / B (内置保守择时)" if key == "PMS_SECTOR_SOURCE" and v not in ("", "gp_hybk", "custom_table", "gp_stock_category"): return "PMS_SECTOR_SOURCE 只能是 空 / gp_hybk / custom_table / gp_stock_category" diff --git a/config/settings.py b/config/settings.py index e0bc39d..bb7bb56 100644 --- a/config/settings.py +++ b/config/settings.py @@ -136,6 +136,20 @@ class Settings(BaseSettings): PMS_EXEC_SLICES: int = 1 # 当日配额分几笔出手 (设计「分笔卖出配额」) PMS_ORDER_TTL_MIN: int = 10 # 单个下发分片的挂单有效期(交易分钟), 到点下游自动撤 + # --- 择时实现A (委托决策系统, 待办 #9; 契约见 BIONIC_PMS_INTERFACE.md) --- + # 决策系统按它凌晨算好的支撑/压力给出执行区间, 盘中只回答「现价在不在区间内」; + # 配额/兜底/停牌/一字板/不追高(当日涨幅) 这些检查始终留在 PMS 本地。 + # 不可用自动退实现B (设计 §13), 绝不因对端故障停出手。 + PMS_EXEC_IMPL: str = "B" # B=内置保守择时(默认) / A=委托决策系统 + PMS_EXEC_API_BASE: str = "" # 实现A接口根地址; 留空=沿用 PMS_JUDGE_API_BASE (同一个 bionic 服务) + PMS_EXEC_PATH: str = "/api/intraday/pms_exec" + PMS_EXEC_TIMEOUT_SEC: int = 8 # 咨询超时(秒)。盘中 tick 等不起长超时, 超了本轮退B + PMS_EXEC_ADVICE_TTL_MIN: int = 10 # 应答有效期(交易分钟), 期内不重复咨询; 对端 valid_min 只缩不放 + PMS_EXEC_FAIL_COOLDOWN_MIN: int = 5 # 咨询失败后的冷却(交易分钟), 冷却内直接走B不再咨询 + PMS_EXEC_LIMIT_BAND: float = 0.10 # 建议价偏离现价超此幅度视为异常数据, 改用本地口径。 + # 建议价通常是执行区间的边缘, 离现价百分之几属正常, 所以放到 10% (A股单日涨跌幅上限); + # 真被这条改写时 reason 里会说明 + # --- QMT WebSocket 直连通道 (协议见 QMT_WS_PROTOCOL.md V1.0) --- # 连接由 pms-ws 常驻进程独占 (app/ws/runner.py); executor 经 pms_qmt_order # 出口表递单, 详见 app/services/dispatcher.py 头部的「进程边界」一节。 diff --git a/scripts/probe_bionic.py b/scripts/probe_bionic.py new file mode 100644 index 0000000..645a596 --- /dev/null +++ b/scripts/probe_bionic.py @@ -0,0 +1,119 @@ +# -*- coding: utf-8 -*- +""" +决策系统 (bionic) 两个 PMS 接口的实机探活 —— 只读, 不落任何表, 不产生任何指令 +============================================================================== +运行 (factorevaluation, 需 .env 能连参数表; --base 显式给地址则连参数表也不用): + + docker compose run --rm --no-deps pms-web python scripts/probe_bionic.py + docker compose run --rm --no-deps pms-web python scripts/probe_bionic.py \ + --base http://192.168.16.178:38000 --code 600000.SH + +探两件事 (待办 #9 / #10 的接通自证, 契约见 BIONIC_PMS_INTERFACE.md): + 1. /api/intraday/pms_exec 择时研判: 送一份手造的现场快照, 看回不回 FIRE/WAIT/UNAVAILABLE + 2. /api/intraday/pms_judge 研判闸: 送一份手造的 DCA 候选, 看回不回 PASS/REJECT + (研判走 LLM, 30~90 秒是正常的; --skip-judge 可只探择时) + +判读: + * exec 回 UNAVAILABLE 时看原因——「无实时价」(收盘后正常)、「没有昨夜结论/已过期」 + (该股不在决策系统每晚分析范围, 换一只探) 都是正常降级不是故障; 盘中探到 FIRE/WAIT 才算全通。 + * judge 回 UNAVAILABLE —— 看 reason: 队列超时说明 brain_queue 拥堵 (考虑独立队列), + 连接拒绝说明地址/网络不通。 + * 两个都通了: 页面把 PMS_JUDGE_API_BASE 填上 (研判闸即生效), + 要切择时A再把 PMS_EXEC_IMPL 改成 A。 +""" +import argparse +import json +import os +import sys +import time + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--base", default="", help="bionic 根地址; 缺省读参数 " + "PMS_EXEC_API_BASE / PMS_JUDGE_API_BASE") + ap.add_argument("--code", default="600000.SH", help="探活用股票代码 (点式)") + ap.add_argument("--skip-judge", action="store_true", help="只探择时, 不打 LLM") + ap.add_argument("--judge-timeout", type=int, default=120) + args = ap.parse_args() + + import requests + + base = args.base.rstrip("/") + if not base: + from app.services import exec_advisor, judge + base = exec_advisor.base_url() or judge.base_url() + if not base: + print("[FAIL] 没有可用地址: --base 未给, PMS_EXEC_API_BASE / " + "PMS_JUDGE_API_BASE 也都为空") + sys.exit(2) + print(f"目标: {base} 探活代码: {args.code}") + + # ---- 1. 择时研判 (纯代码, 应当秒回) ---- + exec_payload = { + "direction": "PMS_EXEC", "ts_code": args.code, "side": "buy", "action": "OPEN", + "qty_left": 500, "is_last_day": False, "tdays_left": 3, + "now": time.strftime("%H:%M"), + "day": {"price": None, "vwap": None, "open": None, "high": None, "low": None, + "day_chg_from_open": None, "bars": 0}, # 现场快照留空 → 由 bionic 自取实时价 + "refs": {"support": None, "pressure": None, "stop": None, "source": None}, + "position": {"total_qty": 0, "avail_qty": 0, "avg_cost": None, "cushion_pct": None}, + } + t0 = time.time() + try: + r = requests.post(f"{base}/api/intraday/pms_exec", json=exec_payload, timeout=15) + body = r.json() + ms = (time.time() - t0) * 1000 + verdict = str(body.get("verdict") or "?") + ok = verdict in ("FIRE", "WAIT", "UNAVAILABLE") + print(f"[{'OK' if ok else 'FAIL'}] pms_exec {r.status_code} {ms:.0f}ms → {verdict}: " + f"{body.get('reason')}") + print(" " + json.dumps({k: body.get(k) for k in + ("limit_price", "valid_min", "confidence", "missing", + "advisor")}, ensure_ascii=False)) + exec_ok = ok + except Exception as e: + print(f"[FAIL] pms_exec 打不通: {type(e).__name__}: {e}") + exec_ok = False + + # ---- 2. 研判闸 (LLM, 慢是正常) ---- + judge_ok = None + if not args.skip_judge: + judge_payload = { + "direction": "PMS_JUDGE", "action": "DCA", "ts_code": args.code, "qty": 200, + "reason": "[探活] 浮亏 -8.5% 触发第一档补仓评估", + "hard_numbers": {"cushion_pct": -0.085, "stage": 1, "price": None}, + "context": {"position": {"ts_code": args.code, "avg_cost": None, + "total_qty": 1000, "cushion_pct": -0.085}, + "recent_ledger": [], "probe": True}, + "must_answer": ["下跌是杀逻辑还是杀情绪"], + } + print(f"研判闸探活中 (走 LLM, 最多等 {args.judge_timeout}s)...") + t0 = time.time() + try: + r = requests.post(f"{base}/api/intraday/pms_judge", json=judge_payload, + timeout=args.judge_timeout) + body = r.json() + sec = time.time() - t0 + verdict = str(body.get("verdict") or "?") + judge_ok = verdict in ("PASS", "REJECT") + tag = "OK" if judge_ok else ("WARN" if verdict == "UNAVAILABLE" else "FAIL") + print(f"[{tag}] pms_judge {r.status_code} {sec:.1f}s → {verdict} " + f"(confidence={body.get('confidence')}): {body.get('reason')}") + except Exception as e: + print(f"[FAIL] pms_judge 打不通: {type(e).__name__}: {e}") + judge_ok = False + + print("-" * 60) + if exec_ok and judge_ok: + print("两个接口都通。下一步: 页面填 PMS_JUDGE_API_BASE=" + base + + " (研判闸即生效); 要切择时A再把 PMS_EXEC_IMPL 改成 A") + elif exec_ok and judge_ok is None: + print("择时接口通 (--skip-judge 未探研判闸)") + sys.exit(0 if exec_ok and judge_ok is not False else 1) + + +if __name__ == "__main__": + main() diff --git a/scripts/run_tests.py b/scripts/run_tests.py index 83bdde4..6523978 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -15,8 +15,9 @@ test_batch8_units.py 榜单变化: 名册指纹/三种语义/尾部闸/落库往返 (60 例) test_batch9_units.py 成本价体检 / 对账按日推进 / 行业闸 / 取整记账 (49 例) test_batch10_units.py 静默失败专项: 关键路径不许丢返回值 + 八条实例 (48 例) + test_batch11_units.py 择时实现A: 本地检查等价/应答折算/缓存冷却/退B (24 例) test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (58 例) - 共 401 例 + 共 425 例 任一子集失败即整体失败 (退出码 1)。 """ import os @@ -28,7 +29,7 @@ ROOT = os.path.dirname(HERE) SUITES = ["test_core_units.py", "test_batch2_units.py", "test_batch3_units.py", "test_batch4_units.py", "test_batch5_units.py", "test_batch6_units.py", "test_batch7_units.py", "test_batch8_units.py", "test_batch9_units.py", - "test_batch10_units.py", "test_wiring.py"] + "test_batch10_units.py", "test_batch11_units.py", "test_wiring.py"] def main(): diff --git a/scripts/test_batch11_units.py b/scripts/test_batch11_units.py new file mode 100644 index 0000000..6275b9a --- /dev/null +++ b/scripts/test_batch11_units.py @@ -0,0 +1,447 @@ +# -*- coding: utf-8 -*- +""" +第十一批模块单测: 择时实现A (委托决策系统, 待办 #9) +================================================================ +运行: 在 tradingSystem 仓库根目录执行 python scripts/test_batch11_units.py +覆盖: + [A] 硬闸拆分后与实现B的等价性 (hard_gate 是从 decide 里拆出来的, 不许拆漂) + [B] apply_advice: 决策系统应答 → 本地决策的折算与建议价保护 + [C] exec_advisor: 档位短路 / 硬闸先行 / 咨询 / 缓存TTL / 失败冷却 / UNAVAILABLE 退B + [D] 参数校验与 run_tick 全链路 (实现A出手、退B出手, 内存桩) +约定同前: 全过输出 "ALL PASS (n cases)" 退出码 0。 + +这批的三条纪律 (与四条项目铁律对齐): + * 拿不到不等于有答案 —— 咨询失败/对端给不出结论, 一律整轮退实现B, 绝不把「拿不到」当 FIRE/WAIT。 + * 配额与收尾兜底始终留在 PMS 本地 —— 本地检查命中时连咨询都不发生 (用 HTTP 计数器钉死)。 + * 默认档位 B 必须与接通前一字不差 —— 实现A是加装不是改装。 +""" +import os +import sys +import traceback +from datetime import datetime + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +HERE = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, HERE) + +from app.core import exec_timing as et # noqa: E402 + +RESULTS = [] + + +def case(name): + def deco(fn): + RESULTS.append((name, fn)) + return fn + return deco + + +PRM = {"sell_avoid_open_min": 30, "buy_halt_dayup": 0.05, + "eod_force_time": "14:45", "eod_force_discount": 0.998, + "advice_limit_band": 0.03} + + +def day(**kw): + d = {"price": 10.0, "vwap": 10.0, "high": 10.5, "low": 9.5, "open": 10.0, + "day_chg_from_open": 0.0, "bars": 60} + d.update(kw) + return d + + +def _dec(side="buy", now="10:30", d=None, quota=1000, fired=0, last=False): + return dict(side=side, now=now, day=d or day(), params=PRM, + is_last_day=last, fired_today=fired, quota=quota) + + +# ================================================================ +# [A] 硬闸拆分后的等价性 +# ================================================================ +@case("[A1] 硬闸命中的场景: hard_gate 与 decide 逐字段一致 (拆分不许拆漂)") +def _(): + scenarios = [ + _dec(d=day(price=0)), # 无价 + _dec(d=day(halted=True)), # 停牌 + _dec(quota=1000, fired=1000), # 配额尽 + _dec(now="09:10"), # 非时段 + _dec(now="12:00"), # 午休 + _dec(side="sell", d=day(limit_down=True)), # 跌停一字板 + _dec(side="buy", d=day(limit_up=True)), # 涨停一字板 + _dec(side="buy", d=day(day_chg_from_open=0.06)), # 不追高 STOP + _dec(side="sell", now="14:50"), # 卖出兜底 + _dec(side="buy", now="14:50", last=True), # 买入末日兜底 + _dec(side="buy", now="14:50", last=False), # 兜底后不新开买单 + _dec(side="hold"), # 非法方向 + ] + for kw in scenarios: + h = et.hard_gate(**kw) + d = et.decide(**kw) + assert h is not None, f"该场景硬闸必须有决定: {kw}" + assert h == d, f"硬闸与 decide 不一致: {kw}\n hard={h}\n decide={d}" + + +@case("[A2] 看法区间 (盘中非兜底、无一字板): hard_gate 放行为 None, decide 给出看法") +def _(): + for kw in (_dec(side="buy", now="10:30"), + _dec(side="sell", now="10:30"), + _dec(side="sell", now="09:40"), # 避开开盘属看法, 不属硬闸 + _dec(side="buy", now="10:30", d=day(price=10.2, vwap=10.0))): + assert et.hard_gate(**kw) is None, kw + assert et.decide(**kw)["action"] in (et.ACT_FIRE, et.ACT_WAIT), kw + + +@case("[A3] 卖出跌停一字板在末日不跳过 (顺延无日可顺), 两个口径同步") +def _(): + kw = _dec(side="sell", d=day(limit_down=True), last=True, now="10:30") + assert et.hard_gate(**kw) is None # 末日不 SKIP → 进看法/研判 + assert et.decide(**kw)["action"] in (et.ACT_FIRE, et.ACT_WAIT) + + +# ================================================================ +# [B] apply_advice: 研判 → 决策 +# ================================================================ +@case("[B1] 对端答 WAIT → WAIT, 理由带 [实现A] 前缀") +def _(): + r = et.apply_advice(side="buy", day=day(), params=PRM, + advice={"verdict": "WAIT", "reason": "现价高于买入区间上沿"}, left=800) + assert r["action"] == et.ACT_WAIT and r["qty_hint"] == 800 + assert r["reason"].startswith("[实现A]") and "买入区间上沿" in r["reason"] + + +@case("[B2] 对端答 FIRE → FIRE, 应答里的建议价直接采用") +def _(): + r = et.apply_advice(side="buy", day=day(price=10.0), params=PRM, + advice={"verdict": "FIRE", "limit_price": 10.05, + "reason": "现价在买入区间内"}, left=600) + assert r["action"] == et.ACT_FIRE and r["limit_price"] == 10.05, r + assert r["forced"] is False + + +@case("[B3] 建议价缺失/为0 → 按本地口径 (买 ×1.002 / 卖 ×0.998)") +def _(): + rb = et.apply_advice(side="buy", day=day(price=10.0), params=PRM, + advice={"verdict": "FIRE"}, left=100) + assert rb["limit_price"] == round(10.0 * 1.002, 2), rb + rs = et.apply_advice(side="sell", day=day(price=10.0), params=PRM, + advice={"verdict": "FIRE", "limit_price": 0}, left=100) + assert rs["limit_price"] == round(10.0 * 0.998, 2), rs + + +@case("[B4] 建议价偏离现价超 band → 本地口径重定并在理由里留痕") +def _(): + r = et.apply_advice(side="buy", day=day(price=10.0), params=PRM, + advice={"verdict": "FIRE", "limit_price": 12.0}, left=100) + assert r["limit_price"] == round(10.0 * 1.002, 2), r + assert "偏离现价" in r["reason"], r["reason"] + # band 可调: 放宽到 25% 后 12.0 (偏 20%) 就该被采用 + r2 = et.apply_advice(side="buy", day=day(price=10.0), + params={**PRM, "advice_limit_band": 0.25}, + advice={"verdict": "FIRE", "limit_price": 12.0}, left=100) + assert r2["limit_price"] == 12.0, r2 + + +@case("[B5] verdict 无法识别 → None (拿不到不等于有答案, 由调用方退实现B)") +def _(): + for v in ("", None, "MAYBE", "PASS", "REJECT"): + assert et.apply_advice(side="buy", day=day(), params=PRM, + advice={"verdict": v}, left=100) is None, v + + +# ================================================================ +# [C] exec_advisor 服务层 (HTTP 打桩) +# ================================================================ +def _advisor(params=None, resp=None, exc=None): + """装桩并返回 (exec_advisor, calls) —— calls 记录每次 HTTP 咨询的 payload。""" + from test_wiring import install_fakes + install_fakes(prices={}, params=params or {}) + from app.services import exec_advisor as ea + calls = [] + + def fake_post(url, payload, timeout): + calls.append({"url": url, "payload": payload, "timeout": timeout}) + if exc: + raise exc + return dict(resp or {}) + ea._post = fake_post + return ea, calls + + +def _adecide(ea, *, side="buy", now="10:30", d=None, quota=1000, fired=0, + last=False, prog=None, pos=None): + return ea.decide(side=side, action="OPEN" if side == "buy" else "TRIM", + ts_code="600000.SH", now=now, day=d or day(), params=PRM, + is_last_day=last, fired_today=fired, quota=quota, + pos=pos or {"support_ref": 9.5, "pressure_ref": 11.0, + "total_qty": 0, "avail_qty": 0}, + tdays_left=3, prog=prog) + + +@case("[C1] 默认档位 B: 决策与 et.decide 一字不差, 零咨询 (实现A是加装不是改装)") +def _(): + ea, calls = _advisor(resp={"verdict": "FIRE"}) + for kw in (dict(side="buy", now="10:30"), dict(side="sell", now="10:30"), + dict(side="sell", now="14:50"), dict(side="buy", now="09:10")): + got = _adecide(ea, **kw) + src = got.pop("source") + want = et.decide(side=kw["side"], now=kw["now"], day=day(), params=PRM, + is_last_day=False, fired_today=0, quota=1000) + assert got == want, (got, want) + assert src == "B" + assert not calls, "档位 B 不该有任何 HTTP 咨询" + + +@case("[C2] 档位A但接口地址全空 → 仍走 B (available=False), 状态一句话说破") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A"}) + r = _adecide(ea) + assert r["source"] == "B" and not calls + st = ea.status() + assert st["available"] is False and "地址为空" in st["note"], st + + +@case("[C3] 档位A: FIRE 研判 → 出手, 建议价采用, 结论写入 prog.exec_advice") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "FIRE", "limit_price": 10.02, + "reason": "现价在买入区间内, 限价挂区间上沿", "valid_min": 6, + "confidence": 70}) + prog = {} + r = _adecide(ea, prog=prog) + assert r["action"] == et.ACT_FIRE and r["limit_price"] == 10.02, r + assert r["source"] == "A" and "[实现A]" in r["reason"] + assert len(calls) == 1 and calls[0]["url"] == "http://b/api/intraday/pms_exec" + adv = prog["exec_advice"] + assert adv["verdict"] == "FIRE" and adv["ymd"] and adv["valid_until_min"], adv + # valid_min=6 < TTL 10 → 有效期按 6 分钟算 + assert adv["valid_until_min"] == et.add_trade_minutes(et.hm_to_min("10:30"), 6), adv + + +@case("[C4] 有效期内复用缓存不再咨询 (source=A缓存), 过期后重新咨询") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "WAIT", "reason": "等回踩", "valid_min": 5}) + prog = {} + r1 = _adecide(ea, now="10:30", prog=prog) + assert r1["action"] == et.ACT_WAIT and r1["source"] == "A" and len(calls) == 1 + r2 = _adecide(ea, now="10:33", prog=prog) # 5 分钟内 + assert r2["source"] == "A缓存" and len(calls) == 1, (r2, len(calls)) + r3 = _adecide(ea, now="10:36", prog=prog) # 过期 + assert r3["source"] == "A" and len(calls) == 2, (r3, len(calls)) + + +@case("[C5] 缓存是昨天的 → 当天首跳直接重新咨询 (隔日不吃旧结论)") +def _(): + from app.core import tradedays as td + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "WAIT", "reason": "x"}) + prog = {"exec_advice": {"ymd": td.ymd() - 1, "verdict": "FIRE", + "valid_until_min": 24 * 60, "reason": "昨日结论"}} + r = _adecide(ea, prog=prog) + assert len(calls) == 1 and r["source"] == "A", (len(calls), r) + assert prog["exec_advice"]["ymd"] == td.ymd() + + +@case("[C6] 咨询异常 → 本轮退实现B并进入冷却, 冷却内不再咨询") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + exc=RuntimeError("connect refused")) + prog = {} + r1 = _adecide(ea, now="10:30", prog=prog) + assert r1["source"].startswith("B(实现A不可用"), r1 + assert "connect refused" in r1["source"] + assert len(calls) == 1 + b = et.decide(side="buy", now="10:30", day=day(), params=PRM, + is_last_day=False, fired_today=0, quota=1000) + assert r1["action"] == b["action"] and r1["reason"] == b["reason"] # 真在走B + assert prog["exec_advice"]["fail_until_min"], prog + r2 = _adecide(ea, now="10:32", prog=prog) # 冷却 (默认5分钟) 内 + assert len(calls) == 1 and "冷却至" in r2["source"], (len(calls), r2) + r3 = _adecide(ea, now="10:40", prog=prog) # 冷却过后恢复咨询 + assert len(calls) == 2, len(calls) + + +@case("[C7] 对端回 UNAVAILABLE (给不出结论) → 退B + 冷却, 不当 FIRE 也不当 WAIT") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "UNAVAILABLE", "reason": "昨夜结论缺失, 请退内置择时"}) + prog = {} + r = _adecide(ea, prog=prog) + assert r["source"].startswith("B(实现A不可用: UNAVAILABLE"), r + assert "UNAVAILABLE" in prog["exec_advice"]["error"], prog + + +@case("[C8] 本地检查先行: 配额尽/兜底/不追高时零咨询 (这些永远不问决策系统)") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "WAIT", "reason": "决策系统说等"}) + r1 = _adecide(ea, quota=1000, fired=1000) # 配额尽 + assert r1["action"] == et.ACT_WAIT and r1["source"] == "guard" and not calls + r2 = _adecide(ea, side="sell", now="14:50") # 兜底: 对端说什么都不算 + assert r2["action"] == et.ACT_FIRE and r2["forced"] and r2["source"] == "guard" + assert not calls + r3 = _adecide(ea, side="buy", d=day(day_chg_from_open=0.08)) # 不追高 + assert r3["action"] == et.ACT_STOP and r3["source"] == "guard" and not calls + + +@case("[C9] valid_min 只缩不放: 对端给 999 分钟按本端 TTL 上限截") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b", + "PMS_EXEC_ADVICE_TTL_MIN": "10"}, + resp={"verdict": "WAIT", "reason": "x", "valid_min": 999}) + prog = {} + _adecide(ea, now="10:30", prog=prog) + assert prog["exec_advice"]["valid_until_min"] == \ + et.add_trade_minutes(et.hm_to_min("10:30"), 10), prog + + +@case("[C10] prog=None (试算语义): 咨询照发但不缓存, 两跳两问") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "WAIT", "reason": "x"}) + _adecide(ea, prog=None) + _adecide(ea, prog=None) + assert len(calls) == 2, len(calls) + + +@case("[C11] 咨询 payload 带齐现场: 行情快照/参考位/持仓/窗口 (bionic 靠它免重取)") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://b"}, + resp={"verdict": "WAIT", "reason": "x"}) + _adecide(ea, side="sell", now="10:31", + pos={"support_ref": 9.5, "pressure_ref": 11.0, "stop_ref": 9.2, + "ref_source": "bionic", "total_qty": 3000, "avail_qty": 3000, + "avg_cost": 9.0, "cushion_pct": 0.11}) + p = calls[0]["payload"] + assert p["direction"] == "PMS_EXEC" and p["side"] == "sell" + assert p["day"]["price"] == 10.0 and p["day"]["vwap"] == 10.0 + assert p["refs"]["support"] == 9.5 and p["refs"]["pressure"] == 11.0 + assert p["position"]["avg_cost"] == 9.0 and p["qty_left"] == 1000 + assert p["now"] == "10:31" and p["tdays_left"] == 3 + + +@case("[C12] PMS_EXEC_API_BASE 独立配置时优先于 PMS_JUDGE_API_BASE") +def _(): + ea, calls = _advisor(params={"PMS_EXEC_IMPL": "A", "PMS_JUDGE_API_BASE": "http://judge", + "PMS_EXEC_API_BASE": "http://exec"}, + resp={"verdict": "WAIT", "reason": "x"}) + _adecide(ea, prog={}) + assert calls[0]["url"].startswith("http://exec/"), calls[0]["url"] + + +# ================================================================ +# [D] 参数校验 + run_tick 全链路 +# ================================================================ +@case("[D1] PMS_EXEC_IMPL 只认 A/B (大小写归一), 其他值被参数中心拒绝") +def _(): + from test_wiring import install_fakes + install_fakes(prices={}) + from app.services import param_store + assert param_store.set_param("PMS_EXEC_IMPL", "A")["ok"] + assert param_store.set_param("PMS_EXEC_IMPL", "b")["ok"] # 归一后合法 + bad = param_store.set_param("PMS_EXEC_IMPL", "C") + assert bad["ok"] is False and "只能是" in bad["error"], bad + from app.services import exec_advisor as ea + param_store.set_param("PMS_EXEC_IMPL", "a") + assert ea.impl() == "A" + + +@case("[D2] run_tick·实现A出手全链路: 研判FIRE → 规则闸 → 下发, exec_advice 落表") +def _(): + from test_wiring import install_fakes + fake = install_fakes(prices={"600000.SH": 10.0}, + params={"PMS_TOTAL_SCALE": "2000000", "PMS_EXEC_IMPL": "A", + "PMS_JUDGE_API_BASE": "http://b"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 9.0}]) + from app.services import exec_advisor as ea, executor + calls = [] + + def fake_post(url, payload, timeout): + calls.append(payload) + return {"verdict": "FIRE", "limit_price": 9.99, "reason": "现价在卖出区间内, 放行分批卖出", + "valid_min": 10, "confidence": 75} + ea._post = fake_post + fake.insert_instruction(instruction_id="INS_A", origin_type="plan", origin_id="P1", + ts_code="600000.SH", action="TRIM", side="sell", qty=3000, + window_tdays=3, status="PROPOSED", + progress={"deadline": "2026-07-29", "is_command": True, + "children": []}) + r = executor.run_tick(now=datetime(2026, 7, 27, 10, 5)) + assert r["ok"] and len(r["fired"]) == 1, r + assert r["fired"][0]["limit"] == 9.99 # 建议价被采用 + assert "[实现A]" in r["fired"][0]["reason"] + ins = fake.instructions["INS_A"] + assert ins["status"] == "DISPATCHED" + assert ins["progress"]["exec_advice"]["verdict"] == "FIRE" # 结论随落表持久化 + assert len(calls) == 1 + # 同日第二跳: 配额已出完 → 硬闸 WAIT, 不再咨询 + r2 = executor.run_tick(now=datetime(2026, 7, 27, 10, 6)) + assert r2["waited"] and len(calls) == 1, (r2, len(calls)) + + +@case("[D3] run_tick·实现A故障不停出手: 咨询挂了照走实现B, 单照下") +def _(): + from test_wiring import install_fakes + fake = install_fakes(prices={"600000.SH": 10.0}, + params={"PMS_TOTAL_SCALE": "2000000", "PMS_EXEC_IMPL": "A", + "PMS_JUDGE_API_BASE": "http://b"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 9.0}]) + from app.services import exec_advisor as ea, executor + + def dead_post(url, payload, timeout): + raise RuntimeError("bionic down") + ea._post = dead_post + fake.insert_instruction(instruction_id="INS_F", origin_type="plan", origin_id="P1", + ts_code="600000.SH", action="TRIM", side="sell", qty=3000, + window_tdays=3, status="PROPOSED", + progress={"deadline": "2026-07-29", "is_command": True, + "children": []}) + # 10:05 价==vwap → 实现B本来就该出手; 实现A故障不能把这一单卡住 + r = executor.run_tick(now=datetime(2026, 7, 27, 10, 5)) + assert r["ok"] and len(r["fired"]) == 1, r + ins = fake.instructions["INS_F"] + assert ins["status"] == "DISPATCHED" + assert ins["progress"]["exec_advice"]["error"], ins["progress"] # 冷却已记 + assert "实现A不可用" in ins["progress"]["last_decision"]["source"] + + +@case("[D4] run_tick·默认档位B下 exec_advice 不出现 (不接通就没有新状态)") +def _(): + from test_wiring import install_fakes + fake = install_fakes(prices={"600000.SH": 10.0}, + params={"PMS_TOTAL_SCALE": "2000000"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 9.0}]) + from app.services import executor + fake.insert_instruction(instruction_id="INS_B0", origin_type="plan", origin_id="P1", + ts_code="600000.SH", action="TRIM", side="sell", qty=3000, + window_tdays=3, status="PROPOSED", + progress={"deadline": "2026-07-29", "is_command": True, + "children": []}) + r = executor.run_tick(now=datetime(2026, 7, 27, 10, 5)) + assert r["ok"] and r["fired"], r + assert "exec_advice" not in fake.instructions["INS_B0"]["progress"] + assert fake.instructions["INS_B0"]["progress"]["last_decision"]["source"] == "B" + + +# ---------------------------------------------------------------- runner +def main(): + passed, failed = 0, 0 + for name, fn in RESULTS: + try: + fn() + print(f" PASS {name}") + passed += 1 + except Exception: + print(f" FAIL {name}") + traceback.print_exc() + failed += 1 + print("-" * 60) + if failed: + print(f"FAILED: {failed} / {passed + failed}") + sys.exit(1) + print(f"ALL PASS ({passed} cases)") + + +if __name__ == "__main__": + main()