chore: 归档文件内容同步为当前canonical版本
This commit is contained in:
@@ -0,0 +1,188 @@
|
||||
{
|
||||
"A_shadow_duplicates": [
|
||||
"300308_monitor.py",
|
||||
"_watchdog_report.py",
|
||||
"accumulation_scanner.py",
|
||||
"batch_reassess.py",
|
||||
"branch_evaluator.py",
|
||||
"candidate_filter.py",
|
||||
"capital_flow_collector.py",
|
||||
"check-prompt-deps.py",
|
||||
"chip_factors.py",
|
||||
"clean_watchlist.py",
|
||||
"cron_health_monitor.py",
|
||||
"data_flow_audit.py",
|
||||
"data_governance.py",
|
||||
"divergence_detector.py",
|
||||
"fix_gateway_port.py",
|
||||
"generate_report.py",
|
||||
"get_realtime_prices.py",
|
||||
"hardcode_scanner.py",
|
||||
"hk_rate.py",
|
||||
"holdings_reconciliation.py",
|
||||
"import_full_stocks.py",
|
||||
"import_holding_xls.py",
|
||||
"intraday_health_check.py",
|
||||
"macro_context_collector.py",
|
||||
"macro_signal_consumer.py",
|
||||
"market_scanner.py",
|
||||
"market_screener.py",
|
||||
"market_watch.py",
|
||||
"meta_growth.py",
|
||||
"mo_alphasift_bridge.py",
|
||||
"mo_bridge.py",
|
||||
"mo_config.py",
|
||||
"mo_data.py",
|
||||
"mo_dsa_opinion.py",
|
||||
"mo_models.py",
|
||||
"mo_provider.py",
|
||||
"mofin_collect.py",
|
||||
"mofin_db.py",
|
||||
"mofin_news.py",
|
||||
"mofin_query.py",
|
||||
"monitor_300308.py",
|
||||
"morning_health_check.py",
|
||||
"multi_timeframe.py",
|
||||
"ocr_client.py",
|
||||
"per_stock_reassess.py",
|
||||
"pre-flight-check.py",
|
||||
"preflight_verify.py",
|
||||
"price_monitor.py",
|
||||
"refresh_macro_context.py",
|
||||
"refresh_mtf_cache.py",
|
||||
"review_needed_watchdog.py",
|
||||
"run_all_tests.py",
|
||||
"self_todo_executor.py",
|
||||
"session_to_cron_bridge.py",
|
||||
"stale_detector.py",
|
||||
"stale_push_wlin.py",
|
||||
"stock_profile.py",
|
||||
"stock_quote.py",
|
||||
"stock_scorer.py",
|
||||
"stock_sector_enrich.py",
|
||||
"strategy-staleness-check.py",
|
||||
"strategy_feedback.py",
|
||||
"strategy_lifecycle.py",
|
||||
"strategy_review.py",
|
||||
"strategy_summary.py",
|
||||
"strategy_tree.py",
|
||||
"sync_cron_prompts.py",
|
||||
"sync_dashboard.py",
|
||||
"sync_decisions_to_db.py",
|
||||
"technical_analysis.py",
|
||||
"trend_detector.py",
|
||||
"vacuum_state_db.py",
|
||||
"verify_reassess_pipeline.py",
|
||||
"watchlist_auto_exit.py"
|
||||
],
|
||||
"B_unref_tools": [
|
||||
"add_hygiene_cron.py",
|
||||
"analyze_health.py",
|
||||
"archive_legacy_json.py",
|
||||
"archive_old_files.py",
|
||||
"audit_data.py",
|
||||
"audit_deadcode.py",
|
||||
"audit_duplication.py",
|
||||
"audit_runtime.py",
|
||||
"backfill_price_events.py",
|
||||
"check_3_dbs.py",
|
||||
"check_97_bug.py",
|
||||
"check_backfill_job.py",
|
||||
"check_candidates_schema.py",
|
||||
"check_current_errors.py",
|
||||
"check_dashboard.py",
|
||||
"check_db_paths.py",
|
||||
"check_db_state.py",
|
||||
"check_default_pool.py",
|
||||
"check_error_freshness.py",
|
||||
"check_fk.py",
|
||||
"check_llm_jobs.py",
|
||||
"check_missing_scripts.py",
|
||||
"check_new_imports.py",
|
||||
"check_price_events.py",
|
||||
"check_srv.py",
|
||||
"check_stocks_table.py",
|
||||
"check_table_schemas.py",
|
||||
"check_todos_schema.py",
|
||||
"check_triggered.py",
|
||||
"classify_diverged.py",
|
||||
"cleanup_disabled_cron.py",
|
||||
"cron_status.py",
|
||||
"db_recovery.py",
|
||||
"dump_health.py",
|
||||
"find_97_rows.py",
|
||||
"find_cron_errors.py",
|
||||
"fix_crontab_cron_to_xmpp.py",
|
||||
"fix_crontab_market.py",
|
||||
"fix_self_todo_path.py",
|
||||
"fix_symlinks.py",
|
||||
"fix_todos_db.py",
|
||||
"full_error_inventory.py",
|
||||
"get_full_errors.py",
|
||||
"get_today_errors.py",
|
||||
"inspect_third_db.py",
|
||||
"inv5.py",
|
||||
"investigate_5_issues.py",
|
||||
"kanban_create.py",
|
||||
"kanban_fix_status.py",
|
||||
"kanban_inspect.py",
|
||||
"key_status.py",
|
||||
"list_tables.py",
|
||||
"merge_third_db.py",
|
||||
"notify_user.py",
|
||||
"notify_user2.py",
|
||||
"notify_user3.py",
|
||||
"notify_user4.py",
|
||||
"notify_user5.py",
|
||||
"notify_user6.py",
|
||||
"notify_user7.py",
|
||||
"pp.py",
|
||||
"register_l134_jobs.py",
|
||||
"remove_cron_watchdog.py",
|
||||
"remove_xiaoguo.py",
|
||||
"remove_xiaoguo2.py",
|
||||
"sc2.py",
|
||||
"sc3.py",
|
||||
"sc5.py",
|
||||
"syntax_check.py",
|
||||
"test_auto_heal.py",
|
||||
"test_bot_import.py",
|
||||
"test_db_only.py",
|
||||
"test_db_write.py",
|
||||
"test_dual_write.py",
|
||||
"test_freshness_query.py",
|
||||
"test_image_pipeline.py",
|
||||
"test_key7_guard.py",
|
||||
"test_llm.py",
|
||||
"test_msg_log.py",
|
||||
"test_multi_profile.py",
|
||||
"test_noping.py",
|
||||
"test_ocr_pipeline_live.py",
|
||||
"test_production.py",
|
||||
"test_raw_insert.py",
|
||||
"test_rerun_action.py",
|
||||
"test_single.py",
|
||||
"test_sn_ocr2.py",
|
||||
"test_xmpp_logger.py",
|
||||
"test_zone_gate.py",
|
||||
"trace_97_source.py",
|
||||
"trigger_weekend_jobs.py",
|
||||
"unify_sources.py",
|
||||
"update_backfill_cron.py",
|
||||
"validate_fixes.py",
|
||||
"verify_300308.py",
|
||||
"verify_deployment.py",
|
||||
"verify_health_fix.py",
|
||||
"verify_health_json.py",
|
||||
"verify_self_check_section.py",
|
||||
"verify_self_todo.py",
|
||||
"xiaoguo_news_processor.py",
|
||||
"xiaoguo_scanner.py",
|
||||
"xiaoguo_sentiment_bridge.py",
|
||||
"xiaoguo_signal_consumer.py"
|
||||
],
|
||||
"keep": [
|
||||
"prepare_report_data.py",
|
||||
"server.py"
|
||||
]
|
||||
}
|
||||
@@ -1,19 +1,30 @@
|
||||
#!/usr/bin/env python3
|
||||
"""batch_reassess.py — 批量补全九维分析(逐只处理,间隔防限流)
|
||||
"""batch_reassess.py — 批量补全12维(九维矩阵)LLM分析(逐只处理,间隔防限流)
|
||||
|
||||
用法: python3 batch_reassess.py [--all] [--code XXXXXX]
|
||||
用法:
|
||||
python3 batch_reassess.py # 所有缺分析/过期的 active 策略
|
||||
python3 batch_reassess.py --type holding # 只处理持仓策略
|
||||
python3 batch_reassess.py --type watchlist # 只处理自选策略
|
||||
python3 batch_reassess.py --type holding --today # 持仓每日刷新(今早未评过的强制重评)
|
||||
python3 batch_reassess.py --code XXXXXX # 单只
|
||||
|
||||
流程:收集最新数据 → 调LLM(gateway)写九维分析+策略 → 保存到DB
|
||||
流程:收集最新数据 → 调LLM(gateway)写12维分析+策略 → 保存到DB
|
||||
"""
|
||||
import sys, json, subprocess, sqlite3, re, time
|
||||
import sys, json, subprocess, sqlite3, re, time, os
|
||||
from datetime import datetime
|
||||
|
||||
# ── 共享 LLM 客户端 + DB 工具(profile-scripts 硬链到同目录)──
|
||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||||
sys.path.insert(0, "/home/hmo/MoFin")
|
||||
from llm_client import call_llm, REASSESS_MODEL, FALLBACK_MODEL, gateway_alive, ocg_alive
|
||||
from mofin_db import snapshot_strategy_history, sync_recommend_tag
|
||||
|
||||
DB = "/home/hmo/MoFin/data/mofin.db"
|
||||
GATEWAY = "http://127.0.0.1:8643/v1/chat/completions"
|
||||
COOLDOWN_HOURS = 1
|
||||
STALE_HOURS = 20 # 分析超过20小时视为过期,需要重评
|
||||
|
||||
def has_llm_analysis(code):
|
||||
"""检查是否为LLM生成的九维分析(>500字)"""
|
||||
"""检查是否为LLM生成的12维分析(>500字)"""
|
||||
conn = sqlite3.connect(DB)
|
||||
r = conn.execute("SELECT LENGTH(full_analysis) FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
conn.close()
|
||||
@@ -33,13 +44,41 @@ def in_cooldown(code):
|
||||
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读策略
|
||||
# 从DB读策略(含 full_analysis / changelog_json / position_advice)
|
||||
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()
|
||||
r = conn.execute("SELECT name, entry_low, entry_high, stop_loss, take_profit, timing_signal, action, rr_ratio, tech_snapshot, sector_context, stock_category, full_analysis, changelog_json, reassessed_at, position_advice FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
if r:
|
||||
data["name"] = r[0]
|
||||
data["entry_low"] = r[1] or 0
|
||||
@@ -52,10 +91,21 @@ def collect_data(code):
|
||||
data["tech_snapshot"] = r[8] or ""
|
||||
data["sector_context"] = r[9] or ""
|
||||
data["stock_category"] = r[10] or ""
|
||||
data["full_analysis"] = r[11] or ""
|
||||
data["changelog_json"] = r[12] or ""
|
||||
data["reassessed_at"] = r[13] or ""
|
||||
data["position_advice"] = r[14] or ""
|
||||
conn.close()
|
||||
|
||||
# 从腾讯API拉最新价和基本面
|
||||
prefix = "sh" if str(code).startswith(("6","9")) else "sz"
|
||||
# 代码前缀: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("~")
|
||||
@@ -80,9 +130,10 @@ def collect_data(code):
|
||||
return data
|
||||
|
||||
def build_prompt(data):
|
||||
"""构建LLM prompt,要求输出完整策略"""
|
||||
cash = 321271 # 可用现金(从DB读取)
|
||||
total = 952879 # 总资产
|
||||
"""构建LLM prompt,先审阅原策略再结合实时数据输出修改判断+九维矩阵分析"""
|
||||
cash, total = get_portfolio()
|
||||
if not total:
|
||||
cash, total = 241330, 929727 # 兜底(DB读不到时)
|
||||
|
||||
# 拉取资金流数据
|
||||
_flow_note = "暂无资金流数据"
|
||||
@@ -122,7 +173,53 @@ def build_prompt(data):
|
||||
except:
|
||||
pass
|
||||
|
||||
return f"""你是一个资深A股分析师。请对{data['code']} {data.get('name','')}做一个完整的九维矩阵分析,并输出策略参数。
|
||||
# ── 构建【原策略全文】section ──
|
||||
_params_parts = []
|
||||
if data.get('action'): _params_parts.append(f"当前策略: {data['action']}")
|
||||
if data.get('timing_signal'): _params_parts.append(f"信号: {data['timing_signal']}")
|
||||
if data.get('entry_low') or data.get('entry_high'):
|
||||
_params_parts.append(f"买入区间: {data.get('entry_low',0)}~{data.get('entry_high',0)}")
|
||||
if data.get('stop_loss'): _params_parts.append(f"止损: {data['stop_loss']}")
|
||||
if data.get('take_profit'): _params_parts.append(f"止盈: {data['take_profit']}")
|
||||
if data.get('position_advice'): _params_parts.append(f"仓位: {data['position_advice']}")
|
||||
_params_str = " | ".join(_params_parts) if _params_parts else "无策略参数"
|
||||
|
||||
# 最近3条变更记录
|
||||
_changelog_str = "无变更记录"
|
||||
try:
|
||||
_cl_raw = data.get('changelog_json', '')
|
||||
if _cl_raw:
|
||||
_cl = json.loads(_cl_raw) if isinstance(_cl_raw, str) else _cl_raw
|
||||
if isinstance(_cl, list) and _cl:
|
||||
_recent = _cl[-3:] if len(_cl) > 3 else _cl
|
||||
_cl_lines = []
|
||||
for i, c in enumerate(_recent):
|
||||
_act = c.get('action', c.get('reason', '')) if isinstance(c, dict) else str(c)
|
||||
_ts = c.get('timestamp', '') if isinstance(c, dict) else ''
|
||||
_cl_lines.append(f" {i+1}. {_ts[:16]} {_act[:80]}")
|
||||
if _cl_lines:
|
||||
_changelog_str = "\n".join(_cl_lines)
|
||||
except:
|
||||
pass
|
||||
|
||||
# 完整分析原文(不截断)
|
||||
_full_analysis = data.get('full_analysis', '') or ''
|
||||
_fa_display = _full_analysis if _full_analysis else '(首次分析,无历史)'
|
||||
|
||||
_orig_strategy_section = f"""当前策略参数: {_params_str}
|
||||
|
||||
变更记录(最近3条):
|
||||
{_changelog_str}
|
||||
|
||||
完整分析原文:
|
||||
{_fa_display}"""
|
||||
|
||||
return f"""你是一个资深A股分析师。请先审阅以下【原策略全文】,判断是否需要修改策略,然后做出完整的九维矩阵分析。
|
||||
|
||||
【原策略全文】
|
||||
{_orig_strategy_section}
|
||||
|
||||
── 以上是已有的策略,以下是当前实时数据,请结合两者做出判断 ──
|
||||
|
||||
⚠️ 重要:以下9个维度不是独立分析的,你必须交叉对比后给出综合结论。
|
||||
例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。
|
||||
@@ -136,11 +233,18 @@ PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿
|
||||
资金流:{_flow_note}(近5日累计)
|
||||
消息面:{_news_note}(最近3条,自动标注抓取时间)
|
||||
当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')}
|
||||
原策略:{(data.get('action','') or '')[:200]}
|
||||
|
||||
我的总资产={total}元,可用现金={cash}元。
|
||||
|
||||
请严格按以下格式输出:
|
||||
请严格按以下格式输出(注意节标题不可省略):
|
||||
|
||||
【维持或修改】明确二选一判断:维持原策略 / 需要修改策略
|
||||
【修改点及理由】
|
||||
如果维持原策略 → 写"无需修改"
|
||||
如果需要修改 → 逐条列出(每条格式:"- 修改点名称:理由说明")
|
||||
【最终新策略】
|
||||
用自然语言输出完整的最终策略全文(200-400字),自包含核心交易逻辑、买入区间价格、止损价、止盈价、仓位比例、风险提示。
|
||||
⚠️ 本段不要使用【综合结论】【买入区间】等标签——用自然语言描述即可。
|
||||
|
||||
【交叉分析】用2-3句话说明哪些维度出现矛盾/共振,最关键的信号是什么
|
||||
① 大盘×基本面 [一句话,说明矛盾关系]
|
||||
@@ -162,13 +266,18 @@ PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿
|
||||
【建议止损】数字
|
||||
【建议止盈】数字
|
||||
|
||||
【建议仓位】只有综合结论为"买入"时才输出此项。仓位计算公式:
|
||||
【建议仓位】⚠️不可省略。综合结论非"买入"时写"不新建仓";为"买入"时按以下公式:
|
||||
基础仓位按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%(理由:一句话说明为什么这个仓位)"""
|
||||
输出格式:"X%(理由:一句话说明为什么这个仓位)"
|
||||
|
||||
⚠️ 输出纪律(必须遵守):
|
||||
1. 直接以【维持或修改】开头,禁止任何寒暄、开场白、分隔线
|
||||
2. 禁止输出 <structured_data> 或任何 XML/JSON/代码块
|
||||
3. 所有【】节标题一个都不能少"""
|
||||
def parse_response(text):
|
||||
"""从LLM回复中提取策略参数"""
|
||||
result = {"signal": "", "entry_low": 0, "entry_high": 0, "stop_loss": 0, "take_profit": 0, "position": ""}
|
||||
@@ -217,28 +326,44 @@ def parse_response(text):
|
||||
return result
|
||||
|
||||
def save_result(code, full_text, parsed):
|
||||
"""保存LLM结果到DB"""
|
||||
"""保存LLM结果到DB(先快照再UPDATE)。空分析拒绝写入。"""
|
||||
if not (full_text or "").strip():
|
||||
print(f" \u274c 拒绝写入空分析(LLM输出为空,保护已有数据)")
|
||||
return
|
||||
conn = sqlite3.connect(DB)
|
||||
now = datetime.now().isoformat()
|
||||
|
||||
# ── 修改前快照 ──
|
||||
snapshot_strategy_history(conn, code, 'batch_12d')
|
||||
|
||||
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:
|
||||
# 区间写入门禁:上下沿都必须为正且 下沿<上沿<下沿x3,否则视为解析错误整体跳过
|
||||
# (防 214.68~2.52 类解析污染,与 GATE_ZONE_SANITY 同级防护)
|
||||
_el, _eh = parsed["entry_low"], parsed["entry_high"]
|
||||
if _el > 0 and _eh > _el and _eh < _el * 3:
|
||||
updates.append("entry_low=?")
|
||||
params.append(parsed["entry_low"])
|
||||
if parsed["entry_high"] > 0:
|
||||
params.append(_el)
|
||||
updates.append("entry_high=?")
|
||||
params.append(parsed["entry_high"])
|
||||
if parsed["stop_loss"] > 0:
|
||||
params.append(_eh)
|
||||
elif _el > 0 or _eh > 0:
|
||||
print(f" ⚠️ 买入区解析异常({_el}~{_eh}),跳过区间写入(保留原值)", flush=True)
|
||||
# 止损/止盈一致性门禁:损>0 时必须在区间下沿之下(0.5x~1.0x),盈>0 时必须在区间上沿之上
|
||||
_sl, _tp = parsed["stop_loss"], parsed["take_profit"]
|
||||
if _sl > 0 and (not _el or _sl < _el) and (not _tp or _sl < _tp):
|
||||
updates.append("stop_loss=?")
|
||||
params.append(parsed["stop_loss"])
|
||||
if parsed["take_profit"] > 0:
|
||||
params.append(_sl)
|
||||
elif _sl > 0:
|
||||
print(f" ⚠️ 止损{_sl}与区间/止盈不一致,跳过写入(保留原值)", flush=True)
|
||||
if _tp > 0 and (not _eh or _tp > _eh) and (not _sl or _tp > _sl):
|
||||
updates.append("take_profit=?")
|
||||
params.append(parsed["take_profit"])
|
||||
params.append(_tp)
|
||||
elif _tp > 0:
|
||||
print(f" ⚠️ 止盈{_tp}与区间/止损不一致,跳过写入(保留原值)", flush=True)
|
||||
if parsed["position"]:
|
||||
updates.append("position_advice=?")
|
||||
params.append(parsed["position"])
|
||||
@@ -247,119 +372,179 @@ def save_result(code, full_text, parsed):
|
||||
sql = f"UPDATE holding_strategies SET {', '.join(updates)} WHERE code=? AND status='active'"
|
||||
conn.execute(sql, params)
|
||||
conn.commit()
|
||||
|
||||
# ── 推荐操作 tag 同步(与 XMPP 动作级信号同源)──
|
||||
sync_recommend_tag(conn, code, parsed.get("signal", ""))
|
||||
|
||||
# 买入信号→推XMPP通知(在conn close前执行)
|
||||
# 买入信号→推XMPP通知(在conn close前执行)——推送质量门禁:
|
||||
# 价格必须>0(live_prices实时价)、区间有效(下沿<上沿<下沿x3)、现价不超过上沿5%、
|
||||
# 损<下沿、盈>上沿、损在(0.5x~1.0x)现价内。任何一项不过 → 不推,只记日志。
|
||||
if parsed.get("signal") == "买入":
|
||||
try:
|
||||
_nr = conn.execute("SELECT name, price FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
_nr = conn.execute("SELECT name FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
_lp = conn.execute("SELECT price FROM live_prices WHERE code=?", (code,)).fetchone()
|
||||
_name = _nr[0] if _nr else code
|
||||
_p = _nr[1] if _nr else 0
|
||||
_p = _lp[0] if _lp and _lp[0] 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]}")
|
||||
_ok, _why = _validate_buy_alert(_p, _el, _eh, _sl, _tp)
|
||||
if _ok:
|
||||
_msg = f"📈 {_name}({code}) 价{_p}→12维分析生成买入信号!区间{_el}~{_eh} 损{_sl} 盈{_tp} 仓位{_pos}"
|
||||
from alert_helper import notify as _notify, ACTION as _ACT
|
||||
_notify("买入信号", _msg, _ACT)
|
||||
print(f" \U0001f4e8 XMPP推送成功: {_msg[:60]}")
|
||||
else:
|
||||
print(f" ⚠️ 买入信号未过推送门禁({_why}),仅记日志不推送", flush=True)
|
||||
except Exception as _e:
|
||||
print(f" ⚠️ XMPP推送失败: {_e}")
|
||||
print(f" \u26a0\ufe0f XMPP推送失败: {_e}")
|
||||
|
||||
conn.close()
|
||||
|
||||
def process_stock(code):
|
||||
|
||||
def _validate_buy_alert(price, el, eh, sl, tp):
|
||||
"""买入信号推送门禁(垃圾信号不发)。
|
||||
返回 (ok, reason)"""
|
||||
if not price or price <= 0:
|
||||
return False, f"无实时价格({price})"
|
||||
if not (el > 0 and eh > el and eh < el * 3):
|
||||
return False, f"区间无效({el}~{eh})"
|
||||
if price > eh * 1.05:
|
||||
return False, f"现价{price}高于区间上沿{eh}超5%(追高信号不推)"
|
||||
if not (sl > 0 and sl < el and price * 0.5 <= sl <= price):
|
||||
return False, f"止损{sl}不合理(需0.5x~1.0x现价且<下沿{el})"
|
||||
if not (tp > eh and tp > sl):
|
||||
return False, f"止盈{tp}需>上沿{eh}且>止损{sl}"
|
||||
return True, ""
|
||||
|
||||
def process_stock(code, force_today=False):
|
||||
"""处理单只股票"""
|
||||
print(f"\n{'='*50}")
|
||||
print(f"处理: {code}")
|
||||
print(f"{'='*50}")
|
||||
|
||||
if has_llm_analysis(code):
|
||||
print(f" ⏭ 已有LLM九维分析,跳过")
|
||||
if in_cooldown(code):
|
||||
print(f" \u23ed 冷却期内,跳过")
|
||||
return False
|
||||
|
||||
if in_cooldown(code):
|
||||
print(f" ⏭ 冷却期内,跳过")
|
||||
# 有分析且未过期 \u2192 跳过(除非 force_today 且今早未评)
|
||||
if has_llm_analysis(code) and not analysis_stale(code, force_today):
|
||||
print(f" \u23ed 已有12维分析且未过期,跳过")
|
||||
return False
|
||||
|
||||
print(f" 收集数据...", flush=True)
|
||||
data = collect_data(code)
|
||||
if not data.get("price"):
|
||||
print(f" ⚠️ 无价格数据,跳过")
|
||||
print(f" \u26a0\ufe0f 无价格数据,跳过")
|
||||
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}")
|
||||
# ── 使用共享 LLM 客户端(替代 curl subprocess)──
|
||||
result = call_llm(prompt, model=REASSESS_MODEL, max_tokens=4096)
|
||||
|
||||
if not result["ok"] or not (result.get("content") or "").strip():
|
||||
print(f" \u274c LLM调用失败或空输出: {result.get('error') or 'empty content'}")
|
||||
return False
|
||||
|
||||
full_text = result["content"]
|
||||
print(f" \u2705 LLM返回({len(full_text)}字, {result['elapsed']:.1f}s, 尝试{result['attempts']}次)", flush=True)
|
||||
|
||||
parsed = parse_response(full_text)
|
||||
|
||||
# ── 截断保护:输出过短且无信号 = 低质输出,升级 pro 重试一次 ──
|
||||
if not parsed.get("signal") and len(full_text) < 1500:
|
||||
print(f" ⚠️ 输出截断({len(full_text)}字)且无信号,升级 {FALLBACK_MODEL} 重试...", flush=True)
|
||||
result2 = call_llm(prompt, model=FALLBACK_MODEL, max_tokens=4096)
|
||||
if result2["ok"] and len((result2.get("content") or "").strip()) > len(full_text):
|
||||
full_text = result2["content"]
|
||||
parsed = parse_response(full_text)
|
||||
print(f" \u2705 升级后({len(full_text)}字)", flush=True)
|
||||
|
||||
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" \u2705 已保存到DB")
|
||||
return True
|
||||
|
||||
def main():
|
||||
# ── 双通道预检:OCG直连 + hermes gateway 兜底,全挂才退出 ──
|
||||
_ocg_ok = ocg_alive()
|
||||
_gw_ok = gateway_alive()
|
||||
if not _ocg_ok and not _gw_ok:
|
||||
print("[FATAL] OCG上游与hermes gateway均不可用,退出")
|
||||
sys.exit(1)
|
||||
if not _ocg_ok:
|
||||
print("[WARN] OCG直连不可用,将使用gateway兜底(agent运行时,较慢)")
|
||||
if not _gw_ok:
|
||||
print("[WARN] hermes gateway不可用,仅使用OCG直连")
|
||||
|
||||
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)
|
||||
rows = conn.execute("SELECT code FROM holding_strategies WHERE status='active' AND decision_type='自选策略' ORDER BY code").fetchall()
|
||||
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)}只")
|
||||
print(f"待处理: {len(codes)}只 (type={dtype or 'all'}, force_today={force_today})")
|
||||
|
||||
ok = 0
|
||||
fail = 0
|
||||
skip = 0
|
||||
failed_codes = []
|
||||
for i, code in enumerate(codes):
|
||||
if has_llm_analysis(code):
|
||||
print(f" [{i+1}/{len(codes)}] ⏭ {code} 已有LLM分析")
|
||||
if has_llm_analysis(code) and not analysis_stale(code, force_today):
|
||||
print(f" [{i+1}/{len(codes)}] \u23ed {code} 已有12维分析且未过期")
|
||||
skip += 1
|
||||
continue
|
||||
|
||||
print(f" [{i+1}/{len(codes)}] ", end="", flush=True)
|
||||
if process_stock(code):
|
||||
if process_stock(code, force_today):
|
||||
ok += 1
|
||||
else:
|
||||
fail += 1
|
||||
failed_codes.append(code)
|
||||
|
||||
# 间隔15秒(防gateway过载)
|
||||
# 间隔8秒(pro model较重但gateway可承受;retry逻辑吸收瞬断)
|
||||
if i < len(codes) - 1:
|
||||
print(f" 等待15秒...", flush=True)
|
||||
time.sleep(15)
|
||||
print(f" 等待8秒...", flush=True)
|
||||
time.sleep(8)
|
||||
|
||||
# ── 失败二轮:主跑结束后休息 60s 让上游恢复,失败股整体重试一次 ──
|
||||
# (凌晨上游空输出高发,二轮可救回大半;仍失败的留给下一轮调度)
|
||||
if failed_codes:
|
||||
print(f"\n{'='*50}")
|
||||
print(f"失败二轮: {len(failed_codes)}只,休息60s后重试...")
|
||||
time.sleep(60)
|
||||
retry_ok = 0
|
||||
for code in failed_codes:
|
||||
print(f" [retry] {code} ", end="", flush=True)
|
||||
if process_stock(code, force_today):
|
||||
retry_ok += 1
|
||||
ok += 1
|
||||
fail -= 1
|
||||
print(f" 等待8秒...", flush=True)
|
||||
time.sleep(8)
|
||||
print(f"失败二轮: {retry_ok}/{len(failed_codes)} 救回")
|
||||
|
||||
print(f"\n{'='*50}")
|
||||
print(f"完成: {ok}成功, {fail}失败, {skip}跳过")
|
||||
|
||||
@@ -17,7 +17,9 @@ DB_PATH = Path("/home/hmo/MoFin/data/mofin.db")
|
||||
UA = "Mozilla/5.0"
|
||||
|
||||
def get_conn():
|
||||
return sqlite3.connect(str(DB_PATH))
|
||||
c = sqlite3.connect(str(DB_PATH), timeout=30)
|
||||
c.execute("PRAGMA busy_timeout=30000")
|
||||
return c
|
||||
|
||||
def log_candidate(conn, code, stage, passed, detail):
|
||||
"""记录过滤日志"""
|
||||
|
||||
@@ -21,33 +21,26 @@ def port_open(port, host="127.0.0.1"):
|
||||
s.close()
|
||||
|
||||
def check_session_health():
|
||||
"""调gateway API,检测session是否卡死。超过15s无响应→不健康"""
|
||||
"""检测 gateway LLM 是否可用——扫 agent.log 最近一次真实调用结果。
|
||||
不再发真实 LLM ping(25s 超时对 20-100s 的冷启动延迟必误报,且每次白烧 22k token)。
|
||||
"""
|
||||
try:
|
||||
payload = json.dumps({
|
||||
"model": "hermes-agent",
|
||||
"messages": [{"role": "user", "content": "ping"}]
|
||||
}).encode()
|
||||
req = urllib.request.Request(GATEWAY_URL, data=payload, method="POST")
|
||||
req.add_header("Content-Type", "application/json")
|
||||
req.add_header("Authorization", f"Bearer {API_KEY}")
|
||||
req.add_header("X-Hermes-Session-Id", SESSION_ID)
|
||||
t0 = time.time()
|
||||
with urllib.request.urlopen(req, timeout=25) as r:
|
||||
data = json.loads(r.read())
|
||||
reply = data.get("choices", [{}])[0].get("message", {}).get("content", "")
|
||||
elapsed = time.time() - t0
|
||||
if reply:
|
||||
print(f"Session {SESSION_ID} 健康 ✓ ({elapsed:.1f}s)")
|
||||
return True
|
||||
else:
|
||||
print(f"Session {SESSION_ID} 返回空", file=sys.stderr)
|
||||
return False
|
||||
except urllib.request.HTTPError as e:
|
||||
print(f"Session {SESSION_ID} HTTP错误: {e.code}", file=sys.stderr)
|
||||
return False
|
||||
sys.path.insert(0, '/home/hmo/MoFin')
|
||||
from xmpp_logger import _scan_agent_log
|
||||
r = _scan_agent_log(time.time(), "zhiwei")
|
||||
if r["status"] == "ok":
|
||||
print(f"Session {SESSION_ID} 健康 ✓ (agent.log: latency={r.get('latency')}, {r.get('age_sec')}s前)")
|
||||
return True
|
||||
# error/unknown:只有近期有明确失败记录才判不健康
|
||||
if r["status"] == "error":
|
||||
print(f"Session {SESSION_ID} 不健康: agent.log 最近调用失败 — {r.get('error','')[:100]}", file=sys.stderr)
|
||||
return False
|
||||
# unknown(无近期调用记录)= 空闲,不算不健康
|
||||
print(f"Session {SESSION_ID} 无近期调用记录(空闲正常)")
|
||||
return True
|
||||
except Exception as e:
|
||||
print(f"Session {SESSION_ID} 不健康: {e}", file=sys.stderr)
|
||||
return False
|
||||
print(f"Session {SESSION_ID} 健康检查异常: {e}(按健康处理)", file=sys.stderr)
|
||||
return True
|
||||
|
||||
def restart_gateway():
|
||||
"""通过systemd重启gateway"""
|
||||
|
||||
@@ -0,0 +1,316 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
mo_data.py — MoFin 统一数据层(纯 DB)
|
||||
|
||||
所有数据从 SQLite 读取。不做 JSON fallback。
|
||||
JSON 文件已弃用,仅保留为历史备份。
|
||||
|
||||
用法:
|
||||
from mo_data import read_portfolio, read_decisions, read_watchlist
|
||||
|
||||
pf = read_portfolio() # 返回和 portfolio.json 一样的 dict 结构
|
||||
dec = read_decisions() # 返回和 decisions.json 一样的 dict 结构
|
||||
wl = read_watchlist() # 返回和 watchlist.json 一样的 dict 结构
|
||||
"""
|
||||
|
||||
import sqlite3, json, sys
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
DB_PATH = '/home/hmo/MoFin/data/mofin.db'
|
||||
SCRIPT_DIR = Path('/home/hmo/MoFin/scripts')
|
||||
|
||||
|
||||
def _get_db():
|
||||
db = sqlite3.connect(DB_PATH)
|
||||
db.row_factory = sqlite3.Row
|
||||
return db
|
||||
|
||||
|
||||
# ── portfolio ─────────────────────────────────────────────────────
|
||||
|
||||
def read_portfolio():
|
||||
"""返回 portfolio.json 等价 dict。纯 DB。"""
|
||||
db = _get_db()
|
||||
rows = db.execute(
|
||||
"SELECT code, name, shares, cost, price, market_value, "
|
||||
"change_pct, currency, position_pct "
|
||||
"FROM holdings WHERE is_active=1"
|
||||
).fetchall()
|
||||
holdings = []
|
||||
for r in rows:
|
||||
h = dict(r)
|
||||
h['_currency'] = h.get('currency', 'CNY')
|
||||
holdings.append(h)
|
||||
|
||||
sum_row = db.execute("SELECT * FROM portfolio_summary WHERE id=1").fetchone()
|
||||
summary = dict(sum_row) if sum_row else {}
|
||||
|
||||
db.close()
|
||||
|
||||
return {
|
||||
"holdings": holdings,
|
||||
"total_assets": summary.get("total_assets", 0),
|
||||
"total_mv": summary.get("total_mv", 0),
|
||||
"stock_value": summary.get("stock_value", summary.get("total_mv", 0)),
|
||||
"cash": summary.get("cash", 0),
|
||||
"frozen_cash": summary.get("frozen_cash", 0),
|
||||
"position_pct": summary.get("position_pct", 0),
|
||||
"currency": summary.get("currency", "CNY"),
|
||||
"updated_at": summary.get("updated_at", ""),
|
||||
}
|
||||
|
||||
|
||||
# ── decisions ─────────────────────────────────────────────────────
|
||||
|
||||
def _parse_json(val, default):
|
||||
if val:
|
||||
try: return json.loads(val)
|
||||
except: pass
|
||||
return default
|
||||
|
||||
|
||||
def read_decisions():
|
||||
"""返回 decisions.json 等价 dict。纯 DB。"""
|
||||
db = _get_db()
|
||||
rows = db.execute(
|
||||
"SELECT code, name, version, price, cost, shares, "
|
||||
"stop_loss, take_profit, entry_low, entry_high, "
|
||||
"currency, strategy_type, action, timing_signal, "
|
||||
"rr_ratio, tech_snapshot, stock_category, sector_context, "
|
||||
"status, trigger_json, changelog_json, source, reason, "
|
||||
"created_at, updated_at, "
|
||||
"avg_price, decision_timestamp, note, quality_check, "
|
||||
"quality_checked_at, quality_issues_json, position_advice, "
|
||||
"signal_factors_json, time_horizon, decision_type, tag "
|
||||
"FROM holding_strategies WHERE status IN ('active','updated') "
|
||||
"ORDER BY code"
|
||||
).fetchall()
|
||||
|
||||
decisions = []
|
||||
for r in rows:
|
||||
d = dict(r)
|
||||
d['trigger'] = _parse_json(r['trigger_json'], {})
|
||||
d['changelog'] = _parse_json(r['changelog_json'], [])
|
||||
d['quality_issues'] = _parse_json(r['quality_issues_json'], {})
|
||||
d['signal_factors'] = _parse_json(r['signal_factors_json'], [])
|
||||
d['timestamp'] = r['decision_timestamp'] or r['created_at'] or ''
|
||||
d['type'] = r['decision_type'] or r['strategy_type'] or '持仓策略'
|
||||
decisions.append(d)
|
||||
|
||||
db.close()
|
||||
|
||||
return {
|
||||
"decisions": decisions,
|
||||
"total": len(decisions),
|
||||
"regenerated_at": datetime.now().strftime('%Y-%m-%d %H:%M'),
|
||||
}
|
||||
|
||||
|
||||
# ── watchlist ─────────────────────────────────────────────────────
|
||||
|
||||
def read_watchlist():
|
||||
"""返回 watchlist 等价 dict。纯 DB。
|
||||
从 holding_strategies(自选策略)读取,watchlist_stocks 已废弃。"""
|
||||
db = _get_db()
|
||||
# 主数据源:holding_strategies 自选策略
|
||||
rows = db.execute(
|
||||
"SELECT code, name, price, entry_low, entry_high, "
|
||||
"stop_loss, currency, updated_at "
|
||||
"FROM holding_strategies WHERE status='active' AND decision_type='自选策略'"
|
||||
).fetchall()
|
||||
|
||||
stocks = []
|
||||
seen = set()
|
||||
for r in rows:
|
||||
code = str(r["code"])
|
||||
if code in seen:
|
||||
continue
|
||||
seen.add(code)
|
||||
stocks.append({
|
||||
"code": code,
|
||||
"name": r["name"] or "",
|
||||
"price": r["price"] or 0,
|
||||
"entry_low": r["entry_low"] or 0,
|
||||
"entry_high": r["entry_high"] or 0,
|
||||
"stop_loss": r["stop_loss"] or 0,
|
||||
"currency": r["currency"] or "CNY",
|
||||
"added_at": r["updated_at"] or "",
|
||||
"analysis": {},
|
||||
})
|
||||
|
||||
return {"stocks": stocks, "total": len(stocks)}
|
||||
|
||||
db.close()
|
||||
|
||||
return {
|
||||
"stocks": stocks,
|
||||
"updated_at": datetime.now().strftime('%Y-%m-%d %H:%M'),
|
||||
}
|
||||
|
||||
|
||||
# ── 便捷别名 ───────────────────────────────────────────────────────
|
||||
|
||||
def read_portfolio_json():
|
||||
return read_portfolio()
|
||||
|
||||
def read_decisions_json():
|
||||
return read_decisions()
|
||||
|
||||
def read_watchlist_json():
|
||||
return read_watchlist()
|
||||
|
||||
|
||||
# ── 统一价格获取(唯一入口,禁止各脚本自拉API)──
|
||||
|
||||
def get_price(code, max_age_minutes=5, use_stale_fallback=True):
|
||||
"""获取单只股票最新价格。
|
||||
|
||||
优先级: live_prices(DB) → stock_quote(API兜底)
|
||||
- live_prices 有且不超过 max_age_minutes → 直接返回
|
||||
- 没有或过期 → 调 stock_quote 拉,写回 live_prices
|
||||
- 都失败 → 返回 (None, None)
|
||||
|
||||
返回 (price, change_pct),两值都是 float 或 None。
|
||||
"""
|
||||
from mofin_db import get_price_from_db
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
# 1. 先读 DB
|
||||
try:
|
||||
db_price, db_chg = get_price_from_db(code)
|
||||
if db_price is not None and db_price > 0:
|
||||
# 检查时效性
|
||||
conn = __import__('sqlite3').connect(str(DB_PATH))
|
||||
row = conn.execute(
|
||||
"SELECT updated_at FROM live_prices WHERE code=?",
|
||||
(str(code).strip(),)
|
||||
).fetchone()
|
||||
conn.close()
|
||||
if row and row[0]:
|
||||
try:
|
||||
updated = datetime.strptime(row[0], "%Y-%m-%d %H:%M:%S")
|
||||
age = (datetime.now() - updated).total_seconds() / 60
|
||||
if age <= max_age_minutes:
|
||||
return (db_price, db_chg)
|
||||
except:
|
||||
pass
|
||||
else:
|
||||
return (db_price, db_chg)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# 2. DB 没有或过期 → 调 stock_quote
|
||||
if not use_stale_fallback:
|
||||
return (None, None)
|
||||
|
||||
try:
|
||||
import subprocess, json
|
||||
r = subprocess.run(
|
||||
[sys.executable, str(SCRIPT_DIR / "stock_quote.py"), str(code)],
|
||||
capture_output=True, text=True, timeout=15
|
||||
)
|
||||
if r.returncode == 0:
|
||||
data = json.loads(r.stdout.strip())
|
||||
price = float(data.get("price", 0))
|
||||
chg = float(data.get("change_pct", 0))
|
||||
if price > 0:
|
||||
# 写回 live_prices
|
||||
try:
|
||||
conn = __import__('sqlite3').connect(str(DB_PATH))
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO live_prices (code, price, change_pct, updated_at)
|
||||
VALUES (?, ?, ?, datetime('now','localtime'))
|
||||
""", (str(code).strip(), price, chg))
|
||||
conn.commit()
|
||||
conn.close()
|
||||
except:
|
||||
pass
|
||||
return (price, chg)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return (None, None)
|
||||
|
||||
|
||||
def get_prices_batch(codes, max_age_minutes=5):
|
||||
"""批量获取价格,返回 {code: (price, change_pct)}"""
|
||||
from mofin_db import get_prices_batch_from_db
|
||||
|
||||
result = {}
|
||||
need_api = []
|
||||
|
||||
# 1. 批量读 DB
|
||||
try:
|
||||
db_prices = get_prices_batch_from_db(codes)
|
||||
for code in codes:
|
||||
cs = str(code).strip()
|
||||
if cs in db_prices:
|
||||
p, c = db_prices[cs]
|
||||
if p and p > 0:
|
||||
result[cs] = (p, c)
|
||||
continue
|
||||
need_api.append(cs)
|
||||
except:
|
||||
need_api = [str(c).strip() for c in codes]
|
||||
|
||||
# 2. 缺失的调 API
|
||||
if need_api:
|
||||
try:
|
||||
import subprocess, json
|
||||
r = subprocess.run(
|
||||
[sys.executable, str(SCRIPT_DIR / "stock_quote.py")] + need_api,
|
||||
capture_output=True, text=True, timeout=30
|
||||
)
|
||||
if r.returncode == 0:
|
||||
for line in r.stdout.strip().split("\n"):
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(line)
|
||||
code = str(data.get("code", "")).strip()
|
||||
price = float(data.get("price", 0))
|
||||
chg = float(data.get("change_pct", 0))
|
||||
if code and price > 0:
|
||||
result[code] = (price, chg)
|
||||
except:
|
||||
pass
|
||||
except:
|
||||
pass
|
||||
|
||||
return result
|
||||
|
||||
|
||||
# ── cash_log 写入 ──────────────────────────────────────────────────
|
||||
|
||||
def write_cash_log(cash_before, cash_after, frozen_before, frozen_after,
|
||||
source, note, verified=0):
|
||||
"""记录现金变更到 cash_log 表。"""
|
||||
change_amount = round(cash_after - cash_before, 2) if cash_after is not None and cash_before is not None else 0
|
||||
db = sqlite3.connect(DB_PATH)
|
||||
try:
|
||||
cur = db.execute(
|
||||
"""INSERT INTO cash_log
|
||||
(timestamp, cash_before, cash_after, frozen_before, frozen_after,
|
||||
change_amount, source, note, verified)
|
||||
VALUES (datetime('now','localtime'), ?, ?, ?, ?, ?, ?, ?, ?)""",
|
||||
(cash_before, cash_after, frozen_before, frozen_after,
|
||||
change_amount, source, note, verified)
|
||||
)
|
||||
db.commit()
|
||||
return cur.lastrowid
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
# ── 自检 ───────────────────────────────────────────────────────────
|
||||
|
||||
if __name__ == "__main__":
|
||||
pf = read_portfolio()
|
||||
print(f"portfolio: {len(pf.get('holdings',[]))} holdings, total_assets={pf.get('total_assets',0)}")
|
||||
|
||||
dec = read_decisions()
|
||||
print(f"decisions: {len(dec.get('decisions',[]))} entries")
|
||||
|
||||
wl = read_watchlist()
|
||||
print(f"watchlist: {len(wl.get('stocks',[]))} stocks")
|
||||
File diff suppressed because it is too large
Load Diff
Executable → Regular
+7
-14
@@ -42,15 +42,12 @@ HERMES_CRON_DIR = Path("/home/hmo/.hermes/profiles/position-analyst/cron")
|
||||
|
||||
def derive_fix_action(detail, msg):
|
||||
"""根据issue信息推导可执行的修复命令"""
|
||||
# 小果扫描 error → 验证脚本是否存在
|
||||
if "xiaoguo_scanner" in msg or "小果扫描" in msg:
|
||||
return f"ls -la /home/hmo/.hermes/profiles/position-analyst/scripts/xiaoguo_scanner.py 2>&1 && echo 'ok'"
|
||||
# system-audit error → 验证拷贝
|
||||
if "system_audit" in msg or "系统审计" in msg:
|
||||
return f"ls -la /home/hmo/.hermes/profiles/position-analyst/scripts/system_audit.py 2>&1"
|
||||
# cron errors(last_status=error)→ 验证文件存在,等下次cron运行自动恢复
|
||||
if "cron" in msg.lower() and "error" in msg.lower() and ("小果" in msg or "系统审计" in msg):
|
||||
return f"ls -la /home/hmo/.hermes/profiles/position-analyst/scripts/xiaoguo_scanner.py /home/hmo/.hermes/profiles/position-analyst/scripts/system_audit.py 2>&1"
|
||||
if "cron" in msg.lower() and "error" in msg.lower() and "系统审计" in msg:
|
||||
return f"ls -la /home/hmo/.hermes/profiles/position-analyst/scripts/system_audit.py 2>&1"
|
||||
# 港股汇率 → 刷新
|
||||
if "港股汇率" in msg:
|
||||
return f"cd {BASE} && python3 hk_rate.py 2>&1"
|
||||
@@ -60,9 +57,9 @@ def derive_fix_action(detail, msg):
|
||||
# delivery目标缺失 → 改为local
|
||||
if "deliver" in msg.lower() or "delivery" in msg.lower():
|
||||
return f"cd {BASE} && echo '需手动设置: cronjob action=update deliver=local'"
|
||||
# 小果→知微桥不通
|
||||
# 小果→知微桥不通(小果已归档,不再自动修复)
|
||||
if "信号桥" in msg:
|
||||
return f"cd {BASE} && python3 scripts/xiaoguo_signal_consumer.py 2>&1"
|
||||
return None # 小果已归档,信号桥不再使用
|
||||
return None
|
||||
|
||||
|
||||
@@ -319,8 +316,6 @@ def check_db_table_count(table, field, value, op="today", threshold=0):
|
||||
cols = [r[1] for r in cur.execute(f"PRAGMA table_info({table})").fetchall()]
|
||||
if "processed" in cols:
|
||||
sql = f"SELECT COUNT(*) FROM {table} WHERE (processed = 0 OR processed IS NULL)"
|
||||
elif "source" in cols:
|
||||
sql = f"SELECT COUNT(*) FROM {table} WHERE source LIKE '%xiaoguo%'"
|
||||
else:
|
||||
sql = f"SELECT COUNT(*) FROM {table}"
|
||||
cur.execute(sql)
|
||||
@@ -400,7 +395,6 @@ def check_cron_paused():
|
||||
"""检查不应暂停的cron是否被误暂停"""
|
||||
should_run = [
|
||||
("3a9fb3300a6a", "价格监控"),
|
||||
("0851c7838ca3", "小果扫描"),
|
||||
("e13323928f3a", "自选提醒"),
|
||||
("b809fcabfa5b", "分支评估"),
|
||||
]
|
||||
@@ -604,10 +598,9 @@ def run_check(item):
|
||||
elif check_spec == "meta:checklist_completeness":
|
||||
ok, detail = check_meta_checklist_completeness()
|
||||
elif check_spec == "pipeline:xiaoguo_signal_flow":
|
||||
today_xiaoguo, d1 = check_db_table_count("signal_news", "created_at", None, "today", 0)
|
||||
unproc, d2 = check_db_table_count("signal_news", None, None, "unprocessed", 30)
|
||||
ok = today_xiaoguo or unproc
|
||||
detail = f"today_xiaoguo={d1}, unprocessed={d2}"
|
||||
# 小果已归档,此管道不再检查
|
||||
ok = True
|
||||
detail = "小果已归档,跳过信号流检查"
|
||||
elif check_spec == "pipeline:registry_audit":
|
||||
ok = True
|
||||
gaps = []
|
||||
|
||||
@@ -34,8 +34,11 @@ def _in_cooldown(code):
|
||||
|
||||
sys.path.insert(0, "/home/hmo/web-dashboard")
|
||||
sys.path.insert(0, "/home/hmo/MoFin")
|
||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) # profile-scripts 硬链目录
|
||||
from strategy_lifecycle import reassess_with_context as reassess_strategy
|
||||
from mo_data import read_decisions, read_portfolio
|
||||
from llm_client import call_llm, REASSESS_MODEL
|
||||
from mofin_db import snapshot_strategy_history
|
||||
|
||||
|
||||
def _build_full_analysis(code, entry, result):
|
||||
@@ -431,18 +434,76 @@ def main():
|
||||
except:
|
||||
pass
|
||||
|
||||
_prompt = f"""你是一个资深股票分析师。请对股票{code}做一个完整的12维矩阵分析(3横×4纵:大盘/行业/个股 × 基本面/消息面/技术面/资金面)。
|
||||
# ── 拉取已有策略全文 + 最近变更 ──
|
||||
_existing_full_analysis = ""
|
||||
_existing_changelog_text = "无变更记录"
|
||||
try:
|
||||
_edb = __import__('sqlite3').connect("/home/hmo/MoFin/data/mofin.db")
|
||||
_er = _edb.execute(
|
||||
"SELECT full_analysis, changelog_json FROM holding_strategies "
|
||||
"WHERE code=? AND status='active'", (code,)
|
||||
).fetchone()
|
||||
if _er:
|
||||
_existing_full_analysis = _er[0] or ""
|
||||
_cl_raw = _er[1] or ""
|
||||
if _cl_raw:
|
||||
_cl = __import__('json').loads(_cl_raw) if isinstance(_cl_raw, str) else _cl_raw
|
||||
if isinstance(_cl, list) and _cl:
|
||||
_recent = _cl[-3:]
|
||||
_existing_changelog_text = "\n".join(
|
||||
[f" [{c.get('timestamp','?')}] {c.get('action','?')}: {c.get('reason','')}"[:120]
|
||||
for c in reversed(_recent)]
|
||||
)
|
||||
_edb.close()
|
||||
except:
|
||||
pass
|
||||
|
||||
_prompt = f"""你是一个资深股票分析师。请对股票{code}评估现有策略是否仍然有效,并输出完整的新策略。
|
||||
|
||||
╔══════════════════════════════════════════════╗
|
||||
║ 📋 第一步:审阅原策略 ║
|
||||
╚══════════════════════════════════════════════╝
|
||||
|
||||
【原策略全文】(上次完整分析):
|
||||
{_existing_full_analysis or '暂无完整策略分析'}
|
||||
|
||||
【当前策略参数】:
|
||||
价格={price} 信号={result.get("timing_signal") or entry.get("timing_signal","")}
|
||||
买入区间={entry.get("entry_low",0)}~{entry.get("entry_high",0)}
|
||||
止损={entry.get("stop_loss",0)} 止盈={entry.get("take_profit",0)}
|
||||
RR={result.get("rr_ratio", entry.get("rr_ratio", 0))}
|
||||
策略={result.get("action") or entry.get("action","")}
|
||||
行业={(result.get("sector_context") or entry.get("sector_context",""))[:50]}(当日实时)
|
||||
技术={(result.get("tech_snapshot") or entry.get("tech_snapshot",""))[:200]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势)
|
||||
|
||||
【最近变更记录】:
|
||||
{_existing_changelog_text}
|
||||
|
||||
╔══════════════════════════════════════════════╗
|
||||
║ 📊 第二步:12维矩阵交叉分析 ║
|
||||
╚══════════════════════════════════════════════╝
|
||||
|
||||
⚠️ 重要:12个维度必须交叉对比,找出矛盾/共振点,给出综合判断。
|
||||
|
||||
当前数据(实时API,每条标注时间窗口,禁止使用模型训练数据):
|
||||
大盘={_macro_desc or "震荡"}(当日实时) | PE/市值={_pe_val} {_pb_val}(最新财报) | 价格={price} 区间={entry.get("entry_low",0)}~{entry.get("entry_high",0)} 止损={entry.get("stop_loss",0)} 止盈={entry.get("take_profit",0)} RR={result.get("rr_ratio",entry.get("rr_ratio",0))} | 信号={result.get("timing_signal") or entry.get("timing_signal","")} | 行业={(result.get("sector_context") or entry.get("sector_context",""))[:50]}(当日实时)
|
||||
策略={(result.get("action") or entry.get("action",""))[:200]}
|
||||
技术={(result.get("tech_snapshot") or entry.get("tech_snapshot",""))[:200]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势)
|
||||
当前实时数据(每条标注时间窗口,禁止使用模型训练数据):
|
||||
大盘={_macro_desc or "震荡"}(当日实时) | PE/市值={_pe_val} {_pb_val}(最新财报)
|
||||
资金流={_flow_note}(近5日累计)
|
||||
消息面={_news_note}(最近3条,自动标注抓取时间)
|
||||
|
||||
格式:
|
||||
╔══════════════════════════════════════════════╗
|
||||
║ 📝 第三步:决策输出 ║
|
||||
╚══════════════════════════════════════════════╝
|
||||
|
||||
请严格按以下顺序输出:
|
||||
|
||||
【维持或修改】判断当前策略是否仍然有效,回答「维持」或「修改」。
|
||||
|
||||
【修改点及理由】(如果维持,写「无需修改」;如果修改,逐条列出):
|
||||
- 修改什么参数/方向
|
||||
- 理由(引用具体维度矛盾或共振)
|
||||
|
||||
【最终新策略】(完整策略全文,self-contained,可直接存入DB)
|
||||
|
||||
【交叉分析】哪些维度矛盾/共振,关键信号
|
||||
① 大盘×基本面 ② 大盘×消息面 ③ 大盘×技术面 ④ 大盘×资金面
|
||||
⑤ 行业×基本面 ⑥ 行业×消息面 ⑦ 行业×技术面 ⑧ 行业×资金面
|
||||
@@ -452,26 +513,42 @@ def main():
|
||||
【综合结论】(买入/关注/观望/卖出)
|
||||
【操作建议】
|
||||
【建议止损】
|
||||
【建议止盈】"""
|
||||
try:
|
||||
_ur = __import__('urllib.request', fromlist=['Request'])
|
||||
_req = _ur.Request("http://127.0.0.1:8643/v1/chat/completions",
|
||||
data=__import__('json').dumps({"model":"deepseek-v4-flash","messages":[{"role":"user","content":_prompt}],"max_tokens":1024}).encode(),
|
||||
headers={"Content-Type":"application/json","Authorization":"Bearer hermes123"})
|
||||
_resp = _ur.build_opener(_ur.ProxyHandler({})).open(_req, timeout=300)
|
||||
_llm_out = __import__('json').loads(_resp.read().decode())["choices"][0]["message"]["content"]
|
||||
_full_analysis_text = _llm_out
|
||||
print(f" ✅ LLM12维分析完成({len(_full_analysis_text)}字)", flush=True)
|
||||
except Exception as _e:
|
||||
print(f" ❌ LLM12维分析失败: {_e}", file=__import__('sys').stderr)
|
||||
_full_analysis_text = None
|
||||
【建议止盈】
|
||||
【建议仓位】⚠️不可省略,非"买入"时写"不新建仓"
|
||||
|
||||
# 保存到DB
|
||||
⚠️ 输出纪律(必须遵守):
|
||||
1. 直接以【维持或修改】开头,禁止任何寒暄、开场白、分隔线
|
||||
2. 禁止输出 <structured_data> 或任何 XML/JSON/代码块
|
||||
3. 所有【】节标题一个都不能少"""
|
||||
_full_analysis_text = None
|
||||
try:
|
||||
_llm_result = call_llm(_prompt, max_tokens=4096, timeout=150, retries=1, backoff=20)
|
||||
if _llm_result["ok"]:
|
||||
_full_analysis_text = _llm_result["content"]
|
||||
print(f" ✅ LLM12维分析完成({len(_full_analysis_text)}字, {_llm_result['elapsed']:.1f}s)", flush=True)
|
||||
else:
|
||||
print(f" ❌ LLM12维分析失败({_llm_result['attempts']}次): {_llm_result['error'][:200]}", flush=True)
|
||||
except Exception as _e:
|
||||
print(f" ❌ LLM12维分析异常: {_e}", flush=True)
|
||||
|
||||
# ── 保存到DB(覆写前先快照)──
|
||||
_fa_conn = __import__('sqlite3').connect("/home/hmo/MoFin/data/mofin.db")
|
||||
_fa_conn.execute("UPDATE holding_strategies SET full_analysis=?, reassessed_at=? WHERE code=? AND status='active'", (_full_analysis_text, __import__('datetime').datetime.now().isoformat(), code))
|
||||
_fa_conn.commit()
|
||||
if _full_analysis_text:
|
||||
# 快照旧策略(使用共享函数)
|
||||
try:
|
||||
snapshot_strategy_history(_fa_conn, code, "per_stock_12d")
|
||||
except Exception as _se:
|
||||
print(f" ⚠️ 快照失败: {_se}", flush=True)
|
||||
|
||||
_fa_conn.execute(
|
||||
"UPDATE holding_strategies SET full_analysis=?, reassessed_at=? WHERE code=? AND status='active'",
|
||||
(_full_analysis_text, __import__('datetime').datetime.now().isoformat(), code))
|
||||
_fa_conn.commit()
|
||||
_fa_conn.close()
|
||||
print(f" ✅ 完整12维分析已保存({len(_full_analysis_text)}字)" if _full_analysis_text else f" ⚠️ 12维分析未完成,跳过保存")
|
||||
if _full_analysis_text:
|
||||
print(f" ✅ 完整12维分析已保存({len(_full_analysis_text)}字)")
|
||||
else:
|
||||
print(f" ⚠️ 12维分析未完成,跳过保存")
|
||||
print(f" [DB] holding_strategies 已更新: {code}")
|
||||
# 从LLM输出提取信号
|
||||
if _full_analysis_text and '【综合结论】' in _full_analysis_text:
|
||||
@@ -480,8 +557,14 @@ def main():
|
||||
if _sig_line:
|
||||
_sig = '买入' if '买入' in _sig_line[0] else '关注' if '关注' in _sig_line[0] else '观望' if '观望' in _sig_line[0] else '卖出' if '卖出' in _sig_line[0] else ''
|
||||
if _sig:
|
||||
__import__('sqlite3').connect('/home/hmo/MoFin/data/mofin.db').execute(
|
||||
"UPDATE holding_strategies SET timing_signal=? WHERE code=? AND status='active'", (_sig, code)).connection.commit()
|
||||
_ts_conn = __import__('sqlite3').connect('/home/hmo/MoFin/data/mofin.db')
|
||||
_ts_conn.execute(
|
||||
"UPDATE holding_strategies SET timing_signal=? WHERE code=? AND status='active'", (_sig, code))
|
||||
_ts_conn.commit()
|
||||
# 推荐操作 tag 同步(与 XMPP 动作级信号同源)
|
||||
from mofin_db import sync_recommend_tag
|
||||
sync_recommend_tag(_ts_conn, code, _sig)
|
||||
_ts_conn.close()
|
||||
print(f" ✅ LLM信号={_sig} 已写入")
|
||||
# 买入信号→推XMPP
|
||||
if _sig == "买入":
|
||||
@@ -490,10 +573,8 @@ def main():
|
||||
"SELECT name, price, entry_low, entry_high, stop_loss, take_profit, position_advice FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
if _nr2:
|
||||
_xm = f"📈 {_nr2[0] or code}({code}) 价{_nr2[1]}→12维买入信号!区间{_nr2[2]}~{_nr2[3]} 损{_nr2[4]} 盈{_nr2[5]} 仓位{_nr2[6] or '-'}"
|
||||
_xr = __import__('urllib.request').Request("http://127.0.0.1:5805/",
|
||||
data=__import__('json').dumps({"body": _xm, "to": "hmo@yoin.fun", "type": "chat"}).encode(),
|
||||
headers={"Content-Type": "application/json"})
|
||||
__import__('urllib.request').urlopen(_xr, timeout=5)
|
||||
from alert_helper import notify as _notify2, ACTION as _ACT2
|
||||
_notify2("买入信号", _xm, _ACT2)
|
||||
print(f" 📨 XMPP推送买入信号")
|
||||
except: pass
|
||||
except: pass
|
||||
|
||||
@@ -608,7 +608,16 @@ def run_once(round_label=""):
|
||||
# === 第三步:买入区偏离检测 + 自动重评 ===
|
||||
reassesed_codes = []
|
||||
# 先做急跌检测(仅持仓,自选股不推送暴跌告警)
|
||||
holdings_codes = {d["code"] for d in active if (d.get("shares") or 0) > 0}
|
||||
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"]
|
||||
# 非持仓跳过
|
||||
|
||||
@@ -142,7 +142,7 @@ def main():
|
||||
reassess_scripts.append(code)
|
||||
print(f"[AUTO_REASSESS] {name}({code}) 价{cur_price:.2f}偏离买入区中心{center:.2f} {drift:+.0f}% → 触发重评")
|
||||
if reassess_scripts:
|
||||
# 调用 per_stock_reassess
|
||||
# 调用 per_stock_reassess(每轮最多5只,防LLM慢导致整批超时;其余下轮继续)
|
||||
reassess_path = None
|
||||
for p in ['/home/hmo/MoFin/scripts/per_stock_reassess.py',
|
||||
'/home/hmo/.hermes/profiles/position-analyst/scripts/per_stock_reassess.py']:
|
||||
@@ -150,12 +150,20 @@ def main():
|
||||
reassess_path = p
|
||||
break
|
||||
if reassess_path:
|
||||
for code in reassess_scripts:
|
||||
r = subprocess.run(['python3', reassess_path, code],
|
||||
capture_output=True, text=True, timeout=60)
|
||||
out = r.stdout.strip()[:200] if r.stdout else ""
|
||||
err = r.stderr.strip()[:200] if r.stderr else ""
|
||||
print(f" → {code}: exited={r.returncode} {out}")
|
||||
MAX_PER_RUN = 5
|
||||
batch = reassess_scripts[:MAX_PER_RUN]
|
||||
if len(reassess_scripts) > MAX_PER_RUN:
|
||||
print(f"[AUTO_REASSESS] 本轮限{MAX_PER_RUN}只,剩余{len(reassess_scripts)-MAX_PER_RUN}只下轮继续")
|
||||
for code in batch:
|
||||
try:
|
||||
# LLM 重评冷启动 20-100s,deepseek-v4-pro 更慢 → 480s
|
||||
r = subprocess.run(['python3', reassess_path, code],
|
||||
capture_output=True, text=True, timeout=480)
|
||||
out = r.stdout.strip()[:200] if r.stdout else ""
|
||||
err = r.stderr.strip()[:200] if r.stderr else ""
|
||||
print(f" → {code}: exited={r.returncode} {out}")
|
||||
except subprocess.TimeoutExpired:
|
||||
print(f" → {code}: 超时480s(LLM仍慢),下轮重试")
|
||||
except Exception as e:
|
||||
print(f"[AUTO_REASSESS FAIL] {e}")
|
||||
# ----- 结束 自选股重评 -----
|
||||
|
||||
@@ -94,7 +94,21 @@ def run():
|
||||
|
||||
# 1. price_monitor 最近运行时间
|
||||
try:
|
||||
conn = sqlite3.connect("/home/hmo/MoFin/data/mofin.db")
|
||||
conn = None
|
||||
last_err = None
|
||||
# malformed 可能是 I/O 风暴下的瞬态 WAL 损坏(2026-07-21 事件):
|
||||
# checkpoint 后自愈。重试一次再告警,避免误报轰炸
|
||||
for _attempt in range(2):
|
||||
try:
|
||||
conn = sqlite3.connect("/home/hmo/MoFin/data/mofin.db")
|
||||
conn.execute("SELECT 1 FROM live_prices LIMIT 1").fetchone()
|
||||
break
|
||||
except Exception as e:
|
||||
last_err = e
|
||||
import time as _t
|
||||
_t.sleep(3)
|
||||
if conn is None:
|
||||
raise last_err
|
||||
lp = conn.execute("SELECT MAX(updated_at) FROM live_prices").fetchone()[0]
|
||||
if lp:
|
||||
lp_dt = datetime.fromisoformat(lp) if isinstance(lp, str) else lp
|
||||
|
||||
Reference in New Issue
Block a user