diff --git a/DEVLOG.md b/DEVLOG.md index 920f4b3..283c296 100644 --- a/DEVLOG.md +++ b/DEVLOG.md @@ -30,6 +30,51 @@ --- +## 2026-08-18 · 决策系统卖出采纳复核:降门槛+全留痕 / 确认即加速 / 清仓完成清场 + +**背景** +用户观测"系统似乎不采纳决策系统的卖出",复核整条链(bionic 生产端 → PMS 消费端)。结论: +接线、代码格式(都是 Tushare 点式,能对上持仓)、置信度量纲(风控 0~100 归一到 0~1)都没 +问题;"不采纳"是门槛标定 + 挂策略路由所致。真机取数印证:bionic 对持仓票的卖出置信度几乎 +全落在 60~84(002015 60/70/70、300474 65、300433/300354 72、688256 70…),而原门槛是 +低于 75 忽略、75~84 落提议、≥85 才自动清仓——于是持仓票的卖出大多被静默忽略或排队等拍板。 + +**做了什么(三件,用户经 AskUserQuestion 拍板)** +一,降门槛 + 全留痕。`PMS_SIGNAL_SELL_CONF_MIN` 默认 0.75 → 0.60(页面可再调),60~84 一律 +浮上来变提议、可见可一键采纳,≥85 仍自动清仓;并在 signal_service 对**持仓票**的"未采纳 +卖出"(低于门槛的 IGNORE)也留痕(action=SIGNAL / verdict=NOTE,当日按来源+股票+动作去重 +一条),不再让一只连日被喊卖的持仓票在账本里查不到。 +二,确认即加速。决策系统卖出命中一只**已有在途清仓指令**的票时,不再简单跳过,而是把该股 +在途的卖出**指令**升级为紧急直通(progress.urgent=True,下一跳 run_tick 按现价直卖、不等 +均价不分桶),每条升级账本留痕。仅挂着提议(还没成指令)的不动。人工清仓 + 风控独立确认 += 最强离场信号,走最快的路。 +三,清仓完成清场。持仓归零(status=CLOSED 且 qty≤0)后,自动撤下该股 策略(网格/做T/跟踪)、 +在途建仓方案与买单、挂着的相关提议;每只闭仓票靠标记(PMS_EXIT_CLEANUP_DONE)**只清一次**, +避免把闭仓后新下的建仓命令误撤(重开的票掉出标记,下次闭仓再清)。挂在 command_poll 每分钟 +一跳、幂等。解决"清仓离场后网格还挂着 ACTIVE、一旦持仓重建又逢跌买入"的隐患。 + +**动了哪些文件** +config/settings.py(门槛 0.60);app/services/signal_service.py(持仓票未采纳留痕 + +`_escalate_inflight_sells` 确认加速);app/services/command_service.py(`cleanup_exited_positions` ++ 纯逻辑 `_cleanup_todo` 标记);app/scheduler.py(command_poll 挂清场); +scripts/test_batch13_units.py(新增 4 例:门槛分档 / 清场标记 / 确认加速 / 闭仓清场,打桩 repo) ++ scripts/run_tests.py(登记 batch13)。 + +**部署与判收** +factorevaluation `make deploy`(动了后端与配置,要重建容器),`make test` 见 ALL SUITES PASS +(已含 batch13)。门槛降了默认值,若表里 pms_runtime_param 已覆盖过该键,需在页面把 +PMS_SIGNAL_SELL_CONF_MIN 也改成 0.60 才生效。真机待看:① 持仓票被喊卖时账本出现"未采纳" +留痕或提议;② 人工清仓一只票后决策系统再喊卖,该清仓指令 progress.urgent 变 True、当即出手; +③ 清仓走完后该股策略变 CANCELLED、在途建仓/提议被撤、账本有 CLEANUP 行。 + +**还欠着什么** +当日去重仍是"先到先赢"(同股同动作一天只消化第一条够格的卖出)——先到的若落成提议、把去重键 +烧掉,后到本可自动清仓的高置信卖出会被当重复丢;本轮未动它(避免过度扩面),若实盘发现高置信 +卖出被前面的弱卖出挡住,再加"更强的卖出可升级已有提议为清仓"。自动清仓线 0.85 未动,按用户 +选择保留,可后续按盘中体感调。 + +--- + ## 2026-08-18 · 上游信号面板改版:删下线框 + 按重要度分档 + 资金异动 z 可读化 + 待触发条 **做了什么** diff --git a/app/scheduler.py b/app/scheduler.py index ed8e705..1c73a0b 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -132,6 +132,13 @@ def command_poll(): from app.services import command_service r = command_service.plan_pending() r.update(command_service.refresh_progress()) + # 清仓完成后清场: 撤该股策略/在途建仓/相关提议 (幂等, 每只闭仓票只清一次)。 + # 挂在这里而非盘中执行位, 是因为它与持仓状态相关、非交易日照样该清干净。 + try: + r["cleanup"] = command_service.cleanup_exited_positions() + except Exception as e: + logger.error("[command_poll] 清场失败 (不影响命令轮询): %s", e) + r["cleanup_error"] = str(e) return r diff --git a/app/services/command_service.py b/app/services/command_service.py index eb05131..3134689 100644 --- a/app/services/command_service.py +++ b/app/services/command_service.py @@ -650,3 +650,143 @@ def _ledger_rejects(rejects: list, command_id: str): reason="命令规划期规则闸拦截") except Exception as e: logger.warning("拒绝留痕失败: %s", e) + + +# ================================================================ 清仓完成后清场 +CLEANUP_MARK_KEY = "PMS_EXIT_CLEANUP_DONE" # 已闭仓且已清场的代码集 (JSON) +BUILD_ACTIONS = ("OPEN", "FILL", "ADD", "DCA") # 建仓类方案/指令 (买入) + + +def _cleanup_todo(closed_codes, done_codes) -> tuple: + """纯逻辑: 算出本轮要清场的代码与新的标记集 (2026-08-18, 便于单测)。 + + 每只票在**一次闭仓周期里只清一次**: todo = 现在闭仓的 − 已清过的。 + 新标记 = 现在仍闭仓的那部分 —— 重新建仓的票会掉出闭仓集、从而掉出标记, + 下次它再闭仓时会被重新清一遍。这条"只清一次"是防呆的关键: 闭仓后你若又下建仓命令, + 持仓状态短暂还是 CLOSED, 标记挡住清场、不会把你刚下的建仓计划误撤。 + """ + closed = set(closed_codes or ()) + done = set(done_codes or ()) + todo = sorted(closed - done) + new_done = closed # 只保留仍闭仓的; 重开的自然移除 + return todo, new_done + + +def _load_cleanup_done() -> set: + try: + raw = pms_repo.get_param(CLEANUP_MARK_KEY) + import json + return set(json.loads(raw)) if raw else set() + except Exception: + return set() + + +def _save_cleanup_done(codes: set): + try: + import json + pms_repo.set_param(CLEANUP_MARK_KEY, json.dumps(sorted(codes)), "system") + except Exception as e: + logger.warning("[清场] 标记写入失败: %s", e) + + +def cleanup_exited_positions(limit: int = 500) -> dict: + """持仓归零(清仓完成)后清场: 撤该股策略 / 在途建仓计划与买单 / 相关提议 (设计 2026-08-18)。 + + 每分钟一跳, 幂等且防呆: + * 只对 status=CLOSED 且 total_qty<=0 的持仓动手; + * 每次闭仓只清一次 (靠 _cleanup_todo 的标记), 避免把闭仓后新下的建仓命令误撤; + * 四类残留一次性清: 策略(网格/做T/跟踪) 撤下、在途建仓方案作废、在途买单撤回、 + 挂着的相关提议驳回。清场动作在账本留一行 (action=CLEANUP)。 + 你清仓离场后, 网格不再空转挂着、也不会因持仓被任何原因重建而复活逢跌买入。 + """ + from app.services import executor, strategy_service + out = {"scanned": 0, "cleaned": [], "errors": [], "note": ""} + try: + positions = pms_repo.list_positions() + except Exception as e: + return {**out, "errors": [f"读持仓失败: {type(e).__name__}: {e}"]} + + closed = {p["ts_code"] for p in positions + if str(p.get("status")) == "CLOSED" and int(p.get("total_qty") or 0) <= 0} + done = _load_cleanup_done() + todo, new_done = _cleanup_todo(closed, done) + out["scanned"] = len(closed) + if not todo: + _save_cleanup_done(new_done) # 重开的票及时移出标记 + out["note"] = "无新闭仓待清场" if closed else "当前无闭仓持仓" + return out + + todo_set = set(todo) + # 各类残留一次性拉取, 再按待清代码过滤 (省得逐票查库) + try: + strategies = [s for s in pms_repo.list_strategies(statuses=["ACTIVE", "PAUSED"], limit=limit) + if s.get("ts_code") in todo_set] + except Exception as e: + strategies = [] + out["errors"].append(f"读策略失败: {e}") + try: + plans = [pl for pl in pms_repo.list_plans(statuses=["PENDING", "GATED", "EXECUTING"], + limit=limit) + if pl.get("ts_code") in todo_set and pl.get("action") in BUILD_ACTIONS] + except Exception as e: + plans = [] + out["errors"].append(f"读方案失败: {e}") + try: + buys = [i for i in pms_repo.list_instructions(statuses=list(LIVE_INSTR), side="buy", + limit=limit) + if i.get("ts_code") in todo_set] + except Exception as e: + buys = [] + out["errors"].append(f"读指令失败: {e}") + try: + props = [pr for pr in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=limit) + if pr.get("ts_code") in todo_set] + except Exception as e: + props = [] + out["errors"].append(f"读提议失败: {e}") + + per_code = {c: {"ts_code": c, "strategies": [], "plans": [], "instructions": [], + "proposals": []} for c in todo} + for s in strategies: + try: + r = strategy_service.set_status(s["strategy_id"], "CANCELLED", by="system") + if r.get("ok") or r.get("status"): + per_code[s["ts_code"]]["strategies"].append(s["strategy_id"]) + except Exception as e: + out["errors"].append(f"{s.get('ts_code')} 撤策略失败: {e}") + for pl in plans: + try: + pms_repo.update_plan(pl["plan_id"], status="CANCELLED") + per_code[pl["ts_code"]]["plans"].append(pl["plan_id"]) + except Exception as e: + out["errors"].append(f"{pl.get('ts_code')} 撤建仓方案失败: {e}") + for ins in buys: + try: + r = executor.cancel_instruction(ins["instruction_id"], reason="清仓完成清场: 撤在途建仓") or {} + if r.get("ok"): + per_code[ins["ts_code"]]["instructions"].append(ins["instruction_id"]) + else: + out["errors"].append(f"{ins.get('ts_code')} 撤买单未成: {r.get('message') or r.get('error')}") + except Exception as e: + out["errors"].append(f"{ins.get('ts_code')} 撤买单失败: {e}") + for pr in props: + try: + if pms_repo.decide_proposal(pr["proposal_id"], "DECLINED"): + per_code[pr["ts_code"]]["proposals"].append(pr["proposal_id"]) + except Exception as e: + out["errors"].append(f"{pr.get('ts_code')} 撤提议失败: {e}") + + for c, acted in per_code.items(): + if any(acted[k] for k in ("strategies", "plans", "instructions", "proposals")): + out["cleaned"].append(acted) + try: + pms_repo.insert_ledger(ts_code=c, action="CLEANUP", arbiter="system", + verdict="PASS", price_at=0, hard_numbers=acted, + reason="清仓完成清场: 撤策略/在途建仓/相关提议") + except Exception as e: + logger.warning("[清场] %s 留痕失败: %s", c, e) + logger.warning("[清场] %s 清仓完成, 已撤 策略%d/方案%d/买单%d/提议%d", + c, len(acted["strategies"]), len(acted["plans"]), + len(acted["instructions"]), len(acted["proposals"])) + _save_cleanup_done(new_done) # 标记本轮闭仓集 (含刚清过的), 只清一次 + return out diff --git a/app/services/signal_service.py b/app/services/signal_service.py index 743b309..041c2f2 100644 --- a/app/services/signal_service.py +++ b/app/services/signal_service.py @@ -145,6 +145,24 @@ def _handle(sig, view, prm, seen, ymd, dry_run, out, strat_codes=frozenset()): if act == sr.ACT_IGNORE: out["ignored"] += 1 + # **持仓票的卖出被忽略也要留痕** (2026-08-18)。digest 对"低于门槛"返回 IGNORE, + # 而 IGNORE 原来一个字不留 —— 一只你正持有、bionic 连日喊卖的票(实测 002015 一天 + # 三条、300474/300433 等), 在账本里查不到, 等于风控在示警而你看不见。未持有的票 + # 不留(会被全市场卖出冲垮), 只给持仓票留, 且当日按(来源,股票,动作)去重一条。 + if (not dry_run and str(sig.get("action") or "").upper() == "SELL" + and pos and int(pos.get("total_qty") or 0) > 0): + k = sr.dedup_key(sig, ymd) + ":IGN" + if k not in seen: + try: + pms_repo.insert_ledger( + ts_code=code, action="SIGNAL", arbiter="rule", verdict="NOTE", + price_at=float(pos.get("price") or 0), + hard_numbers={**d["hard_numbers"], "msg_id": sig.get("msg_id")}, + reason="决策系统卖出未采纳(置信度低于门槛), 持仓票请留意: " + d["reason"][:400]) + seen.add(k) + out.setdefault("held_sell_noted", []).append(code) + except Exception as e: + logger.warning("[信号消化] 持仓票卖出留痕失败 %s: %s", code, e) return if act == sr.ACT_RECORD: # 只给持有的票留痕, 否则全市场广播会把评审账本冲垮 @@ -187,8 +205,16 @@ def _handle(sig, view, prm, seen, ymd, dry_run, out, strat_codes=frozenset()): out["ignored"] += 1 return if _has_inflight(code): + # **确认即加速** (2026-08-18): 走到这里 act 已是 EXIT/PROPOSE —— 决策系统确认要卖。 + # 若这只票已有在途清仓指令(多为你人工下的清仓命令), 不再简单跳过, 而是把它升级成 + # 紧急直通(马上按现价卖、不等均价不分桶): 人工意图 + 风控独立确认 = 最强离场信号。 + # 只升级已下发的卖出**指令**; 仅挂着提议(还没成指令)的不动, 那本就等你拍板。 + esc = [] if dry_run else _escalate_inflight_sells(code, d["reason"]) out["ignored"] += 1 - out.setdefault("skipped_inflight", []).append(code) + if esc: + out.setdefault("escalated", []).append({"ts_code": code, "instructions": esc}) + else: + out.setdefault("skipped_inflight", []).append(code) return brief = {"ts_code": code, "qty": d["qty"], "confidence": d["hard_numbers"]["confidence"], @@ -328,6 +354,37 @@ def _has_inflight(code: str) -> bool: return False +def _escalate_inflight_sells(code: str, reason: str) -> list: + """把该股在途的卖出**指令**升级为紧急直通 (确认即加速, 2026-08-18)。 + + 只动已下发的卖出指令 (progress.urgent 置 True, 下一跳 run_tick 就按现价直卖); + 已是紧急的不重复动。每条升级在账本留一行, 供复盘"为什么这条清仓突然加速了"。 + 仅挂着提议(还没成指令)的不在此列 —— 那本就等用户拍板, 不该被信号自动推成紧急。 + """ + ids = [] + try: + for ins in pms_repo.list_instructions(statuses=list(executor.LIVE), ts_code=code, + limit=50): + if str(ins.get("side") or "").lower() != "sell": + continue + prog = dict(ins.get("progress") or {}) + if prog.get("urgent"): + continue + prog["urgent"] = True + prog["escalated_by_signal"] = str(reason)[:200] + pms_repo.update_instruction(ins["instruction_id"], progress=prog) + ids.append(ins["instruction_id"]) + pms_repo.insert_ledger( + ts_code=code, action="EXIT", arbiter="rule", verdict="PASS", price_at=0, + ref_id=ins["instruction_id"], + reason="决策系统卖出确认, 在途清仓升级为紧急直通: " + str(reason)[:300]) + logger.warning("[信号消化] %s 决策系统确认卖出, 在途清仓 %s 升级紧急直通", + code, ins["instruction_id"]) + except Exception as e: + logger.warning("[信号消化] 升级在途清仓失败 %s: %s", code, e) + return ids + + def _pos_of(view: dict, ts_code: str): for x in view["positions"]: if x["ts_code"] == ts_code: diff --git a/config/settings.py b/config/settings.py index 091cba6..42f6ace 100644 --- a/config/settings.py +++ b/config/settings.py @@ -244,7 +244,10 @@ class Settings(BaseSettings): PMS_SIGNAL_CONSUMER: str = "pms_1" PMS_SIGNAL_STREAM_INTRADAY: str = "intraday_signals:{ymd}" # db2, 每日一条流 PMS_SIGNAL_STREAM_SELL: str = "bionic:signals:llm_sell_actions" # db3, 固定 key - PMS_SIGNAL_SELL_CONF_MIN: float = 0.75 # 低于此置信度的卖出信号不消化 + # 2026-08-18 实测: bionic 对持仓票的卖出置信度(0~100 证据一致性)多落在 60~84, + # 原门槛 0.75 把大量持仓票的卖出静默忽略。降到 0.60: 60~84 一律浮上来变提议(可见、 + # 一键采纳), ≥85 仍自动清仓; 加上 signal_service 对持仓票的"未采纳也留痕", 不再丢信号。 + PMS_SIGNAL_SELL_CONF_MIN: float = 0.60 # 低于此置信度的卖出信号不消化(仅持仓票留痕) PMS_SIGNAL_AUTO_EXIT_CONF: float = 0.85 # 高于此置信度直接转清仓指令, 之间则落提议 PMS_SIGNAL_TRIM_RATIO: float = 0.3333 # 中等置信度时的减仓比例 diff --git a/scripts/run_tests.py b/scripts/run_tests.py index ebb51ad..bfb86dc 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -33,7 +33,7 @@ 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_batch11_units.py", "test_batch12_units.py", - "test_wiring.py"] + "test_batch13_units.py", "test_wiring.py"] def main(): diff --git a/scripts/test_batch13_units.py b/scripts/test_batch13_units.py new file mode 100644 index 0000000..077bd25 --- /dev/null +++ b/scripts/test_batch13_units.py @@ -0,0 +1,192 @@ +# -*- coding: utf-8 -*- +""" +第十三批模块单测 (实机运行, 零外部依赖) +======================================== +运行: 在 tradingSystem 仓库根目录执行 python scripts/test_batch13_units.py +覆盖 (2026-08-18 决策系统卖出采纳复核后的三处改动): + * 消化门槛降到 0.60 后 digest 的分档 (纯逻辑); + * 清仓完成清场的"每次闭仓只清一次"标记 (_cleanup_todo 纯逻辑); + * 确认即加速: 决策系统卖出命中在途清仓 → 只把卖出**指令**升级紧急 (打桩 repo); + * 闭仓清场: 撤策略/在途建仓/买单/提议, 幂等第二遍无动作 (打桩 repo/service)。 +""" +import os +import sys +import traceback + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from app.core import signal_rules as sr # noqa: E402 +from app.services import command_service as csvc # noqa: E402 +from app.services import signal_service as ssvc # noqa: E402 +from app.repo import pms_repo # noqa: E402 + +RESULTS = [] + + +def case(name): + def deco(fn): + RESULTS.append((name, fn)) + return fn + return deco + + +HELD = {"ts_code": "300474.SZ", "total_qty": 3000, "avail_qty": 3000, "price": 8.0} +PRM60 = {"sell_conf_min": 0.60, "auto_exit_conf": 0.85, "trim_ratio": 1 / 3} + + +# ================================================================ 门槛分档 +@case("门槛 0.60·持仓票卖出: <60 忽略 / 60~84 提议 / ≥85 清仓") +def _(): + def sell(conf): + return {"source": "risk_sell", "ts_code": "300474.SZ", "action": "SELL", + "confidence": conf} + assert sr.digest(sell(0.55), HELD, PRM60)["action"] == sr.ACT_IGNORE + assert sr.digest(sell(0.60), HELD, PRM60)["action"] == sr.ACT_PROPOSE # 恰好到线, 采纳 + assert sr.digest(sell(0.72), HELD, PRM60)["action"] == sr.ACT_PROPOSE # 实测这一带最多 + assert sr.digest(sell(0.84), HELD, PRM60)["action"] == sr.ACT_PROPOSE + assert sr.digest(sell(0.85), HELD, PRM60)["action"] == sr.ACT_EXIT + assert sr.digest(sell(0.92), HELD, PRM60)["action"] == sr.ACT_EXIT + # 未持有仍一律忽略 (不因降门槛就对没持仓的票动手) + assert sr.digest(sell(0.95), {"total_qty": 0}, PRM60)["action"] == sr.ACT_IGNORE + + +# ================================================================ 清场标记 (纯逻辑) +@case("清场标记·每次闭仓只清一次, 重开的票掉出标记下次再清") +def _(): + # 首轮: A/B 刚闭仓, 都要清 + todo, done = csvc._cleanup_todo({"A", "B"}, set()) + assert todo == ["A", "B"] and done == {"A", "B"} + # 次轮: 还是 A/B 闭仓且已清过 → 不再清 + todo, done = csvc._cleanup_todo({"A", "B"}, {"A", "B"}) + assert todo == [] and done == {"A", "B"} + # A 重新建仓(掉出闭仓集), 新来 C 闭仓 → 只清 C; A 从标记移除 + todo, done = csvc._cleanup_todo({"B", "C"}, {"A", "B"}) + assert todo == ["C"] and done == {"B", "C"} and "A" not in done + # A 又闭仓 → 因已移出标记, 会被重新清 + todo, done = csvc._cleanup_todo({"A", "B", "C"}, {"B", "C"}) + assert todo == ["A"] and done == {"A", "B", "C"} + + +# ---------------------------------------------------------------- 打桩工具 +class _Rec: + """记录型假 repo/service: 存住调用, 供断言。""" + def __init__(self): + self.instr = [] # 在途指令 + self.updates = {} # instruction_id -> progress + self.ledger = [] + self.plans = [] + self.strategies = [] + self.proposals = [] + self.positions = [] + self.plan_cancelled = [] + self.prop_declined = [] + self.strat_cancelled = [] + self.instr_cancelled = [] + self.param = {} + + +def _patch_signal(rec): + pms_repo.list_instructions = lambda **kw: [i for i in rec.instr + if (kw.get("ts_code") in (None, i["ts_code"])) + and (kw.get("side") in (None, i.get("side")))] + + def _upd(iid, **kw): + if "progress" in kw: + rec.updates[iid] = kw["progress"] + return 1 + pms_repo.update_instruction = _upd + pms_repo.insert_ledger = lambda **kw: rec.ledger.append(kw) + + +@case("确认加速·只升级在途卖出指令, 买单与已紧急的不动") +def _(): + rec = _Rec() + rec.instr = [ + {"instruction_id": "S1", "ts_code": "600150.SH", "side": "sell", "progress": {}}, + {"instruction_id": "S2", "ts_code": "600150.SH", "side": "sell", + "progress": {"urgent": True}}, # 已紧急, 不重复 + {"instruction_id": "B1", "ts_code": "600150.SH", "side": "buy", "progress": {}}, + ] + _patch_signal(rec) + ids = ssvc._escalate_inflight_sells("600150.SH", "决策系统确认卖出") + assert ids == ["S1"], ids # 只 S1 被升级 + assert rec.updates["S1"]["urgent"] is True + assert "escalated_by_signal" in rec.updates["S1"] + assert "B1" not in rec.updates and "S2" not in rec.updates + assert any(l.get("action") == "EXIT" and "紧急直通" in l.get("reason", "") + for l in rec.ledger) + + +def _patch_cleanup(rec): + import app.services.executor as ex + import app.services.strategy_service as strat + pms_repo.list_positions = lambda **kw: rec.positions + pms_repo.list_strategies = lambda **kw: [s for s in rec.strategies + if s.get("status") in (kw.get("statuses") or [s.get("status")])] + pms_repo.list_plans = lambda **kw: rec.plans + pms_repo.list_instructions = lambda **kw: [i for i in rec.instr + if kw.get("side") in (None, i.get("side"))] + pms_repo.list_proposals = lambda **kw: rec.proposals + pms_repo.update_plan = lambda pid, **kw: rec.plan_cancelled.append(pid) or 1 + pms_repo.decide_proposal = lambda pid, st: rec.prop_declined.append((pid, st)) or 1 + pms_repo.insert_ledger = lambda **kw: rec.ledger.append(kw) + pms_repo.get_param = lambda k: rec.param.get(k) + pms_repo.set_param = lambda k, v, by=None: rec.param.__setitem__(k, v) + strat.set_status = lambda sid, st, by=None: (rec.strat_cancelled.append((sid, st)) + or {"ok": True}) + ex.cancel_instruction = lambda iid, reason=None: (rec.instr_cancelled.append(iid) + or {"ok": True}) + + +@case("闭仓清场·撤策略/建仓方案/买单/提议, 幂等第二遍无动作") +def _(): + rec = _Rec() + rec.positions = [ + {"ts_code": "300474.SZ", "status": "CLOSED", "total_qty": 0}, # 刚闭仓, 要清 + {"ts_code": "600150.SH", "status": "HOLDING", "total_qty": 1700}, # 在持, 不碰 + ] + rec.strategies = [{"strategy_id": "G1", "ts_code": "300474.SZ", "status": "ACTIVE"}, + {"strategy_id": "G2", "ts_code": "600150.SH", "status": "ACTIVE"}] + rec.plans = [{"plan_id": "P1", "ts_code": "300474.SZ", "action": "OPEN"}, + {"plan_id": "P2", "ts_code": "300474.SZ", "action": "EXIT"}] # 卖出方案不撤 + rec.instr = [{"instruction_id": "IB", "ts_code": "300474.SZ", "side": "buy"}] + rec.proposals = [{"proposal_id": "PR1", "ts_code": "300474.SZ", "action": "ADD"}] + _patch_cleanup(rec) + + r = csvc.cleanup_exited_positions() + assert len(r["cleaned"]) == 1 and r["cleaned"][0]["ts_code"] == "300474.SZ", r + assert rec.strat_cancelled == [("G1", "CANCELLED")] # 只撤闭仓票的策略 + assert rec.plan_cancelled == ["P1"] # 只撤建仓方案, 不撤 EXIT + assert rec.instr_cancelled == ["IB"] + assert rec.prop_declined == [("PR1", "DECLINED")] + assert any(l.get("action") == "CLEANUP" for l in rec.ledger) + assert "300474.SZ" in rec.param.get(csvc.CLEANUP_MARK_KEY, "") + + # 幂等: 第二遍 (标记已记 300474) → 不再动任何东西 + rec.strat_cancelled.clear(); rec.plan_cancelled.clear() + rec.instr_cancelled.clear(); rec.prop_declined.clear() + r2 = csvc.cleanup_exited_positions() + assert r2["cleaned"] == [] and not rec.strat_cancelled and not rec.plan_cancelled, r2 + + +# ================================================================ runner +def main(): + passed = failed = 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()