# -*- 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 = 1500 # 分页每页条数 (全市场约 5000, 三四页取完) 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 < page_size: # 最后一页 break if matched and offset >= matched: # 够了 (服务端若忽略 limit 也在这里收住) 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