436 lines
20 KiB
Python
436 lines
20 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""技术面取数、落表、映射与对外状态 (2026-09-11 技术面接入)。
|
||
|
||
职责三件:
|
||
一, 分页拉决策系统 GET /api/v1/market/technical 的全市场读数, 宽松解析后落进
|
||
pms_tech_daily。**状态非 OK 不落表**; 分页里任一页失败不推进 data_date, 整轮算失败,
|
||
绝不落半份 (先记账后动作, 宁可这轮没有也不要落一份残缺的)。
|
||
二, 每早对相关票 (在持 + 计划主榜观察 + 待拍板) 用 core.tech_rules 合成技术面立场, 写进
|
||
运行参数 PMS_TECH_STATE_MAP。与 logic_state_service 的映射同一手法: 读不到留空带原因、
|
||
超 FRESH_DAYS 天按无读数、绝不折成看空 (设计原则二)。
|
||
三, 给页面出状态 (开关/新鲜度/行数/映射覆盖/错误) 与按票的研究面读数 (点击时调)。
|
||
|
||
取数纯函数 (fetch_page / fetch_all_pages / parse_item) 不读参数中心, 以便脱库单测。
|
||
总开关 PMS_TECH_ENABLED 关掉即整模块歇工 (不拉、不建映射、状态标注已停, 行为逐字如旧)。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
from collections import Counter
|
||
from datetime import date, datetime
|
||
|
||
from app.core import tech_rules, tradedays
|
||
from app.core.command_spec import normalize_code
|
||
from app.repo import tech_repo
|
||
from app.services import param_store
|
||
|
||
logger = logging.getLogger("pms.tech")
|
||
|
||
MAP_KEY = "PMS_TECH_STATE_MAP"
|
||
FRESH_DAYS = 3 # 映射超过这么多自然日没刷新就当没读数 (调度断了不拿旧态拦人)
|
||
DEFAULT_PATH = "/api/v1/market/technical"
|
||
PAGE_SIZE = 1000 # 分页每页条数 (接口 limit 上限实测 1000, 设成一样避免误判最后一页)
|
||
MAX_PAGES = 20 # 防呆: 分页硬上限
|
||
|
||
|
||
class TechFeedError(RuntimeError):
|
||
"""拉技术面失败。调用方按无读数降级, 不产生任何拦截。"""
|
||
|
||
|
||
# ================================================================ 参数
|
||
def _getf(key: str, d: float) -> float:
|
||
v = param_store.get(key, None)
|
||
try:
|
||
return float(v) if v not in (None, "") else float(d)
|
||
except (TypeError, ValueError):
|
||
return float(d)
|
||
|
||
|
||
def _params() -> dict:
|
||
g, gi, gb = param_store.get, param_store.get_int, param_store.get_bool
|
||
base = (g("PMS_TECH_API_BASE", "") or "").strip().rstrip("/")
|
||
if not base: # 留空沿用研判接口根地址 (同一台 bionic)
|
||
base = (g("PMS_JUDGE_API_BASE", "") or "").strip().rstrip("/")
|
||
return {
|
||
"enabled": gb("PMS_TECH_ENABLED", True),
|
||
"base": base,
|
||
"path": (g("PMS_TECH_API_PATH", DEFAULT_PATH) or DEFAULT_PATH).strip(),
|
||
"timeout": gi("PMS_TECH_TIMEOUT", 20),
|
||
"keep_days": gi("PMS_TECH_KEEP_DAYS", 40),
|
||
"stale_tdays": gi("PMS_TECH_STALE_TDAYS", 2),
|
||
"rules": {
|
||
"squeeze_lookback": gi("PMS_TECH_SQUEEZE_LOOKBACK", 5),
|
||
"open_bw_growth": _getf("PMS_TECH_OPEN_BW_GROWTH", 0.20),
|
||
"choppy_flips": gi("PMS_TECH_CHOPPY_FLIPS", 4),
|
||
"flip_fresh_days": gi("PMS_TECH_FLIP_FRESH_DAYS", 2),
|
||
},
|
||
}
|
||
|
||
|
||
def enabled() -> bool:
|
||
return param_store.get_bool("PMS_TECH_ENABLED", True)
|
||
|
||
|
||
# ================================================================ 解析 (纯函数)
|
||
def _f(v):
|
||
try:
|
||
return None if v in (None, "") else float(v)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
|
||
|
||
def _i(v):
|
||
try:
|
||
return None if v in (None, "") else int(float(v))
|
||
except (TypeError, ValueError):
|
||
return None
|
||
|
||
|
||
def _s(v):
|
||
if v is None:
|
||
return None
|
||
s = str(v).strip()
|
||
return s or None
|
||
|
||
|
||
def _tinyint(v):
|
||
"""接口布尔 → 1/0; None 保留 None (缺字段不误判成 0)。"""
|
||
if v is None:
|
||
return None
|
||
if isinstance(v, str):
|
||
return 1 if v.strip() in ("1", "true", "True", "是") else 0
|
||
return 1 if v else 0
|
||
|
||
|
||
def parse_item(it: dict, data_date: int, algo_version=None):
|
||
"""接口一只票 → 落表行 dict。代码归一 (SH600000 → 600000.SH); 认不出返回 None。"""
|
||
if not isinstance(it, dict):
|
||
return None
|
||
code = normalize_code(str(it.get("stock_code") or it.get("code") or it.get("ts_code") or ""))
|
||
if not code or "." not in code:
|
||
return None
|
||
boll = it.get("boll") if isinstance(it.get("boll"), dict) else {}
|
||
bbi = it.get("bbiboll") if isinstance(it.get("bbiboll"), dict) else {}
|
||
sar = it.get("sar") if isinstance(it.get("sar"), dict) else {}
|
||
return {
|
||
"data_date": int(data_date), "ts_code": code,
|
||
"stock_name": _s(it.get("stock_name") or it.get("name")),
|
||
"boll_upper": _f(boll.get("upper")), "boll_mid": _f(boll.get("mid")),
|
||
"boll_lower": _f(boll.get("lower")), "boll_bw_pct": _f(boll.get("bandwidth_pct")),
|
||
"boll_pos": _f(boll.get("pos")), "boll_state": _s(boll.get("state")),
|
||
"boll_squeeze": _tinyint(boll.get("squeeze")),
|
||
"bbi": _f(bbi.get("bbi")), "bbi_upper": _f(bbi.get("upper")),
|
||
"bbi_lower": _f(bbi.get("lower")), "bbi_pos": _f(bbi.get("pos")),
|
||
"bbi_state": _s(bbi.get("state")), "bbi_dist_pct": _f(bbi.get("dist_pct")),
|
||
"sar_value": _f(sar.get("value")), "sar_side": _s(sar.get("side")),
|
||
"sar_flip_days": _i(sar.get("flip_days")), "sar_dist_pct": _f(sar.get("dist_pct")),
|
||
"reanchored": _tinyint(it.get("reanchored")), "in_pool": _tinyint(it.get("in_pool")),
|
||
"quality": _s(it.get("quality")), "bars_used": _i(it.get("bars_used")),
|
||
"last_bar_date": _i(it.get("last_bar_date")), "bars_lag": _i(it.get("bars_lag")),
|
||
"algo_version": _s(algo_version),
|
||
}
|
||
|
||
|
||
# ================================================================ 取数 (纯函数, 可注入 fetch)
|
||
def _http_get(url, params, timeout):
|
||
import requests
|
||
r = requests.get(url, params=(params or None), timeout=timeout)
|
||
r.raise_for_status()
|
||
return r.json()
|
||
|
||
|
||
def fetch_page(*, base, path, timeout, limit, offset, fetch=None) -> dict:
|
||
"""拉一页。fetch 可注入 (单测)。失败或应答非对象抛 TechFeedError。"""
|
||
base = (base or "").strip().rstrip("/")
|
||
if not base:
|
||
raise TechFeedError("技术面接口地址为空(PMS_TECH_API_BASE 与 PMS_JUDGE_API_BASE 都没填)")
|
||
path = (path or DEFAULT_PATH).strip()
|
||
if not path.startswith("/"):
|
||
path = "/" + path
|
||
url = base + path
|
||
try:
|
||
payload = (fetch or _http_get)(url, {"limit": int(limit), "offset": int(offset)}, int(timeout))
|
||
except Exception as e:
|
||
logger.warning("拉技术面失败 %s offset=%s: %s: %s", url, offset, type(e).__name__, e)
|
||
raise TechFeedError(f"{type(e).__name__}: {e}") from e
|
||
if not isinstance(payload, dict):
|
||
raise TechFeedError(f"应答不是 JSON 对象: {type(payload).__name__}")
|
||
return payload
|
||
|
||
|
||
def fetch_all_pages(*, base=None, path=None, timeout=None, page_size=PAGE_SIZE, fetch=None) -> dict:
|
||
"""分页拉全市场。返回 {status, data_date, algo_version, rows, pages, total, matched}。
|
||
任一页 status 非 OK 或抛错 → 整轮失败 (抛 TechFeedError), 不返回半份。"""
|
||
if base is None or path is None or timeout is None:
|
||
p = _params()
|
||
base = p["base"] if base is None else base
|
||
path = p["path"] if path is None else path
|
||
timeout = p["timeout"] if timeout is None else timeout
|
||
rows, offset, pages, matched = [], 0, 0, None
|
||
data_date = algo = None
|
||
while pages < MAX_PAGES:
|
||
payload = fetch_page(base=base, path=path, timeout=timeout,
|
||
limit=page_size, offset=offset, fetch=fetch)
|
||
status = str(payload.get("status") or "")
|
||
if status != "OK":
|
||
raise TechFeedError(f"接口状态非 OK: {status or '缺'}")
|
||
if data_date is None:
|
||
data_date = _i(payload.get("data_date"))
|
||
algo = _s(payload.get("algo_version"))
|
||
matched = _i(payload.get("matched"))
|
||
items = payload.get("items") if isinstance(payload.get("items"), list) else []
|
||
for it in items:
|
||
row = parse_item(it, data_date, algo)
|
||
if row:
|
||
rows.append(row)
|
||
pages += 1
|
||
got = len(items)
|
||
offset += got
|
||
if got == 0: # 空页, 到底了
|
||
break
|
||
if matched is not None and offset >= matched: # 拉够 matched 只
|
||
break
|
||
# 不能用 got < page_size 判最后一页: 接口 limit 上限可能小于 page_size, 那样每页都返回
|
||
# 上限条数、got 恒小于 page_size, 第一页就误判成最后一页 (2026-09-11 真机: 上限 1000、
|
||
# matched 4991, 旧逻辑只落了 1000/4991)。matched 缺失时才退回用 got 判, 另有 MAX_PAGES 兜底。
|
||
if matched is None and got < page_size:
|
||
break
|
||
return {"status": "OK", "data_date": data_date, "algo_version": algo,
|
||
"rows": rows, "pages": pages, "total": len(rows), "matched": matched}
|
||
|
||
|
||
# ================================================================ 落表
|
||
def pull(*, fetch=None) -> dict:
|
||
"""拉一轮全市场 + 落表 + 删旧。返回读数供调度日志。失败不落表、不推进日期。"""
|
||
out = {"ok": True, "at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||
"data_date": None, "rows": 0, "pruned": 0, "error": None}
|
||
p = _params()
|
||
if not p["enabled"]:
|
||
out.update(ok=False, error="PMS_TECH_ENABLED 关着, 本轮不拉取")
|
||
return out
|
||
try:
|
||
got = fetch_all_pages(base=p["base"], path=p["path"], timeout=p["timeout"], fetch=fetch)
|
||
except TechFeedError as e:
|
||
out.update(ok=False, error=str(e)[:300])
|
||
return out
|
||
dd, rows = got["data_date"], got["rows"]
|
||
if not dd or not rows:
|
||
out.update(ok=False, error="接口未给 data_date 或无读数行")
|
||
return out
|
||
try:
|
||
out["rows"] = tech_repo.upsert_daily(rows)
|
||
except Exception as e: # noqa: BLE001
|
||
out.update(ok=False, error=f"落表失败: {type(e).__name__}: {e}")
|
||
return out
|
||
out["data_date"] = dd
|
||
out["pages"], out["matched"] = got["pages"], got["matched"]
|
||
# 删旧: 保留 keep_days 个交易日, 门槛 = data_date 往前第 keep_days 个交易日
|
||
try:
|
||
before = tradedays.ymd(tradedays.prev_trade_day(dd, int(p["keep_days"])))
|
||
out["pruned"] = tech_repo.prune(before)
|
||
except Exception as e: # noqa: BLE001
|
||
out["prune_error"] = f"{type(e).__name__}: {e}"
|
||
return out
|
||
|
||
|
||
# ================================================================ 映射 (照 logic_state_service)
|
||
def load_map() -> dict:
|
||
raw = param_store.get(MAP_KEY, "") or ""
|
||
if not raw:
|
||
return {}
|
||
try:
|
||
m = json.loads(raw)
|
||
except (TypeError, ValueError):
|
||
return {}
|
||
return m if isinstance(m, dict) else {}
|
||
|
||
|
||
def _save_map(m: dict) -> dict:
|
||
r = param_store.set_param(MAP_KEY, json.dumps(m, ensure_ascii=False), "system")
|
||
if not r.get("ok"):
|
||
logger.error("[技术面] 映射写入失败: %s —— 扫描层沿用上一份 (最多 %d 天后自动失效)",
|
||
r.get("error"), FRESH_DAYS)
|
||
return r
|
||
|
||
|
||
def state_map(m: dict | None = None) -> dict:
|
||
"""{ts_code: 立场 dict} —— 只在映射新鲜时给; 超 FRESH_DAYS 自然日一律按无读数返回空。"""
|
||
m = load_map() if m is None else (m or {})
|
||
at = str(m.get("at") or "")[:10]
|
||
try:
|
||
age = (date.today() - date.fromisoformat(at)).days
|
||
except ValueError:
|
||
return {}
|
||
if age > FRESH_DAYS:
|
||
return {}
|
||
states = m.get("states")
|
||
return dict(states) if isinstance(states, dict) else {}
|
||
|
||
|
||
def _compact(st: dict) -> dict:
|
||
"""synthesize 的结果压成映射里存的紧凑态 (方案附录丁「逐票紧凑状态」)。"""
|
||
return {"stance": st.get("stance"), "strength": st.get("strength"),
|
||
"phase": st.get("phase"), "confirm": st.get("confirm"),
|
||
"sar_side": st.get("sar_side"), "sar_value": st.get("sar_value"),
|
||
"sar_flip_days": st.get("sar_flip_days"), "choppy": st.get("choppy"),
|
||
"reason": st.get("reason"), "data_date": st.get("data_date")}
|
||
|
||
|
||
def _tdays_elapsed(from_ymd, to_ymd) -> int:
|
||
"""from_ymd 到 to_ymd 之间经过的交易日数 (不含 from, 含 to); to<=from 返回 0。"""
|
||
n = tradedays.trade_days_left(to_ymd, from_ymd) # 含两端的交易日数
|
||
return max(0, n - 1) if n else 0
|
||
|
||
|
||
def _relevant_codes() -> list:
|
||
"""映射要覆盖的票: 在持 + 计划主榜观察 + 待拍板提议。任一来源失败只跳过该来源、记日志。"""
|
||
from app.repo import pms_repo
|
||
from app.services import plan_feed
|
||
codes = []
|
||
try:
|
||
codes += [r.get("ts_code") for r in pms_repo.list_positions(only_open=True)]
|
||
except Exception as e: # noqa: BLE001
|
||
logger.warning("[技术面] 取持仓失败: %s", e)
|
||
try:
|
||
plan = plan_feed.get_plan()
|
||
codes += [r.get("ts_code") for r in (plan.get("main") or [])]
|
||
codes += [r.get("ts_code") for r in (plan.get("observe") or [])]
|
||
except Exception as e: # noqa: BLE001
|
||
logger.warning("[技术面] 取计划失败: %s", e)
|
||
try:
|
||
codes += [r.get("ts_code") for r in
|
||
pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200)]
|
||
except Exception as e: # noqa: BLE001
|
||
logger.warning("[技术面] 取待拍板提议失败: %s", e)
|
||
return [c for c in dict.fromkeys(normalize_code(c) for c in codes if c) if c]
|
||
|
||
|
||
def build_map(*, now=None) -> dict:
|
||
"""对相关票算技术面立场并写映射。与 logic_state_service.pull_for_held 同构:
|
||
任何一步失败都不抛, 取不到写空映射带原因 (扫描层据此按无读数)。"""
|
||
now = now or datetime.now()
|
||
stamp = now.strftime("%Y-%m-%d %H:%M:%S")
|
||
out = {"ok": True, "at": stamp, "data_date": None, "codes": 0, "states": 0,
|
||
"by_stance": {}, "error": None}
|
||
p = _params()
|
||
if not p["enabled"]:
|
||
_save_map({"at": stamp, "data_date": None, "states": {}, "note": "PMS_TECH_ENABLED 关着"})
|
||
out.update(ok=False, error="PMS_TECH_ENABLED 关着")
|
||
return out
|
||
try:
|
||
dd = tech_repo.latest_date()
|
||
except Exception as e: # noqa: BLE001
|
||
_save_map({"at": stamp, "data_date": None, "states": {}, "error": f"读最新日失败: {e}"})
|
||
out.update(ok=False, error=f"读最新日失败: {type(e).__name__}: {e}")
|
||
return out
|
||
if not dd:
|
||
_save_map({"at": stamp, "data_date": None, "states": {}, "note": "库里还没有技术面读数"})
|
||
out.update(ok=False, error="库里还没有技术面读数")
|
||
return out
|
||
out["data_date"] = dd
|
||
# 读数陈旧 (最新日距今超 stale_tdays 个交易日) → 整份按无读数, 只记原因
|
||
if _tdays_elapsed(dd, tradedays.ymd(now.date())) > int(p["stale_tdays"]):
|
||
_save_map({"at": stamp, "data_date": dd, "stale": True, "states": {},
|
||
"note": f"技术面读数 {dd} 距今超过 {p['stale_tdays']} 个交易日, 按无读数"})
|
||
out.update(ok=False, error="读数陈旧")
|
||
return out
|
||
codes = _relevant_codes()
|
||
out["codes"] = len(codes)
|
||
if not codes:
|
||
_save_map({"at": stamp, "data_date": dd, "states": {}, "note": "没有需要技术面的票"})
|
||
return out
|
||
since = tradedays.ymd(tradedays.prev_trade_day(dd, tech_rules.CHOPPY_WINDOW + 5))
|
||
try:
|
||
hist = tech_repo.history_multi(codes, since=since)
|
||
except Exception as e: # noqa: BLE001
|
||
_save_map({"at": stamp, "data_date": dd, "states": {}, "error": f"取历史失败: {e}"})
|
||
out.update(ok=False, error=f"取历史失败: {type(e).__name__}: {e}")
|
||
return out
|
||
states = {}
|
||
for c in codes:
|
||
rows = hist.get(c) or []
|
||
if not rows or int(rows[-1].get("data_date") or 0) != dd:
|
||
continue # 今日没有这只票的读数 → 无读数, 不拿旧行硬合成
|
||
states[c] = _compact(tech_rules.synthesize(rows, params=p["rules"]))
|
||
r = _save_map({"at": stamp, "data_date": dd, "states": states})
|
||
if not r.get("ok"):
|
||
out["error"] = f"映射写入失败: {r.get('error')}"
|
||
out["states"] = len(states)
|
||
out["by_stance"] = dict(Counter(s.get("stance") or "空" for s in states.values()))
|
||
return out
|
||
|
||
|
||
def pull_and_map(*, fetch=None) -> dict:
|
||
"""拉取 + 落表 + 重建映射, 一把梭 (手动拉取接口与 08:40 补拉用)。落表失败就不建映射。"""
|
||
pr = pull(fetch=fetch)
|
||
mp = build_map() if pr.get("ok") else {"ok": False, "error": "未落表, 跳过建映射"}
|
||
return {"ok": bool(pr.get("ok")), "pull": pr, "map": mp}
|
||
|
||
|
||
# ================================================================ 页面挂载与状态
|
||
def attach(rows, states=None):
|
||
"""把技术面立场挂到行上 (原地改 r['tech'])。无读数的票挂 None, 页面按无读数显示。"""
|
||
states = state_map() if states is None else (states or {})
|
||
for r in rows or []:
|
||
r["tech"] = states.get(normalize_code(r.get("ts_code") or "")) or None
|
||
return rows
|
||
|
||
|
||
def _row_view(r: dict) -> dict:
|
||
"""一行落表读数 → 页面友好的三块 (布林/多空布林/SAR)。"""
|
||
return {
|
||
"data_date": r.get("data_date"),
|
||
"boll": {"upper": r.get("boll_upper"), "mid": r.get("boll_mid"), "lower": r.get("boll_lower"),
|
||
"bw_pct": r.get("boll_bw_pct"), "pos": r.get("boll_pos"),
|
||
"state": r.get("boll_state"), "squeeze": bool(r.get("boll_squeeze"))},
|
||
"bbi": {"bbi": r.get("bbi"), "pos": r.get("bbi_pos"),
|
||
"state": r.get("bbi_state"), "dist_pct": r.get("bbi_dist_pct")},
|
||
"sar": {"value": r.get("sar_value"), "side": r.get("sar_side"),
|
||
"flip_days": r.get("sar_flip_days"), "dist_pct": r.get("sar_dist_pct")},
|
||
"quality": r.get("quality"), "reanchored": bool(r.get("reanchored")),
|
||
}
|
||
|
||
|
||
def research_feed(ts_code: str) -> dict:
|
||
"""单票研究面的技术面一块: 最新读数 + 合成立场 + 近日翻向次数。点击时直接读库算。"""
|
||
code = normalize_code(ts_code or "")
|
||
if not code:
|
||
return {"tech": None, "error": "代码认不出"}
|
||
if not enabled():
|
||
return {"tech": None, "note": "PMS_TECH_ENABLED 关着"}
|
||
p = _params()
|
||
try:
|
||
dd = tech_repo.latest_date()
|
||
since = tradedays.ymd(tradedays.prev_trade_day(dd, tech_rules.CHOPPY_WINDOW + 5)) if dd else 0
|
||
rows = tech_repo.history(code, since=since, limit=tech_rules.CHOPPY_WINDOW + 5)
|
||
except Exception as e: # noqa: BLE001
|
||
return {"tech": None, "error": f"{type(e).__name__}: {e}"}
|
||
if not rows:
|
||
return {"tech": None, "note": "这只票没有技术面读数"}
|
||
st = tech_rules.synthesize(rows, params=p["rules"])
|
||
return {"tech": {**_compact(st), "latest": _row_view(rows[-1]),
|
||
"flips": tech_rules.count_sar_flips(rows), "bars": len(rows)}}
|
||
|
||
|
||
def status() -> dict:
|
||
"""技术面模块状态: 开关/最新读数日/今日行数/映射覆盖与新鲜度/错误。页面顶栏芯片与运维看。"""
|
||
p = _params()
|
||
out = {"enabled": p["enabled"], "base_set": bool(p["base"]), "data_date": None,
|
||
"rows_today": 0, "map": {}, "error": None}
|
||
if not p["enabled"]:
|
||
out["note"] = "PMS_TECH_ENABLED 关着"
|
||
return out
|
||
try:
|
||
dd = tech_repo.latest_date()
|
||
out["data_date"] = dd
|
||
out["rows_today"] = tech_repo.count_on(dd) if dd else 0
|
||
except Exception as e: # noqa: BLE001
|
||
out["error"] = f"读库失败: {type(e).__name__}: {e}"
|
||
m = load_map()
|
||
if m:
|
||
out["map"] = {"at": m.get("at"), "data_date": m.get("data_date"),
|
||
"states": len(m.get("states") or {}), "fresh": bool(state_map(m)),
|
||
"note": m.get("note"), "error": m.get("error"), "stale": m.get("stale")}
|
||
return out
|