@@ -990,7 +990,7 @@ body.dock-r:not(.r-fold) .side-r .strip{display:none;}
@@ -1741,7 +1741,7 @@ body.dock-r:not(.r-fold) .side-r .strip{display:none;}
-
+
强刷 + 灌行业映射
只重读
@@ -2043,7 +2043,9 @@ createApp({
return Math.abs(g) > 0.10 ? g : 0;
});
- const wsc = computed(() => wsRaw.value.channel || {});
+ // 通道状态优先取 ws-channel (管理员接口), 取不到回落 dispatch-mode 里的 channel ——
+ // 纯交易员没有 ws-channel 权限, 原来「通道没连上」这类告警对他们永远不亮 (2026-08-28)。
+ const wsc = computed(() => wsRaw.value.channel || (dm.value || {}).channel || {});
const wst = computed(() => wsc.value.stat || {});
const wsMode = computed(() => wsRaw.value.mode || '');
const wsOrders = computed(() => wsRaw.value.orders || []);
@@ -2166,6 +2168,15 @@ createApp({
const heldPositions = computed(() => (positions.value || []).filter(p => (p.total_qty || 0) > 0));
// ---- 今日在办分档 + 候选处置 (2026-08-12 重构新增) ----
const todayYmd = computed(() => (String(health.value.now || '').slice(0, 10).replace(/-/g, '')) || '');
+ const todayDateStr = computed(() => String(health.value.now || '').slice(0, 10));
+ // 「今天发生了什么 / 这只票今天的记录」只收 decided_at 是今天的行 (2026-08-28):
+ // /api/ledger 按 id 倒序取近 60 条不带日期条件, 当天记录少时列表被前几天的行填满,
+ // 主面板又只显示时分 —— 昨天 14:31 的驳回看起来就是今天 14:31 发生的。
+ const ledgerToday = computed(() => {
+ const t = todayDateStr.value;
+ if (!t) return ledger.value || [];
+ return (ledger.value || []).filter(l => String(l.decided_at || '').slice(0, 10) === t);
+ });
// 是否处于连续交易时段 (09:30-11:30 / 13:00-15:00 的交易日)。缺实时价在盘中才是问题,
// 休市/非交易日本就没有实时行情, 不该报警 —— 2026-08-12 用户提。
const marketOpen = computed(() => {
@@ -2179,7 +2190,10 @@ createApp({
const _LIVE_INS = ['PROPOSED', 'RULE_PASSED', 'JUDGE_PASSED', 'DISPATCHED'];
const _isArch = r => !!(r && r.archived_at); // 软归档: 有 archived_at = 已从视图移除 (行仍在库)
const _sameDayIns = r => { const y = todayYmd.value;
- return !y || String(r.instruction_id || '').indexOf('INS_' + y + '_') === 0; };
+ const id = String(r.instruction_id || '');
+ // 策略指令的前缀是 STR{ymd} (strategy_runner), 不是 INS_ —— 原来它们一进终态就从
+ // 「今日成交/已经结束」和消息栏同时消失, 交易员眼看着自己的单子凭空不见 (2026-08-28)
+ return !y || id.indexOf('INS_' + y + '_') === 0 || id.indexOf('STR' + y) === 0; };
// 已归档的一律不进「今日在办」三档 —— 归档=从所有面板隐藏。
const insLive = computed(() => (instructions.value || []).filter(r => _LIVE_INS.includes(r.status) && !_isArch(r)));
const insDone = computed(() => (instructions.value || []).filter(r => r.status === 'CONFIRMED' && _sameDayIns(r) && !_isArch(r)));
@@ -2198,7 +2212,8 @@ createApp({
const stratArchivable = r => ['CANCELLED', 'DONE'].includes(r.status);
// 每只候选今天为什么下单/没下单: 已落地的事实(指令/提议)优先, 其余看 /api/open-scan 的只读处置。
function dispOf(code) {
- const inst = (instructions.value || []).find(i => i.ts_code === code && i.action === 'OPEN');
+ const inst = (instructions.value || []).find(i => i.ts_code === code && i.action === 'OPEN'
+ && (_LIVE_INS.includes(i.status) || _sameDayIns(i))); // 旧的终态单不算「已下单」(2026-08-28)
if (inst) return { cls: 'done', label: '已下单', why: '在『今日在办』(' + tx('insStatus', inst.status) + ')' };
const prop = (proposals.value || []).find(p => p.ts_code === code && p.action === 'OPEN');
if (prop) return { cls: 'queue', label: '待你确认', why: '已生成提议,在『等我拍板』等你点头' };
@@ -2266,20 +2281,35 @@ createApp({
return n;
} catch (e) { return null; }
}
+ const quickBusy = ref(false); // 防重复下达: 双击会发出两条相同的资金命令 (2026-08-28)
async function quickCmd(cmd_type, params, confirmMsg) {
- if (confirmMsg) {
- try { await ElementPlus.ElMessageBox.confirm(confirmMsg, '确认', { type: 'warning' }); }
- catch (e) { return; }
- }
- let d = await call('post', '/api/commands', { cmd_type, params: params || {}, note: '交易员视图' });
- if (!d.ok && (d.conflicts || []).length) {
- try {
- await ElementPlus.ElMessageBox.confirm('和还没做完的命令有冲突,还是要下达吗?', '命令冲突', { type: 'warning' });
- d = await call('post', '/api/commands', { cmd_type, params: params || {}, note: '交易员视图', force: true });
- } catch (e) { return; }
- }
- if (d.ok) { ElementPlus.ElMessage.success('已下达 ' + (d.command_id || '')); await loadAll(); }
- else ElementPlus.ElMessage.error((d.errors || [d.error]).join('; '));
+ if (quickBusy.value) { ElementPlus.ElMessage.info('上一条命令还在下达,请稍候'); return; }
+ quickBusy.value = true;
+ try {
+ if (confirmMsg) {
+ try { await ElementPlus.ElMessageBox.confirm(confirmMsg, '确认', { type: 'warning' }); }
+ catch (e) { return; }
+ }
+ let d = await call('post', '/api/commands', { cmd_type, params: params || {}, note: '交易员视图' });
+ if (!d.ok && (d.conflicts || []).length) {
+ try {
+ await ElementPlus.ElMessageBox.confirm('和还没做完的命令有冲突,还是要下达吗?', '命令冲突', { type: 'warning' });
+ d = await call('post', '/api/commands', { cmd_type, params: params || {}, note: '交易员视图', force: true });
+ } catch (e) { return; }
+ }
+ if (d.ok && String(d.status || '') === 'CANCELLED') {
+ // 后端对「该产方案却一条没产」的命令回 ok:true + status=CANCELLED ——
+ // 原来只看 ok 报「已下达」, 用户以为风险在削减, 实际命令已自我作废 (2026-08-28)
+ const why = ((d.plan || {}).reject_summary) || ((d.plan || {}).notes || []).join(';')
+ || '没有产出任何可执行的方案';
+ ElementPlus.ElMessage.warning('命令没有生效(已自动作废):' + why);
+ await loadAll();
+ } else if (d.ok) {
+ ElementPlus.ElMessage.success('已下达 ' + (d.command_id || '')); await loadAll();
+ } else {
+ ElementPlus.ElMessage.error((d.errors || [d.error]).join('; '));
+ }
+ } finally { quickBusy.value = false; }
}
// 个股
function actExitStock(row) { quickCmd('EXIT_STOCK', { ts_code: row.ts_code }, '清仓 ' + nm(row.ts_code) + '?系统会在允许下单的时间段里挑时机把它卖光。'); }
@@ -2455,6 +2485,12 @@ createApp({
.filter(p => String(p.value) !== String(paramsRaw.value[p.key]))
.map(p => ({ key: p.key, value: p.value }));
const d = await call('post', '/api/params', { items });
+ if (d.ok === false && !d.results) {
+ // 传输层失败 (超时/断网): 没有 results 不等于全成功 —— 原来这里报「已保存」
+ // 还立刻用服务器旧值把你刚改的输入冲掉 (2026-08-28)。保留输入, 由你重试。
+ ElementPlus.ElMessage.error('保存失败(网络或服务不可用):' + (d.error || '') + ',你的修改还在输入框里,请重试');
+ return;
+ }
const bad = (d.results || []).filter(r => !r.ok);
if (bad.length) ElementPlus.ElMessage.error(bad.map(b => b.error).join('; '));
else ElementPlus.ElMessage.success('已保存 ' + items.length + ' 项');
@@ -2726,7 +2762,7 @@ createApp({
// 会把日志和连接数刷得很难看, 排查问题时反而碍事。
function wsPollSync() {
const want = tab.value === 'ws' && wsAuto.value;
- if (want && !wsTimer) { loadWs(); wsTimer = setInterval(loadWs, 3000); }
+ if (want && !wsTimer) { loadWs(); wsTimer = setInterval(() => { if (authed.value) loadWs(); }, 3000); }
if (!want && wsTimer) { clearInterval(wsTimer); wsTimer = null; }
}
watch([tab, wsAuto], wsPollSync);
@@ -2736,7 +2772,11 @@ createApp({
loading.value = true; err.value = '';
await Promise.all([loadOverview(), loadParams(), loadCatalog(), loadCommands(),
loadPositions(), loadInstructions(), loadLedger(), loadProposals(),
- loadDispatchMode(), loadWs(), loadStrategies(), loadOpLog(), loadOpenScan(), loadUpstream(),
+ loadDispatchMode(),
+ // ws-channel 是管理员专属接口: 交易员不调它, 免得每次刷新都弹一条「无权限」
+ // (通道告警对交易员改由 dispatch-mode 里的 channel 兜底, 见 wsc 的注释)
+ (isAdmin.value ? loadWs() : Promise.resolve()),
+ loadStrategies(), loadOpLog(), loadOpenScan(), loadUpstream(),
loadMacro()]);
await loadNames();
lastRefresh.value = new Date().toTimeString().slice(0, 8);
@@ -2816,8 +2856,11 @@ createApp({
}
async function decide(row, decision) {
const d = await call('post', '/api/proposals/' + row.proposal_id + '/decide', { decision });
- if (d.ok) ElementPlus.ElMessage.success(decision === 'ACCEPTED'
- ? ('已采纳, 指令 ' + d.instruction_id) : '已驳回');
+ if (d.ok) {
+ ElementPlus.ElMessage.success(decision === 'ACCEPTED'
+ ? ('已采纳, 指令 ' + d.instruction_id) : '已驳回');
+ if (d.warning) ElementPlus.ElMessage.warning(d.warning); // 一次性纪律失效等提示 (2026-08-28)
+ }
else ElementPlus.ElMessage.error(d.error);
await Promise.all([loadProposals(), loadInstructions(), loadLedger()]);
}
@@ -2973,7 +3016,7 @@ createApp({
if (refreshTimer) return;
// 固定 30s 自动刷新(开关开时)。不再按 marketOpen 分档 —— 交易日历对模拟/未来日期可能
// 误判成非交易日, 那样会把刷新降到 5 分钟、体感像"卡住不刷"。休市时多刷几次也无害。
- refreshTimer = setInterval(() => { if (autoRefresh.value) refreshLive(); }, 30000);
+ refreshTimer = setInterval(() => { if (autoRefresh.value && authed.value) refreshLive(); }, 30000);
}
onUnmounted(() => { if (refreshTimer) clearInterval(refreshTimer); });
onMounted(async () => {
@@ -3154,7 +3197,7 @@ createApp({
});
return { tab, loading, err, health, ov, params, catalog, commands, plans, plansOf,
- positions, lots, lotsOf, instructions, ledger, proposals, report, reportDrawer,
+ positions, lots, lotsOf, instructions, ledger, ledgerToday, proposals, report, reportDrawer,
opsDrawer, opsResult, opsLoading, issuing, form, curSpec, dirtyCount, dm,
money, pct, groupLabel, fieldLabel, stTag, cuTag, canCancel, progPct, insPct,
scaleGap,
diff --git a/app/ws/runner.py b/app/ws/runner.py
index a4d815d..9f1424e 100644
--- a/app/ws/runner.py
+++ b/app/ws/runner.py
@@ -173,7 +173,6 @@ class WsRunner:
logger.error("通道表自检失败 (每 60 秒重试): %s", msg)
return False
- self._db_ready = True
try:
st = await _db(qmt_repo.get_state)
stored = int(st.get("last_seq") or 0)
@@ -189,8 +188,13 @@ class WsRunner:
# 只回 ack{duplicate:true} 带当前状态。幂等键就是为这一刻准备的。
logger.warning("有 %s 张委托上次卡在 SENDING, 已退回队列重发 (幂等键兜底)", n)
except Exception as e:
- logger.error("启动自检失败, 水位按 0 起算 (下轮重连会重试): %s", _brief_err(e))
+ # _db_ready 必须保持 False —— 原来在恢复块之前就置 True, 这里失败后
+ # 「下轮重连会重试」是句空话 (232 行按 _db_ready 短路), 进程会带着
+ # last_seq=0 上线, 把真缺口误判成冷启动 (2026-08-28 审查修)。
+ self._db_ready = False
+ logger.error("启动自检失败, 水位恢复未完成 (60 秒后重试): %s", _brief_err(e))
return False
+ self._db_ready = True
return True
def _operable(self) -> bool:
@@ -209,7 +213,9 @@ class WsRunner:
try:
await self._refresh_params()
if self._operable():
- await _db(qmt_repo.beat, self._stat)
+ # 传浅拷贝快照: qmt_repo.beat 在工作线程里 json.dumps, 而 reader 协程
+ # 同时在往 _stat 插新键 —— 迭代中改字典会炸掉当轮心跳 (2026-08-28 修)
+ await _db(qmt_repo.beat, dict(self._stat))
except Exception as e:
logger.warning("心跳写入失败 (库不可用?): %s", _brief_err(e))
with contextlib.suppress(asyncio.TimeoutError):
@@ -236,6 +242,7 @@ class WsRunner:
continue
url = self._p("url", settings.PMS_QMT_WS_URL)
+ t_started = time.monotonic()
try:
await _db(qmt_repo.set_conn, "CONNECTING")
await self._session(url, seed, peer)
@@ -243,6 +250,11 @@ class WsRunner:
except asyncio.CancelledError:
raise
except Exception as e:
+ # 退避按**连续失败**计 (2026-08-28 审查修): 一段健康跑了 5 分钟以上的
+ # 会话断掉, 说明不是连不上, 计数归零从最短退避重来 —— 原来 attempt 按
+ # 进程生命周期只增不减, 常驻几天后每次普通抖动都要干等满 30 秒。
+ if time.monotonic() - t_started > 300:
+ attempt = 0
self._stat["reconnects"] += 1
# 断开原因同样存一份: set_conn 的 last_error 会在下次连上时被覆盖,
# 断了又连的场景里那条线索活不过 30 秒, 而这正是最需要它的场景。
@@ -304,6 +316,13 @@ class WsRunner:
for t in tasks + [stopper]:
t.cancel()
await asyncio.gather(*tasks, stopper, return_exceptions=True)
+ # 停机路径的最后一次 ack 要在**这里**发 (2026-08-28 审查修): _shutdown 里
+ # 那两个 `if self._ws is not None` 块永远走不到 —— 到那时本 finally 早已把
+ # _ws 置 None、连接也关了, 是死代码。趁连接还在把待 ack 的水位刷出去,
+ # QMT 侧就不会在每次停机后多留一段已收未确认的消息。
+ if self._stop.is_set() and self._ws is not None:
+ with contextlib.suppress(Exception):
+ await asyncio.wait_for(self._flush_ack(), timeout=3)
self._ws = None
async def _handshake(self, seed: str, peer: str):
@@ -341,7 +360,11 @@ class WsRunner:
logger.error(self._warn)
else:
self._warn = ""
- await _db(qmt_repo.set_conn, "CONNECTING", server_seq=server_seq, resync=resync)
+ # resync 标记**只置位, 不在这里清零** (2026-08-28 审查修): 传 None 表示不动该列。
+ # 原来每次干净握手都把 resync 写回 0 —— 缺口置位后几分钟一次正常重连, "须人工
+ # 全量对账"的持久标记就没了, 丢失区间的成交无人补账。清零只走页面的人工接口。
+ await _db(qmt_repo.set_conn, "CONNECTING", server_seq=server_seq,
+ resync=(True if resync else None))
if resync:
# §5.1 / §6.2: 对端补不齐我们要的区间 (日志已滚动)。此时**不能**装作没事 ——
# 中间那段成交我们永远拿不到了, 必须走全量快照对账, 对不齐就停一切自主动作。
@@ -559,10 +582,20 @@ class WsRunner:
iid = pl.get("instruction_id") or env.get("corr_id")
try:
if type_ == wsc.T_ACK:
- n = await _db(qmt_repo.update_order, iid,
- status=str(pl.get("status") or wsc.ST_ACCEPTED).upper(),
- broker_order_id=pl.get("broker_order_id"))
- self._note_orphan(n, type_, iid)
+ # 终态保护 (2026-08-28 审查修, 与 _on_order_update 同一条纪律):
+ # _boot 重发换来的 ack{duplicate:true} 是新 seq 的新消息, 不带 status 时
+ # 这里的 ACCEPTED 兜底会把已 FILLED 的委托改回在途 —— 幻影单从此常驻
+ # queue_depth 与撤单清单。已终态的委托, ack 一律只记日志不动状态。
+ row = await _db(qmt_repo.get_order, iid) or {}
+ cur = str(row.get("status") or "").upper()
+ if cur and wsc.is_final(cur):
+ logger.info("[ack] %s 已是终态 %s, ack(duplicate=%s) 不回退状态",
+ iid, cur, bool(pl.get("duplicate")))
+ else:
+ n = await _db(qmt_repo.update_order, iid,
+ status=str(pl.get("status") or wsc.ST_ACCEPTED).upper(),
+ broker_order_id=pl.get("broker_order_id"))
+ self._note_orphan(n, type_, iid)
if pl.get("duplicate"):
logger.info("[ack] %s 幂等命中 (对端已受理过), 当前状态 %s",
iid, pl.get("status"))
diff --git a/scripts/backtest_churn.py b/scripts/backtest_churn.py
index 2fae1c5..2d94010 100644
--- a/scripts/backtest_churn.py
+++ b/scripts/backtest_churn.py
@@ -40,6 +40,19 @@ from datetime import datetime, timedelta
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from app.db.session import fetch_all, DBUnavailable # noqa: E402 只读单表, 过单表守卫
+from app.core import tradedays as _td # noqa: E402 真交易日计数
+
+
+def tdays_between(d0, d1) -> int:
+ """(d0, d1] 内的交易日数 (同日买卖 = 0)。护栏档位 N 按交易日定义, 持有天数必须同尺 ——
+ 自然日近似会把周五买周一卖 (1 个交易日) 算成 3 天, N=1 档直接漏掉 (2026-08-28 修)。"""
+ n, cur, guard = 0, d0, 0
+ while cur < d1 and guard < 400:
+ cur += timedelta(days=1)
+ guard += 1
+ if _td.is_trade_day(cur):
+ n += 1
+ return n
# ------------------------------------------------------------------ 可调常量
MIN_SAMPLE = 5 # 少于这个数只报数不下结论
@@ -219,7 +232,7 @@ def trade_cost(notional: float, is_sell: bool) -> float:
# ------------------------------------------------------------------ 来回配对
def pair_round_trips(rows_by_code: dict) -> list:
"""同一只票: 把"一次买"和其后"第一次卖"配成一个来回(粗配, 不做逐笔 FIFO)。
- 返回每个来回: 入场/出场时间价、持有交易日近似(自然日近似)、来回毛收益、是否决策系统驱动。
+ 返回每个来回: 入场/出场时间价、持有交易日(按交易日历真算, 2026-08-28 起)、来回毛收益、是否决策系统驱动。
说明: 这是方向与量级的体检, 不是会计账; 精算逐笔在 pms_lot, 要精算另说。"""
trips = []
for code, rows in rows_by_code.items():
@@ -234,12 +247,15 @@ def pair_round_trips(rows_by_code: dict) -> list:
bt, st = b["decided_at"], s["decided_at"]
bt = bt if isinstance(bt, datetime) else datetime.fromisoformat(str(bt))
st = st if isinstance(st, datetime) else datetime.fromisoformat(str(st))
- hold_days = (st.date() - bt.date()).days
+ hold_days = tdays_between(bt.date(), st.date()) # 真交易日 (2026-08-28)
gross = _pct(b["price_at"], s["price_at"])
+ hn_b, hn_s = (b.get("hn") or {}), (s.get("hn") or {})
trips.append({
"ts_code": code, "buy_at": bt, "sell_at": st,
"buy_price": _f(b["price_at"]), "sell_price": _f(s["price_at"]),
"hold_days": hold_days, "gross_ret": gross,
+ "buy_amt": _f(hn_b.get("amount")) or None,
+ "sell_amt": _f(hn_s.get("amount")) or None,
"signal_driven": tag["signal_driven"], "sell_reason": s.get("reason"),
"cohort": tag.get("cohort"),
})
@@ -334,15 +350,17 @@ def section_cost(trips):
if not sig:
print(" 无样本。")
return
- # 无法拿到每笔真实金额时, 以"单位名义1"估相对成本率; 有 hard_numbers.amount 时用真实额
- total_rate = 0.0
- for t in sig:
- rt_cost_rate = (COMMISSION_RATE * 2 + STAMP_RATE + TRANSFER_RATE * 2
- + SLIPPAGE_BPS / 10000.0 * 2)
- total_rate += rt_cost_rate
print(f" 每个来回的往返摩擦约 {(_rt_cost_rate()):.3%} (佣金双边+印花+过户+滑点双边)")
- print(f" {len(sig)} 组来回累计摩擦 ≈ 名义规模的 {total_rate:.2%} "
- f"(即平均毛收益要先跑赢这条线才算真挣到)")
+ # 有真实金额 (账本 hard_numbers.amount) 的来回按真实额算, 含每边最低 5 元佣金;
+ # 没有的只报口径, 不再打印"N×费率"那个等名义假设的合计 (它不是钱, 是比率和, 2026-08-28 修)
+ real = [t for t in sig if t.get("buy_amt") and t.get("sell_amt")]
+ if real:
+ yuan = sum(trade_cost(t["buy_amt"], False) + trade_cost(t["sell_amt"], True)
+ for t in real)
+ print(f" 其中 {len(real)} 组带真实金额: 累计摩擦 {yuan:,.0f} 元 (含每边最低 5 元佣金)")
+ if len(real) < len(sig):
+ print(f" 其余 {len(sig) - len(real)} 组账本未记金额, 只能按费率口径读: "
+ f"每组来回的毛收益要先跑赢 {_rt_cost_rate():.3%} 才算真挣到")
def _rt_cost_rate():
@@ -359,13 +377,15 @@ def _guard_sweep(label, affected):
for t in affected:
if t["hold_days"] is None or t["hold_days"] > N:
continue # 护栏只管"刚建仓 N 日内"的卖出
- if (t["gross_ret"] or 0) <= HARD_RISK_DROP:
+ if t["gross_ret"] is None:
+ continue # 卖价缺失: 不能当 0 收益混进样本 (2026-08-28)
+ if t["gross_ret"] <= HARD_RISK_DROP:
continue # 硬止损放行, 不受护栏拦
seq = _fwd_closes_after(t["ts_code"], t["sell_at"], n_max=GUARD_HOLD_TO)
if len(seq) < GUARD_HOLD_TO:
continue
# 实际(卖了): 拿到 gross_ret, 并付了一次卖出摩擦; 之后空仓 = 0
- actual = (t["gross_ret"] or 0) - _rt_cost_rate() / 2
+ actual = t["gross_ret"] - _rt_cost_rate() / 2
# 反事实(没卖, 持有到 T+H 再看): 用 T+H 相对入场的收益, 只付了买入侧摩擦
held = _pct(t["buy_price"], seq[GUARD_HOLD_TO - 1])
if held is None:
diff --git a/scripts/backtest_entry.py b/scripts/backtest_entry.py
index 3e54a1b..489b5bb 100644
--- a/scripts/backtest_entry.py
+++ b/scripts/backtest_entry.py
@@ -44,6 +44,18 @@ from datetime import datetime, timedelta
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from app.db.session import fetch_all, DBUnavailable # noqa: E402 只读单表, 过单表守卫
+from app.core import tradedays as _td # noqa: E402 真交易日计数
+
+
+def tdays_between(d0, d1) -> int:
+ """(d0, d1] 内的交易日数 (同日买卖 = 0) —— 与 backtest_churn 同口径 (2026-08-28)。"""
+ n, cur, guard = 0, d0, 0
+ while cur < d1 and guard < 400:
+ cur += timedelta(days=1)
+ guard += 1
+ if _td.is_trade_day(cur):
+ n += 1
+ return n
# ------------------------------------------------------------------ 可调常量
MIN_SAMPLE = 5 # 少于这个数只报数不下结论
@@ -65,7 +77,9 @@ TAKE_PROFIT_DOM = "take_profit" # 卖出 dominant_signal 是这个 = 止盈
QUICK_DAYS = 10
FWD_HORIZONS = (1, 5, 20) # 入场后前向看这几个交易日
-# 入场过滤规则的阈值 (⑤ 反事实扫的候选闸门, 都只用买入日及之前的信息, 无未来函数)
+# 入场过滤规则的阈值 (⑤ 反事实扫的候选闸门)。口径 (2026-08-28 理清):
+# 逆势买 / 放量阴线 只用买入**前一交易日**及更早的数据 —— 实盘上午就拿得到, 无前视;
+# 追高 用买入日全天高低 —— 收盘才知道, 属事后画像口径, 它的反事实读数是乐观上限。
CHASE_TOP_FRAC = 0.70 # 追高: 买价落在当日振幅的顶部这个比例以上 (1.0=买在最高)
DOWNTREND_5D = -0.03 # 逆势: 买入日相对前5个交易日收盘已跌超这个幅度
REDVOL_VOL_MULT = 1.5 # 放量阴线: 买入日是阴线且量能≥前5日均量的这个倍数
@@ -224,10 +238,12 @@ def forward_returns(ts_code: str, dt: datetime, base_price: float) -> dict:
def entry_context(ts_code: str, buy_dt: datetime, buy_price: float) -> dict:
"""重建买入时点的价格情形。返回 {chase, day_chg, mom5, red_vol}, 缺数据的项为 None。
- chase 追高度 = (买价-当日最低)/(当日最高-当日最低), 1.0=买在最高, 越高越追。
- day_chg 买入日涨跌 = 当日 percent(优先) 或 收盘/前收-1, 正=买在红盘。
- mom5 前5日动量 = 买入日收盘/前5个交易日收盘-1, 负=买在已经下跌的票上(逆势)。
- red_vol 放量阴线 = 买入日收盘<开盘 且 量能≥前5日均量×倍数 (True/False/None)。
+ chase 追高度 = (买价-当日最低)/(当日最高-当日最低), 1.0=买在最高。**事后口径**
+ (全天高低要收盘才定), 只作画像与乐观上限, 不是实盘可复刻的闸门。
+ day_chg 买入日全天涨跌 —— 同为事后口径。
+ mom5 前一交易日收盘 / 再前5个交易日收盘 - 1 —— **盘中可得** (2026-08-28 改),
+ 负=买在已经下跌的票上(逆势)。
+ red_vol 前一交易日是阴线且量能≥更早5日均量×倍数 —— **盘中可得** (2026-08-28 改)。
"""
bars = _bars(ts_code, buy_dt)
ctx = {"chase": None, "day_chg": None, "mom5": None, "red_vol": None}
@@ -236,28 +252,31 @@ def entry_context(ts_code: str, buy_dt: datetime, buy_price: float) -> dict:
bd = buy_dt.date()
bar = bars.get(bd)
if bar is None:
- # 账本时间戳那天没有行情(极少见), 用其后第一根近似
- later = sorted(d for d in bars if d >= bd)
- if not later:
- return ctx
- bd = later[0]
- bar = bars[bd]
+ # 买入日没有行情 (极少见): **直接跳过该样本** —— 原来用"其后第一根"近似,
+ # 那是买入之后那天的数据, 与"无未来函数"的承诺相悖 (2026-08-28 修)
+ return ctx
hi, lo = bar["high"], bar["low"]
if buy_price > 0 and hi > lo:
+ # 追高度用当日全天高低点: 这是**事后口径** (买入时刻只知道到那一刻的高低),
+ # 结果是过滤效果的乐观上限, 读数时要打这个折 (2026-08-28 说明)
ctx["chase"] = max(0.0, min(1.0, (buy_price - lo) / (hi - lo)))
if bar["pct"]:
- ctx["day_chg"] = bar["pct"] / 100.0
+ ctx["day_chg"] = bar["pct"] / 100.0 # 同为事后口径 (全天涨跌)
elif bar["pre_close"] > 0:
ctx["day_chg"] = bar["close"] / bar["pre_close"] - 1.0
prior = sorted((d, b) for d, b in bars.items() if d < bd)
- if len(prior) >= 5:
- c5 = prior[-5][1]["close"]
- if c5 > 0:
- ctx["mom5"] = bar["close"] / c5 - 1.0
- vols = [b["vol"] for _, b in prior[-5:] if b["vol"] > 0]
- if vols and bar["vol"] > 0:
+ # mom5 与 red_vol 改用**买入前一交易日及更早**的数据 (2026-08-28 修): 原来用买入日
+ # 收盘和全天量 —— 那要收盘才知道, 实盘过滤器在上午拿不到。改后这两条规则是真正
+ # 可实现的 (前一日收盘 / 前一日阴线放量), 反事实结果不再靠日内前视撑着。
+ if len(prior) >= 6:
+ y = prior[-1][1] # 买入前一交易日
+ c5 = prior[-6][1]["close"] # 前一日再往前 5 个交易日
+ if c5 > 0 and y["close"] > 0:
+ ctx["mom5"] = y["close"] / c5 - 1.0
+ vols = [b["vol"] for _, b in prior[-6:-1] if b["vol"] > 0]
+ if vols and y["vol"] > 0:
avgv = sum(vols) / len(vols)
- ctx["red_vol"] = (bar["close"] < bar["open"]) and (bar["vol"] >= REDVOL_VOL_MULT * avgv)
+ ctx["red_vol"] = (y["close"] < y["open"]) and (y["vol"] >= REDVOL_VOL_MULT * avgv)
return ctx
@@ -290,7 +309,7 @@ def pair_round_trips(rows_by_code: dict) -> list:
bt, st = b["decided_at"], s["decided_at"]
bt = bt if isinstance(bt, datetime) else datetime.fromisoformat(str(bt))
st = st if isinstance(st, datetime) else datetime.fromisoformat(str(st))
- hold_days = (st.date() - bt.date()).days
+ hold_days = tdays_between(bt.date(), st.date()) # 真交易日 (2026-08-28)
gross = _pct(b["price_at"], s["price_at"])
trips.append({
"ts_code": code, "buy_at": bt, "sell_at": st,
@@ -407,11 +426,12 @@ def _filters():
"""候选入场闸门: 每条是 (名称, 说明, 判定函数(trip,ctx)->bool 命中即"该拦")。
只用买入日及之前的信息, 无未来函数。"""
return [
- ("追高", f"买价落在当日振幅顶部{(1-CHASE_TOP_FRAC):.0%}以内(追高度≥{CHASE_TOP_FRAC})",
+ ("追高[事后口径]", f"买价落在当日振幅顶部{(1-CHASE_TOP_FRAC):.0%}以内(追高度≥{CHASE_TOP_FRAC}; "
+ f"全天高低收盘才定, 此条是乐观上限)",
lambda t, c: c["chase"] is not None and c["chase"] >= CHASE_TOP_FRAC),
- ("逆势买", f"买入日已较前5日跌超{abs(DOWNTREND_5D):.0%}(前5日动量≤{DOWNTREND_5D:.0%})",
+ ("逆势买", f"前一交易日已较其前5日跌超{abs(DOWNTREND_5D):.0%}(盘中可得)",
lambda t, c: c["mom5"] is not None and c["mom5"] <= DOWNTREND_5D),
- ("放量阴线", f"买入日是阴线且量能≥前5日均量×{REDVOL_VOL_MULT}",
+ ("放量阴线", f"前一交易日是阴线且量能≥更早5日均量×{REDVOL_VOL_MULT}(盘中可得)",
lambda t, c: c["red_vol"] is True),
]
diff --git a/scripts/check_db.py b/scripts/check_db.py
index becdd80..86c9521 100644
--- a/scripts/check_db.py
+++ b/scripts/check_db.py
@@ -79,8 +79,16 @@ def check_ws_channel():
else:
line("OK", f"对端公钥已配置 ({peer[:12]}...)")
- if param_store.get_bool("PMS_QMT_SIGN_SEED_HEX") or param_store.get("PMS_QMT_SIGN_SEED_HEX"):
- fail("私钥出现在 ParamStore 可读路径 —— 违反协议 §10.1.1, 请检查 SECRET_KEYS")
+ # 自证必须**直查参数表** (2026-08-28 修): param_store.get 对 SECRET_KEYS 无条件回
+ # 空串 (那正是堵读取的闸), 经它去验"表里有没有私钥"是永真检查 —— 谁把 seed 写进
+ # 参数表, 这里照样报 OK。
+ try:
+ from app.repo import pms_repo as _pr
+ _raw = _pr.all_params()
+ if any(_raw.get(k) for k in param_store.SECRET_KEYS):
+ fail("私钥出现在参数表 pms_runtime_param —— 违反协议 §10.1.1, 请删除该行并轮换密钥")
+ except Exception as _e:
+ warn(f"密钥落库自证跳过 (参数表读不了): {type(_e).__name__}: {_e}")
try:
ch = dispatcher.channel_status()
diff --git a/scripts/gen_keys.py b/scripts/gen_keys.py
index 0f86b29..43b46ff 100644
--- a/scripts/gen_keys.py
+++ b/scripts/gen_keys.py
@@ -177,12 +177,19 @@ def check():
print(" 那一行 base64 粘进 .env 即可 (整段 PEM 有换行, 放不进 .env)")
print("\n[4] 密钥不落库自证 (协议 §10.1.1)")
+ # 直查参数表 (2026-08-28 修): param_store.get 对密钥键恒回空串, 经它验证是永真检查
from app.services import param_store
- leaked = [k for k in param_store.SECRET_KEYS if param_store.get(k)]
+ try:
+ from app.repo import pms_repo as _pr
+ _raw = _pr.all_params()
+ leaked = [k for k in param_store.SECRET_KEYS if _raw.get(k)]
+ except Exception as _e:
+ leaked = []
+ print(f" (参数表读不了, 本项跳过: {type(_e).__name__}: {_e})")
if leaked:
- bad(f"密钥可经 ParamStore 读出: {leaked}")
+ bad(f"密钥出现在参数表 pms_runtime_param: {leaked} —— 请删除该行并轮换密钥")
else:
- ok("ParamStore 读不到密钥, 页面参数列表里也不会出现")
+ ok("参数表里没有密钥行; ParamStore 读取路径也已封死")
print("\n" + "-" * 66)
if failed:
diff --git a/scripts/migrate_archived_at.py b/scripts/migrate_archived_at.py
index d250ebe..67df6b3 100644
--- a/scripts/migrate_archived_at.py
+++ b/scripts/migrate_archived_at.py
@@ -96,8 +96,11 @@ def main():
print(f"以下 {len(failed)} 张表加列失败, 完整语句如下 —— 可直接拿到物理库 (my_quant_db) 执行:")
for t, s, err in failed:
print(f"\n### {t} ({err.splitlines()[0]})\n{s};")
- if failed or miss:
- print(f"\nFAILED: {len(failed) + len(miss)} 张表仍未就绪")
+ if failed or miss or unknown:
+ # unknown (查列失败, 多为表不存在) 也算未就绪 (2026-08-28 修): 原来 todo 非空时
+ # 它被忘掉, 四张表迁了三张也报 ALL OK —— 缺列的表要等页面报 SQL 错才暴露
+ print(f"\nFAILED: {len(failed) + len(miss) + len(unknown)} 张表仍未就绪"
+ + (f" (含查列失败 {len(unknown)} 张: {', '.join(unknown)})" if unknown else ""))
sys.exit(1)
print(f"ALL OK: {len(todo)} 张表已加 {COL} 列。页面的「移除 / 显示已完成」现在可用。")
diff --git a/scripts/probe_strategy_signals.py b/scripts/probe_strategy_signals.py
index f101b74..2dfb620 100644
--- a/scripts/probe_strategy_signals.py
+++ b/scripts/probe_strategy_signals.py
@@ -284,6 +284,11 @@ def main():
print("\n【四】候选池合格票与吸筹的重合 (拍板②「资格不放宽」的实测依据)")
try:
from app.services import command_service, plan_feed
+ # 只读承诺自证 (2026-08-28 修): plan_feed.get_plan 默认会落一行 pms_plan_snapshot
+ # (还可能 upsert 行业映射) —— 探测脚本提前落库会吞掉当天正式链路的榜单变化提示
+ # (probe_plan_api 早有注释点过这个坑)。本进程内把两个写函数替换成空操作。
+ plan_feed._snapshot_quiet = lambda *a, **kw: {"skipped": "probe 只读, 不落快照"}
+ plan_feed._sync_themes_quiet = lambda *a, **kw: {"skipped": "probe 只读, 不灌映射"}
try:
black = command_service.blacklist()
except Exception:
diff --git a/scripts/report_strategy_score.py b/scripts/report_strategy_score.py
index 9e3cdd0..3ed3784 100644
--- a/scripts/report_strategy_score.py
+++ b/scripts/report_strategy_score.py
@@ -156,7 +156,17 @@ def main():
continue
n_traded += 1
buy_avg = b["amt"] / b["qty"] if b["qty"] else 0.0
- realized = sl["amt"] - sl["qty"] * buy_avg if sl["qty"] else 0.0
+ if sl["qty"] and not b["qty"]:
+ # 只有卖出没有买入 (买入腿痕迹缺失/被对账冲销): 按零成本算会把整笔卖出额
+ # 报成利润, 判分虚高误导阈值决策 (2026-08-28 修) —— 这条不算差价, 单独点名。
+ print(f" {code} {sid} [{s.get('status')}]: 卖 {sl['qty']} 股/{sl['amt']:,.0f} 元, "
+ f"但**查无买入腿成交** —— 差价无法计算, 不计入合计, 请核对该策略的成交归属")
+ continue
+ sell_q = min(sl["qty"], b["qty"])
+ if sl["qty"] > b["qty"]:
+ print(f" {code} {sid}: 卖出 {sl['qty']} 股 > 买入 {b['qty']} 股, "
+ f"超出部分 {sl['qty'] - b['qty']} 股不属本策略买入, 差价只按配对部分算")
+ realized = (sl["amt"] * (sell_q / sl["qty"]) - sell_q * buy_avg) if sl["qty"] else 0.0
net_qty = b["qty"] - sl["qty"]
cur = _f(px.get(code))
floating = net_qty * (cur - buy_avg) if (net_qty > 0 and cur > 0 and buy_avg > 0) else None
diff --git a/scripts/reset_ledger.py b/scripts/reset_ledger.py
index fe9f03a..f9f64e1 100644
--- a/scripts/reset_ledger.py
+++ b/scripts/reset_ledger.py
@@ -205,8 +205,10 @@ def main():
if args.reset_ws:
try:
+ from datetime import datetime as _dt
execute("UPDATE pms_ws_state SET last_seq = 0, acked_seq = 0, server_seq = 0, "
- "resync_flag = 0, conn_state = 'INIT', updated_at = NOW() WHERE id = 1")
+ "resync_flag = 0, conn_state = 'INIT', updated_at = :ts WHERE id = 1",
+ {"ts": _dt.now()}) # 库端 NOW() 是 UTC, 会写出倒退 8 小时的时间戳 (2026-08-28)
print(" OK pms_ws_state 水位归零 —— 重连后对端会从 seq 1 全量补发")
except Exception as e:
print(f" WARN pms_ws_state 归零失败 (表可能还没建): {type(e).__name__}: {e}")
diff --git a/scripts/run_tests.py b/scripts/run_tests.py
index 5ea5117..1784a0f 100644
--- a/scripts/run_tests.py
+++ b/scripts/run_tests.py
@@ -34,8 +34,12 @@
判分脚本聚合与对照分组 (41 例)
test_batch18_units.py 登录与权限: 角色判定/会话票签验/bshop 返回解析/
接口鉴权(维护类归管理员) (16 例)
- test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (58 例)
- 共 560 例
+ test_batch19_units.py 2026-08-28 审查修复回归: 科创板最小申报统一口径/
+ 取整与部分卖/同轮买卖互斥/信号百分制契约/除权核销
+ 缩放/日历按年降级/网格中枢与止盈闩锁/买入暂停
+ 按来源分记/宏观失败路径保留留痕 (17 例)
+ test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (66 例)
+ 共 585 例
任一子集失败即整体失败 (退出码 1)。
"""
import os
@@ -50,7 +54,7 @@ SUITES = ["test_core_units.py", "test_batch2_units.py", "test_batch3_units.py",
"test_batch10_units.py", "test_batch11_units.py", "test_batch12_units.py",
"test_batch13_units.py", "test_batch14_units.py", "test_batch15_units.py",
"test_batch16_units.py", "test_batch17_units.py",
- "test_batch18_units.py", "test_wiring.py"]
+ "test_batch18_units.py", "test_batch19_units.py", "test_wiring.py"]
def main():
diff --git a/scripts/test_batch10_units.py b/scripts/test_batch10_units.py
index 7324841..07ebb8a 100644
--- a/scripts/test_batch10_units.py
+++ b/scripts/test_batch10_units.py
@@ -456,8 +456,9 @@ def _():
from app.services import signal_service
from app.core import signal_rules as sr
install_fakes(prices={"600000.SH": 10.0})
- seen, out = set(), {"ignored": 0, "recorded": 0, "exits": [], "proposals": [],
- "errors": []}
+ # seen 自 2026-08-28 起是**有序 dict** (裁剪时裁最旧的, 不再按字典序), 用法同集合
+ seen, out = {}, {"ignored": 0, "recorded": 0, "exits": [], "proposals": [],
+ "errors": []}
sig = {"msg_id": "M1", "ts_code": "600000.SH", "action": "SELL", "confidence": 0.95,
"source": "test", "reason": "风控"}
view = {"positions": [{"ts_code": "600000.SH", "total_qty": 1000, "avail_qty": 1000,
@@ -488,8 +489,8 @@ def _():
from app.services import signal_service
from app.core import signal_rules as sr
install_fakes(prices={"600000.SH": 10.0})
- seen, out = set(), {"ignored": 0, "recorded": 0, "exits": [], "proposals": [],
- "errors": []}
+ seen, out = {}, {"ignored": 0, "recorded": 0, "exits": [], "proposals": [],
+ "errors": []}
sig = {"msg_id": "M1", "ts_code": "600000.SH", "action": "SELL", "confidence": 0.95,
"source": "test", "reason": "风控"}
view = {"positions": [{"ts_code": "600000.SH", "total_qty": 1000, "avail_qty": 1000,
@@ -621,7 +622,6 @@ def _():
@case("[J3] 自主提议给规则闸的 day 必须是真行情, 不是拿 price 拼出来的空壳")
def _():
- import inspect
from app.services import proposal_service
# 曾经是 {"vwap": price, "day_chg_from_open": None} —— 当日涨幅恒 None, 于是
# "不追高(涨幅)"这一项对所有自主买入从来没有真正跑过。这里直接验行为: 造一只
diff --git a/scripts/test_batch11_units.py b/scripts/test_batch11_units.py
index 269d4aa..72ca176 100644
--- a/scripts/test_batch11_units.py
+++ b/scripts/test_batch11_units.py
@@ -261,6 +261,10 @@ def _():
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)
+ # 2026-08-28 补: 第二次咨询同样失败, r3 必须仍是"退实现B"的有效决策且重新挂上冷却
+ assert r3 and r3.get("action"), r3
+ assert "B" in str(r3.get("source") or ""), r3
+ assert prog["exec_advice"]["fail_until_min"] > 10 * 60 + 40 - 1, prog["exec_advice"]
@case("[C7] 对端回 UNAVAILABLE (给不出结论) → 退B + 冷却, 不当 FIRE 也不当 WAIT")
diff --git a/scripts/test_batch17_units.py b/scripts/test_batch17_units.py
index 6fccf41..56485b4 100644
--- a/scripts/test_batch17_units.py
+++ b/scripts/test_batch17_units.py
@@ -483,6 +483,9 @@ def _():
out2 = _scan_stubbed(rec2, held=held, accum=accum, heat={}, dry_run=False)
assert len(rec2.attach_calls) == 2, rec2.attach_calls
assert any("名额已满" in x.get("reason", "") for x in rec2.ledger), rec2.ledger
+ # 2026-08-28 补: 真跑那半段的返回值也要锁住 (原来 out2 算了没断言)
+ assert out2["ok"], out2
+ assert len(out2["blocked"]) == 1 and "名额已满" in out2["blocked"][0]["why"], out2["blocked"]
@case("[冒烟] 排除项: 冻结票与热度停更都不挂; 词表外定性浮到 unknown_states")
diff --git a/scripts/test_batch19_units.py b/scripts/test_batch19_units.py
new file mode 100644
index 0000000..80ac151
--- /dev/null
+++ b/scripts/test_batch19_units.py
@@ -0,0 +1,398 @@
+# -*- coding: utf-8 -*-
+"""
+第十九批: 2026-08-28 全库审查修复的纯逻辑回归 —— 不联网、不连库
+=====================================================================
+这一批钉住的都是当次审查改掉的真伤, 每条用例头上写清"原来错在哪":
+ 1. 科创板最小申报统一口径 (sizer.lot_of / lot_qty 浮点容差 / split_batches 可行性线);
+ 2. planner 取整与部分卖修正 (ceil_lot 小数截断 / _min_sell 200 股线);
+ 3. 规则闸科创板买卖申报校验 (买 <200 拒单 / 部分卖 <200 拒、全清放行);
+ 4. 动作引擎同轮买卖互斥 (TRIM 触发时买入侧让路, 不再自动对倒);
+ 5. 信号口径 (_norm_conf_pct 百分制契约: 1 = 1% 不是 100%; digest 科创板卖量修正);
+ 6. 除权调整已核销批次同步缩放 (recon.apply_ex_right);
+ 7. 交易日历按年探测降级 (chinesecalendar 装了但没有当年数据);
+ 8. 网格只买中枢下方 / 跟踪止盈部分卖一次性闩锁 (strategy_runner);
+ 9. 策略买入暂停按来源分记 (strategy_service, accum 与 signal 互不误伤);
+ 10. 宏观失败路径保留当日动作留痕 (macro_service._upsert_unavailable)。
+运行: python scripts/test_batch19_units.py
+"""
+import os
+import sys
+import traceback
+from datetime import date
+
+sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
+
+RESULTS = []
+
+
+def case(name):
+ def deco(fn):
+ RESULTS.append((name, fn))
+ return fn
+ return deco
+
+
+# ================================================================
+# [A] sizer: 最小申报数量的唯一出处
+# ================================================================
+@case("[A1] lot_of: 688/689 开头 200 股, 其余 100; 空值与 '含688不在开头' 不误判")
+def _():
+ from app.core.sizer import lot_of
+ assert lot_of("688802.SH") == 200 and lot_of("689009.SH") == 200
+ assert lot_of("600000.SH") == 100 and lot_of("000001.SZ") == 100
+ assert lot_of("300688.SZ") == 100 # 688 在中间不算科创板
+ assert lot_of(None) == 100 and lot_of("") == 100
+
+
+@case("[A2] lot_qty 浮点容差: 407 元买 4.07 元的票**恰好一手**, 不许被浮点误差算成零手")
+def _():
+ from app.core.sizer import lot_qty
+ # 407/4.07 在浮点里是 99.999…, 原实现直接 int() 会把"恰好买得起一手"算成 0 手
+ assert lot_qty(407, 4.07) == 100, lot_qty(407, 4.07)
+ assert lot_qty(814, 4.07) == 200
+ assert lot_qty(1500, 10.0) == 100 and lot_qty(999, 10.0) == 0
+ assert lot_qty(0, 10.0) == 0 and lot_qty(1000, 0) == 0 and lot_qty(None, 5) == 0
+ assert lot_qty(4000, 10.0, 200) == 400 # lot=200 时按 200 的整数倍
+
+
+@case("[A3] split_batches min_lot=200: 可行性线抬到科创板 200 股, 降档与失败话术都说清")
+def _():
+ from app.core.sizer import split_batches
+ # 50/25/25 里 25% 批只有 100 股 (<200) → 自动降档到 60/40 (300/200 股, 都合法)
+ r = split_batches(6000, 10.0, min_lot=200)
+ assert r["ok"] and r["scheme"] == (0.6, 0.4), r
+ assert [b["qty"] for b in r["batches"]] == [300, 200], r
+ # 全部阶梯都买不足 200 股 → 失败, 原因里点名科创板
+ r2 = split_batches(1500, 10.0, min_lot=200)
+ assert not r2["ok"] and "科创板最少 200 股" in r2["reason"], r2
+ # 主板行为一字不变
+ r3 = split_batches(6000, 10.0)
+ assert r3["ok"] and r3["scheme"] == (0.5, 0.25, 0.25), r3
+
+
+# ================================================================
+# [B] planner: 取整与科创板部分卖
+# ================================================================
+@case("[B1] ceil_lot 小数向上取整: 100.5 → 200 (原实现先截断再进位, 缺口永远盖不掉)")
+def _():
+ from app.core.planner import ceil_lot, floor_lot
+ assert ceil_lot(100.5) == 200, ceil_lot(100.5)
+ assert ceil_lot(101) == 200 and ceil_lot(100) == 100 and ceil_lot(0) == 0
+ assert ceil_lot(100.0000001) == 100 # 1e-9 容差防浮点噪声顶成 200
+ assert floor_lot(199) == 100 and floor_lot(-5) == 0
+
+
+@case("[B2] _min_sell: 科创板部分卖 <200 时, 可卖够就抬到 200, 不够就本轮不切; 主板原样")
+def _():
+ from app.core.planner import _min_sell
+ assert _min_sell(100, 1000, "688802.SH") == 200 # 抬到 200 (偏保守多卖一点)
+ assert _min_sell(100, 150, "688802.SH") == 0 # 可卖不足 200, 这票本轮不切部分卖
+ assert _min_sell(200, 1000, "688802.SH") == 200 # 已合法, 原样
+ assert _min_sell(100, 1000, "600000.SH") == 100 # 主板不动
+ assert _min_sell(0, 1000, "688802.SH") == 0
+
+
+@case("[B3] plan_sector_exit 停牌票不许静默消失: 无价也下整票卖单 (金额 0, 留痕说明)")
+def _():
+ from app.core.planner import plan_sector_exit
+ r = plan_sector_exit(sector="半导体", positions=[
+ {"ts_code": "600000.SH", "total_qty": 6000, "price": 10.0, "sector": "半导体"},
+ # 停牌票: 行情缺失, price_ok=False (positions_view 拿摊薄成本顶的价)
+ {"ts_code": "688111.SH", "total_qty": 3000, "price": 8.0, "price_ok": False,
+ "sector": "半导体"},
+ ])
+ assert r["ok"], r
+ ex = {i["ts_code"]: i for i in r["items"] if i["action"] == "EXIT"}
+ assert set(ex) == {"600000.SH", "688111.SH"}, ex # 停牌票也在
+ assert ex["688111.SH"]["amount"] == 0.0 and ex["688111.SH"].get("need_price"), ex
+ assert ex["688111.SH"]["qty"] == 3000
+ assert any("取不到现价" in n for n in r["notes"]), r["notes"]
+
+
+# ================================================================
+# [C] 规则闸: 科创板申报数量
+# ================================================================
+@case("[C1] 规则闸买入: 科创板 <200 股必拒 (交易所会拒单); 主板整百照旧")
+def _():
+ from app.core import rule_gate as rg
+ ctx = {"ts_code": "688802.SH", "position": {"total_qty": 0, "avail_qty": 0},
+ "caps": None, "params": {}, "flags": {},
+ "day": {"price": 10.0, "day_chg_from_open": 0.0, "ma5": 10.0}}
+ r = rg.check(side="buy", action="OPEN", qty=100, price=10.0, ctx=ctx)
+ assert any("科创板买入申报最少 200" in f for f in r["failed"]), r["failed"]
+ r2 = rg.check(side="buy", action="OPEN", qty=200, price=10.0, ctx=ctx)
+ assert not any("科创板" in f for f in r2["failed"]), r2["failed"]
+ ctx3 = dict(ctx, ts_code="600000.SH")
+ r3 = rg.check(side="buy", action="OPEN", qty=100, price=10.0, ctx=ctx3)
+ assert not any("LOT_INVALID" in f for f in r3["failed"]), r3["failed"]
+
+
+@case("[C2] 规则闸卖出: 科创板部分卖 <200 拒; 余额不足 200 一次性清出放行 (零股全清)")
+def _():
+ from app.core import rule_gate as rg
+
+ def sell(code, qty, total):
+ return rg.check(side="sell", action="TRIM", qty=qty, price=10.0,
+ ctx={"ts_code": code, "caps": None, "params": {}, "flags": {},
+ "position": {"total_qty": total, "avail_qty": total},
+ "day": {"price": 10.0}})
+ r = sell("688802.SH", 100, 1000)
+ assert any("科创板部分减持最少 200" in f for f in r["failed"]), r["failed"]
+ assert sell("688802.SH", 200, 1000)["passed"], "200 股部分卖合法"
+ assert sell("688802.SH", 150, 150)["passed"], "余额 150 一次性清出是交易所允许的例外"
+ assert sell("600000.SH", 100, 1000)["passed"], "主板 100 股照旧"
+
+
+# ================================================================
+# [D] 动作引擎: 同轮买卖互斥
+# ================================================================
+@case("[D1] 同一只票同轮 TRIM+ADD 同时成立 → 买入侧让路, 不再自动对倒空耗手续费")
+def _():
+ from app.core import action_engine as ae
+ params = {"scale": 2000000, "cushion_solid": 0.03, "trim_peak": 0.06,
+ "trim_giveback": 0.5, "stock_target_default": 0.06}
+ # 峰值 8% 回吐到 3% (过半) → TRIM 成立; 垫 3% 且创 5 日新高 → ADD 也成立
+ # (峰值是全时段只增不减, 加仓看近 5 日窗口 —— 两套时间基准可以同时为真)
+ pos = {"ts_code": "600000.SH", "total_qty": 6000, "avail_qty": 6000,
+ "cushion_pct": 0.03, "cushion_peak": 0.08, "price": 10.3,
+ "market_value": 61800.0, "target_pct": 0.06}
+ mkt = {"600000.SH": {"high5": 10.3, "ma5": 10.3, "tdays_since_open": None,
+ "tdays_since_last_add": None}}
+ r = ae.scan(positions=[pos], params=params, market=mkt)
+ acts = [c["action"] for c in r["candidates"]]
+ assert acts == ["TRIM"], r["candidates"] # 只留减仓, 买入侧让路
+ assert any(s["action"] == "ADD" and "互斥" in s["why"] for s in r["skipped"]), r["skipped"]
+ # 对照: 峰值不足、TRIM 不触发时, ADD 照常产出 (互斥只在同轮同票双触发时生效)
+ pos2 = dict(pos, cushion_peak=0.04)
+ r2 = ae.scan(positions=[pos2], params=params, market=mkt)
+ assert [c["action"] for c in r2["candidates"]] == ["ADD"], r2["candidates"]
+
+
+# ================================================================
+# [E] 信号口径
+# ================================================================
+@case("[E1] _norm_conf_pct 百分制契约: 1 = 1% (原启发式把 1 当 100% 直接触发自动清仓)")
+def _():
+ from app.core.signal_rules import _norm_conf_pct
+ assert abs(_norm_conf_pct(1) - 0.01) < 1e-12, _norm_conf_pct(1)
+ assert abs(_norm_conf_pct(0.9) - 0.009) < 1e-12 # 0~1% 噪声级, 不再漏缩放
+ assert abs(_norm_conf_pct(92) - 0.92) < 1e-12
+ assert _norm_conf_pct(150) == 1.0 and _norm_conf_pct(-5) == 0.0
+ assert _norm_conf_pct("abc") == 0.0 # 解析不了按 0, 落"低于门槛"档
+
+
+@case("[E2] digest 科创板中置信减持: 量抬到 200 / 持仓不足 200 退化全卖; 主板不变")
+def _():
+ from app.core import signal_rules as sr
+ prm = {"sell_conf_min": 0.75, "auto_exit_conf": 0.85, "trim_ratio": 1 / 3}
+
+ def d(code, held, conf=0.8):
+ sig = {"source": "risk_sell", "ts_code": code, "action": "SELL", "confidence": conf}
+ return sr.digest(sig, {"total_qty": held, "avail_qty": held}, prm)
+ r = d("688111.SH", 300)
+ assert r["action"] == sr.ACT_PROPOSE and r["qty"] == 200, r # 100 → 抬到 200
+ r2 = d("688111.SH", 150)
+ assert r2["qty"] == 150, r2 # 不足 200: 一次性全清是合法例外
+ r3 = d("600000.SH", 3000)
+ assert r3["qty"] == 1000, r3 # 主板 1/3 照旧
+ r4 = d("688111.SH", 3000)
+ assert r4["qty"] == 1000, r4 # 量本来就 ≥200, 不动
+
+
+# ================================================================
+# [F] 除权 / 日历
+# ================================================================
+@case("[F1] apply_ex_right: 已部分核销的批次, closed_qty 与核销均价同比例调 (单位不混算)")
+def _():
+ from app.core.recon import apply_ex_right
+ lots = [{"id": 1, "qty": 500, "open_price": 20.0,
+ "closed_qty": 500, "close_avg_price": 22.0},
+ {"id": 2, "qty": 1000, "open_price": 18.0, "closed_qty": 0,
+ "close_avg_price": None}]
+ out = apply_ex_right(lots, 2.0) # 10 送 10
+ a, b = out[0], out[1]
+ assert a["qty"] == 1000 and abs(a["open_price"] - 10.0) < 1e-9, a
+ # 原来只调剩余数量: 剩余是新股数单位、已核销还是旧单位, 摊薄成本照样错
+ assert a["closed_qty"] == 1000 and abs(a["close_avg_price"] - 11.0) < 1e-9, a
+ assert b["qty"] == 2000 and b["closed_qty"] == 0 and b["close_avg_price"] is None, b
+ assert "除权调整" in a["note"]
+
+
+@case("[F2] 交易日历按年探测: 库装了但没当年数据 → degraded=True, 工作日放行不静默跳")
+def _():
+ import app.core.tradedays as td0
+ orig_has, orig_fn = td0._HAS_CAL, td0._is_workday
+ orig_cache = dict(td0._YEAR_OK)
+ try:
+ td0._HAS_CAL = True
+
+ def fake_workday(d):
+ if d.year >= 2027: # 模拟: 库只有 2026 及以前的数据
+ raise NotImplementedError("no data for 2027")
+ return True
+ td0._is_workday = fake_workday
+ td0._YEAR_OK.clear()
+ # 原来 degraded 只看"装没装": 年初库没升级时, 全年法定节假日都被当交易日,
+ # 页面却显示一切正常 —— 这正是最常见的降级场景。
+ assert td0.calendar_degraded(date(2026, 8, 28)) is False
+ assert td0.calendar_degraded(date(2027, 1, 15)) is True
+ assert td0.is_trade_day(date(2027, 1, 15)) is True # 周五: 降级按工作日放行
+ assert td0.is_trade_day(date(2027, 1, 16)) is False # 周六照样拦
+ assert td0._YEAR_OK.get(2027) is False and td0._YEAR_OK.get(2026) is True
+ finally:
+ td0._HAS_CAL, td0._is_workday = orig_has, orig_fn
+ td0._YEAR_OK.clear()
+ td0._YEAR_OK.update(orig_cache)
+
+
+# ================================================================
+# [G] 策略运行侧: 网格中枢 / 止盈闩锁
+# ================================================================
+@case("[G1] 网格只买中枢下方: 上半区回落一档不接盘, 档位照常推进 (不再高买低不买)")
+def _():
+ from app.services import strategy_runner as srun
+ prm = {"lower": 9.0, "upper": 11.0, "center": 10.0, "step_pct": 0.02,
+ "per_lot": 100, "max_capital": 50000}
+ lv = srun._grid_levels(prm)
+ k = srun._band(lv, 10.6) # 中枢上方的一档
+ assert lv[k] >= 10.0, (k, lv)
+ st = {"last_band": k + 1, "filled_levels": {}}
+ d = srun._eval_grid({"ts_code": "600000.SH", "params": prm},
+ {"avail_qty": 0, "add_qty": 0}, {"price": 10.6}, None,
+ {"state": st, "notes": [], "buy_paused": False})
+ assert d is None and st["last_band"] == k, (d, st) # 不买, 但档位随价下移
+ # 中枢下方照常接 (与 batch17 科创板用例同一条路, 这里钉主板+显式中枢)
+ b0 = srun._band(lv, 9.5)
+ st2 = {"last_band": b0 + 1, "filled_levels": {}}
+ d2 = srun._eval_grid({"ts_code": "600000.SH", "params": prm},
+ {"avail_qty": 0, "add_qty": 0}, {"price": 9.5}, None,
+ {"state": st2, "notes": [], "buy_paused": False})
+ assert d2 and d2["side"] == "buy", d2
+ # 没显式配 center 时用 (下界+上界)/2 兜底, 行为一致
+ prm2 = {"lower": 9.0, "upper": 11.0, "step_pct": 0.02, "per_lot": 100}
+ st3 = {"last_band": k + 1, "filled_levels": {}}
+ d3 = srun._eval_grid({"ts_code": "600000.SH", "params": prm2},
+ {"avail_qty": 0, "add_qty": 0}, {"price": 10.6}, None,
+ {"state": st3, "notes": [], "buy_paused": False})
+ assert d3 is None and st3["last_band"] == k, (d3, st3)
+
+
+@case("[G2] 跟踪止盈部分卖一次性闩锁: 同一高水位只卖一次, 创新高后才许再卖; 全清不上锁")
+def _():
+ from app.services import strategy_runner as srun
+
+ def trail(avail, price, state, ratio=0.5):
+ # start_line 显式给, 别让缺省值走 param_store (这批单测不连库)
+ st = {"ts_code": "600000.SH", "params": {"giveback": 0.05, "sell_ratio": ratio,
+ "start_line": 0.03}}
+ pos = {"avg_cost": 8.0, "avail_qty": avail, "total_qty": avail * 2,
+ "cushion_pct": price / 8.0 - 1}
+ return srun._eval_trail(st, pos, {"price": price}, None,
+ {"state": state, "notes": []})
+ state = {"armed": True, "high_water": 12.0}
+ d1 = trail(1000, 10.0, state)
+ assert d1 and d1["qty"] == 500, d1 # 第一次回落: 卖一半
+ assert state.get("trail_fired_hw") == 12.0, state # 闩锁记下高水位
+ # 原 bug: 卖完条件仍成立, 每隔一单再卖剩余一半, 几何级联直到卖光
+ assert trail(500, 10.0, state) is None, "同一高水位不许再卖"
+ trail(500, 13.0, state) # 创新高 → 高水位抬到 13
+ assert state["high_water"] == 13.0, state
+ d3 = trail(500, 12.3, state) # 新一轮回落 ≥5% → 允许再卖
+ assert d3 and d3["qty"] == 200, d3
+ # 全清路径不上锁: 清仓意图失败了就该重试
+ state4 = {"armed": True, "high_water": 12.0, "trail_fired_hw": 12.0}
+ d4 = trail(1000, 10.0, state4, ratio=1.0)
+ assert d4 and d4["action"] == srun.A_EXIT, d4
+
+
+# ================================================================
+# [H] 策略买入暂停按来源分记
+# ================================================================
+@case("[H1] buypause 多来源并存: accum 解除只摘自己的, 不放开风控停的; 旧格式条目自动迁移")
+def _():
+ from app.repo import pms_repo
+ from app.services import strategy_service as svc
+ store = {}
+ orig_get, orig_set = pms_repo.get_param, pms_repo.set_param
+ orig_ls = pms_repo.list_strategies
+ try:
+ pms_repo.get_param = lambda k: store.get(k)
+ pms_repo.set_param = lambda k, v, by="user": store.__setitem__(k, str(v)) or 1
+ pms_repo.list_strategies = (
+ lambda *, ts_code=None, statuses=None, limit=500, include_archived=False:
+ [{"strategy_id": "S1", "ts_code": ts_code}])
+ # 旧格式条目 (没有 sources 子表) 先躺在表里 → pause_buy 迁移成 sources
+ import json
+ store[svc.BUYPAUSE_KEY] = json.dumps(
+ {"600000.SH": {"reason": "旧风控", "source": "signal", "at": "2026-08-27"}})
+ assert svc.pause_buy("600000.SH", reason="定性失效", source="accum") == ["S1"]
+ m = svc.buypause_map()
+ assert set(m["600000.SH"]["sources"]) == {"signal", "accum"}, m
+ # 原 bug: 先写先赢、后来的来源被吞 —— advisor 按 accum 解除时把风控停的也放开了
+ r = svc.clear_buypause("600000.SH", only_source="accum")
+ assert r["ok"] and r["cleared"] and r["still_paused_by"] == ["signal"], r
+ m2 = svc.buypause_map()
+ assert "600000.SH" in m2 and set(m2["600000.SH"]["sources"]) == {"signal"}, m2
+ # 来源不匹配: 不动, 不算错
+ r2 = svc.clear_buypause("600000.SH", only_source="accum")
+ assert r2["ok"] and r2["cleared"] is False, r2
+ # 无 only_source: 整条解除
+ r3 = svc.clear_buypause("600000.SH")
+ assert r3["cleared"] and "600000.SH" not in svc.buypause_map(), r3
+ finally:
+ pms_repo.get_param, pms_repo.set_param = orig_get, orig_set
+ pms_repo.list_strategies = orig_ls
+
+
+# ================================================================
+# [I] 宏观失败路径不冲留痕
+# ================================================================
+@case("[I1] _upsert_unavailable: 重扫失败保留当日 CMD_ISSUED/建议/周期, 不再整行重写")
+def _():
+ from app.repo import macro_repo
+ from app.core import macro_rules as mr
+ from app.services import macro_service as ms
+ got = {}
+ orig_get, orig_up = macro_repo.get_signal, macro_repo.upsert_signal
+ try:
+ macro_repo.get_signal = lambda key, d: {
+ "action": "CMD_ISSUED", "ref_id": "CMD_MACRO_1", "note": "已下减仓命令",
+ "detail": {"cycle": {"done_shift": 1}, "exit_acted": True, "advice": "减"}}
+ macro_repo.upsert_signal = lambda **kw: got.update(kw) or 1
+ ms._upsert_unavailable("stock_fx_hedge", 20260828, "数据源超时")
+ # 原来失败路径按默认值整行覆盖: 上午的 CMD_ISSUED 被冲掉 → "当日不重复下命令"
+ # 判据失效, 次日取昨日周期也拿不到
+ assert got["action"] == "CMD_ISSUED" and got["ref_id"] == "CMD_MACRO_1", got
+ assert got["zone"] == mr.Z_UNAVAILABLE and got["value"] is None, got
+ assert got["detail"]["cycle"] == {"done_shift": 1}, got["detail"]
+ assert got["detail"]["exit_acted"] is True and got["detail"]["advice"] == "减", got
+ assert "重扫失败" in got["note"] and "已下减仓命令" in got["note"], got["note"]
+ # 当日无留痕 (action=NONE) 时: 正常落 UNAVAILABLE, note 就是失败原因本身
+ got.clear()
+ macro_repo.get_signal = lambda key, d: None
+ ms._upsert_unavailable("stock_fx_hedge", 20260828, "数据源超时")
+ assert got["action"] == "NONE" and got["note"] == "数据源超时", got
+ finally:
+ macro_repo.get_signal, macro_repo.upsert_signal = orig_get, orig_up
+
+
+def main():
+ passed, failed = 0, []
+ for name, fn in RESULTS:
+ try:
+ fn()
+ passed += 1
+ print(f" ✓ {name}")
+ except Exception as e:
+ failed.append((name, e))
+ print(f" ✗ {name}: {type(e).__name__}: {e}")
+ traceback.print_exc()
+ print()
+ if failed:
+ print(f"FAILED {len(failed)}/{len(RESULTS)}")
+ sys.exit(1)
+ print(f"ALL PASS ({passed} cases)")
+
+
+if __name__ == "__main__":
+ main()
diff --git a/scripts/test_batch2_units.py b/scripts/test_batch2_units.py
index 9f202ae..9983d14 100644
--- a/scripts/test_batch2_units.py
+++ b/scripts/test_batch2_units.py
@@ -202,6 +202,7 @@ def _():
for i in r["items"]:
if i["action"] in (pl.A_EXIT, pl.A_TRIM):
used[i["ts_code"]] = used.get(i["ts_code"], 0) + i["qty"]
+ assert used, "35 万档必须真的排出减持动作 (2026-08-28 补: 空计划曾能零断言通过)"
for c, q in used.items():
assert q <= hold[c], (c, q, hold[c])
diff --git a/scripts/test_batch3_units.py b/scripts/test_batch3_units.py
index 48932ae..36c6653 100644
--- a/scripts/test_batch3_units.py
+++ b/scripts/test_batch3_units.py
@@ -38,14 +38,21 @@ def day(**kw):
# ================================================================ exec_timing
-@case("分日配额·整除/上取整到一手/最后一日全出/零股尾巴并入")
+@case("分日配额·整除/上取整到一手/最后一日全出/零股尾巴只对清仓并入")
def _():
assert et.daily_quota(6000, 3) == 2000
assert et.daily_quota(5000, 3) == 1700 # 1666.7 → 上取整到一手
assert et.daily_quota(1000, 1) == 1000 # 最后一日全出
- assert et.daily_quota(150, 3) == 150 # 尾巴不足一手 → 一次出完
+ # 2026-08-28 口径修正: 部分减持 (allow_odd_tail=False) **不再并零股尾巴** ——
+ # 并进去整片变成非整百, 会被规则闸按 LOT_INVALID 整片拒掉, 连整数部分都卖不出。
+ # 150 剩量的部分减持: 本次出 100, 尾巴 50 留给窗口收口报部分完成。
+ assert et.daily_quota(150, 3) == 100
+ # 整票清仓 (allow_odd_tail=True) 才允许把零股尾巴并进本片一次出完
+ assert et.daily_quota(150, 3, allow_odd_tail=True) == 150
assert et.daily_quota(0, 3) == 0
assert et.daily_quota(100, 5) == 100
+ # 科创板: 按 lot=200 取整
+ assert et.daily_quota(1000, 3, lot=200) == 400
# 整票清仓允许零股
assert et.daily_quota(14050, 3, allow_odd_tail=True) == 4700
assert et.daily_quota(14050, 1, allow_odd_tail=True) == 14050
diff --git a/scripts/test_batch5_units.py b/scripts/test_batch5_units.py
index 5dd3acc..7cecaca 100644
--- a/scripts/test_batch5_units.py
+++ b/scripts/test_batch5_units.py
@@ -50,10 +50,18 @@ def _():
assert s["source"] == sr.SRC_RISK_SELL and s["action"] == "SELL"
assert abs(s["confidence"] - 0.88) < 1e-9, s # 88 → 0.88, 两条流尺度不同
assert s["dominant_signal"] == "破位" and "支撑" in s["reason"]
- # 已经是 0~1 的也不会被再除一次
+ # 2026-08-28 口径修正: 风控流按文档**固定 0~100 制, 一律除以 100**, 不再做
+ # "大于 1 才除"的猜测 —— 原口径下 0~100 制里的 1 (即 1%) 会被当成 100% 直接清仓。
assert abs(sr.parse_risk_sell({"data": json.dumps({"ts_code": "x", "action": "SELL",
"confidence": 0.9})})["confidence"]
- - 0.9) < 1e-9
+ - 0.009) < 1e-9
+ # 危险边界: 1 是 1%, 绝不能被解释成 100%
+ assert abs(sr.parse_risk_sell({"data": json.dumps({"ts_code": "x", "action": "SELL",
+ "confidence": 1})})["confidence"]
+ - 0.01) < 1e-9
+ # 盘中流 (0~1 制) 的口径不变: 0.9 就是 90%
+ assert abs(sr.parse_intraday({"ts_code": "x", "action": "SELL",
+ "confidence": "0.9"})["confidence"] - 0.9) < 1e-9
# 坏 JSON → 明确标记, 不抛异常
bad = sr.parse_risk_sell({"data": "{不是JSON"})
assert bad["ts_code"] == "" and "解析失败" in bad["parse_error"]
diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py
index b09b124..4a522d4 100644
--- a/scripts/test_batch6_units.py
+++ b/scripts/test_batch6_units.py
@@ -621,6 +621,9 @@ def run():
@case("密钥不进 ParamStore (协议 §10.1.1)")
def _():
from app.services import param_store
+ # 2026-08-28 补: 名单本身要钉死 —— 原断言拿实现自己的 SECRET_KEYS 当预期,
+ # 有人把键从名单里删掉, 保护和测试会一起消失
+ assert {"PMS_QMT_SIGN_SEED_HEX", "PMS_QMT_PEER_PUBKEY_B64"} <= set(param_store.SECRET_KEYS), param_store.SECRET_KEYS
snap_keys = {p["key"] for p in [{"key": k} for k in param_store._editable_keys()]}
for k in param_store.SECRET_KEYS:
assert k not in snap_keys, f"{k} 不该出现在可调参数里"
diff --git a/scripts/test_batch7_units.py b/scripts/test_batch7_units.py
index 68bd8f3..3434d10 100644
--- a/scripts/test_batch7_units.py
+++ b/scripts/test_batch7_units.py
@@ -224,6 +224,8 @@ def _():
pf.assert_fresh(p, max_stale_tdays=1, today="2026-08-05")
except pf.PlanFeedError as e:
assert "2026-07-29" in str(e) and "交易日" in str(e), str(e)
+ else:
+ raise AssertionError("超期计划没有抛 PlanFeedError (2026-08-28 补: 原来不抛也静默通过)")
# ================================================================ 筛选
diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py
index 1e16135..2daa62d 100644
--- a/scripts/test_wiring.py
+++ b/scripts/test_wiring.py
@@ -34,7 +34,9 @@ class FakeRepo:
self.positions, self.lots, self.instructions = {}, [], {}
self.proposals, self.ledger, self.reports, self.industry = {}, [], {}, {}
self.cash_flows = []
+ self.strategies = {}
self._lot_id = 0
+ self._strategy_id = 0
# --- runtime param ---
def all_params(self):
@@ -74,10 +76,13 @@ class FakeRepo:
and (not cmd_class or c["cmd_class"] == cmd_class)]
return sorted(out, key=lambda c: -c["id"])[:limit]
- def update_command(self, cid, *, status=None, progress=None, done_at=None, note=None):
+ def update_command(self, cid, *, status=None, progress=None, done_at=None, note=None,
+ only_if_status=None):
c = self.commands.get(cid)
if not c:
return 0
+ if only_if_status is not None and c.get("status") not in list(only_if_status):
+ return 0 # 条件更新没抢到 (与真 repo 的 WHERE status IN 同语义)
if status is not None:
c["status"] = status
if progress is not None:
@@ -338,6 +343,62 @@ class FakeRepo:
self.industry[r["ts_code"]] = r["industry"]
return len(rows)
+ # --- pms_strategy (个股交易方案) ---
+ # 2026-08-28 审查补: 桩里原来**一个策略函数都没有**, 于是所有走
+ # pms_repo.list_strategies / active_strategy_codes 的服务代码在单测里都命中
+ # 真 repo → 连库异常 → 被各自的 try/except 按"空集"吞掉 —— 策略相关的接线
+ # (清仓撤策略 / 风控暂停买入 / 动作引擎排除策略票) 从来没被装配自检真正跑过。
+ def list_strategies(self, *, ts_code=None, statuses=None, limit=500,
+ include_archived=False):
+ out = [dict(s) for s in self.strategies.values()
+ if (not ts_code or s["ts_code"] == ts_code)
+ and (not statuses or s["status"] in list(statuses))
+ and (include_archived or not s.get("archived_at"))]
+ return sorted(out, key=lambda s: -s["id"])[:int(limit)]
+
+ def get_strategy(self, strategy_id):
+ s = self.strategies.get(strategy_id)
+ return dict(s) if s else None
+
+ def active_strategies(self):
+ return self.list_strategies(statuses=["ACTIVE"], limit=1000)
+
+ def active_strategy_codes(self):
+ return {s["ts_code"] for s in self.strategies.values()
+ if s["status"] == "ACTIVE" and s.get("ts_code")}
+
+ def insert_strategy(self, *, strategy_id, ts_code, stype, autonomy="auto", params=None,
+ state=None, status="ACTIVE", note=None):
+ if strategy_id in self.strategies: # 真表 strategy_id 唯一键
+ raise Exception(f"Duplicate entry '{strategy_id}' for key 'strategy_id'")
+ self._strategy_id += 1
+ self.strategies[strategy_id] = {
+ "id": self._strategy_id, "strategy_id": strategy_id, "ts_code": ts_code,
+ "type": stype, "status": status, "autonomy": autonomy,
+ "params": dict(params or {}), "state": dict(state or {}), "note": note,
+ "archived_at": None, "created_at": datetime.now(), "updated_at": datetime.now()}
+ return 1
+
+ def update_strategy(self, strategy_id, **fields):
+ s = self.strategies.get(strategy_id)
+ if not s:
+ return 0
+ upd = dict(fields)
+ # 与真 repo 同义: params/state 传 dict 落 *_json 列; 其余按白名单更新
+ touched = False
+ for k in ("params", "state"):
+ if k in upd:
+ s[k] = dict(upd.pop(k) or {})
+ touched = True
+ for k in ("status", "autonomy", "note"):
+ if k in upd:
+ s[k] = upd[k]
+ touched = True
+ if not touched:
+ return 0
+ s["updated_at"] = datetime.now()
+ return 1
+
class FakeQmtRepo:
"""ws 通道三表的内存替身 (pms_qmt_order / pms_qmt_inbox / pms_ws_state)。
@@ -376,10 +437,11 @@ class FakeQmtRepo:
def enqueue_order(self, *, instruction_id, parent_id, ts_code, side, qty, limit_price,
valid_until, intent="OPEN", note=None):
self.orders[instruction_id] = {
- "instruction_id": instruction_id, "parent_id": parent_id,
- # 真表的列名是 parent_instruction_id, consume_ws_trades 反查用的是它。
- # 桩里两个键都放, 少一个的话 SMOKE_ 那道闸在单测里永远"看起来没生效"。
- "parent_instruction_id": parent_id, "ts_code": ts_code,
+ # 真表的列名是 parent_id (ddl_pms_v1.sql / qmt_repo.enqueue_order)。
+ # 2026-08-28 修: 桩里原来多放了一个不存在的 parent_instruction_id 键,
+ # 还配了一句写反的注释 —— 恰好掩护了 consume_ws_trades 读错列名的真 bug
+ # (SMOKE 闸从未生效)。桩必须与真表同形, 只留 parent_id。
+ "instruction_id": instruction_id, "parent_id": parent_id, "ts_code": ts_code,
"side": side, "qty": int(qty), "limit_price": float(limit_price),
"valid_until": int(valid_until), "intent": intent, "note": note,
"status": "QUEUED", "cancel_state": "NONE", "cancel_id": None}
@@ -505,11 +567,12 @@ def install_fakes(prices=None, positions=None, params=None, high5=None, prev_clo
industry_repo.primary_industry_map = lambda codes, level="l3": {}
industry_repo.invalidate = lambda: None
downstream_repo.latest_filled_order_id = lambda: "ANCHOR_0"
- # 回放游标预置成非空 —— 否则 replay_fills 会走「冷启动只对齐不追认」那条路 (见
- # ledger_service._seed_cursor), 下面那几个回放用例就测不到入账。冷启动本身另有专门用例。
- # 用 "0" 而不是随便一个字符串: next_cursor 只进不退, 且非数字 id 会退化成字典序比较,
- # 占位值若比真实 order_id 大 (比如 "SEED"), 游标就永远推不动了。
- fake.params.setdefault("PMS_REPLAY_CURSOR", "0")
+ downstream_repo.latest_filled_time = lambda: "2026-07-27 09:00:00"
+ # 回放游标预置成 v2 时间游标 (从 epoch 起 = 什么都算新) —— 否则 replay_fills 会走
+ # 「冷启动只对齐不追认」那条路, 下面那几个回放用例就测不到入账。冷启动另有专门用例。
+ # 2026-08-28 起游标是 {"v":2,"t":时间,"seen":{order_id:时间}}, 不再按 order_id 字典序。
+ fake.params.setdefault("PMS_REPLAY_CURSOR",
+ '{"v": 2, "t": "1970-01-01 00:00:00", "seen": {}}')
# 上游选股计划接口在单测里一律停用 (base 为空 → plan_feed 立即抛 PlanFeedError,
# 不会发出任何 HTTP 请求)。要测候选池的用例自己 stub plan_feed.candidates。
fake.params.setdefault("PMS_PLAN_API_BASE", "")
@@ -790,11 +853,13 @@ def _():
assert fake.positions["600000.SH"]["base_qty"] == 6000
assert abs(float(fake.positions["600000.SH"]["cushion_pct"]) - 0.10) < 1e-4
assert fake.instructions["INS_A"]["status"] == "CONFIRMED"
- assert fake.params["PMS_REPLAY_CURSOR"] == "101"
- # 幂等: 游标已推进, 同一批不再重复入账
- downstream_repo.fetch_filled_orders = lambda **kw: []
+ import json as _json
+ cur = _json.loads(fake.params["PMS_REPLAY_CURSOR"])
+ assert cur["t"] == "2026-07-27 09:40:00" and "101" in cur["seen"], cur
+ # 幂等: **同一批再喂一遍**也不重复入账 (seen 集合按 order_id 去重 ——
+ # 这才是真幂等; 旧口径只测了"喂空批不入账", 测不住重复消费)
r2 = ls.replay_fills()
- assert r2["fills"] == 0 and fake.positions["600000.SH"]["total_qty"] == 6000
+ assert r2["fills"] == 0 and fake.positions["600000.SH"]["total_qty"] == 6000, r2
finally:
downstream_repo.fetch_filled_orders = orig
@@ -839,14 +904,15 @@ def _():
orig = downstream_repo.fetch_filled_orders
try:
downstream_repo.fetch_filled_orders = lambda **kw: hist
- downstream_repo.latest_filled_order_id = lambda: "SELL_ZZZ_999"
+ downstream_repo.latest_filled_time = lambda: "2026-07-28 15:00:00"
r = ls.replay_fills()
# 关键: 一条都不能入账。旧系统多年的历史成交若被并入 BASE, 摊薄成本与安全垫全错,
# 而补仓/加仓/保垫减仓都挂在安全垫上 —— 一错就是整条纪律链。
assert r.get("seeded") and r["fills"] == 0 and r["actions"] == 0, r
- assert r["cursor"] == "SELL_ZZZ_999", r
+ assert r["cursor"] == "2026-07-28 15:00:00", r
assert not fake.lots and not fake.positions, "冷启动不该产生任何批次或持仓"
- assert fake.params["PMS_REPLAY_CURSOR"] == "SELL_ZZZ_999"
+ import json as _json
+ assert _json.loads(fake.params["PMS_REPLAY_CURSOR"])["t"] == "2026-07-28 15:00:00"
# 游标就位后, 新成交照常入账
new = [{"order_id": "ZZZ_NEW", "ts_code": "600000.SH", "side": "buy", "qty": 100,
@@ -1041,19 +1107,36 @@ def _():
assert fake.reports[rep["ymd"]]["ymd"] == rep["ymd"]
-@case("账本服务·除权检测走通 (10送10 → 批次按比例调整)")
+@case("账本服务·除权检测走通 (10送10 → 批次按比例调整; 下游为检测源)")
def _():
+ # 2026-08-28 口径重写: 送转发生在券商账户, 账本数量在对账之前不会变 ——
+ # 检测必须比「昨日快照 vs 下游当前数量」, 而不是账本自己比自己 (那样永远检不出,
+ # 红股会被对账当普通差异吸收, 老批次成本价不除权, 摊薄成本虚高)。
from app.services import ledger_service as ls
fake = install_fakes(prices={"600000.SH": 5.0}, params={"PMS_TOTAL_SCALE": "2000000"},
- positions=[{"ts_code": "600000.SH", "total_qty": 2000,
+ positions=[{"ts_code": "600000.SH", "total_qty": 1000,
"avg_cost": 10.0}])
- fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=2000, open_price=10.0,
+ fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=1000, open_price=10.0,
open_date="2026-07-01")
fake.upsert_report(20260726, {"snapshot": {"600000.SH": {"qty": 1000, "price": 10.0}}})
- r = ls.detect_and_apply_ex_right()
- assert r["ex_rights"] and abs(r["ex_rights"][0]["ratio"] - 2.0) < 1e-6, r
- lot = fake.lots[0]
- assert lot["qty"] == 4000 and abs(lot["open_price"] - 5.0) < 1e-6, lot
+ orig_src = ls.positions_source
+ try:
+ ls.positions_source = lambda: {"source": "table", "mode": "test", "alerts": [],
+ "rows": [{"ts_code": "600000.SH", "qty": 2000}],
+ "columns": {"qty": "q"}, "raw_count": 1,
+ "as_of": 0, "age_sec": 0}
+ r = ls.detect_and_apply_ex_right()
+ assert r["ex_rights"] and abs(r["ex_rights"][0]["ratio"] - 2.0) < 1e-6, r
+ lot = fake.lots[0]
+ assert lot["qty"] == 2000 and abs(lot["open_price"] - 5.0) < 1e-6, lot
+
+ # 正常加仓日不误报: 账本数量已因今日成交变动 → 跳过, 不产生 MISMATCH
+ fake2_pos = fake.positions["600000.SH"]
+ fake2_pos["total_qty"] = 3000 # 账本已变 (今天有成交), 快照仍是 1000
+ r2 = ls.detect_and_apply_ex_right()
+ assert not r2["mismatches"] and not r2["ex_rights"], r2
+ finally:
+ ls.positions_source = orig_src
@case("盘前准备·T+1 可卖重置")
@@ -1616,22 +1699,39 @@ def _():
class FakeRedis:
- """Redis Stream 的最小替身 (消费组 + xreadgroup + ack)。"""
+ """Redis Stream 的最小替身, 语义对齐真 Redis 消费组 (2026-08-28 扩):
+ - ">" 只投递**从未投递过**的消息, 读到即进本消费者的待确认清单 (delivered);
+ - "0" 只回自己名下**读了没确认**的消息 —— 消费组的"不确认会重投"仅指这条路;
+ - xack 从待确认清单移除;
+ - xrevrange 是只读取样 (试算用), 不产生任何投递痕迹。
+ 原桩把"读了不确认"做成了"下次照常再读", 恰好掩护了 dry_run 吞消息的真 bug。"""
def __init__(self, msgs=None):
self.msgs = dict(msgs or {})
+ self.delivered = {} # key -> [(id, fields)] 已投递未确认
self.acked, self.groups = [], []
def xgroup_create(self, key, group, id="$", mkstream=False):
self.groups.append((key, group))
def xreadgroup(self, group, consumer, streams, count=10, block=0):
- key = list(streams)[0]
- m = self.msgs.pop(key, [])
+ key, start = list(streams)[0], list(streams.values())[0]
+ if str(start) == "0": # 回捞自己名下待确认的
+ m = list(self.delivered.get(key) or [])
+ return [(key, m)] if m else []
+ m = self.msgs.pop(key, []) # ">": 只投新消息, 且立刻进待确认清单
+ if m:
+ self.delivered.setdefault(key, []).extend(m)
return [(key, m)] if m else []
+ def xrevrange(self, key, count=10):
+ rows = list(self.delivered.get(key) or []) + list(self.msgs.get(key) or [])
+ return list(reversed(rows))[:count]
+
def xack(self, key, group, msg_id):
self.acked.append(msg_id)
+ self.delivered[key] = [(i, f) for i, f in (self.delivered.get(key) or [])
+ if i != msg_id]
def xlen(self, key):
return len(self.msgs.get(key, []))
@@ -1720,6 +1820,12 @@ def _():
by = _install_signal_fakes(fake2, sell_msgs=[_sell_msg("3-9", "600000.SH", 95)])
r2 = ss.consume()
assert "skipped" in r2 and not fake2.instructions, r2
+ # 2026-08-28 补: 开关关着必须**碰都不碰**消费组 —— 读了不 ACK 的消息会永久滞留,
+ # 开关重开后再也补不回来 (原来 by 拿了没断言, 恰是这条最该断)
+ from config.settings import settings as _st
+ _r3c = by[_st.SIGNAL_REDIS_DB_ACTIONS]
+ assert _r3c.acked == [] and not _r3c.delivered, (_r3c.acked, _r3c.delivered)
+ assert _r3c.msgs.get("bionic:signals:llm_sell_actions"), "消息必须原封不动留在流里"
fake3 = install_fakes(prices={"600000.SH": 10.0}, params={"PMS_TOTAL_SCALE": "2000000"},
positions=[{"ts_code": "600000.SH", "total_qty": 6000,
@@ -2054,6 +2160,222 @@ def _():
assert by["601111.SH"]["src"] == "buy_plan" and by["601111.SH"]["price"] == 9.0
+# ================================================================ 2026-08-28 审查补:
+# 撤销链路 / 进度口径 / 策略接线 / 消化幂等 —— 这批用例钉住的都是当次审查修掉的真伤
+@case("命令撤销·在途指令必须经下游撤回; 下游拒撤 → 命令保持在途可重撤")
+def _():
+ from app.services import command_service as csvc, dispatcher
+ fake = install_fakes(
+ prices={"600000.SH": 10.0, "000001.SZ": 8.0},
+ params={"PMS_TOTAL_SCALE": "2000000"},
+ positions=[{"ts_code": "600000.SH", "total_qty": 14000, "base_qty": 7000,
+ "avg_cost": 8.93},
+ {"ts_code": "000001.SZ", "total_qty": 10000, "base_qty": 10000,
+ "avg_cost": 8.8}])
+ r = csvc.issue("REDUCE_EXPOSURE", {"pct": "5%", "window_tdays": 3})
+ assert r["ok"], r
+ cid = r["command_id"]
+ plans = fake.list_plans(command_id=cid)
+ assert plans, "方案未落表"
+ # 手工把一条方案物化成在途指令 (绕过择时, 只测「撤销必须过下游」这条链)
+ fake.insert_instruction(instruction_id="INS_CXL_1", origin_type="plan",
+ origin_id=plans[0]["plan_id"], ts_code=plans[0]["ts_code"],
+ action=plans[0]["action"], side="sell", qty=plans[0]["qty"],
+ limit_price=8.0, status="DISPATCHED")
+ orig = dispatcher.cancel
+ try:
+ dispatcher.cancel = lambda **kw: {"ok": False, "error": "ws 断连, 撤单没发出去"}
+ r2 = csvc.cancel(cid)
+ # 2026-07-31 教训的第五条路: 原来这里只把本端行标 CANCELLED, 下游委托继续挂着
+ # 继续成交, 页面却回"已撤销"。现在: 没撤成 → 命令**保持在途**继续被跟踪。
+ assert r2["ok"] is False and r2["failed"], r2
+ assert "仍在下游挂着" in r2["message"], r2["message"]
+ assert fake.commands[cid]["status"] == "EXECUTING", fake.commands[cid]
+ assert fake.instructions["INS_CXL_1"]["status"] == "DISPATCHED", \
+ "下游拒撤时指令不许在本端标终态"
+ finally:
+ dispatcher.cancel = orig
+ # 下游恢复后再点一次撤销: 方案已 CANCELLED 也要能匹配到在途指令 (匹配用**全部**方案,
+ # 不带状态过滤 —— 带了的话上一次撤到一半的在途指令会漏成孤儿)
+ r3 = csvc.cancel(cid)
+ assert r3["ok"] and "下游已确认" in r3["message"], r3
+ assert fake.instructions["INS_CXL_1"]["status"] == "CANCELLED"
+ assert fake.commands[cid]["status"] == "CANCELLED"
+
+
+@case("命令进度·目标金额 0 但方案未出清 (盘前清仓无现价) 不判 DONE; 出清才完结")
+def _():
+ from app.services import command_service as csvc
+ fake = install_fakes(
+ prices={}, # 盘前: 一只现价都取不到
+ params={"PMS_TOTAL_SCALE": "2000000"},
+ positions=[{"ts_code": "600000.SH", "total_qty": 6000, "avail_qty": 6000,
+ "avg_cost": 10.0},
+ {"ts_code": "000001.SZ", "total_qty": 3000, "avail_qty": 3000,
+ "avg_cost": 8.0}])
+ r = csvc.issue("LIQUIDATE_ALL", {"confirm": "YES"})
+ assert r["ok"] and r["status"] == "EXECUTING", r
+ cid = r["command_id"]
+ assert fake.commands[cid]["progress"]["target_amount"] == 0.0 # 全部票取不到价
+ exit_p = [p for p in fake.list_plans(command_id=cid) if p["action"] == "EXIT"]
+ assert len(exit_p) == 2 and all(float(p["amount"]) == 0.0 for p in exit_p), exit_p
+ assert all("取不到现价" in (p.get("reason") or "") for p in exit_p), exit_p
+ # 原来 settle 见 target<=0 直接 DONE —— 一股没卖、方案永不物化、页面显示"已完成"。
+ r2 = csvc.refresh_progress(cid)
+ assert r2["commands"][0]["status"] == "EXECUTING", r2
+ assert fake.commands[cid]["status"] == "EXECUTING", fake.commands[cid]
+ for p in exit_p:
+ fake.update_plan(p["plan_id"], status="DONE", filled_qty=p["qty"])
+ csvc.refresh_progress(cid)
+ assert fake.commands[cid]["status"] == "DONE", fake.commands[cid]
+
+
+@case("命令进度·GATED 批金额不进完成分母 (BASE 出清即 DONE, 不再拖满窗口判 PARTIAL)")
+def _():
+ from app.core import command_spec as cspec
+ from app.services import command_service as csvc
+ fake = install_fakes(params={"PMS_TOTAL_SCALE": "2000000"})
+ fake.insert_command(command_id="CMD_G1", cmd_class=cspec.CLS_TASK,
+ cmd_type="INCREASE_EXPOSURE", ts_code=None, params={},
+ status="EXECUTING",
+ progress={"target_amount": 100000.0, "deadline": "2099-01-01"})
+ fake.insert_plans([
+ {"plan_id": "PL_G1", "command_id": "CMD_G1", "ts_code": "600000.SH",
+ "action": "OPEN", "qty": 6000, "amount": 60000.0, "status": "DONE",
+ "filled_qty": 6000, "priority": 1, "deadline": "2099-01-01"},
+ {"plan_id": "PL_G2", "command_id": "CMD_G1", "ts_code": "600000.SH",
+ "action": "FILL", "qty": 4000, "amount": 40000.0, "status": "GATED",
+ "filled_qty": 0, "priority": 2, "deadline": "2099-01-01"},
+ ])
+ csvc.refresh_progress("CMD_G1")
+ c = fake.commands["CMD_G1"]
+ # GATED 批没有解锁机制, 算进分母的话每条建仓命令必然拖满窗口被判 PARTIAL
+ assert c["status"] == "DONE", c
+ assert c["progress"]["done_amount"] == 60000.0, c["progress"]
+ assert c["progress"]["gated_amount"] == 40000.0, c["progress"]
+
+
+@case("信号消化·挂着策略的票: 高置信风控卖出不清仓, 落等拍板提议 + 暂停策略买入")
+def _():
+ from app.services import signal_service as ss, strategy_service
+ 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": 10.0}])
+ fake.insert_strategy(strategy_id="STR_1", ts_code="600000.SH", stype="GRID",
+ params={"lower": 9.0, "upper": 11.0})
+ by = _install_signal_fakes(fake, sell_msgs=[_sell_msg("9-1", "600000.SH", 92)])
+ r = ss.consume()
+ assert r["ok"], r
+ # 强制离场会推翻你特意设的策略 —— 只提示、不自动清仓
+ assert not r["exits"], r
+ assert not any(i["action"] == "EXIT" for i in fake.instructions.values()), \
+ fake.instructions
+ assert r["proposals"] and r["proposals"][0]["on_strategy"] is True, r
+ assert r.get("strategy_buy_paused") == ["STR_1"], r
+ prop = list(fake.proposals.values())[0]
+ assert prop["action"] == "EXIT" and prop["qty"] == 6000, prop # 采纳 = 撤策略并**全**清
+ assert prop["status"] == "WAIT_USER" and prop["judge_verdict"] == "STRATEGY_RISK", prop
+ assert "策略" in prop["hard_numbers"]["reason"], prop
+ assert len(by[3].acked) == 1, by[3].acked
+ # 买入暂停表按来源落了 risk_sell 这一条 (卖出/平回不受影响, 页面可恢复)
+ m = strategy_service.buypause_map()
+ assert "600000.SH" in m, m
+ assert "risk_sell" in (m["600000.SH"].get("sources") or {}), m["600000.SH"]
+
+
+@case("信号消化·试算不吞消息: 试算后正式消费, 同一条风控卖出照常转指令")
+def _():
+ from app.services import signal_service as ss
+ 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": 10.0}])
+ by = _install_signal_fakes(fake, sell_msgs=[_sell_msg("7-1", "600000.SH", 95)])
+ key = "bionic:signals:llm_sell_actions"
+ r = ss.consume(dry_run=True)
+ assert r["exits"] and r["exits"][0].get("dry_run") is True, r
+ # 原 bug: 试算走消费组读了不 ACK → 消息进待确认清单**永不再投递**, 点一次试算就把
+ # 未消化的风控卖出永久吞掉。现在试算走只读 XREVRANGE, 不留任何投递痕迹。
+ assert not fake.instructions and by[3].acked == [] and not by[3].delivered, \
+ (fake.instructions, by[3].acked, by[3].delivered)
+ assert by[3].msgs.get(key), "试算后消息必须原封不动留在流里"
+ r2 = ss.consume() # 正式消费: 同一条消息真正入账
+ assert r2["exits"] and r2["exits"][0]["ts_code"] == "600000.SH", r2
+ ins = [i for i in fake.instructions.values() if i["action"] == "EXIT"]
+ assert ins and ins[0]["qty"] == 6000, ins
+ assert by[3].acked == ["7-1"], by[3].acked
+
+
+@case("信号消化·上一跳读了没确认的消息 (中途崩溃), 下一跳自动捞回重消化")
+def _():
+ from app.services import signal_service as ss
+ 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": 10.0}])
+ by = _install_signal_fakes(fake, sell_msgs=[_sell_msg("8-1", "600000.SH", 95)])
+ key = "bionic:signals:llm_sell_actions"
+ r3 = by[3]
+ # 模拟上一跳: 消息按 ">" 被读走 (进本消费者的待确认清单) 后处理中途崩溃, 没 ACK
+ got = r3.xreadgroup("g", "c", {key: ">"}, count=10)
+ assert got and not r3.msgs.get(key) and r3.delivered.get(key), (got, r3.delivered)
+ # 本跳 consume: _read 先按 "0" 回捞自己名下待确认的 → 正常消化 + ACK。
+ # "读了不确认 = 下次还在" 是对消费组的误解 —— 不回捞的话这条消息永远不会再投递。
+ r = ss.consume()
+ assert r["exits"] and r["exits"][0]["ts_code"] == "600000.SH", r
+ assert r3.acked == ["8-1"] and not r3.delivered.get(key), (r3.acked, r3.delivered)
+ ins = [i for i in fake.instructions.values() if i["action"] == "EXIT"]
+ assert ins and ins[0]["progress"]["from_signal"] is True, ins
+
+
+@case("自主提议·当日涨幅超上限的加仓被 NO_CHASE_DAYUP 拦在整条链路上; 减持不受累")
+def _():
+ from app.services import market, proposal_service as ps
+ fake = _prop_fakes(params={"PMS_AUTONOMY": "full"})
+ orig = market.day_snapshot
+ try:
+ # 造当日大涨 10% (> 默认上限 5%)。batch10 [J3] 钉的是 _market_ctx 取真快照,
+ # 这里钉整条行为链: scan_and_route → _route_one → 规则闸真的拿到涨幅并拦下。
+ market.day_snapshot = lambda c: {"price": 11.0, "vwap": 10.8, "open": 10.0,
+ "high": 11.1, "low": 10.0,
+ "day_chg_from_open": 0.10, "bars": 60}
+ r = ps.scan_and_route()
+ finally:
+ market.day_snapshot = orig
+ rej = {(x["ts_code"], x["action"], x["by"]) for x in r["rejected"]}
+ assert ("600000.SH", "ADD", "rule") in rej, r
+ assert any("NO_CHASE_DAYUP" in f for x in r["rejected"] for f in x["failed"]), \
+ r["rejected"]
+ assert not any(i["action"] == "ADD" for i in fake.instructions.values())
+ led = [x for x in fake.ledger
+ if x["verdict"] == "REJECT" and x["ts_code"] == "600000.SH"]
+ assert led and any("NO_CHASE_DAYUP" in f for f in led[0]["failed_checks"]), led
+ # 减持是离场保护, 不追高闸只拦买入侧
+ assert ("000001.SZ", "TRIM") in {(x["ts_code"], x["action"]) for x in r["executed"]}, r
+
+
+@case("清仓覆盖冲突·一键清仓当场撤该票 ACTIVE 策略 (不等持仓归零) 并留痕")
+def _():
+ from app.services import command_service as csvc
+ 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": 10.0}])
+ fake.insert_strategy(strategy_id="STR_LQ", ts_code="600000.SH", stype="GRID")
+ fake.insert_proposal(proposal_id="PRP_LQ", ts_code="600000.SH", action="ADD",
+ qty=1000, hard_numbers={}, expire_at=None)
+ r = csvc.issue("LIQUIDATE_ALL", {"confirm": "YES"})
+ assert r["ok"], r
+ # 2026-08-19 实机: 清仓命令在追一个自己还在被网格买进的持仓。命令下达即撤, 不等归零。
+ assert fake.strategies["STR_LQ"]["status"] == "CANCELLED", fake.strategies
+ assert fake.proposals["PRP_LQ"]["status"] == "DECLINED", fake.proposals
+ prog = fake.commands[r["command_id"]]["progress"]
+ stopped = prog["stopped_buyside"]
+ assert stopped["strategies"] == ["STR_LQ"], stopped
+ assert any("清仓覆盖冲突" in n for n in prog["notes"]), prog["notes"]
+ assert any(x["action"] == "CLEANUP" and x["ref_id"] == r["command_id"]
+ for x in fake.ledger), "撤策略必须在评审账本留痕"
+
+
def main():
# 静音日志: 本套里有好几条用例**故意**触发异常与告警来验证「守成」行为
# (调度守卫吞异常、外部成交告警、连续对账升级 ERROR、窗口耗尽告警),
diff --git a/scripts/watch.py b/scripts/watch.py
index 3390851..7c5f255 100644
--- a/scripts/watch.py
+++ b/scripts/watch.py
@@ -94,7 +94,8 @@ def head(now):
flag += " ⛔ 全局暂停执行"
elif halt_b:
flag += " ⛔ 全局暂停买入"
- if brake:
+ from datetime import datetime as _dt
+ if brake and int(_dt.now().strftime("%Y%m%d")) < brake: # 到期的刹车不再常驻标题 (2026-08-28)
flag += f" ⛔ 刹车至 {brake}"
print(f" PMS 联测监视 · {now} · 下发 {mode} · 自主 {auto}{flag}")
_rule()
@@ -147,18 +148,35 @@ def account():
print(f" ⚠ 取不到现价: {', '.join(v['price_missing'][:6])}")
-def _left_today(ins: dict, used: int):
- """本日剩余额度 = 当日配额 − 今日已投放。用的是执行器那两个纯函数, 口径一致。"""
+def _left_today(ins: dict, kids: list, ymd: int):
+ """今日 (投, 余, 废) 三个数, **与执行器逐字同口径** (2026-08-28 修):
+ 紧急单与窗口末日单的「已投放」按在途量算 (_inflight_today), 普通多日单按
+ 已成交+在途算 (_consumed_today) —— 原来一律按后者, 紧急清仓当天有成交就显示
+ 「余 0」, 而执行器实际还会持续补单, 同一句话两种真相。"""
from app.core import exec_timing as et, tradedays as td
+ from app.core.sizer import lot_of
+ from app.services import executor
try:
remaining = max(0, int(ins.get("qty") or 0) - int(ins.get("exec_qty") or 0))
- dl = (ins.get("progress") or {}).get("deadline")
+ prog = ins.get("progress") or {}
+ dl = prog.get("deadline")
left_days = td.trade_days_left(dl, None) if dl else 1
- quota = et.daily_quota(remaining, left_days,
+ is_urgent = bool(prog.get("urgent"))
+ quota = et.daily_quota(remaining, left_days, lot=lot_of(ins.get("ts_code")),
allow_odd_tail=(ins.get("action") == "EXIT"))
- return max(0, quota - int(used))
+ try:
+ used = (executor._inflight_today(kids, ymd) if (is_urgent or left_days <= 1)
+ else executor._consumed_today(kids, ymd))
+ except Exception:
+ used = sum(int(c.get("qty") or 0) for c in kids)
+ try:
+ consumed = executor._consumed_today(kids, ymd)
+ except Exception:
+ consumed = used
+ void = sum(int(c.get("qty") or 0) for c in kids) - consumed
+ return used, max(0, quota - int(used)), void
except Exception:
- return "?"
+ return "?", "?", "?"
def instructions():
@@ -188,16 +206,8 @@ def instructions():
why = _short(d.get("reason") or "(本轮还没轮到它)", 24)
at = d.get("at") or ""
kids = [c for c in (prog.get("children") or []) if int(c.get("ymd") or 0) == ymd]
- try:
- used = executor._consumed_today(kids, ymd)
- except Exception:
- used = sum(int(c.get("qty") or 0) for c in kids)
- void = sum(int(c.get("qty") or 0) for c in kids) - used
- # 「余」必须**现算**, 不能读 last_decision 里的 qty_hint —— 那是执行器上一跳
- # 存下来的快照, 而投/废是此刻现算的。两者取自不同时刻, 会拼出自相矛盾的一行:
- # 2026-08-03 实机 002518.SZ 显示 `0/0/200` —— 废了 200 却还说余 0, 因为执行器
- # 那一跳时第二张分片还挂着(算已投), 等 watch 来看时它刚到期。三个数必须同源。
- left = _left_today(r, used)
+ # 「余」必须**现算**且三个数同源 (2026-08-03 教训), 口径委托给 _left_today
+ used, left, void = _left_today(r, kids, ymd)
q = f"{used}/{left}/{void}"
print(" " + _pad(r["ts_code"], 11)
+ _pad("买" if r["side"] == "buy" else "卖", 5)
@@ -266,7 +276,14 @@ def crosscheck():
return
try:
orders = qmt_repo.list_orders(limit=200)
- ins = {r["instruction_id"]: r for r in pms_repo.list_instructions(limit=200)}
+ # 父指令全集要**含已归档**且窗口放大 (2026-08-28 修): 用户把已完成指令「移除」
+ # (软归档) 后, 其委托仍在近 200 张出口行里 —— 原来按未归档近 200 条查, 正常
+ # 归档就触发「本端找不到父指令」的重号误报, 真告警反被当噪音。
+ try:
+ ins = {r["instruction_id"]: r
+ for r in pms_repo.list_instructions(limit=1000, include_archived=True)}
+ except TypeError:
+ ins = {r["instruction_id"]: r for r in pms_repo.list_instructions(limit=1000)}
except Exception:
return
bad = []
diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py
index 84efe8a..70ca82e 100644
--- a/scripts/ws_smoke.py
+++ b/scripts/ws_smoke.py
@@ -71,16 +71,19 @@ def cmd_status(args):
+ (f" · 挂起待人工 {ch['orphan_held']} 笔" if ch.get("orphan_held") else ""))
print(f" 出口队列 {ch['queue'] or '空'}")
# 快照是对账的事实源 (协议 §6.2), 拿不到就意味着账本建不起来 —— 必须一眼能看见
+ snap_err = False
try:
snap = qmt_repo.latest_snapshot("positions")
except Exception as e:
- snap = None
- print(f" 持仓快照 读取失败: {type(e).__name__}: {e}")
+ snap, snap_err = None, True
+ print(f" 持仓快照 读取失败: {type(e).__name__}: {e} —— 先修本端库/表, "
+ f"别按「对端没发」去排查")
if snap:
items = (snap["payload"] or {}).get("items")
print(f" 持仓快照 {len(items) if isinstance(items, list) else '?'} 个条目 · "
f"{snap['age_sec']:.0f}s 前 (seq={snap['seq']})")
- elif snap is None:
+ elif snap is None and not snap_err:
+ # 读取失败时不打这句 (2026-08-28 修): "读不到"与"对端从未发过"排查方向完全不同
print(f" 持仓快照 **从未收到** (对账没有事实源, 账本建不起来)")
s = st.get("stat") or {}
if s: