537 lines
30 KiB
Python
537 lines
30 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""
|
||
自主提议: 扫描 → 规则闸 → 研判闸 → 按自主档位分流 (设计 §6 / §7)
|
||
==================================================================
|
||
分流规则 (设计原文):
|
||
full 闸门与研判通过即执行 → 直接落指令
|
||
propose_only 增持类全部待用户确认 (一期默认) → 落 pms_proposal 队列
|
||
off 不扫描
|
||
|
||
两条无条件覆盖档位的规矩:
|
||
* **减持方向不设确认门槛** —— TRIM 保垫减仓属纯规则自动执行, 任何档位都直接落指令。
|
||
* **−15% 及更深的补仓永远需用户确认** —— 即便档位是 full, 也强制入队。
|
||
|
||
研判闸不可用时 (决策系统未接通/超时), 按设计自动降级为 propose_only + ERROR 告警,
|
||
**绝不把「研判拿不到」当成「研判通过」**。
|
||
|
||
2026-08-06 加了第五类动作「新建仓 OPEN」: 除了管已有持仓, 也从上游候选池里挑新票建底仓,
|
||
走的是同一条 规则闸 → 研判闸 → 档位分流。三处与已有四类不同的地方, 都在本文件里:
|
||
* 档位走自己的 `PMS_OPEN_AUTONOMY` (默认 full), 不跟随 `PMS_AUTONOMY` —— 「新建仓要不要
|
||
人点头」和「加仓要不要人点头」是两个决定。但 `PMS_AUTONOMY=off` 仍然是总闸, 它一关,
|
||
连新建仓一起停: off 的意思就是「自主动作全部停下」, 不该有子开关能绕过它。
|
||
* 研判驳回对新建仓做当日去重 (加仓类**有意**不去重, 那条纪律不动), 理由见
|
||
`pms_repo.judge_rejected_today`。
|
||
* 单轮扫描给研判一个时间预算, 用尽的候选留到下一跳 —— 见 `_judge_budget_left` 的说明。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import time
|
||
from datetime import datetime, timedelta
|
||
|
||
from app.core import action_engine as ae
|
||
from app.core import command_spec as cs
|
||
from app.core import rule_gate
|
||
from app.core import tradedays as td
|
||
from app.repo import pms_repo
|
||
from app.services import (command_service, executor, industry, judge, market, param_store,
|
||
plan_feed, portfolio)
|
||
|
||
logger = logging.getLogger("pms.proposal")
|
||
|
||
AUTONOMY_FULL, AUTONOMY_PROPOSE, AUTONOMY_OFF = "full", "propose_only", "off"
|
||
|
||
|
||
def scan_and_route(*, now=None, dry_run: bool = False) -> dict:
|
||
"""自主提议扫描一轮。dry_run=True 只出候选与判定, 不落任何表。"""
|
||
now = now or datetime.now()
|
||
out = {"ok": True, "autonomy": None, "open_autonomy": None, "candidates": 0,
|
||
"open_candidates": 0, "executed": [], "queued": [],
|
||
"rejected": [], "skipped": [], "errors": [], "degraded": False, "dry_run": dry_run}
|
||
|
||
autonomy = param_store.get("PMS_AUTONOMY", AUTONOMY_PROPOSE)
|
||
out["autonomy"] = autonomy
|
||
if autonomy == AUTONOMY_OFF:
|
||
out["skipped"].append({"why": "自主档位 off, 不扫描 (新建仓一并停 —— off 是总闸)"})
|
||
return out
|
||
if param_store.get_bool("PMS_GLOBAL_EXEC_HALT", False):
|
||
out["skipped"].append({"why": "全局暂停执行 (休假模式)"})
|
||
return out
|
||
open_autonomy = param_store.get("PMS_OPEN_AUTONOMY", AUTONOMY_FULL)
|
||
out["open_autonomy"] = open_autonomy
|
||
|
||
try:
|
||
view = portfolio.positions_view()
|
||
params = _scan_params(view)
|
||
mkt = _market_ctx(view["held"], now)
|
||
params["_mkt"] = mkt # 规则闸要用同一份 MA5, 不再重取
|
||
# 跳过三类: ①已有在途提议或指令的 ②今天已被规则闸拒过的 ③今天已被研判闸驳回的新建仓。
|
||
# ② 是 2026-07-29 的教训 —— 组合已超总仓上限时, 16 只深亏票的补仓候选每分钟被拒
|
||
# 一次, 一天往评审账本灌几千行一模一样的记录。闸门结论当天基本不会变, 记一次就够。
|
||
# ③ 只挡新建仓: 加仓类对研判驳回**有意**不去重 (契约里那句「研判结论会变」), 那条
|
||
# 纪律不动; 而候选池是几十只的量级, 不挡的话 ② 那个洞会原样从研判闸重来一遍。
|
||
skip = _inflight_keys() | _rejected_today_keys() | _judge_rejected_open_keys()
|
||
scanned = ae.scan(positions=view["held"], params=params, market=mkt, skip=skip)
|
||
except Exception as e:
|
||
logger.exception("提议扫描失败")
|
||
return {**out, "ok": False, "errors": [f"扫描失败: {type(e).__name__}: {e}"]}
|
||
|
||
out["candidates"] = len(scanned["candidates"])
|
||
out["skipped"].extend(scanned["skipped"])
|
||
stock_params = command_service.effective_stock_params()
|
||
brake_active = td.ymd() < param_store.get_int("PMS_BRAKE_UNTIL", 0)
|
||
|
||
# ---- 新建仓: 候选池那一路 (取不到候选池不影响上面四类, 反之亦然) ----
|
||
open_cands = []
|
||
if open_autonomy == AUTONOMY_OFF:
|
||
out["skipped"].append({"action": ae.A_OPEN, "why": "新建仓档位 off, 不扫描候选池"})
|
||
else:
|
||
try:
|
||
open_cands = _scan_open(view, params, stock_params, skip, mkt, out)
|
||
except Exception as e:
|
||
logger.exception("新建仓扫描失败")
|
||
out["errors"].append(f"新建仓扫描失败: {type(e).__name__}: {e}")
|
||
out["open_candidates"] = len(open_cands)
|
||
|
||
# 研判的时间预算从这一刻起算。**先跑已有持仓的四类, 再跑新建仓** —— 预算真用尽时,
|
||
# 被推到下一跳的一定是新建仓, 已有持仓的动作行为与 2026-08-06 之前一致。
|
||
deadline = time.monotonic() + max(
|
||
0, param_store.get_int("PMS_JUDGE_TICK_BUDGET_SEC", 150))
|
||
|
||
for c in list(scanned["candidates"]) + open_cands:
|
||
try:
|
||
_route_one(c, view, params, stock_params, brake_active, now, dry_run, out,
|
||
deadline=deadline)
|
||
except Exception as e:
|
||
logger.exception("提议分流失败 %s", c.get("ts_code"))
|
||
out["errors"].append(f"{c.get('ts_code')} {c.get('action')}: "
|
||
f"{type(e).__name__}: {e}")
|
||
out["ok"] = not out["errors"]
|
||
return out
|
||
|
||
|
||
# ================================================================ 新建仓的取数与筛选
|
||
def _scan_open(view, params, stock_params, skip, mkt, out) -> list:
|
||
"""候选池 → 新建仓候选。产出的候选数天生不超过剩余名额 (名额与金额在纯逻辑里边走边扣)。
|
||
|
||
与命令驱动那条路 (command_service._candidates) 有意不同的两点:
|
||
1. **只认上游计划接口这一个来源。** buy_plan 那张旧表在目标架构下没有明确的写入方,
|
||
白名单是给人下命令用的 —— 无人值守地建仓, 事实源必须单一。
|
||
2. **价格只认实时价。** 那边用 market.plan_price, 拿不到实时价会回落昨收; 它自己的
|
||
注释就写着「拿昨收当现价去做不追高这类判断会出错, 只给规划期定量用」。
|
||
这条路是要真下单的, 所以走 market.get_prices, 取不到就整只跳过并留痕。
|
||
"""
|
||
src = (param_store.get("PMS_CANDIDATE_SOURCE", plan_feed.SRC_PLAN_API)
|
||
or plan_feed.SRC_PLAN_API).strip()
|
||
if src == plan_feed.SRC_BUY_PLAN:
|
||
out["skipped"].append({"action": ae.A_OPEN,
|
||
"why": f"候选池来源是 {src}, 自主新建仓只认上游计划接口"})
|
||
return []
|
||
|
||
black = {c for c, d in (stock_params or {}).items() if d.get("black")}
|
||
held = [x["ts_code"] for x in view["held"]]
|
||
try:
|
||
sel = plan_feed.candidates(held=held, black=black)
|
||
except plan_feed.PlanFeedError as e:
|
||
# 拿不到 ≠ 今天没票可买。显式留痕, 本轮不产新建仓候选, 已有持仓的四类照常。
|
||
logger.error("[新建仓] 候选池取数失败, 本轮不建仓: %s", e)
|
||
out["skipped"].append({"action": ae.A_OPEN, "why": f"候选池取不到, 本轮不建仓: {e}"})
|
||
return []
|
||
items = list(sel.get("items") or [])
|
||
if not items:
|
||
out["skipped"].append({"action": ae.A_OPEN,
|
||
"why": f"候选池过滤后为空 (考察 {sel.get('considered')} 只, "
|
||
f"落选明细 {sel.get('dropped')})"})
|
||
return []
|
||
|
||
codes = [x["ts_code"] for x in items]
|
||
prices = market.get_prices(codes)
|
||
# 行业名一次批量取, 本轮复用。滚动扣减要算行业集中度, 所以必须先有它 ——
|
||
# caps_ctx 对**从没持仓过**的票带不出行业名, 不显式传的话 sizer.check_caps 遇到
|
||
# sector 为空会整段跳过, 行业集中度那道硬拦截就静默失效了。
|
||
# (industry.get_many 走 gp_stock_category 时是逐只查库、没有缓存; 给它加按日缓存能
|
||
# 省掉这几十次往返, 但那会改到既有函数的时序行为, 单独提、单独拍板, 这次不夹带。)
|
||
sectors = industry.get_many(codes) if codes else {}
|
||
# 今天被决策系统盘中判过转多的票 (signal_service 落的留痕)。**只用来排序, 不改资格**:
|
||
# 不在候选池里的票不会因为有信号就被建仓, 候选层那一整套过滤一道都不绕。
|
||
# 它的增量是实打实的 —— 候选按分数降序取, 名额只剩两个时第 25 名永远轮不上;
|
||
# 而「此刻转多」是盘中才有的新信息, 昨夜算出来的分数与买入区间都表达不了它。
|
||
sig_buy = _buy_signals_today()
|
||
cands = [{**x, "price": prices.get(x["ts_code"]), "sector": sectors.get(x["ts_code"]),
|
||
"sig_buy": sig_buy.get(x["ts_code"])} for x in items]
|
||
hit = [c["ts_code"] for c in cands if c.get("sig_buy")]
|
||
if hit:
|
||
logger.info("[新建仓] 候选池里今天被决策系统判过转多的: %s", hit)
|
||
|
||
t = view["totals"]
|
||
p = view["params"]
|
||
slots = int(p["max_names"] or 0) - int(t["names_count"] or 0)
|
||
room = float(p["portfolio_cap"] or 0) * float(p["scale"] or 0) - float(t["portfolio_mv"] or 0)
|
||
# 真实可用资金封顶。**这一条比规则闸严一档, 是刻意的**: 规则闸那条「拿不到 ws 资金快照
|
||
# 只告警不拦」是为**已经排好的命令**设计的 —— 通道故障不该升级成业务停摆。而无人值守地
|
||
# 从零建仓完全可以等一等, 07-30 那次 (scale 200 万 / 账户实际 98 万, 方案一路放行到
|
||
# 下游才被拒, 而拒了不自动重发) 不该换个入口重演。
|
||
if str(t.get("cash_source") or "") == portfolio.CASH_WS and t.get("cash_avail") is not None:
|
||
room = min(room, float(t["cash_avail"]))
|
||
elif param_store.get_bool("PMS_OPEN_REQUIRE_WS_CASH", True):
|
||
out["skipped"].append({"action": ae.A_OPEN,
|
||
"why": f"拿不到 ws 资金快照 ({t.get('cash_why') or '原因未知'}), "
|
||
f"本轮不自动新建仓 (PMS_OPEN_REQUIRE_WS_CASH=True)"})
|
||
return []
|
||
|
||
res = ae.scan_open(candidates=cands, params=params, caps=portfolio.caps_ctx(view),
|
||
room_amt=room, slots=slots, skip=skip)
|
||
out["skipped"].extend(res["skipped"])
|
||
|
||
# 当日行情快照只对**真的产出了候选**的票取 —— 候选数已被名额与金额扣到很小,
|
||
# 不必为整池几十只票各拉一遍分钟线。规则闸的「不追高」靠的就是这一份。
|
||
for c in res["candidates"]:
|
||
code = c["ts_code"]
|
||
d = {"ma5": market.get_ma5(code)}
|
||
try:
|
||
d["day"] = market.day_snapshot(code) or {}
|
||
except Exception as e:
|
||
logger.warning("[新建仓] 取当日快照失败 %s: %s", code, e)
|
||
d["day"] = {}
|
||
mkt[code] = d
|
||
return res["candidates"]
|
||
|
||
|
||
def _route_one(c, view, params, stock_params, brake_active, now, dry_run, out,
|
||
*, deadline=None):
|
||
code, action, side = c["ts_code"], c["action"], c["side"]
|
||
is_open = (action == ae.A_OPEN)
|
||
pos = _pos_of(view, code)
|
||
# 新建仓的现价取候选自带的那份。**不能取 pos 的**: 从没持仓过的票在账本里没有行,
|
||
# `_pos_of` 回的是空壳、没有 price, 取它会是 0, 规则闸一上来就判 PRICE_MISSING。
|
||
price = float(c.get("price") or 0) if is_open else float(pos.get("price") or 0)
|
||
|
||
# ---- 一级: 规则闸 (自主动作受刹车约束, is_command=False) ----
|
||
gate = rule_gate.check(
|
||
side=side, action=action, qty=c["qty"], price=price,
|
||
ctx={"ts_code": code, "position": pos,
|
||
# 当日行情走 _market_ctx 取到的真实快照。原来这里是
|
||
# {"vwap": price, "day_chg_from_open": None}
|
||
# 两个都是编的: vwap 拿现价顶等于"现价恰好等于均价", day_chg 给 None 让
|
||
# 「不追高」整道闸静默跳过。取不到就留空, 让规则闸记 *_MISSING 警告。
|
||
"day": {**((params.get("_mkt") or {}).get(code, {}).get("day") or {}),
|
||
"price": price,
|
||
"ma5": (params.get("_mkt") or {}).get(code, {}).get("ma5")},
|
||
"params": {"no_chase_ma5": params.get("no_chase_ma5"),
|
||
"buy_halt_dayup": params.get("buy_halt_dayup"),
|
||
"sector_source_ready": view["sector_ready"]},
|
||
# 新建仓必须**显式**把 is_new_name 与行业名传进去。caps_ctx 只在
|
||
# view["positions"] 里找得到该票时才带得出行业, 而新票根本不在里面 ——
|
||
# 不传的话 sizer.check_caps 遇到 sector 为空会整段跳过, 行业集中度那道硬拦截
|
||
# 就静默失效了 (不报错、不少数据, 就是不生效)。
|
||
"caps": (portfolio.caps_ctx(view, ts_code=code, is_new_name=True,
|
||
sector=c.get("sector")) if is_open
|
||
else (portfolio.caps_ctx(view, ts_code=code) if side == "buy" else None)),
|
||
"flags": {"buy_halt": params.get("buy_halt"), "exec_halt": params.get("exec_halt"),
|
||
"brake_active": brake_active,
|
||
"blacklisted": bool(stock_params.get(code, {}).get("black")),
|
||
"is_command": False}})
|
||
if not gate["passed"]:
|
||
out["rejected"].append({"ts_code": code, "action": action, "by": "rule",
|
||
"failed": gate["failed"]})
|
||
if not dry_run:
|
||
pms_repo.insert_ledger(ts_code=code, action=action, arbiter="rule", verdict="REJECT",
|
||
price_at=price, hard_numbers=c["hard_numbers"],
|
||
failed_checks=gate["failed"], reason="自主提议未过规则闸")
|
||
return
|
||
|
||
# ---- 二级: 研判闸 (补足/加仓/补仓/新建仓; 减持不送研判) ----
|
||
verdict = {"verdict": judge.PASS, "reason": "", "degraded": False}
|
||
if c.get("judge_required"):
|
||
if not _judge_budget_left(deadline):
|
||
# 本轮研判时间用尽 —— 整条跳过, 下一跳重来。
|
||
# **不当成「研判不可用」入人工队列**: 那会在自动档位下凭空造出一个人工确认队列,
|
||
# 与「不用每天靠人」这件事正好相反。规则闸白跑一次的成本是毫秒级。
|
||
out["skipped"].append({"ts_code": code, "action": action,
|
||
"why": "本轮研判时间预算用尽, 下一跳继续 "
|
||
f"(PMS_JUDGE_TICK_BUDGET_SEC="
|
||
f"{param_store.get_int('PMS_JUDGE_TICK_BUDGET_SEC', 150)}s)"})
|
||
return
|
||
verdict = judge.request(c, context={"position": _judge_ctx(pos),
|
||
"recent_ledger": _recent_ledger(code)})
|
||
if verdict["verdict"] == judge.REJECT:
|
||
out["rejected"].append({"ts_code": code, "action": action, "by": "judge",
|
||
"failed": [verdict.get("reason") or "研判驳回"]})
|
||
if not dry_run:
|
||
pms_repo.insert_ledger(ts_code=code, action=action, arbiter="judge",
|
||
verdict="REJECT", price_at=price,
|
||
hard_numbers=c["hard_numbers"],
|
||
reason=verdict.get("reason") or "研判驳回")
|
||
return
|
||
if verdict.get("degraded"):
|
||
out["degraded"] = True
|
||
|
||
# ---- 三级: 按档位分流 ----
|
||
# 新建仓走自己的档位 (PMS_OPEN_AUTONOMY), 不跟随全局 —— 「新建仓要不要人点头」和
|
||
# 「加仓要不要人点头」是两个不同的决定。全局 off 已经在 scan_and_route 入口拦掉了,
|
||
# 走到这里说明总闸是开的。三条无条件覆盖档位的规矩对新建仓一样有效: 减持自动、
|
||
# 深档补仓强制确认、研判不可用一律入队。
|
||
autonomy = (out.get("open_autonomy") or AUTONOMY_FULL) if is_open else out["autonomy"]
|
||
force_queue = bool(c.get("needs_user_confirm")) or verdict.get("degraded")
|
||
auto_exec = (side == "sell") or (autonomy == AUTONOMY_FULL and not force_queue)
|
||
|
||
if dry_run:
|
||
(out["executed"] if auto_exec else out["queued"]).append(
|
||
{**_brief(c), "route": "auto" if auto_exec else "queue",
|
||
"judge": verdict["verdict"], "dry_run": True})
|
||
return
|
||
|
||
if auto_exec:
|
||
iid = _make_instruction(c, price, now)
|
||
pms_repo.insert_ledger(ts_code=code, action=action,
|
||
arbiter="judge" if c.get("judge_required") else "rule",
|
||
verdict="PASS", price_at=price, hard_numbers=c["hard_numbers"],
|
||
ref_id=iid,
|
||
reason=(verdict.get("reason") or c["reason"])[:500])
|
||
out["executed"].append({**_brief(c), "instruction_id": iid,
|
||
"why": ("减持方向自动执行" if side == "sell"
|
||
else ("新建仓档位 full" if is_open else "档位 full"))})
|
||
else:
|
||
pid = _make_proposal(c, price, verdict)
|
||
why = ("深档补仓强制确认" if c.get("needs_user_confirm")
|
||
else ("研判不可用, 降级人工确认" if verdict.get("degraded")
|
||
else ("新建仓档位 propose_only" if is_open else "档位 propose_only")))
|
||
out["queued"].append({**_brief(c), "proposal_id": pid, "why": why})
|
||
|
||
|
||
def _judge_budget_left(deadline) -> bool:
|
||
"""这一轮还够不够再送一次研判。
|
||
|
||
**这不是节流, 是让一次心跳做得完。** judge.request 是同步阻塞的, 单次上限
|
||
PMS_JUDGE_TIMEOUT (默认 90 秒), 而 scheduler 给所有调度任务设的软超时是 240 秒 ——
|
||
三只票送研判就顶破了, 任务被 celery 打死在中途, 而且是在已经落了一部分表之后。
|
||
以前送研判的只有已有持仓那几只、多数轮次还被去重挡掉, 所以一直没撞上; 新建仓上线后
|
||
冷启动那天会有十来条候选, 第一跳就会捅穿。
|
||
|
||
预算不够时调用方整条跳过并留痕, 下一分钟的心跳接着做 —— 一条候选都不丢, 也没有任何
|
||
按天计的上限。决策系统那侧的裁决按「日期+股票+动作」缓存半小时, 所以真正慢的只有
|
||
缓存过期后的第一跳, 之后同一批候选都是秒回。
|
||
"""
|
||
if deadline is None:
|
||
return True
|
||
need = max(1, param_store.get_int("PMS_JUDGE_TIMEOUT", 90))
|
||
return (time.monotonic() + need) <= deadline
|
||
|
||
|
||
# ================================================================ 落表
|
||
def _make_instruction(c, price, now) -> str:
|
||
ymd = td.ymd(now)
|
||
seq = int(now.strftime("%H%M%S"))
|
||
iid = cs.make_instruction_id(ymd, c["ts_code"], c["action"], seq % 1000)
|
||
window = param_store.get_int("PMS_EXEC_WINDOW_TDAYS", 3)
|
||
pms_repo.insert_instruction(
|
||
instruction_id=iid, origin_type="proposal", origin_id=None, ts_code=c["ts_code"],
|
||
action=c["action"], side=c["side"], qty=c["qty"], limit_price=None,
|
||
window_tdays=window, status=executor.ST_PROPOSED,
|
||
progress={"deadline": str(td.window_deadline(now.date(), window)),
|
||
"is_command": False, "children": [], "auto": True,
|
||
"reason": c["reason"]})
|
||
g = bump_once_guards(c["ts_code"], c["action"], c.get("hard_numbers"), now) or {}
|
||
if not g.get("ok"):
|
||
# 指令已经落表了, 不回滚; 但要在评审账本上留一条痕, 否则这条纪律失效没有任何记录
|
||
try:
|
||
pms_repo.insert_ledger(
|
||
ts_code=c["ts_code"], action=c["action"], arbiter="rule", verdict="WARN",
|
||
price_at=float(price or 0), ref_id=iid,
|
||
hard_numbers={"once_guard_fields": {k: str(v) for k, v in
|
||
(g.get("fields") or {}).items()}},
|
||
reason=f"一次性守卫计数器未写入 ({g.get('error')}) —— "
|
||
f"该票 {c['action']} 的「只做一次」本轮失效, 留意重复出手")
|
||
except Exception:
|
||
logger.exception("一次性守卫失败留痕也没写上 %s", c["ts_code"])
|
||
return iid
|
||
|
||
|
||
def bump_once_guards(ts_code: str, action: str, hard_numbers=None, now=None):
|
||
"""把「只做一次」的三个计数器写上。
|
||
|
||
设计 §6 写了三条一次性约束, 但它们的计数器此前**只被读、从没被写过** —— 也就是说这三条
|
||
纪律一直是失效的, 只是被「同一票同一动作有在途提议就不重复提」这条兜底遮住了, 而那条兜底
|
||
恰好在规则闸拒绝时失灵 (没生成提议 → 没东西可去重), 于是同一个候选每分钟重来一次:
|
||
|
||
fill_count 回踩补足「每票 1 次」
|
||
last_add_date 盈利加仓「距上次 ≥2 交易日」—— 恒 None 时永远算作"很久没加过"
|
||
dca_count 补仓「各档评估一次」—— 恒 0 时每轮都当第一次评估
|
||
|
||
写入时机取「动作真的要落地」这一刻 (落指令), 而不是产出候选那一刻: 候选被闸门拦下不算
|
||
做过, 用掉一次名额不合理。dca_count 记的是**档位**而不是次数, 所以取本次触及的档序。
|
||
"""
|
||
now = now or datetime.now()
|
||
fields = {}
|
||
if action == "FILL":
|
||
fields["fill_count"] = 1
|
||
elif action == "ADD":
|
||
fields["last_add_date"] = now.date()
|
||
elif action == "DCA":
|
||
fields["dca_count"] = int((hard_numbers or {}).get("stage") or 1)
|
||
if not fields:
|
||
return {"ok": True, "fields": {}}
|
||
try:
|
||
n = pms_repo.update_position(ts_code, **fields)
|
||
except Exception as e:
|
||
# 写不上不该把已落表的指令带崩, 但**必须回成败**: 计数器没写上, 这条一次性纪律
|
||
# 就退回失效状态, 同一票同一动作在首条指令离开 LIVE 之后会再来一次。
|
||
# (2026-07-31 静默失败专项: 原来这里 return None, 调用方拿不到任何信号)
|
||
logger.error("一次性守卫计数器写入失败 %s %s: %s —— 该票这条一次性纪律本轮失效",
|
||
ts_code, action, e)
|
||
return {"ok": False, "fields": fields, "error": f"{type(e).__name__}: {e}"}
|
||
if n == 0:
|
||
logger.error("一次性守卫计数器没写到任何行 %s %s (持仓行不存在?) —— "
|
||
"该票这条一次性纪律本轮失效", ts_code, action)
|
||
return {"ok": False, "fields": fields, "error": "持仓行不存在, 影响 0 行"}
|
||
return {"ok": True, "fields": fields}
|
||
|
||
|
||
def _make_proposal(c, price, verdict) -> str:
|
||
ttl = param_store.get_int("PMS_PROPOSAL_TTL_HOURS", 24)
|
||
pid = f"PRP_{td.ymd()}_{c['ts_code'].replace('.', '')}_{c['action']}"
|
||
hn = {**(c.get("hard_numbers") or {}), "price": price, "reason": c["reason"],
|
||
"needs_user_confirm": c.get("needs_user_confirm", False)}
|
||
pms_repo.insert_proposal(
|
||
proposal_id=pid, ts_code=c["ts_code"], action=c["action"], qty=c["qty"],
|
||
hard_numbers=hn, expire_at=datetime.now() + timedelta(hours=ttl),
|
||
judge_verdict=verdict.get("verdict"),
|
||
judge_reason=(verdict.get("reason") or "")[:500])
|
||
return pid
|
||
|
||
|
||
# ================================================================ 上下文
|
||
def _scan_params(view: dict) -> dict:
|
||
p = dict(view["params"])
|
||
p.update({
|
||
"cushion_solid": param_store.get_float("PMS_CUSHION_SOLID", 0.03),
|
||
"trim_peak": param_store.get_float("PMS_TRIM_PEAK", 0.06),
|
||
"trim_giveback": param_store.get_float("PMS_TRIM_GIVEBACK", 0.5),
|
||
"dca_triggers": param_store.get_tuple_floats("PMS_DCA_TRIGGERS", (-0.08, -0.15)),
|
||
"dca_deep_confirm": param_store.get_float("PMS_DCA_DEEP_CONFIRM", -0.15),
|
||
"dca_max_ratio": param_store.get_float("PMS_DCA_MAX_RATIO", 0.5),
|
||
"no_chase_ma5": param_store.get_float("PMS_NO_CHASE_MA5", 0.06),
|
||
"buy_halt_dayup": param_store.get_float("PMS_BUY_HALT_DAYUP", 0.05),
|
||
"build_window_tdays": param_store.get_int("PMS_BUILD_WINDOW_TDAYS", 10),
|
||
"fill_max_loss": param_store.get_float("PMS_FILL_MAX_LOSS", -0.03),
|
||
"open_signal_priority": param_store.get_bool("PMS_OPEN_SIGNAL_PRIORITY", True),
|
||
})
|
||
return p
|
||
|
||
|
||
def _market_ctx(held: list, now) -> dict:
|
||
"""每票的 MA5 / 5日高点 / 建仓天数 / 距上次加仓天数 (交易日口径)。"""
|
||
out = {}
|
||
today = now.date() if hasattr(now, "date") else now
|
||
for p in held or []:
|
||
code = p["ts_code"]
|
||
d = {"ma5": market.get_ma5(code), "high5": market.get_high5(code)}
|
||
# 当日行情快照 —— **规则闸的「不追高(当日涨幅)」全靠它**。
|
||
# 原来这里没取, `_route_one` 给规则闸硬编码 `day_chg_from_open: None`, 而规则闸对
|
||
# None 是整道跳过 → 动作引擎产出的每一笔自主买单 (FILL/ADD/DCA) 从来没过过这道闸,
|
||
# 账本里也查不到痕迹。executor 那条路一直是拿 day_snapshot 的, 只有这里漏了。
|
||
# 取不到就留空, 由规则闸记 DAYUP_MISSING 警告 —— 不再无声无息。
|
||
try:
|
||
d["day"] = market.day_snapshot(code) or {}
|
||
except Exception as e: # 行情不可用不该让整轮扫描崩掉
|
||
logger.warning("[提议] 取当日快照失败 %s: %s", code, e)
|
||
d["day"] = {}
|
||
opened, last_add = p.get("opened_date"), p.get("last_add_date")
|
||
d["tdays_since_open"] = _tdays_between(opened, today)
|
||
d["tdays_since_last_add"] = _tdays_between(last_add, today)
|
||
out[code] = d
|
||
return out
|
||
|
||
|
||
def _tdays_between(start, today):
|
||
"""start(含) 到 today 的交易日数; start 为空返回 None (调用方按「无约束」处理)。"""
|
||
if not start:
|
||
return None
|
||
try:
|
||
return max(0, td.trade_days_left(today, start) - 1)
|
||
except Exception:
|
||
return None
|
||
|
||
|
||
def _rejected_today_keys() -> set:
|
||
"""今天已被规则闸拒过的 (代码, 动作)。读不到就返回空集 —— 去重是降噪, 不是纪律,
|
||
读失败时宁可多记几行日志, 也不能因此漏扫一个本该评估的动作。"""
|
||
try:
|
||
return pms_repo.rule_rejected_today(datetime.now().replace(
|
||
hour=0, minute=0, second=0, microsecond=0))
|
||
except Exception as e:
|
||
logger.warning("读当日规则闸拒绝记录失败 (按未拒过继续扫描): %s", e)
|
||
return set()
|
||
|
||
|
||
def _buy_signals_today() -> dict:
|
||
"""今天的盘中转多留痕。读不到就返回空 —— 它只影响候选的先后, 不影响资格,
|
||
所以读失败时降级成「按分数排」就够了, 不该因此让整轮新建仓停摆。"""
|
||
try:
|
||
return pms_repo.buy_signals_today(datetime.now().replace(
|
||
hour=0, minute=0, second=0, microsecond=0))
|
||
except Exception as e:
|
||
logger.warning("[新建仓] 读当日转多留痕失败 (本轮按纯分数排序): %s", e)
|
||
return {}
|
||
|
||
|
||
def _judge_rejected_open_keys() -> set:
|
||
"""今天已被研判闸驳回的**新建仓** (代码, 动作)。读不到就返回空集 (同上: 去重是降噪)。
|
||
|
||
只取 OPEN 那些键。加仓类对研判驳回**有意**不做当日去重 —— 契约里那句「研判结论会变」,
|
||
节流责任放在决策系统那侧的半小时缓存上, 这条纪律一个字不动。
|
||
"""
|
||
try:
|
||
keys = pms_repo.judge_rejected_today(datetime.now().replace(
|
||
hour=0, minute=0, second=0, microsecond=0))
|
||
except Exception as e:
|
||
logger.warning("读当日研判驳回记录失败 (按未驳回过继续扫描): %s", e)
|
||
return set()
|
||
return {(code, act) for code, act in keys if act == ae.A_OPEN}
|
||
|
||
|
||
def _inflight_keys() -> set:
|
||
"""已有在途提议或在途指令的 (代码, 动作) —— 同一件事不重复提。"""
|
||
keys = set()
|
||
try:
|
||
for p in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200):
|
||
keys.add((p["ts_code"], p["action"]))
|
||
except Exception as e:
|
||
logger.warning("读提议队列失败: %s", e)
|
||
try:
|
||
for i in pms_repo.list_instructions(statuses=list(executor.LIVE), limit=300):
|
||
keys.add((i["ts_code"], i.get("action")))
|
||
except Exception as e:
|
||
logger.warning("读在途指令失败: %s", e)
|
||
return keys
|
||
|
||
|
||
def _pos_of(view: dict, ts_code: str) -> dict:
|
||
for x in view["positions"]:
|
||
if x["ts_code"] == ts_code:
|
||
return x
|
||
return {"ts_code": ts_code, "total_qty": 0, "avail_qty": 0, "frozen_reason": "NONE"}
|
||
|
||
|
||
def _judge_ctx(pos: dict) -> dict:
|
||
"""送研判的账本快照 (设计 §7: PMS 备齐硬数字与账本流水作为 context)。"""
|
||
keep = ("ts_code", "price", "avg_cost", "total_qty", "base_qty", "add_qty", "dca_qty",
|
||
"cushion_pct", "cushion_peak", "pct_of_scale", "target_pct", "support_ref",
|
||
"pressure_ref", "stop_ref", "ref_source", "sector", "neg_cushion_days")
|
||
return {k: pos.get(k) for k in keep}
|
||
|
||
|
||
def _recent_ledger(ts_code: str, limit: int = 10) -> list:
|
||
try:
|
||
rows = pms_repo.list_ledger(ts_code=ts_code, limit=limit)
|
||
except Exception:
|
||
return []
|
||
return [{"at": str(r.get("decided_at")), "action": r.get("action"),
|
||
"arbiter": r.get("arbiter"), "verdict": r.get("verdict"),
|
||
"price": r.get("price_at"), "reason": r.get("reason")} for r in rows]
|
||
|
||
|
||
def _brief(c: dict) -> dict:
|
||
return {"ts_code": c["ts_code"], "action": c["action"], "side": c["side"],
|
||
"qty": c["qty"], "reason": c["reason"]}
|