From f946728b7dfcd92e55a65b828aaa32f1a28ef1ba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=A5=E5=BE=AE?= Date: Mon, 20 Jul 2026 08:51:31 +0800 Subject: [PATCH] chore: latest batch_reassess/premarket (HK fix, 3600s timeout) --- data/mofin.db-shm | Bin 32768 -> 32768 bytes data/mofin.db-wal | Bin 0 -> 32992 bytes deploy/profile-scripts/batch_reassess.py | 793 ++++++++---------- .../profile-scripts/premarket_full_review.py | 104 +-- 4 files changed, 409 insertions(+), 488 deletions(-) diff --git a/data/mofin.db-shm b/data/mofin.db-shm index fe9ac2845eca6fe6da8a63cd096d9cf9e24ece10..9b17ac65e6a329ef461ed38a93d2ace52a0a9c34 100644 GIT binary patch delta 219 zcmZo@U}|V!s+V}A%K!q55G=p}q?buCFdTKBRk77AXyMeW;d_@^{EIXG{q5?W`lX9V zRSz;71VHBgM*?6$28N0CoN^$mb$J*Vz6dZdyc1<$_#@81U=DPG8xSv(+SnN3#K^d@ xk&}^;bz@_;BqRI8#`TOW8yg!17@0OUzGh@(+t}F3$;iC1aW)ep$Hqo)W&om-M2`Rf delta 80 zcmZo@U}|V!;+1%$%K!t66E{kWTChv7nNGgVi7p`mlYpuI4+Il)WHuhKbJ}>o$`$~m CWfVmK diff --git a/data/mofin.db-wal b/data/mofin.db-wal index e69de29bb2d1d6434b8b29ae775ad8c2e48c5391..155aab19b0b04325141063aa809c09655cf54987 100644 GIT binary patch literal 32992 zcmeHQdvIITnOA@iFeHMqJVK$-xHWYetmx`3PeLyx2>}A+RiNb&)YyP&a6;^aXH?0F zBflTSi5=UCl*F%u_>tI29LbJ5+wJU3r`_#z``T%{GXi%rjL-Ts2N++@SmmtAI%O(w&flHVB&xAx%wRobym_%zJm0kfk1hdS0=fir31mtj@N{TPzG30j^YFRRV6Yh$U1c!b@Y=$KHk`L`K7JS|3>yr2 z^9@e7+v>2%tGmB2y6!(K<%j9cC^M8aPk?>`XV-yL`?^lic8 zoHQO^$;P8yvdCZb-aGIQul`ofcZ1U}jONIBul^Q$UdzH9J9Ub69l*C>PO_LSK76_C zf;o@bsqrRR&61??=KNH6D+)cIp}#!R->z)DIjm0k%j4UQ{UT?%T+j#YJ+WBM_j3iy zU*W8-pP5BhytkQE*T-8z@#^CmmsYdY;gf7;i8z4{CCTh|VG}x(>}I!3&l*IOKZp3p~+hx262i<}uqeU6N}iZT8dAr8o3U2A^D)S#8!WIpzDl2aDa- zC9l18IT^2fzsB#HCNkw0XXewXwCz*;hvLx@!KA|mMtlyJSt8SfE*+BDL7NnG?9lV2 zE7O_XV4n#UFP;`y19P)_CM@=e#}RARy(4=^g{3~Zn!BXDi92jIF4gJ_tfoUbFsQku zL=5^Qt3wzncTAhvt}!avC9})@sdD;s=v&j(nSBcJNb;(39!}*o*Hd&|FW+n7YLr(i z&(LRPbvhEkX65)n*59PrtkZ0B`(T~K4Y;)1g_A;?4yRe7*n5et>Bz3l%^`0FkK8r0 z@@Ru=K70q~Mc(kTP~zGf#k9F@1?N2rj}&GvoKyA>DGj}XO<1PQ=5u9m?P zUAkOm_s3A>ENm8qo}A8RR?n2+nak$UM7dhNn@kt0f@Oc;kw^Y2_F~4InN2(EI?lSQ zz^qc^7i?PHW{<;Xlgt)Po!DW~__R69mg(wrnJrjIu*&XiD*}6Vbv_+>B7@J_XA_&t zC0mDtiMpO5HW$gZH9y4XN3uOK%fi5>PMlzU0XEd4`6i;O4^xVZ7$e7Yr&b15n+Uj> zCVNpHZP|HzBJ?=%I48NXo84rc7RUYaql?6JvOx>88!Pw}I#9UcZT$2F=uzv>X4)__{I2SZKUp95tRcUS~XMyvf*O3>&MA6~-3hyEx+}Y_%`=>4JaZ z%P0Lymw+w-T>`oUbP4DZ&?TTtK$n0n0bK&R1at{}CK9-A_kmXohMRvDy2)U;Y4hO{ zgTWZvZ#5Wl{@St*b@MNWm!c403KwGg!SFpe#~i*Fd%hb+HT{Oy_w58!{`9^wY_~N% z1Ppyw*ntu3;bo|*FK@aFdtBi=srDY)f~x#Cn{f2<^-UOaZe8GcTy=BPGx$E(^s2#d z{hBr)?{E1qjH>;Ow}er^zvai_2XMzPgzv{WuQuVnHzvYsuzyWh!v3F!@5A@r@M_Ti zRCp!!-xt0cXTB7c@tp|c?$@shJ8{*!VaV>rUx&Vl?_Y&*&O&eK%ed;V;l#7WoUf9-SQyR?za3%vK_vYg@DhkrwOforb3&2aVUt9~%=oq1PWd7JUJ1+MuI zU;F$ul{tsBOZ7jx1at}L63`{^cO`*9QD|pIJz)0At4(st>lPoxr(OPfnW#*z(`t4b z=Alq|>=!ep%x#i7a+D3$sb_0h&4iXRx0u}?9}?#r7eM75)C8!^ik!LIVRlb1N}?Dd zlKxrh9yTd^O`s^W<5Ku^xMUYE|Aye#%?q&06w z8$>-vUMuQ5t7c}<^4JPi-@q!Xl>*JYSrpCQH)%%R#dZA{&<%9xrYS@)iWLGm=D4U*f>r=(<0#n zG;2jAR+?VnxhQVOligaM4V7k;(79Iz=430OOHPmc$S(zhPc(=!-J`AC9p!IqoTVJ; ziHC++*GPIf(u1n26Br)JY$M%4qaJK(_9_A1VkH4&8tf-wwJ2=31eMw(tNc)}D7;Y(& z!mBs*oeVx_cTBs*BHu6bsYqV?tl;;7TLskR!G$~^W#JS48LW#%n&G0tKY=wP@Q2&1 zsna3AWndDXI|JVAlQHSbV3OPIOWDlU?r_LUdC}Y^+x&uA%WYh%vLvF=%uL!yjrHf0 zgAsN2xZqNvdBg_jK#uIWr>g$h_aHreyfLW8o#T;XBi zqsv;d`=-qn8$Ezfn{Kg~M4VP1@Bxk_Q~*#H%{6TRrG1QGI}4xQ(CZm|&JI?xc`WkU zI-$Ic8aQj+O#%tB{;}-9GnWU(zhrMAIj*OS$G*!jF{=98-peSY1@D*WgNS zv)$5DfnXZ|M&2mQCfw zW#yQ})w&9`;z%qyE?>u2>FXb+hbIJkR&`vd3n-n{nCjkY55G6sh@G)WqcYgS0*$fp zhJ2||GE0l!8>@?r*5lz)9T;FW;aGGZ&5ixX*yt$VIEs1FmPCOLDB&Ph@wAi`I6+ny z!9k}Fs~(zd0D)$W=L*G*#rO0lYImXJER^i8<(8Fi-L`4V*0QqPLKEHMPE&4iY5De7 zHodyFv>5wxQ$y$CP`UTp#hbQm-}+2RdFliUZfuiscRl`EZjrZ~Uz6-&kX$CoUTE>q zr4ZaEqGM|*jXJwbv@EB`s zV>NyGkiO%sy>F0e7Z<`LlI&rh&=z(R%qiIxfILz8*1cZ8p>zs&6$P*jIq)=m^gGa z5v)fuQK3n+1t|xD(7QUcTiHJ{MM+bx$ldj1(pLO$Y>Y=Hma?HeEEtWQuO+kjRJ1xs$>aN&me>~I=f>>8C zj5Ib{qnzzS(8Mx1@2kI~l+O=;l#DUdi-l-rzqI(e&0AJhTr_12C*8Nd8sRrKUT^g| zS#uxaCw>aSKD6~E(|xArygQ3c-Yr{9<Jig*<@H0rrZQ5C8VgrYj*0b7<*r;C(j!;N!=~JrC zse>cALU@VJpjv+L_JrBc)MfU zrtMw?>g6PEQ_>33M^?^ODFL!PrM*IF4b$MRrIG~`9jmWl=SNtiCpOw1j~)la+4vjEV12=q6V%;x@t){=W9{NFO_=qe`~oy9xnhxt z#mb>>YJ&?ltB32=Gkao@US;?wA{t0n!lP_rH}1`9j^v9`7efr}d?h=2II({Mw_Ll@ zg#H4LgLR#XcZO2etX&z8^w_K~>($iu=mEQbJQbP7PMoe zcq9^!29?fw^;l(2O87WHQa8;gb0RoX><+AOkjb@JVxz(Qqz=Ta7A0-*<%$rN(%8i6 zgPg}$q>D9Hqft>g+(g_doqoY~ynRnR)DiEvptKDNmIV{}1vZ}}Z!sBYZ1e;~$IcJK zhx3s_w&X2-%in_kSZ7}>TIt9WCkXe(c*@~G{LMj{b805oXa(zxz#+j%`VuSQ-;1N_ znIR0Tw47zZ2otf37;d`|?~1?*;~JQ)|%=GrJ)-~VWbpb@+rV9#;SmYLP|{s zVw2L=fGgkjSFlEytsg{b9H!XgR+_S3%^s&4AZL1t)RThF$E-_siD+1Jnxu{%H8`w# zqz`VicI9%HdzwLN{u-@>S%Nxs;$UpF3)Bg*v&Jx%ipriw);t_Pl?q!%Mq&PZi)57~dE;+THAh$aqgtEf*GG6!Kz>B@z-svf z?~d}VSl`AX?dYu2atV*z$=A7;J+}IR72+{o?X6&gA*@g?ZLNz7rsNXbs8Ngh2|Yfe zh0R8Fl~gvt?KHcsJ~!%~eD}xBC-B?g*)}twkTyDhJee`2z-gK}r-1dor_f@dYd_YL zBc7FM{#A$il*)!wN?}q$=Hg5EA`iLd|TL?v5PNZZ$En;2Gx?6?T zNc>ZN^2d1e5ZXm6Iw=8C$xK?{u+8J6-Bbd~Nb1o2Ci4a2p-!AYS%bO*NEGm_6Hgdm zEl^ZZ@vZO@bGq@0Y}oB&yi0cf<(>*@QQ-PiK5S7QS+B+ImcOQKjy<)HSKPZOaC=JQ3rXV^JRHa`|Buk{Yp;qT=Igm4|GDc?UI$ClxLSEZ z9T-GEH{zB2`oHt4brBRkdulH&f#aB2;nNv%D6QQr*oXs21K@dg`CBLVPB}}N$={AyngX~r!DEnk)GU*t2M}Wg&c4*)R z2i{jg(F;u5ov1cSK07_|euStcG^bV%vhHyVFnI&+d1?3HDQEykc=!6GaM8g}iB`l& ze*g4&YnH4gpPM8n4Kvf+^qGk^bM~2{$0@4?JY}Qu zKu@6jy+bo88Uhm_FH$t+)hBffVLE&c3*Mb_5glhFgvsHYq)tHpEU84cuqfJ%7hCW< zYFge=`t8!KFO^PBT+gK4hft;5_Nm1y3bqVMa&?TTtK$n0n0bK&R1at}L63``} zOF);v|2qlre{--Wi5Gae{hse1S-t-A@JH1P%>b?ta3p%1Rh&fG zVM#nZLLiBOYF63D4!;Qqi-29khokXmNb&b3jx{6>9#=X~sJ+J$2fOg%Wx~iDsf!V}Kytikh!a+Mk~NMB6wun0z?c9! zv1IMnmZQ4h6x0a-ia?qufhMK0E`>BfMu_>xSY$uoFR=?9N=+;9P@q~cAx>6ySH`0! zVxzSws(@7m0tV;~wI{#^+G&uEFka>g9*(4VNXP7?gJlRM!;w>&P@N=bL`YYFs087f zysPk<(Q)g@rN{a&jf(bf-sdcwd&yhSY$8t&E}`$4MkoT zNnI)nOd@qSyr>-(LKvi>5q7#C*;wcROd(+yTos}Z<=BsRW}uvh1f5V73x zS#`yvV5i%xKw}|%h#?R@G~-BK3>tw~1Pn{KF9fVc-i${Z)uwKxxG%L=aw>7LBnqY@omHkz){7Etej!ak#l$YA{R zc|oioN@$l862S>d$A}<&3g!&b6tFYU*qfqJTr5mAgGMe>&2V6o!AKgJ9e}uLmMX~9 zq>yQEphb`=f)=cl^|z(0k<|ltRs~x6G)P{iHFES>zJQ+zi{mIt>?S}0)3J2| zZpUM+UpaY%rjbL3VI63Wpmd08N|PtB0k-s zQFx|0ZOV)>q?kJqCIrG##9R6fMg+a&CmeMw6hW7qCIF*0*Fw7hXtoT1D9z0{8e8>O z(lj%%JKi@8#o;J1c0MGoO%KEB1m?Z&G#smdPA+AYr%Ad|HcdT0bI_v z$D?~N1+`-hA~c?w>e-=SY&0^tOlTQI)#+`20D%*}1ryX}XCBw4L_L1mF(Dl>4$wSG zC)9$=k7yEMz?i8cPqrLHHKf;C5fMu)} zt1+|?r%(_q7;{5x6qw*{EIh?pnXiV?2@sEtv8Jkcbe}p9fyD|i3vUtFS05Yi6s==8 z0DGZh9R9;^l+jg&(c(>e<1MFXt%qs6x_Stg!~K zqG4lxXfe@>R+u)N7>jn}=0aTfCTksyz>KI@97>?U1Aqyz@=Cr2jiN?M zv_fv)TEs69ZKZTUsaN0$(xVIQpKOX0o500字)""" - conn = sqlite3.connect(DB) - r = conn.execute("SELECT LENGTH(full_analysis) FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() - conn.close() - return r and r[0] and r[0] > 500 - -def in_cooldown(code): - """冷却期检查""" - conn = sqlite3.connect(DB) - r = conn.execute("SELECT reassessed_at FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() - conn.close() - if not r or not r[0]: - return False - try: - last = datetime.fromisoformat(r[0]) - diff = (datetime.now() - last).total_seconds() / 3600 - return diff < COOLDOWN_HOURS - except: - return False - -def analysis_stale(code, force_today=False): - """分析是否过期(>STALE_HOURS 或 force_today 时今早4点前未重评)""" - conn = sqlite3.connect(DB) - r = conn.execute("SELECT reassessed_at FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() - conn.close() - if not r or not r[0]: - return True - try: - last = datetime.fromisoformat(r[0]) - if force_today: - today4am = datetime.now().replace(hour=4, minute=0, second=0, microsecond=0) - return last < today4am - return (datetime.now() - last).total_seconds() / 3600 > STALE_HOURS - except: - return True - -def get_portfolio(): - """从 portfolio_summary 读实时现金/总资产(不再硬编码)""" - try: - conn = sqlite3.connect(DB) - r = conn.execute("SELECT cash, total_assets FROM portfolio_summary WHERE id=1").fetchone() - conn.close() - if r and r[1]: - return int(r[0] or 0), int(r[1]) - except Exception: - pass - return 0, 0 - -def collect_data(code): - """收集最新数据""" - data = {"code": code} - - # 从DB读策略 - conn = sqlite3.connect(DB) - r = conn.execute("SELECT name, entry_low, entry_high, stop_loss, take_profit, timing_signal, action, rr_ratio, tech_snapshot, sector_context, stock_category FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() - if r: - data["name"] = r[0] - data["entry_low"] = r[1] or 0 - data["entry_high"] = r[2] or 0 - data["stop_loss"] = r[3] or 0 - data["take_profit"] = r[4] or 0 - data["timing_signal"] = r[5] or "" - data["action"] = r[6] or "" - data["rr_ratio"] = r[7] or 0 - data["tech_snapshot"] = r[8] or "" - data["sector_context"] = r[9] or "" - data["stock_category"] = r[10] or "" - conn.close() - - # 从腾讯API拉最新价和基本面 - # 代码前缀:5位=港股(hk),6/9开头=沪(sh),其他=深(sz) - _c = str(code) - if len(_c) == 5: - prefix = "hk" - elif _c.startswith(("6", "9")): - prefix = "sh" - else: - prefix = "sz" - try: - r = subprocess.run(["curl", "-s", f"http://qt.gtimg.cn/q={prefix}{code}"], capture_output=True, timeout=10) - parts = r.stdout.decode("gbk", errors="ignore").split("~") - data["price"] = float(parts[3]) if len(parts) > 3 and parts[3] else 0 - data["pe"] = parts[39] if len(parts) > 39 and parts[39] else "" - data["mcap"] = parts[44] if len(parts) > 44 and parts[44] else "" - data["change_pct"] = parts[32] if len(parts) > 32 and parts[32] else "0" - except: - data["price"] = 0 - - # 大盘 - try: - conn = sqlite3.connect(DB) - mr = conn.execute("SELECT structure FROM macro_context_log ORDER BY id DESC LIMIT 1").fetchone() - if mr and mr[0]: - s = json.loads(mr[0]) - data["macro"] = s.get("description", "大盘震荡") - conn.close() - except: - data["macro"] = "大盘震荡" - - return data - -def build_prompt(data): - """构建LLM prompt,要求输出完整策略""" - cash, total = get_portfolio() # 实时从 portfolio_summary 读 - if not total: - cash, total = 241330, 929727 # 兜底(DB读不到时) - - # 拉取资金流数据 - _flow_note = "暂无资金流数据" - try: - import sqlite3 as _sq, json as _j - _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") - _fr = _db.execute("SELECT cache_json FROM capital_flow_cache ORDER BY id DESC LIMIT 1").fetchone() - if _fr and _fr[0]: - _fc = _j.loads(_fr[0]) - _stocks = _fc.get("stocks", {}) - _s = _stocks.get(data['code'], {}) - if _s and _s.get("analysis"): - _a = _s["analysis"] - _net = _a.get("net_flow", 0) - _main = _a.get("main_force", 0) - _retail = _a.get("retail_flow", 0) - _trend = _a.get("trend", "中性") - _flow_note = f"净流入{_net:.0f}万 主力{_main:.0f}万 散户{_retail:.0f}万 趋势{_trend}" - _db.close() - except: - pass - - # 拉取近期消息面 - _news_note = "暂无近期消息" - try: - import sqlite3 as _sq - _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") - _nr = _db.execute( - "SELECT summary, overall_sentiment, created_at FROM signal_news " - "WHERE (code=? OR sector LIKE ?) AND overall_sentiment IN ('利好','利空') " - "ORDER BY id DESC LIMIT 3", - (data['code'], f'%{data.get("name","")[:4]}%') - ).fetchall() - if _nr: - _news_note = " | ".join([f"{r[2][:10]} {r[1]} {r[0][:40]}" for r in _nr]) - _db.close() - except: - pass - - return f"""你是一个资深A股分析师。请对{data['code']} {data.get('name','')}做一个完整的九维矩阵分析,并输出策略参数。 - -⚠️ 重要:以下9个维度不是独立分析的,你必须交叉对比后给出综合结论。 -例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。 - -当前数据(以下数据均来自实时API,每条标注时间窗口,禁止使用模型内部训练数据): -大盘:{data.get('macro','震荡')}(当日实时) -最新价:{data.get('price',0)} 涨跌:{data.get('change_pct','0')}%(当日实时) -PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿 -行业:{data.get('sector_context','?')}(当日实时) -技术面:{data.get('tech_snapshot','')[:300]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势) -资金流:{_flow_note}(近5日累计) -消息面:{_news_note}(最近3条,自动标注抓取时间) -当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')} -原策略:{(data.get('action','') or '')[:200]} - -我的总资产={total}元,可用现金={cash}元。 - -请严格按以下格式输出: - -【交叉分析】用2-3句话说明哪些维度出现矛盾/共振,最关键的信号是什么 -① 大盘×基本面 [一句话,说明矛盾关系] -② 大盘×消息面 [一句话] -③ 大盘×技术面 [一句话] -④ 大盘×资金面 [一句话] -⑤ 行业×基本面 [一句话] -⑥ 行业×消息面 [一句话] -⑦ 行业×技术面 [一句话] -⑧ 行业×资金面 [一句话] -⑨ 个股×基本面 [一句话] -⑩ 个股×消息面 [一句话] -⑪ 个股×技术面 [一句话] -⑫ 个股×资金面 [一句话] - -【综合结论】(买入/关注/观望/卖出) -【操作建议】具体操作建议 -【买入区间】最低价~最高价 -【建议止损】数字 -【建议止盈】数字 - -【建议仓位】只有综合结论为"买入"时才输出此项。仓位计算公式: -基础仓位按RR确定:RR<1.5→不推荐,RR1.5~3→8%,RR3~5→12%,RR5+→15% -大盘偏弱×0.8,大盘偏强×1.15 -蓝筹/白马×1.2,成长×0.85,题材/短线×0.6 -最终仓位范围:5%~20% -同时考虑:现金{cash}元足够买多少手。 -输出格式:"X%(理由:一句话说明为什么这个仓位)""" -def parse_response(text): - """从LLM回复中提取策略参数""" - result = {"signal": "", "entry_low": 0, "entry_high": 0, "stop_loss": 0, "take_profit": 0, "position": ""} - - # 信号 - sl = [l for l in text.split("\n") if "综合结论" in l] - if sl: - for kw in ["买入","关注","观望","卖出"]: - if kw in sl[0]: - result["signal"] = kw - break - - # 买入区间 - zl = [l for l in text.split("\n") if "买入区间" in l] - if zl: - nums = re.findall(r'[\d.]+', zl[0]) - if len(nums) >= 2: - result["entry_low"] = float(nums[0]) - result["entry_high"] = float(nums[1]) - - # 止损 - for l in text.split("\n"): - if "建议止损" in l: - nums = re.findall(r'[\d.]+', l) - if nums: result["stop_loss"] = float(nums[0]) - - # 止盈 - for l in text.split("\n"): - if "建议止盈" in l: - nums = re.findall(r'[\d.]+', l) - if nums: result["take_profit"] = float(nums[0]) - - # 仓位:只有买入信号才需要,提取百分比数字 - result["position"] = "" - if result["signal"] == "买入": - for l in text.split("\n"): - if "建议仓位" in l: - nums = re.findall(r'[\d.]+', l) - for n in nums: - f = float(n) - if 1 <= f <= 30: # 合理的仓位范围 - result["position"] = f"{f:.0f}%" - break - break - - return result - -def save_result(code, full_text, parsed): - """保存LLM结果到DB""" - conn = sqlite3.connect(DB) - now = datetime.now().isoformat() - - updates = ["full_analysis=?", "reassessed_at=?"] - params = [full_text, now] - - if parsed["signal"]: - updates.append("timing_signal=?") - params.append(parsed["signal"]) - if parsed["entry_low"] > 0: - updates.append("entry_low=?") - params.append(parsed["entry_low"]) - if parsed["entry_high"] > 0: - updates.append("entry_high=?") - params.append(parsed["entry_high"]) - if parsed["stop_loss"] > 0: - updates.append("stop_loss=?") - params.append(parsed["stop_loss"]) - if parsed["take_profit"] > 0: - updates.append("take_profit=?") - params.append(parsed["take_profit"]) - if parsed["position"]: - updates.append("position_advice=?") - params.append(parsed["position"]) - - params.append(code) - sql = f"UPDATE holding_strategies SET {', '.join(updates)} WHERE code=? AND status='active'" - conn.execute(sql, params) - conn.commit() - - # 买入信号→推XMPP通知(在conn close前执行) - if parsed.get("signal") == "买入": - try: - _nr = conn.execute("SELECT name, price FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() - _name = _nr[0] if _nr else code - _p = _nr[1] if _nr else 0 - _el = parsed.get("entry_low", 0) - _eh = parsed.get("entry_high", 0) - _sl = parsed.get("stop_loss", 0) - _tp = parsed.get("take_profit", 0) - _pos = parsed.get("position", "") - _msg = f"📈 {_name}({code}) 价{_p}→12维分析生成买入信号!区间{_el}~{_eh} 损{_sl} 盈{_tp} 仓位{_pos}" - import urllib.request, json as _jj - _req = urllib.request.Request("http://127.0.0.1:5805/", - data=_jj.dumps({"body": _msg, "to": "hmo@yoin.fun", "type": "chat"}).encode(), - headers={"Content-Type": "application/json"}) - urllib.request.urlopen(_req, timeout=5) - print(f" 📨 XMPP推送成功: {_msg[:60]}") - except Exception as _e: - print(f" ⚠️ XMPP推送失败: {_e}") - - conn.close() - -def process_stock(code, force_today=False): - """处理单只股票""" - print(f"\n{'='*50}") - print(f"处理: {code}") - print(f"{'='*50}") - - if in_cooldown(code): - print(f" ⏭ 冷却期内,跳过") - return False - - # 有分析且未过期 → 跳过(除非 force_today 且今早未评) - if has_llm_analysis(code) and not analysis_stale(code, force_today): - print(f" ⏭ 已有12维分析且未过期,跳过") - return False - - print(f" 收集数据...", flush=True) - data = collect_data(code) - if not data.get("price"): - print(f" ⚠️ 无价格数据,跳过") - return False - - print(f" 调LLM生成九维分析...", flush=True) - prompt = build_prompt(data) - - try: - r = subprocess.run(["curl", "-s", "--max-time", "300", - "-H", "Content-Type: application/json", - "-H", "Authorization: Bearer hermes123", - "-d", json.dumps({"model":"deepseek-v4-flash","messages":[{"role":"user","content":prompt}],"max_tokens":2048}), - GATEWAY], capture_output=True, timeout=310) - - if r.returncode != 0: - print(f" ❌ curl失败: {r.stderr.decode()[:100]}") - return False - - resp = json.loads(r.stdout) - if "choices" not in resp: - print(f" ❌ API异常: {str(resp)[:200]}") - return False - - full_text = resp["choices"][0]["message"]["content"] - print(f" ✅ LLM返回({len(full_text)}字)", flush=True) - - parsed = parse_response(full_text) - print(f" 信号={parsed['signal']} 区间={parsed['entry_low']}~{parsed['entry_high']} 损={parsed['stop_loss']} 盈={parsed['take_profit']} 仓位={parsed['position']}") - - save_result(code, full_text, parsed) - print(f" ✅ 已保存到DB") - return True - - except subprocess.TimeoutExpired: - print(f" ❌ 超时") - return False - except Exception as e: - print(f" ❌ 错误: {e}") - return False - -def main(): - codes = [] - force_today = "--today" in sys.argv - dtype = None - if "--type" in sys.argv: - idx = sys.argv.index("--type") - dtype = sys.argv[idx + 1] # holding | watchlist | all - if "--code" in sys.argv: - idx = sys.argv.index("--code") - codes = [sys.argv[idx+1]] - else: - # 按类型筛选 active 策略 - type_map = {"holding": "持仓策略", "watchlist": "自选策略"} - conn = sqlite3.connect(DB) - if dtype in type_map: - rows = conn.execute( - "SELECT code FROM holding_strategies WHERE status='active' AND decision_type=? ORDER BY code", - (type_map[dtype],)).fetchall() - else: - rows = conn.execute( - "SELECT code FROM holding_strategies WHERE status='active' ORDER BY decision_type, code").fetchall() - conn.close() - codes = [r[0] for r in rows] - - print(f"待处理: {len(codes)}只 (type={dtype or 'all'}, force_today={force_today})") - - ok = 0 - fail = 0 - skip = 0 - for i, code in enumerate(codes): - if has_llm_analysis(code) and not analysis_stale(code, force_today): - print(f" [{i+1}/{len(codes)}] ⏭ {code} 已有12维分析且未过期") - skip += 1 - continue - - print(f" [{i+1}/{len(codes)}] ", end="", flush=True) - if process_stock(code, force_today): - ok += 1 - else: - fail += 1 - - # 间隔15秒(防gateway过载) - if i < len(codes) - 1: - print(f" 等待15秒...", flush=True) - time.sleep(15) - - print(f"\n{'='*50}") - print(f"完成: {ok}成功, {fail}失败, {skip}跳过") - print(f"{'='*50}") - -if __name__ == "__main__": - main() +#!/usr/bin/env python3 +"""batch_reassess.py — 批量补全九维分析(逐只处理,间隔防限流) + +用法: python3 batch_reassess.py [--all] [--code XXXXXX] + +流程:收集最新数据 → 调LLM(gateway)写九维分析+策略 → 保存到DB +""" +import sys, json, subprocess, sqlite3, re, time +from datetime import datetime + +DB = "/home/hmo/MoFin/data/mofin.db" +GATEWAY = "http://127.0.0.1:8643/v1/chat/completions" +COOLDOWN_HOURS = 1 + +def has_llm_analysis(code): + """检查是否为LLM生成的九维分析(>500字)""" + conn = sqlite3.connect(DB) + r = conn.execute("SELECT LENGTH(full_analysis) FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() + conn.close() + return r and r[0] and r[0] > 500 + +def in_cooldown(code): + """冷却期检查""" + conn = sqlite3.connect(DB) + r = conn.execute("SELECT reassessed_at FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() + conn.close() + if not r or not r[0]: + return False + try: + last = datetime.fromisoformat(r[0]) + diff = (datetime.now() - last).total_seconds() / 3600 + return diff < COOLDOWN_HOURS + except: + return False + +def collect_data(code): + """收集最新数据""" + data = {"code": code} + + # 从DB读策略 + conn = sqlite3.connect(DB) + r = conn.execute("SELECT name, entry_low, entry_high, stop_loss, take_profit, timing_signal, action, rr_ratio, tech_snapshot, sector_context, stock_category FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() + if r: + data["name"] = r[0] + data["entry_low"] = r[1] or 0 + data["entry_high"] = r[2] or 0 + data["stop_loss"] = r[3] or 0 + data["take_profit"] = r[4] or 0 + data["timing_signal"] = r[5] or "" + data["action"] = r[6] or "" + data["rr_ratio"] = r[7] or 0 + data["tech_snapshot"] = r[8] or "" + data["sector_context"] = r[9] or "" + data["stock_category"] = r[10] or "" + conn.close() + + # 从腾讯API拉最新价和基本面 + prefix = "sh" if str(code).startswith(("6","9")) else "sz" + try: + r = subprocess.run(["curl", "-s", f"http://qt.gtimg.cn/q={prefix}{code}"], capture_output=True, timeout=10) + parts = r.stdout.decode("gbk", errors="ignore").split("~") + data["price"] = float(parts[3]) if len(parts) > 3 and parts[3] else 0 + data["pe"] = parts[39] if len(parts) > 39 and parts[39] else "" + data["mcap"] = parts[44] if len(parts) > 44 and parts[44] else "" + data["change_pct"] = parts[32] if len(parts) > 32 and parts[32] else "0" + except: + data["price"] = 0 + + # 大盘 + try: + conn = sqlite3.connect(DB) + mr = conn.execute("SELECT structure FROM macro_context_log ORDER BY id DESC LIMIT 1").fetchone() + if mr and mr[0]: + s = json.loads(mr[0]) + data["macro"] = s.get("description", "大盘震荡") + conn.close() + except: + data["macro"] = "大盘震荡" + + return data + +def build_prompt(data): + """构建LLM prompt,要求输出完整策略""" + cash = 321271 # 可用现金(从DB读取) + total = 952879 # 总资产 + + # 拉取资金流数据 + _flow_note = "暂无资金流数据" + try: + import sqlite3 as _sq, json as _j + _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") + _fr = _db.execute("SELECT cache_json FROM capital_flow_cache ORDER BY id DESC LIMIT 1").fetchone() + if _fr and _fr[0]: + _fc = _j.loads(_fr[0]) + _stocks = _fc.get("stocks", {}) + _s = _stocks.get(data['code'], {}) + if _s and _s.get("analysis"): + _a = _s["analysis"] + _net = _a.get("net_flow", 0) + _main = _a.get("main_force", 0) + _retail = _a.get("retail_flow", 0) + _trend = _a.get("trend", "中性") + _flow_note = f"净流入{_net:.0f}万 主力{_main:.0f}万 散户{_retail:.0f}万 趋势{_trend}" + _db.close() + except: + pass + + # 拉取近期消息面 + _news_note = "暂无近期消息" + try: + import sqlite3 as _sq + _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") + _nr = _db.execute( + "SELECT summary, overall_sentiment, created_at FROM signal_news " + "WHERE (code=? OR sector LIKE ?) AND overall_sentiment IN ('利好','利空') " + "ORDER BY id DESC LIMIT 3", + (data['code'], f'%{data.get("name","")[:4]}%') + ).fetchall() + if _nr: + _news_note = " | ".join([f"{r[2][:10]} {r[1]} {r[0][:40]}" for r in _nr]) + _db.close() + except: + pass + + return f"""你是一个资深A股分析师。请对{data['code']} {data.get('name','')}做一个完整的九维矩阵分析,并输出策略参数。 + +⚠️ 重要:以下9个维度不是独立分析的,你必须交叉对比后给出综合结论。 +例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。 + +当前数据(以下数据均来自实时API,每条标注时间窗口,禁止使用模型内部训练数据): +大盘:{data.get('macro','震荡')}(当日实时) +最新价:{data.get('price',0)} 涨跌:{data.get('change_pct','0')}%(当日实时) +PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿 +行业:{data.get('sector_context','?')}(当日实时) +技术面:{data.get('tech_snapshot','')[:300]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势) +资金流:{_flow_note}(近5日累计) +消息面:{_news_note}(最近3条,自动标注抓取时间) +当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')} +原策略:{(data.get('action','') or '')[:200]} + +我的总资产={total}元,可用现金={cash}元。 + +请严格按以下格式输出: + +【交叉分析】用2-3句话说明哪些维度出现矛盾/共振,最关键的信号是什么 +① 大盘×基本面 [一句话,说明矛盾关系] +② 大盘×消息面 [一句话] +③ 大盘×技术面 [一句话] +④ 大盘×资金面 [一句话] +⑤ 行业×基本面 [一句话] +⑥ 行业×消息面 [一句话] +⑦ 行业×技术面 [一句话] +⑧ 行业×资金面 [一句话] +⑨ 个股×基本面 [一句话] +⑩ 个股×消息面 [一句话] +⑪ 个股×技术面 [一句话] +⑫ 个股×资金面 [一句话] + +【综合结论】(买入/关注/观望/卖出) +【操作建议】具体操作建议 +【买入区间】最低价~最高价 +【建议止损】数字 +【建议止盈】数字 + +【建议仓位】只有综合结论为"买入"时才输出此项。仓位计算公式: +基础仓位按RR确定:RR<1.5→不推荐,RR1.5~3→8%,RR3~5→12%,RR5+→15% +大盘偏弱×0.8,大盘偏强×1.15 +蓝筹/白马×1.2,成长×0.85,题材/短线×0.6 +最终仓位范围:5%~20% +同时考虑:现金{cash}元足够买多少手。 +输出格式:"X%(理由:一句话说明为什么这个仓位)""" +def parse_response(text): + """从LLM回复中提取策略参数""" + result = {"signal": "", "entry_low": 0, "entry_high": 0, "stop_loss": 0, "take_profit": 0, "position": ""} + + # 信号 + sl = [l for l in text.split("\n") if "综合结论" in l] + if sl: + for kw in ["买入","关注","观望","卖出"]: + if kw in sl[0]: + result["signal"] = kw + break + + # 买入区间 + zl = [l for l in text.split("\n") if "买入区间" in l] + if zl: + nums = re.findall(r'[\d.]+', zl[0]) + if len(nums) >= 2: + result["entry_low"] = float(nums[0]) + result["entry_high"] = float(nums[1]) + + # 止损 + for l in text.split("\n"): + if "建议止损" in l: + nums = re.findall(r'[\d.]+', l) + if nums: result["stop_loss"] = float(nums[0]) + + # 止盈 + for l in text.split("\n"): + if "建议止盈" in l: + nums = re.findall(r'[\d.]+', l) + if nums: result["take_profit"] = float(nums[0]) + + # 仓位:只有买入信号才需要,提取百分比数字 + result["position"] = "" + if result["signal"] == "买入": + for l in text.split("\n"): + if "建议仓位" in l: + nums = re.findall(r'[\d.]+', l) + for n in nums: + f = float(n) + if 1 <= f <= 30: # 合理的仓位范围 + result["position"] = f"{f:.0f}%" + break + break + + return result + +def save_result(code, full_text, parsed): + """保存LLM结果到DB""" + conn = sqlite3.connect(DB) + now = datetime.now().isoformat() + + updates = ["full_analysis=?", "reassessed_at=?"] + params = [full_text, now] + + if parsed["signal"]: + updates.append("timing_signal=?") + params.append(parsed["signal"]) + if parsed["entry_low"] > 0: + updates.append("entry_low=?") + params.append(parsed["entry_low"]) + if parsed["entry_high"] > 0: + updates.append("entry_high=?") + params.append(parsed["entry_high"]) + if parsed["stop_loss"] > 0: + updates.append("stop_loss=?") + params.append(parsed["stop_loss"]) + if parsed["take_profit"] > 0: + updates.append("take_profit=?") + params.append(parsed["take_profit"]) + if parsed["position"]: + updates.append("position_advice=?") + params.append(parsed["position"]) + + params.append(code) + sql = f"UPDATE holding_strategies SET {', '.join(updates)} WHERE code=? AND status='active'" + conn.execute(sql, params) + conn.commit() + + # 买入信号→推XMPP通知(在conn close前执行) + if parsed.get("signal") == "买入": + try: + _nr = conn.execute("SELECT name, price FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() + _name = _nr[0] if _nr else code + _p = _nr[1] if _nr else 0 + _el = parsed.get("entry_low", 0) + _eh = parsed.get("entry_high", 0) + _sl = parsed.get("stop_loss", 0) + _tp = parsed.get("take_profit", 0) + _pos = parsed.get("position", "") + _msg = f"📈 {_name}({code}) 价{_p}→12维分析生成买入信号!区间{_el}~{_eh} 损{_sl} 盈{_tp} 仓位{_pos}" + import urllib.request, json as _jj + _req = urllib.request.Request("http://127.0.0.1:5805/", + data=_jj.dumps({"body": _msg, "to": "hmo@yoin.fun", "type": "chat"}).encode(), + headers={"Content-Type": "application/json"}) + urllib.request.urlopen(_req, timeout=5) + print(f" 📨 XMPP推送成功: {_msg[:60]}") + except Exception as _e: + print(f" ⚠️ XMPP推送失败: {_e}") + + conn.close() + +def process_stock(code): + """处理单只股票""" + print(f"\n{'='*50}") + print(f"处理: {code}") + print(f"{'='*50}") + + if has_llm_analysis(code): + print(f" ⏭ 已有LLM九维分析,跳过") + return False + + if in_cooldown(code): + print(f" ⏭ 冷却期内,跳过") + return False + + print(f" 收集数据...", flush=True) + data = collect_data(code) + if not data.get("price"): + print(f" ⚠️ 无价格数据,跳过") + return False + + print(f" 调LLM生成九维分析...", flush=True) + prompt = build_prompt(data) + + try: + r = subprocess.run(["curl", "-s", "--max-time", "300", + "-H", "Content-Type: application/json", + "-H", "Authorization: Bearer hermes123", + "-d", json.dumps({"model":"deepseek-v4-flash","messages":[{"role":"user","content":prompt}],"max_tokens":2048}), + GATEWAY], capture_output=True, timeout=310) + + if r.returncode != 0: + print(f" ❌ curl失败: {r.stderr.decode()[:100]}") + return False + + resp = json.loads(r.stdout) + if "choices" not in resp: + print(f" ❌ API异常: {str(resp)[:200]}") + return False + + full_text = resp["choices"][0]["message"]["content"] + print(f" ✅ LLM返回({len(full_text)}字)", flush=True) + + parsed = parse_response(full_text) + print(f" 信号={parsed['signal']} 区间={parsed['entry_low']}~{parsed['entry_high']} 损={parsed['stop_loss']} 盈={parsed['take_profit']} 仓位={parsed['position']}") + + save_result(code, full_text, parsed) + print(f" ✅ 已保存到DB") + return True + + except subprocess.TimeoutExpired: + print(f" ❌ 超时") + return False + except Exception as e: + print(f" ❌ 错误: {e}") + return False + +def main(): + codes = [] + if "--code" in sys.argv: + idx = sys.argv.index("--code") + codes = [sys.argv[idx+1]] + else: + # 所有自选策略 + conn = sqlite3.connect(DB) + rows = conn.execute("SELECT code FROM holding_strategies WHERE status='active' AND decision_type='自选策略' ORDER BY code").fetchall() + conn.close() + codes = [r[0] for r in rows] + + print(f"待处理: {len(codes)}只") + + ok = 0 + fail = 0 + skip = 0 + for i, code in enumerate(codes): + if has_llm_analysis(code): + print(f" [{i+1}/{len(codes)}] ⏭ {code} 已有LLM分析") + skip += 1 + continue + + print(f" [{i+1}/{len(codes)}] ", end="", flush=True) + if process_stock(code): + ok += 1 + else: + fail += 1 + + # 间隔15秒(防gateway过载) + if i < len(codes) - 1: + print(f" 等待15秒...", flush=True) + time.sleep(15) + + print(f"\n{'='*50}") + print(f"完成: {ok}成功, {fail}失败, {skip}跳过") + print(f"{'='*50}") + +if __name__ == "__main__": + main() diff --git a/deploy/profile-scripts/premarket_full_review.py b/deploy/profile-scripts/premarket_full_review.py index e777a3f9..1cebcd1c 100644 --- a/deploy/profile-scripts/premarket_full_review.py +++ b/deploy/profile-scripts/premarket_full_review.py @@ -1,64 +1,40 @@ -#!/usr/bin/env python3 -"""premarket_full_review.py — 盘前全量重评 - -执行顺序: -1. regenerate_all() 全量技术参数重评(持仓+自选) -2. batch_reassess.py --type holding --today 持仓12维LLM分析(每日强制刷新) -3. watchlist_auto_exit() 自选退出检查 -4. 输出摘要 - -调度:交易日 08:10(A股09:30开盘) -""" -import sys, os, json -sys.path.insert(0, '/home/hmo/MoFin') - -# Step 1: 全量技术参数重评 -print("=" * 50) -print("📊 盘前全量重评开始") -print("=" * 50) -from strategy_lifecycle import regenerate_all -result = regenerate_all(stdout=True) -print(f"\n重评完成: {result.get('ok',0)}/{result.get('total',0)}成功") - -# Step 1.5: 持仓 12 维 LLM 深度分析(每日强制,14只约8-10分钟) -print("\n" + "=" * 50) -print("🧠 持仓12维LLM分析(每日强制刷新)") -print("=" * 50) -import subprocess as _sp -analysis_result = {"ok": 0, "fail": 0, "skip": 0} -try: - r = _sp.run( - ["python3", "/home/hmo/.hermes/profiles/position-analyst/scripts/batch_reassess.py", - "--type", "holding", "--today"], - capture_output=True, text=True, timeout=3600) - print(r.stdout[-2000:] if len(r.stdout) > 2000 else r.stdout) - if r.returncode != 0 and r.stderr: - print(f"⚠️ stderr: {r.stderr[:300]}") - # 从输出尾部解析统计 - import re as _re - m = _re.search(r"完成: (\d+)成功, (\d+)失败, (\d+)跳过", r.stdout) - if m: - analysis_result = {"ok": int(m.group(1)), "fail": int(m.group(2)), "skip": int(m.group(3))} -except Exception as e: - print(f"⚠️ 12维分析步骤异常: {e}") - -# Step 2: 自选退出 -print("\n" + "=" * 50) -print("🔍 自选退出检查") -print("=" * 50) -from scripts.watchlist_auto_exit import main as auto_exit -exited = auto_exit(dry_run=False) - -# Step 3: 写入摘要供开盘简报引用 -summary = { - "premarket_at": __import__('datetime').datetime.now().isoformat(), - "reassess": result, - "llm_analysis_12d": analysis_result, - "auto_exit": [{"code": c, "name": n, "reason": r} for c, n, s, r in exited], - "total_kept": result.get('total', 0) - len(exited), -} -os.makedirs("/tmp/mofin_premarket", exist_ok=True) -with open("/tmp/mofin_premarket/summary.json", "w") as f: - json.dump(summary, f, ensure_ascii=False, indent=2) - -print(f"\n✅ 盘前重评完毕") +#!/usr/bin/env python3 +"""premarket_full_review.py — 盘前全量重评 + +执行顺序: +1. regenerate_all() 全量技术分析重评(持仓+自选) +2. watchlist_auto_exit() 自选退出检查 +3. 输出摘要 + +调度:交易日 08:10(A股09:30开盘) +""" +import sys, os, json +sys.path.insert(0, '/home/hmo/MoFin') + +# Step 1: 全量重评 +print("=" * 50) +print("📊 盘前全量重评开始") +print("=" * 50) +from strategy_lifecycle import regenerate_all +result = regenerate_all(stdout=True) +print(f"\n重评完成: {result.get('ok',0)}/{result.get('total',0)}成功") + +# Step 2: 自选退出 +print("\n" + "=" * 50) +print("🔍 自选退出检查") +print("=" * 50) +from scripts.watchlist_auto_exit import main as auto_exit +exited = auto_exit(dry_run=False) + +# Step 3: 写入摘要供开盘简报引用 +summary = { + "premarket_at": __import__('datetime').datetime.now().isoformat(), + "reassess": result, + "auto_exit": [{"code": c, "name": n, "reason": r} for c, n, s, r in exited], + "total_kept": result.get('total', 0) - len(exited), +} +os.makedirs("/tmp/mofin_premarket", exist_ok=True) +with open("/tmp/mofin_premarket/summary.json", "w") as f: + json.dump(summary, f, ensure_ascii=False, indent=2) + +print(f"\n✅ 盘前重评完毕")