#!/usr/bin/env python3 """price_monitor.py — 高频价格监控脚本(批量版) 规则:进入区间报一次,离开区间报一次,中间不重复。 每次运行时一次性刷新所有持仓+自选股的实时价。 """ import urllib.request import os, sys, time, json import sqlite3 from datetime import datetime from mo_data import read_decisions from alert_logger import record_alert, clear_alert # ── 温区感知 + 信号溯源(2026-08-13)── try: from regime_gate import _load_weights except Exception: _load_weights = None def _load_active_strats(): """当前温区激活策略集合(strategy_weights.json matched)""" if _load_weights is None: return None try: w = _load_weights() if w and w.get("weights"): return {k for k, v in w["weights"].items() if v.get("matched")} except Exception: pass return None def _is_strategy_active(d_entry): """该持仓策略是否在当前温区激活。未标记策略默认激活(兼容持仓管理/风控)""" act = _load_active_strats() if act is None: return True for key in ("tag", "version", "strategy_type", "decision_type"): v = d_entry.get(key) if v and str(v) in act: return True # 该策略在权重表中但不在激活集 → 非激活 try: w = _load_weights() or {} wmap = w.get("weights") or {} for key in ("tag", "version", "strategy_type", "decision_type"): v = d_entry.get(key) if v and v in wmap: return v in act except Exception: pass return True def _record_signal(code, name, reason, strategy=""): """信号溯源:记录到 signal_ledger""" try: from signal_ledger import record_signal from regime_gate import get_current_regime from temp_band import get_market_temp rg = get_current_regime() tp = get_market_temp() record_signal(code=code, name=name, strategy=strategy, version=strategy, regime=rg.get("regime", ""), temp_band=tp.get("band", ""), reason=reason, source_module="price_monitor") except Exception: pass BREACH_PATH = "/home/hmo/.hermes/zone_breach.json" STATE_PATH = "/home/hmo/.hermes/price_trigger_state.json" # DB 模块(同步实时价到 mofin.db) sys.path.insert(0, "/home/hmo/MoFin") try: from mofin_db import get_conn, DB_PATH from mo_models import calc_total_mv, calc_total_assets HAS_DB = True except ImportError: HAS_DB = False # 策略重评依赖(技术面驱动,非机械百分比) sys.path.insert(0, "/home/hmo/web-dashboard") try: from strategy_lifecycle import reassess_strategy, reassess_with_context HAS_REASSESS = True except ImportError: HAS_REASSESS = False UA = "Mozilla/5.0" # ── XMPP推送 ────────────────────────────────────────────────────────── XMPP_USER = "hmo@yoin.fun" XMPP_BRIDGE = "http://127.0.0.1:5805/" def push_to_xmpp(text): """原始直推(已废弃直用)——保留给极少数必须原样的场景。 新代码请用 _push_action/_push_digest。""" if not text.strip(): return try: payload = json.dumps({ "to": XMPP_USER, "body": text.strip(), "type": "chat", }).encode("utf-8") req = urllib.request.Request(XMPP_BRIDGE, data=payload, headers={"Content-Type": "application/json"}) urllib.request.urlopen(req, timeout=5) except Exception as e: print(f"[XMPP推送失败] {e}", file=sys.stderr) # ── 分级推送(2026-07-21 信噪比纪律,红线#12)── # ACTION: 破止损/重评确认的操作信号 — 直通不限速 # INFO: 未确认的进区提示 — 聚合成摘要,30min 限 1 条 def _push_action(category, text): try: from alert_helper import notify, ACTION notify(category, text, ACTION) except Exception as e: print(f"[ACTION推送失败] {e}", file=sys.stderr) def _push_digest(category, text): try: from alert_helper import notify, INFO notify(category, text, INFO) except Exception as e: print(f"[INFO推送失败] {e}", file=sys.stderr) # ── 批量拉取价格 ────────────────────────────────────────────────────────── def fetch_all_prices(codes): """腾讯批量行情API:一次请求拉取所有股票(A股+港股) A股:sh600110 / sz000001 港股:hk00700 返回 {code: (price, change, change_pct)} """ if not codes: return {} # 构建批量查询串 symbols = [] code_map = {} # symbol -> original_code for code in codes: code_s = str(code).strip() if len(code_s) == 6: # A股:沪市以5/6/9开头,深市以0/3开头 if code_s.startswith(('5', '6', '9')): sym = f"sh{code_s}" else: sym = f"sz{code_s}" else: sym = f"hk{code_s}" symbols.append(sym) code_map[sym] = code_s url = f"http://qt.gtimg.cn/q={','.join(symbols)}" try: req = urllib.request.Request(url, headers={"User-Agent": UA}) with urllib.request.urlopen(req, timeout=10) as r: text = r.read().decode("gbk") except Exception as e: print(f"⚠️ 批量拉取失败: {e}", file=sys.stderr) try: record_alert(level="warning", source="price_monitor", title="批量拉取失败", detail=str(e)[:200]) except Exception: pass return {} results = {} for line in text.strip().split("\n"): line = line.strip() if not line or "=" not in line: continue try: # 格式: v_sh600110="1~诺德股份~600110~11.84~11.90~..." raw_value = line.split("=", 1)[1].strip().strip('"').strip(";") fields = raw_value.split("~") if len(fields) < 6: continue sym = line.split("=", 1)[0].strip().lstrip("v_") orig_code = code_map.get(sym) if not orig_code: continue price = float(fields[3]) if fields[3] else 0 prev_close = float(fields[4]) if fields[4] else 0 change = price - prev_close if prev_close > 0 else 0 change_pct = fields[32] if len(fields) > 32 and fields[32] else "0" results[orig_code] = (price, change, change_pct) except (ValueError, IndexError): continue return results def refresh_data_prices(): """一次性刷新所有持仓+自选股的实时价(完全DB版,不写JSON)""" all_codes = set() # 从DB读所有需要拉取价格的代码 try: conn = get_conn() for r in conn.execute("SELECT code FROM holdings WHERE is_active=1"): all_codes.add(r['code']) for r in conn.execute("SELECT code FROM watchlist_stocks"): all_codes.add(r['code']) for r in conn.execute("SELECT code FROM holding_strategies WHERE status='active'"): all_codes.add(r['code']) conn.close() except Exception as e: print(f"⚠️ 从DB读代码失败: {e}", file=sys.stderr) try: record_alert(level="warning", source="price_monitor", title="读取持仓失败", detail=str(e)[:200]) except Exception: pass return 0 if not all_codes: return 0 # 一次性批量拉取 prices = fetch_all_prices(list(all_codes)) updated = len(prices) # === 弹性同步实时价到 mofin.db === # 防死锁策略(经2026-07-14 WAL死锁复盘改进): # ① 启动时 checkpoint WAL(清理残留事务) # ② 统一 BEGIN IMMEDIATE 包裹整个写操作 # ③ 5次重试 + 指数退避: 1s → 2s → 4s → 8s → 16s(共~31s) # ④ get_conn() 的 busy_timeout=30000 保证等待上限 # ⑤ 每个写操作检查返回值,任一失败立即 rollback + 重试 # ⑥ try/finally 确保连接始终释放 if HAS_DB and prices: # 先checkpoint一次,清理上次被kill残留的WAL try: c = get_conn() c.execute("PRAGMA wal_checkpoint(TRUNCATE)") c.close() except Exception: pass max_tries = 5 conn = None for db_attempt in range(max_tries): try: conn = get_conn() # BEGIN IMMEDIATE 立即获取写锁——失败则等 busy_timeout(30s) conn.execute("BEGIN IMMEDIATE") # ── 构建 holdings 更新数据 ── db_holdings = [] for r in conn.execute("SELECT * FROM holdings WHERE is_active=1"): h = dict(r) code = str(h.get('code', '')) if code in prices: price_val, _, change_pct = prices[code] if price_val > 0: h['price'] = round(price_val, 2) h['change_pct'] = float(change_pct) if change_pct else 0 db_holdings.append(h) # ── 写 holdings 表 ── for h in db_holdings: currency = str(h.get('currency', 'CNY')).upper() if currency not in ('CNY', 'HKD'): raise ValueError(f"非法币种: {currency}") conn.execute(""" INSERT INTO holdings (code, name, shares, cost, price, market_value, change_pct, currency, position_pct, added_at, is_active) VALUES (?,?,?,?,?,?,?,?,?,datetime('now','localtime'),1) ON CONFLICT(code) DO UPDATE SET name=excluded.name, shares=excluded.shares, cost=excluded.cost, price=excluded.price, market_value=excluded.market_value, change_pct=excluded.change_pct, currency=excluded.currency, position_pct=excluded.position_pct """, ( h.get('code'), h.get('name'), h.get('shares', 0), h.get('cost'), h.get('price'), h.get('market_value'), h.get('change_pct'), h.get('currency', 'CNY'), h.get('position_pct'), )) # ── 写 portfolio_summary ── mv = calc_total_mv(db_holdings) existing = conn.execute( 'SELECT cash, frozen_cash FROM portfolio_summary WHERE id=1' ).fetchone() db_cash = existing['cash'] if existing else 0.0 db_frozen = existing['frozen_cash'] if existing else 0.0 assets = calc_total_assets({'holdings': db_holdings, 'cash': db_cash, 'frozen_cash': db_frozen}) position_pct = round(mv / assets * 100, 2) if assets > 0 else 0 conn.execute(""" INSERT INTO portfolio_summary (id, total_assets, total_mv, stock_value, cash, frozen_cash, position_pct, total_pnl, currency, updated_at) VALUES (1,?,?,?,?,?,?,?,?,datetime('now','localtime')) ON CONFLICT(id) DO UPDATE SET total_assets=excluded.total_assets, total_mv=excluded.total_mv, stock_value=excluded.stock_value, cash=excluded.cash, frozen_cash=excluded.frozen_cash, position_pct=excluded.position_pct, total_pnl=excluded.total_pnl, currency=excluded.currency, updated_at=datetime('now','localtime') """, ( assets, mv, mv, db_cash, db_frozen, position_pct, 0, 'CNY', )) # ── 写 live_prices ── for h in db_holdings: code = h.get('code', '') if code: p = h.get('price', 0) cp = h.get('change_pct', 0) conn.execute( "INSERT OR REPLACE INTO live_prices (code, price, change_pct, updated_at) " "VALUES (?,?,?,datetime('now','localtime'))", (code, p, cp) ) # 补充策略股/自选股的价格(不在holdings中的) for code, pdata in prices.items(): if code not in {h.get('code') for h in db_holdings}: price_val = pdata[0] if isinstance(pdata, (list, tuple)) else pdata.get('price', 0) cp_val = pdata[1] if isinstance(pdata, (list, tuple)) else pdata.get('change_pct', 0) conn.execute( "INSERT OR REPLACE INTO live_prices (code, price, change_pct, updated_at) " "VALUES (?,?,?,datetime('now','localtime'))", (code, price_val, cp_val) ) conn.commit() conn.close() conn = None if db_attempt > 0: print(f"DB同步成功(第{db_attempt+1}次重试)") break # success except (sqlite3.OperationalError, sqlite3.DatabaseError) as e: if conn: try: conn.rollback() except Exception: pass try: conn.close() except Exception: pass conn = None err_str = str(e) if "locked" in err_str or "cannot commit" in err_str or "busy" in err_str: if db_attempt < max_tries - 1: wait = 2 ** db_attempt # 1, 2, 4, 8, 16 print(f"⏳ DB锁(尝试{db_attempt+1}/{max_tries}): {e} → {wait}s后重试", file=sys.stderr) time.sleep(wait) else: print(f"❌ DB锁(重试{max_tries}次耗尽): {e}", file=sys.stderr) else: print(f"❌ DB错误: {e}", file=sys.stderr) break except Exception as e: if conn: try: conn.rollback() except Exception: pass try: conn.close() except Exception: pass conn = None print(f"⚠️ DB同步异常: {e}", file=sys.stderr) try: record_alert(level="error", source="price_monitor", title="DB同步异常", detail=str(e)[:200]) except Exception: pass break else: # for-else: loop exhausted without break print("❌ DB同步失败(所有重试耗尽)", file=sys.stderr) # 尝试紧急 WAL checkpoint(释放死锁) try: c = sqlite3.connect(str(DB_PATH), timeout=1) c.execute("PRAGMA wal_checkpoint(TRUNCATE)") c.close() print(" ↪ 紧急WAL checkpoint完成", file=sys.stderr) except Exception as we: print(f" ↪ WAL checkpoint也失败: {we}", file=sys.stderr) return updated # ── 区间偏离检测 ────────────────────────────────────────────────────────── def load_state(): try: with open(STATE_PATH) as f: return json.load(f) except: return {} def save_state(state): os.makedirs(os.path.dirname(STATE_PATH), exist_ok=True) with open(STATE_PATH, 'w') as f: json.dump(state, f, ensure_ascii=False, indent=2) def load_breaches(): try: with open(BREACH_PATH) as f: return json.load(f) except: return {} def save_breaches(data): os.makedirs(os.path.dirname(BREACH_PATH), exist_ok=True) with open(BREACH_PATH, 'w') as f: json.dump(data, f, ensure_ascii=False, indent=2) def record_event(code, name, event_type, price, trigger_value, event_label=""): """记录一次价格触发事件到 DB price_events 表(唯一权威存储,JSON 已退役)。 price_events.code 有 FK -> stocks(code),未注册的股票(新候选/港股) 先注册再写事件,否则 FK 失败事件丢失。 """ now = datetime.now().isoformat() if HAS_DB: try: from mofin_db import get_conn, write_price_event _c = get_conn() _exch, _typ = ("HK", "H") if len(str(code)) == 5 else (("SH", "A") if str(code).startswith(("6", "9")) else ("SZ", "A")) _c.execute("INSERT OR IGNORE INTO stocks (code, name, exchange, type, updated_at) VALUES (?,?,?,?,?)", (str(code), name or str(code), _exch, _typ, now)) _c.commit() write_price_event(_c, code=code, name=name, event_type=event_type, price=round(price, 2), trigger_value=str(trigger_value), event_label=event_label) _c.close() except Exception as e: print(f"[price_events DB写入失败] {e}", file=sys.stderr) def get_trigger_zones(trigger): """返回该trigger所有可监控的区间列表,跳过已执行的batch""" zones = [] for key, label in [ ("entry_zone", "加仓区间"), ("batch1_price", "试仓区间"), ("batch2_price", "加仓区间"), ("take_profit_zone", "止盈区间"), ("watch_low", "关注区间"), ("watch_high", "减仓区间"), ("watch_break", "止损区间") ]: status_key = key.replace("_price", "_status") if status_key in trigger and trigger[status_key] == "executed": continue val = trigger.get(key, "") if val and "~" in val: try: parts = val.split("~") lo, hi = float(parts[0]), float(parts[1]) zones.append((key, label, lo, hi)) except: pass sl = trigger.get("stop_loss", "") if sl: try: sl_price = float(sl) if isinstance(sl, (int, float)) else float(sl) zones.append(("stop_loss", "止损", 0, sl_price)) except: pass return zones def _cleanup_lock(): """清理进程锁文件""" try: os.remove("/tmp/price_monitor.lock") except Exception: pass def _handle_sigterm(signum, frame): """收到SIGTERM时清理锁文件后退出""" _cleanup_lock() sys.exit(0) def _handle_sigalrm(signum, frame): """收到SIGALRM强制超时时清理锁文件后退出""" _cleanup_lock() print(f"[TIMEOUT] 本轮执行超时({signum}s),已清理锁文件", file=sys.stderr, flush=True) sys.exit(0) def run_once(round_label=""): """执行一轮完整的监控流程""" import os, signal # 必须在开头import,否则os变量会被后面的局部import绑定覆盖 signal.signal(signal.SIGTERM, _handle_sigterm) signal.signal(signal.SIGALRM, _handle_sigalrm) os.nice(10) # 降低优先级,避免与DB其他写操作抢占 # ── 进程锁:同一时间只跑一个实例 ── _lk = "/tmp/price_monitor.lock" _pid = None try: with open(_lk) as _f: _pid = int(_f.read().strip()) os.kill(_pid, 0) print(f"[LOCK] 已有实例(PID {_pid})在运行,跳过本轮", file=sys.stderr, flush=True) return except (FileNotFoundError, ProcessLookupError, ValueError): pass with open(_lk, "w") as _f: _f.write(str(os.getpid())) signal.alarm(120) # 硬上限120s,超时自动清理锁退出 label = f" [{round_label}]" if round_label else "" start = time.time() TIME_BUDGET = 90 # 预留30s给输出和清理,90s内必须完成核心逻辑 # === 第一步:一次性刷新所有价格 === refreshed = refresh_data_prices() # === 第二步:检查触发条件 === try: dec = read_decisions() except: print(f"❌{label} 无法读取decisions(DB)", file=sys.stderr) return active = [d for d in dec.get("decisions", []) if d.get("status") == "active"] state = load_state() outputs = [] state_updated = False # 时间冷却:同股同区间30分钟内不重复推 _push_cooldown = {} _cooldown_file = "/home/hmo/.hermes/.price_push_cooldown.json" try: import os if os.path.exists(_cooldown_file): with open(_cooldown_file) as _f: _push_cooldown = json.load(_f) except Exception: _push_cooldown = {} def _can_push(code, zone_key): now = time.time() key = f"{code}_{zone_key}" last = _push_cooldown.get(key, 0) if now - last < 1800: # 30分钟 return False _push_cooldown[key] = now # 持久化写入 try: with open(_cooldown_file, "w") as _f: json.dump(_push_cooldown, _f) except Exception: pass return True # 收集所有需要检查的代码 check_codes = set() for d in active: trig = d.get("trigger", {}) if trig: check_codes.add(d["code"]) # 批量拉取这些股票的价格 prices = fetch_all_prices(list(check_codes)) # 本轮进区事件收集(聚合成一条摘要推送,替代逐条轰炸) _zone_entries = [] for d in active: code = d["code"] trig = d.get("trigger", {}) if not trig: continue zones = get_trigger_zones(trig) if not zones: continue price_info = prices.get(code) if not price_info: continue price, _, _ = price_info if price == 0: continue name = d.get("name", code) if code not in state: state[code] = {} # ── v_combo波段出场(2026-08-13 已移除硬编码)── # 老莫原则:k线形态判断在LLM提示词上体现(batch_reassess 提示词已含波段出场形态), # 不在代码里硬编码"连续2日收破MA10"特定规则。 # swing持仓的 stop_loss/entry_zone/take_profit_zone 已由通用 zones 循环触发 LLM 重评, # timing_signal 含"卖出/止盈"才发推荐 → 统一"进区→重评→维持才发"。 # 时间预算检查:如果超时,跳过重评只做状态记录 _budget_low = (time.time() - start) > TIME_BUDGET for key, label, lo, hi in zones: in_zone = lo <= price <= hi prev_in_zone = state[code].get(key, None) if in_zone and prev_in_zone != True: if key == "stop_loss": outputs.append(f"⚠️ {name}({code}) {price} → 跌破止损{hi}!") record_event(code, name, "stop_loss", price, str(hi)) # 止损触发 → 立即重评并推送给Dad(时间不够则直接推原始告警) if _budget_low: outputs.append(f" 📨 止损触发(超时跳过重评)→已推送Dad") if _can_push(code, "stop_loss"): _push_action("止损告警", f"⚠️ {name}({code}) {price} → 跌破止损{hi}!") else: try: cost = d.get("cost", 0) or 0 shares = d.get("shares", 0) or 0 current_action = d.get("action", "") result = reassess_with_context(code, name, price, cost, shares, current_action) if result: timing_signal = result.get("timing_signal", "") action = result.get("action", "") if "买入" in timing_signal or "加仓" in timing_signal or timing_signal in ("卖出","止盈"): buy_lo = d.get("entry_low", 0) buy_hi = d.get("entry_high", 0) rr = result.get("rr_ratio", 0) # 分级(老爸规则:只有持仓风控动作才ACTION直推;买入/加仓机会进摘要) if timing_signal in ("卖出","止盈") and (d.get("shares") or 0) > 0: if _can_push(code, "stop_loss"): msg = f"🔔 {name}({code}) 价{price}→触发操作区间{max(buy_lo,0):.2f}~{buy_hi:.2f},已触发重评|RR={rr}" _push_action("操作信号", msg) else: _zone_entries.append(f"{name}({code}) {price}→入区+重评{timing_signal}|RR={rr}") outputs.append(f" 📨 止损重评→已推送Dad: {action}") except Exception as e: outputs.append(f" ⚠️ 止损重评失败: {e}") try: record_alert(level="error", source="price_monitor", title="止损重评失败", detail=str(e)[:200], code=code) except Exception: pass else: extra = "" if "_price" in key: batch_shares = trig.get(key.replace("_price", "_shares"), "") action = trig.get(key.replace("_price", "_action"), "") if batch_shares: extra = f" {action}{batch_shares}股" if action else f" {batch_shares}股" elif key in ("take_profit_zone",): act = trig.get("take_profit_action", "") if act: extra = f"({act})" outputs.append(f"⚡ {name}({code}) {price} → 进入{label}{lo}~{hi}{extra}") record_event(code, name, "entry_zone", price, f"{lo}~{hi}", label) # 进入区间 → 立即重评并推送给Dad(时间不够则记入摘要,不逐条轰炸) if _budget_low: if _can_push(code, key): _zone_entries.append(f"{name}({code}) {price}→{label}{lo}~{hi}") outputs.append(f" 📨 区间触发(超时)→记入摘要") else: try: cost = d.get("cost", 0) or 0 shares = d.get("shares", 0) or 0 current_action = d.get("action", "") result = reassess_with_context(code, name, price, cost, shares, current_action) if result: timing_signal = result.get("timing_signal", "") action = result.get("action", "") # 格式化区间描述(止盈区lo=0时美化显示) if key == "take_profit_zone" and lo == 0: zone_desc = f"止盈监控(目标{hi:.0f})" else: zone_desc = f"操作区间{lo}~{hi}" if "买入" in timing_signal or "加仓" in timing_signal or timing_signal in ("卖出","止盈"): rr = result.get("rr_ratio", 0) # 分级(老爸规则:只有持仓风控动作才ACTION直推;买入/加仓机会进摘要) if timing_signal in ("卖出","止盈") and (d.get("shares") or 0) > 0: if _can_push(code, key): msg = f"🔔 {name}({code}) 价{price}→触发{zone_desc},已触发重评|RR={rr}" _push_action("操作信号", msg) outputs.append(f" 📨 区间触发重评→已推送Dad: {action}") # 2026-08-13 信号溯源:风控动作记录 _record_signal(code, name, f"风控:{timing_signal} {zone_desc}", strategy=d.get("tag") or d.get("version") or "") else: # 2026-08-13 温区感知:非激活策略的买入/加仓机会抑制(不进摘要) if _is_strategy_active(d): _zone_entries.append(f"{name}({code}) {price}→{zone_desc}+重评{timing_signal}|RR={rr}") outputs.append(f" 📋 机会记入摘要: {timing_signal} RR={rr}") # 2026-08-13 信号溯源:买入机会记录 _record_signal(code, name, f"买入机会:{zone_desc} {timing_signal}", strategy=d.get("tag") or d.get("version") or "") else: outputs.append(f" 🧊 温区未激活跳过: {name}({code}) {timing_signal} RR={rr}") else: reason = f"重评结果:{timing_signal},不构成操作建议" outputs.append(f" 📋 本地日志(不推): {reason}") except Exception as e: outputs.append(f" ⚠️ 区间重评失败: {e}") try: record_alert(level="error", source="price_monitor", title="区间重评失败", detail=str(e)[:200], code=code) except Exception: pass state[code][key] = True state_updated = True elif not in_zone and prev_in_zone == True: if key != "stop_loss": outputs.append(f"📌 {name}({code}) {price} → 离开{label}{lo}~{hi}") state[code][key] = False state_updated = True # === 第二步收尾:进区事件只记日志不推送(2026-07-22 老爸:进区提示=噪音,关停)=== # 正式推荐由 12维→tag→摘要队列 通道负责;破止损/卖出/止盈(持仓)仍走 ACTION 直推。 if _zone_entries: outputs.append(f"📋 进区事件{len(_zone_entries)}只(仅记日志不推送): {'; '.join(_zone_entries[:5])}") # === 第三步:买入区偏离检测 + 自动重评 === # 2026-08-11 修复:告警/重评只在 A股连续竞价时段执行。 # 9:15-9:25 集合竞价,API 返回虚拟撮合价(非真实成交)——9:17 曾误报 # 688002 暴跌-20%(竞价期 API 误返回 52周低 135.2),造成假告警。 # 连续竞价 = 9:30-11:30 + 13:00-15:00;非此时段(竞价/午休/盘后)跳过告警+重评。 _now_hm = datetime.now().strftime("%H:%M") _trading = (_now_hm >= "09:30" and _now_hm <= "11:30") or (_now_hm >= "13:00" and _now_hm <= "15:00") if not _trading: print(f"[时段] {_now_hm} 非连续竞价时段(集合竞价/午休/盘后),跳过告警与重评(价格已刷新)", flush=True) # 仍提交状态(防止告警标记滞留) if state_updated: save_state(state) return reassesed_codes = [] # 先做急跌检测(仅持仓,自选股不推送暴跌告警) holdings_codes = set() for d in active: shares = d.get("shares", 0) if isinstance(shares, (int, float)): if shares > 0: holdings_codes.add(d["code"]) else: # 非数值shares(如被错误写入的字符串),兜底处理 holdings_codes.add(d["code"]) print(f" [WARN] {d.get('code')} shares为非数值({shares!r}),视为持仓处理", flush=True) for d in active: code = d["code"] # 非持仓跳过 if code not in holdings_codes: continue name = d.get("name", code) price_info = prices.get(code) if not price_info: continue price, _, change_pct = price_info if price == 0: continue # 单日跌幅>7%告警(不依赖zone边界,盘中急跌即触发) try: cp = float(change_pct) if change_pct else 0 except: cp = 0 if cp <= -7: prev_alert = state.get(code, {}).get("__sharp_decline_triggered", False) if not prev_alert: stop_loss = d.get("stop_loss", 0) sl_note = f" 止损{stop_loss}" if stop_loss else "" msg = f"🔻 {name}({code}) {price} 暴跌{cp:.1f}%!{sl_note}" _push_action("急跌告警", msg) outputs.append(msg) state.setdefault(code, {})["__sharp_decline_triggered"] = True state_updated = True # 立即持久化,防止后续超时导致状态丢失而重复推送 save_state(state) elif cp > -5: # 反弹后清除告警标记,下次再跌还能报 state.setdefault(code, {}).pop("__sharp_decline_triggered", None) for d in active: code = d["code"] name = d.get("name", code) price_info = prices.get(code) if not price_info: continue price, _, _ = price_info if price == 0: continue # 从 decisions (DB holding_strategies) 中读取 analysis 的买入区 entry_low = d.get("entry_low", 0) entry_high = d.get("entry_high", 0) if not entry_low or not entry_high: continue in_buy_zone = entry_low <= price <= entry_high prev_in_buy_zone = state.get(code, {}).get("__buy_zone", None) # 状态变化时才触发 if in_buy_zone and prev_in_buy_zone == False: # 重新进入买入区 → 重评确认区间是否仍然有效 outputs.append(f"🔄 {name}({code}) {price} → 重新进入买入区{entry_low}~{entry_high},触发技术面重评") do_reassess = True elif not in_buy_zone and prev_in_buy_zone == True: # 离开买入区 → 立即重评,更新止损/止盈/区间 outputs.append(f"🔄 {name}({code}) {price} → 离开买入区{entry_low}~{entry_high},立即技术面重评") do_reassess = True else: do_reassess = False if do_reassess and HAS_REASSESS: try: cost = d.get("cost", 0) or 0 shares = d.get("shares", 0) or 0 profit_pct = (price - cost) / cost * 100 if cost else 0 is_deep_loss = profit_pct < -20 sentiment = "neutral" if d.get("tech_snapshot"): if "bearish" in d["tech_snapshot"]: sentiment = "bearish" elif "bullish" in d["tech_snapshot"]: sentiment = "bullish" # 调用技术面驱动重评(非机械百分比) result = reassess_strategy( code, name, price, cost, shares, current_action=d.get("action", ""), volume_signal="中性", sentiment=sentiment, ) outputs.append(f" 📊 新策略: 损{result['stop_loss']} 盈{result['take_profit']} 区{result['entry_low']}~{result['entry_high']} RR={result['rr_ratio']}") reassesed_codes.append(code) except Exception as e: outputs.append(f" ⚠️ 重评失败: {e}") # 更新买入区状态 if "__buy_zone" not in state.get(code, {}): if code not in state: state[code] = {} state[code]["__buy_zone"] = in_buy_zone state_updated = True # 如果有重评过的股票,更新 DB holding_strategies(此前写入 decisions.json,已废弃) if reassesed_codes and HAS_REASSESS: # ── 5分钟冷却:regenerate_all 开销太大,不每2分钟跑一次 ── _regen_marker = "/tmp/price_monitor_regen_at" _skip_regen = False try: if os.path.exists(_regen_marker): with open(_regen_marker) as _f: _last_regen = float(_f.read().strip()) if time.time() - _last_regen < 300: _skip_regen = True except: pass if _skip_regen: outputs.append(f" ⏭ 跳过全量重评(距上次<5min),下次再跑") else: try: from strategy_lifecycle import regenerate_all r = regenerate_all(stdout=False) outputs.append(f" ✅ 策略已全量重评: {r.get('ok',0)}/{r.get('total',0)}成功") outputs.append(f" 📌 触发股票: {', '.join(reassesed_codes)}") try: with open(_regen_marker, "w") as _f: _f.write(str(time.time())) except: pass except Exception as e: outputs.append(f" ⚠️ 全量重评失败: {e}") # === 第四步:输出 === now_str = datetime.now().strftime("%H:%M:%S") elapsed = time.time() - start if outputs: print(f"\n🔔 {now_str}{label}") for o in outputs: print(o) print(f"\n{json.dumps({'type':'价格监控','time':now_str,'triggers':outputs}, ensure_ascii=False)}") else: # 无触发时 SILENT(中继不推送) print(f"[SILENT]{label} 价格正常 | {refreshed}只已刷新 | {elapsed:.1f}s") if state_updated: save_state(state) # 输出耗时 print(f"⏱{label} {elapsed:.1f}s", flush=True) # 取消超时定时器(正常完成) signal.alarm(0) # ── 策略追踪评估(2026-07-27 老爸:检查推荐是否触发止盈/止损)── try: import mofin_db _conn = mofin_db.get_conn() tracked = mofin_db.check_strategy_outcomes(_conn) _conn.close() if tracked: print(f" [TRACK] {tracked}条推荐触发止盈/止损", flush=True) except Exception as e: print(f" [TRACK] 检查失败: {e}", flush=True) # 清理进程锁 try: os.remove("/tmp/price_monitor.lock") except Exception: pass def _singleton_guard(max_age_sec=300): """自愈式单例守卫(2026-08-05 进程堆积事故修复) 背景:cron每2分钟拉起本脚本,外网抖动时实例挂起不退,timeout(仅SIGTERM)杀不动, 5个并发烧~155% CPU。规则: - 有其他"新鲜"实例(存活max_age)→ SIGKILL 接管(自愈) """ import subprocess as _sp, os as _os, sys as _sys my_pid = _os.getpid() try: out = _sp.run(["ps", "-C", "python3", "-o", "pid,etimes,cmd"], capture_output=True, text=True, timeout=10).stdout for line in out.splitlines(): if "price_monitor.py" not in line: continue parts = line.split(None, 2) if len(parts) < 3: continue try: pid = int(parts[0]); age = int(parts[1]) except ValueError: continue if pid == my_pid: continue if age > max_age_sec: try: _os.kill(pid, 9) print(f"[guard] SIGKILL卡死实例 pid={pid} age={age}s", flush=True) except ProcessLookupError: pass else: print(f"[guard] 已有新鲜实例 pid={pid} age={age}s 在跑, 本实例退出", flush=True) _sys.exit(0) except Exception as _e: print(f"[guard] 守卫异常(放行): {_e}", flush=True) def main(): """每cron触发跑一轮""" _singleton_guard(max_age_sec=300) run_once() if __name__ == "__main__": main()