diff --git a/README.md b/README.md
index 9530e91..7601a74 100644
--- a/README.md
+++ b/README.md
@@ -389,6 +389,7 @@ IS_MOCK = False
| v1.6.0 | 目标跟踪:概念/主题/股票(watch_targets),概念智能体分析+受益个股,持仓管理与目标管理,报告支持双类型 |
| v1.7.0 | 每日定时报告:工作日 9:00 盘前 / 15:30 盘后,简版正文+详细版附件,全球市场数据,cron 定时+手动触发 |
| v1.8.0 | 🌙 静默期:舆情监控/持仓跟踪/定时报告三大自动化任务均可独立配置静默时段(多段/跨午夜),命中时段不发送邮件(日志标记 quiet) |
+| v1.9.0 | ①自动化页开关/参数全部**即时保存**(切换即生效,无需再点保存按钮,刷新不回弹);②持仓跟踪新增**重大变化推送方式**:聚合一起发(默认,一轮全部分析完统一发一封汇总邮件)/ 分散发(每只出结果立即单独发);③**行情数据自动延续**:新增 extend_daily 增量扩展(行情/新闻/指数从库内最新日期延续到最新交易日,幂等),服务每小时自动检查+数据管理页「📅 更新行情到最新交易日」手动触发,数据不再停留在建库日 |
---
diff --git a/app.py b/app.py
index 1ef8564..f01fe3d 100644
--- a/app.py
+++ b/app.py
@@ -619,6 +619,8 @@ def api_settings_save():
set_setting("tracking_enabled", _bool_str(track["tracking_enabled"]))
if "tracking_notify" in track:
set_setting("tracking_notify", _bool_str(track["tracking_notify"]))
+ if "tracking_push_mode" in track:
+ set_setting("tracking_push_mode", str(track["tracking_push_mode"]).strip().lower())
# 静默期(各自动化任务独立配置)
quiet = body.get("quiet") or {}
for prefix in ("monitor", "tracking", "report"):
@@ -786,6 +788,8 @@ def api_admin_stats():
tables = ("stocks", "stock_daily", "news", "institutions", "inst_ratings",
"fund_holdings", "watchlist", "analysis_cache", "analysis_history",
"strategy_backtests", "market_index")
+ latest_daily = query_one("SELECT MAX(date) d FROM stock_daily")
+ latest_news = query_one("SELECT MAX(publish_date) d FROM news")
return jsonify({
"tables": {t: table_count(t) for t in tables},
"vector": {
@@ -794,6 +798,8 @@ def api_admin_stats():
},
"is_mock": IS_MOCK,
"db": "stock_advisor.db",
+ "latest_daily": (latest_daily or {}).get("d"),
+ "latest_news": (latest_news or {}).get("d"),
})
@@ -814,6 +820,29 @@ def api_admin_reseed():
return jsonify({"ok": True, "msg": "重灌任务已启动,可在数据管理页刷新查看进度"})
+@app.route("/api/admin/update-data", methods=["POST"])
+def api_admin_update_data():
+ """把行情/新闻/指数增量扩展到最新交易日(模拟数据延续),后台执行"""
+ import threading
+
+ def run():
+ try:
+ from seed_data import extend_daily
+ with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f:
+ f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 开始更新数据\n")
+ res = extend_daily()
+ with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f:
+ f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 完成: {res}\n")
+ log.info("data extend done: %s", res)
+ except Exception as e:
+ log.exception("data extend fail")
+ with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f:
+ f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 失败: {e}\n")
+
+ threading.Thread(target=run, daemon=True).start()
+ return jsonify({"ok": True, "msg": "行情更新已启动(增量扩展到最新交易日),稍后刷新数据管理页查看"})
+
+
@app.route("/api/admin/healthcheck")
def api_admin_healthcheck():
"""外部依赖连通性检查"""
@@ -835,10 +864,35 @@ def api_admin_healthcheck():
if __name__ == "__main__":
+ import threading
init_db()
from engine.notifier import start_monitor
from engine.agent import start_tracking
start_monitor()
start_tracking()
+ # 行情自动更新器:每小时检查,若行情最新日期落后于最新交易日则增量扩展(模拟数据延续)
+ class DataUpdater(threading.Thread):
+ def __init__(self):
+ super().__init__(daemon=True, name="data-updater")
+ self._stop = threading.Event()
+
+ def run(self):
+ log.info("行情自动更新器启动(每小时检查)")
+ while not self._stop.is_set():
+ try:
+ import datetime as _dt
+ last = query_one("SELECT MAX(date) d FROM stock_daily")
+ today = _dt.date.today()
+ while today.weekday() >= 5:
+ today -= _dt.timedelta(days=1)
+ if not last or last["d"] < today.isoformat():
+ from seed_data import extend_daily
+ res = extend_daily()
+ log.info("行情自动扩展: %s", res)
+ except Exception as e:
+ log.warning("data updater error: %s", e)
+ self._stop.wait(3600)
+
+ DataUpdater().start()
print(f"✅ {SERVICE_NAME} 启动: http://0.0.0.0:{SERVICE_PORT}")
app.run(host=SERVICE_HOST, port=SERVICE_PORT, threaded=True)
diff --git a/engine/agent.py b/engine/agent.py
index fddc6f3..52681c1 100644
--- a/engine/agent.py
+++ b/engine/agent.py
@@ -222,8 +222,8 @@ def strip_json_block(text):
# ===================================================================== 执行
-def track_stock(code, focus=""):
- """执行一次跟踪,返回 {ok, report_id, meta, ...}"""
+def track_stock(code, focus="", notify=True):
+ """执行一次跟踪,返回 {ok, report_id, meta, ...}。notify=False 时不单独发邮件(聚合推送模式)"""
with _track_lock:
seg = collect_chain(code)
if not seg:
@@ -267,7 +267,8 @@ def track_stock(code, focus=""):
(code, s["name"], s["industry"], clean_report, json.dumps(meta, ensure_ascii=False),
json.dumps(sources, ensure_ascii=False)))
rid = query_one("SELECT MAX(id) id FROM tracking_reports")["id"]
- _notify_if_significant(rid, s, meta)
+ if notify:
+ _notify_if_significant(rid, s, meta)
return {"ok": True, "report_id": rid, "meta": meta}
except Exception as e:
log.exception("track %s fail", code)
@@ -323,8 +324,9 @@ def _notify_if_significant(rid, stock, meta):
# ===================================================================== 概念/主题跟踪
-def track_concept(name, keywords="", focus=""):
- """跟踪一个概念/主题:采集相关新闻 → 受益个股梳理 → LLM 深度分析 → 影响度判定"""
+def track_concept(name, keywords="", focus="", notify=True):
+ """跟踪一个概念/主题:采集相关新闻 → 受益个股梳理 → LLM 深度分析 → 影响度判定。
+ notify=False 时不单独发邮件(聚合推送模式)"""
kw = keywords.strip() or name
kw_list = [k.strip() for k in kw.replace(",", ",").split(",") if k.strip()]
kws = " ".join(kw_list)
@@ -411,7 +413,8 @@ def track_concept(name, keywords="", focus=""):
(cid, name, "概念/主题", clean_report, json.dumps(meta, ensure_ascii=False),
json.dumps(sources, ensure_ascii=False)))
rid = query_one("SELECT MAX(id) id FROM tracking_reports")["id"]
- _notify_if_significant(rid, {"name": name, "code": cid}, meta)
+ if notify:
+ _notify_if_significant(rid, {"name": name, "code": cid}, meta)
return {"ok": True, "report_id": rid, "meta": meta}
except Exception as e:
log.exception("track concept %s fail", name)
@@ -437,9 +440,69 @@ def delete_target(tid):
return {"ok": True}
+def _notify_aggregate(items):
+ """聚合推送:一轮跟踪全部完成后,统一发一封汇总邮件(默认推送方式)"""
+ if not items:
+ return
+ cfg = tracking_config()
+ try:
+ if not cfg["notify"]:
+ return
+ qc = quiet_config("tracking")
+ if qc["enabled"] and in_quiet_period(qc["ranges"]):
+ log.info("tracking aggregate notify suppressed by quiet period (%d items)", len(items))
+ return
+ from engine.notifier import send_email
+ mc = mail_config()
+ rows = ""
+ for it in items:
+ meta = it.get("meta") or {}
+ is_concept = it.get("type") == "concept"
+ impact = int(meta.get("impact_score", 0))
+ level = "🔴 重大" if impact >= 65 else "🟡 关注"
+ name = escape_html(it.get("name", ""))
+ code = it.get("code", "")
+ rows += (
+ f"
"
+ f"| {'🎯 概念' if is_concept else '💼 股票'} | "
+ f"{name}"
+ f"{' ' + str(code) + ' ' if not is_concept and code else ''} | "
+ f"{meta.get('change_kind', '')} | "
+ f"{impact}/100 ({level}) | "
+ f"{escape_html(meta.get('summary', ''))} | "
+ f"
"
+ )
+ send_email(
+ f"[持仓跟踪] 聚合报告:{len(items)} 个目标出现重大动态(影响度≥{int(cfg['impact_threshold'])})",
+ f"""
+
+
🧭 持仓跟踪聚合报告 · {time.strftime('%Y-%m-%d %H:%M')}
+
+
本轮共跟踪 {items[0].get('_total', len(items))} 个目标,其中 {len(items)} 个达到重大变化判定阈值(影响度 ≥ {int(cfg['impact_threshold'])}):
+
+
详细分析报告请在系统「🧭 持仓跟踪智能体 → 最近跟踪报告」中查看。
+
""",
+ cfg=mc)
+ set_tracking_state(last_alert=int(items[0]["meta"].get("impact_score", 0)))
+ log.info("tracking aggregate notify sent: %d items", len(items))
+ except Exception as e:
+ log.warning("tracking aggregate notify fail: %s", e)
+
+
+def escape_html(s):
+ return str(s or "").replace("&", "&").replace("<", "<").replace(">", ">").replace('"', """)
+
+
# ===================================================================== 批量与调度
def track_all(progress=None):
- """跟踪全部目标:持仓(自选股)+ 目标(概念/主题/股票)"""
+ """跟踪全部目标:持仓(自选股)+ 目标(概念/主题/股票)。
+ 推送方式(tracking_push_mode):
+ aggregate(默认) —— 一轮全部分析完后,汇总重大变化统一发一封邮件;
+ scattered —— 每只股票分析完立即单独发邮件。
+ """
global _cycle_running
if _cycle_running:
return {"tracked": 0, "msg": "已有跟踪任务进行中,请稍后再试"}
@@ -465,16 +528,33 @@ def track_all(progress=None):
uniq.append((typ, key, name))
if not uniq:
return {"tracked": 0, "msg": "暂无跟踪目标:请添加持仓/自选股或概念主题目标"}
+ cfg = tracking_config()
+ scattered = cfg.get("push_mode", "aggregate") == "scattered"
+ significant = [] # 聚合推送模式收集重大变化
results = []
for i, (typ, key, name) in enumerate(uniq):
- if typ == "stock":
- r = track_stock(key)
- else:
- r = track_concept(key, name)
+ try:
+ if typ == "stock":
+ r = track_stock(key, notify=scattered)
+ else:
+ r = track_concept(key, name, notify=scattered)
+ except Exception as e:
+ r = {"error": str(e)}
results.append({"type": typ, "name": name, **r})
+ if r.get("ok"):
+ meta = r.get("meta") or {}
+ if int(meta.get("impact_score", 0)) >= int(cfg["impact_threshold"]):
+ significant.append({"type": typ, "name": name, "code": key,
+ "meta": meta, "_total": len(uniq)})
set_tracking_state(last_run=time.strftime("%Y-%m-%d %H:%M:%S"), last_stock=name)
if progress:
progress(i + 1, len(uniq))
+ # 聚合推送:一轮分析完统一发一封汇总邮件
+ if not scattered and significant:
+ try:
+ _notify_aggregate(significant)
+ except Exception as e:
+ log.warning("aggregate notify error: %s", e)
return {"tracked": len(results), "results": results}
finally:
_cycle_running = False
diff --git a/seed_data.py b/seed_data.py
index c033e98..6cf6ee4 100644
--- a/seed_data.py
+++ b/seed_data.py
@@ -21,7 +21,7 @@ import random
import sys
from config import (CHROMA_NEWS_COLLECTION, CHROMA_PROFILE_COLLECTION, LOG_DIR)
-from database import init_db, executemany, query_one, wipe_all
+from database import init_db, executemany, query_one, query, execute, wipe_all
from rag import vector_store as vs
random.seed(42)
@@ -427,6 +427,123 @@ def build_vectors(news, stocks):
print(f" 概况索引条数: {vs.collection_count(CHROMA_PROFILE_COLLECTION)}")
+# ===================================================================== 增量更新
+def extend_daily(target_date=None):
+ """把模拟行情/指数/新闻从库里最新日期增量扩展到最新交易日(幂等,可重复调用)。
+ 返回 (起始日期, 结束日期, 新增交易日数, 新增新闻数) 或 (None, None, 0, 0)。
+ """
+ init_db()
+ last = query_one("SELECT MAX(date) d FROM stock_daily")
+ last_date = last["d"] if last else None
+ today = dt.date.today()
+ if target_date is None:
+ target_date = today
+ if isinstance(target_date, str):
+ target_date = dt.date.fromisoformat(target_date)
+ # 最新交易日:今天若是周末则回退到周五
+ while target_date.weekday() >= 5:
+ target_date -= dt.timedelta(days=1)
+ if last_date:
+ start = dt.date.fromisoformat(last_date) + dt.timedelta(days=1)
+ else:
+ start = target_date - dt.timedelta(days=179)
+ dates = []
+ d = start
+ while d <= target_date:
+ if d.weekday() < 5:
+ dates.append(d.isoformat())
+ d += dt.timedelta(days=1)
+ if not dates:
+ return None, None, 0, 0
+ print(f">>> 增量扩展行情 {dates[0]} ~ {dates[-1]}({len(dates)} 个交易日)...")
+
+ stocks = query("SELECT code,name,industry,board,total_shares,float_shares FROM stocks")
+ daily = []
+ index = {}
+ price = {}
+ for s in stocks:
+ rows = query("SELECT date,close,volume FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 30",
+ (s["code"],))
+ if not rows:
+ continue
+ last_close = rows[0]["close"]
+ base_vol = rows[0]["volume"] or 1
+ closes = [r["close"] for r in reversed(rows)]
+ rets = [(closes[i + 1] - closes[i]) / closes[i] for i in range(len(closes) - 1)]
+ mean_r = sum(rets) / max(len(rets), 1)
+ vol = (sum((r - mean_r) ** 2 for r in rets) / max(len(rets), 1)) ** 0.5 if rets else 0.02
+ vol = max(0.005, min(0.05, vol))
+ drift = 0.0002
+ p = prev = last_close
+ for dd in dates:
+ r = random.gauss(drift, vol)
+ if random.random() < 0.02:
+ r += random.gauss(0, vol * 1.6)
+ prev = p
+ p = max(0.5, p * (1 + r))
+ open_p = prev * (1 + random.gauss(0, vol * 0.5))
+ high = max(open_p, p) * (1 + abs(random.gauss(0, vol * 0.35)))
+ low = min(open_p, p) * (1 - abs(random.gauss(0, vol * 0.35)))
+ volume = base_vol * (1 + 1.5 * abs(r) / max(vol, 1e-6)) * random.uniform(0.6, 1.4)
+ amount = volume * (open_p + p) / 2
+ chg = (p - prev) / prev * 100
+ daily.append((s["code"], dd, round(open_p, 2), round(high, 2), round(low, 2),
+ round(p, 2), round(volume, 0), round(amount, 0), round(chg, 2)))
+ price[s["code"]] = p
+
+ # 指数延续
+ idx_rows = query("SELECT date,sh,sz,cy FROM market_index ORDER BY date DESC LIMIT 1")
+ if idx_rows:
+ sh, sz, cy = idx_rows[0]["sh"], idx_rows[0]["sz"], idx_rows[0]["cy"]
+ else:
+ sh, sz, cy = 3245.0, 10580.0, 2120.0
+ for dd in dates:
+ sh_r = sum(random.gauss(0.0004, 0.008) for _ in range(6)) / 6
+ sz_r = sh_r + random.gauss(0, 0.004)
+ cy_r = sh_r + random.gauss(0, 0.006)
+ sh *= (1 + sh_r); sz *= (1 + sz_r); cy *= (1 + cy_r)
+ index[dd] = {"sh": round(sh, 2), "sz": round(sz, 2), "cy": round(cy, 2)}
+
+ executemany(
+ "INSERT OR REPLACE INTO stock_daily(code,date,open,high,low,close,volume,amount,change_pct) "
+ "VALUES(?,?,?,?,?,?,?,?,?)", daily)
+ executemany("INSERT OR REPLACE INTO market_index(date,sh,sz,cy) VALUES(?,?,?,?)",
+ [(d, v["sh"], v["sz"], v["cy"]) for d, v in index.items()])
+ gm = _gen_global(dates)
+ executemany("INSERT OR REPLACE INTO global_markets(date,data) VALUES(?,?)",
+ [(d, json.dumps(v, ensure_ascii=False)) for d, v in gm.items()])
+ # 回填市值
+ for code in price:
+ execute("UPDATE stocks SET market_cap=ROUND((SELECT close FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 1)*total_shares,2) WHERE code=?",
+ (code, code))
+
+ # 新增新闻(仅保留落在新增日期内的)
+ news = gen_news(dates, price)
+ date_set = set(dates)
+ new_news = [n for n in news if n["publish_date"] in date_set]
+ if new_news:
+ executemany(
+ "INSERT INTO news(title,content,source,category,publish_date,related_stocks,sentiment,is_positive) "
+ "VALUES(?,?,?,?,?,?,?,?)",
+ [(n["title"], n["content"], n["source"], n["category"], n["publish_date"],
+ n["related"], n["sentiment"], n["is_positive"]) for n in new_news])
+ # 向量追加
+ ids, docs, metas = [], [], []
+ for n in new_news:
+ for code in n["related"].split(","):
+ ids.append(f"news-{n['publish_date']}-{n['title']}-{code}")
+ docs.append(f"{n['title']}\n{n['content']}")
+ metas.append({"code": code, "title": n["title"], "date": n["publish_date"],
+ "category": n["category"], "sentiment": n["sentiment"], "news_id": 0})
+ try:
+ for i in range(0, len(ids), 16):
+ vs.add_documents(ids[i:i + 16], docs[i:i + 16], metas[i:i + 16], CHROMA_NEWS_COLLECTION)
+ except Exception as e:
+ print("新闻向量追加失败(可忽略,下次重灌会重建):", e)
+ print(f">>> 新增新闻 {len(new_news)} 条(向量已追加)")
+ return dates[0], dates[-1], len(dates), len(new_news)
+
+
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--skip-vector", action="store_true", help="跳过向量索引重建")
diff --git a/settings.py b/settings.py
index edc93a9..4017196 100644
--- a/settings.py
+++ b/settings.py
@@ -115,6 +115,8 @@ def tracking_config():
"interval_min": max(15, int(get_setting("tracking_interval", TRACKING_DEFAULTS["tracking_interval"]))),
"notify": _truthy(get_setting("tracking_notify", TRACKING_DEFAULTS["tracking_notify"])),
"impact_threshold": float(get_setting("tracking_impact_threshold", TRACKING_DEFAULTS["tracking_impact_threshold"])),
+ # 推送方式: aggregate=聚合一起发(默认, 一轮全部分析完统一发一封汇总) / scattered=分散发(每只出结果立即单独发)
+ "push_mode": (get_setting("tracking_push_mode", "aggregate") or "aggregate").strip().lower(),
}
diff --git a/static/js/admin.js b/static/js/admin.js
index 81c3fe1..4345ef0 100644
--- a/static/js/admin.js
+++ b/static/js/admin.js
@@ -22,6 +22,12 @@ async function refreshStats() {
🏢 公司概况索引 (stock_profiles_v1)
${d.vector.profiles} 条
+
+ 行情最新日期${d.latest_daily || '—'}
+
+
+ 新闻最新日期${d.latest_news || '—'}
+
数据模式
${d.is_mock ? '模拟数据' : '真实数据'}
@@ -98,6 +104,17 @@ async function reseed() {
}
}
+async function updateData() {
+ try {
+ const r = await api('/api/admin/update-data', { method: 'POST' });
+ toast(r.msg || '行情更新已启动');
+ setTimeout(refreshStats, 8000);
+ setTimeout(refreshStats, 20000);
+ } catch (e) {
+ toast('启动失败:' + e.message);
+ }
+}
+
async function health() {
$('#healthBox').innerHTML = '
检查中…
';
try {
diff --git a/static/js/automation.js b/static/js/automation.js
index 96c2e39..f2a797a 100644
--- a/static/js/automation.js
+++ b/static/js/automation.js
@@ -30,6 +30,7 @@ async function loadAuto() {
$('#trackingInterval').value = tr.interval_min;
$('#trackingImpact').value = tr.impact_threshold;
$('#trackingNotify').checked = !!tr.notify;
+ $('#trackingPushMode').value = (tr.push_mode === 'scattered') ? 'scattered' : 'aggregate';
renderTrackingState(tr);
// 静默期(跟踪)
const tq = (autoSettings.quiet && autoSettings.quiet.tracking) || {};
@@ -100,7 +101,8 @@ async function saveTracking() {
tracking_enabled: $('#trackingEnabled').checked,
tracking_interval: $('#trackingInterval').value,
tracking_impact_threshold: $('#trackingImpact').value,
- tracking_notify: $('#trackingNotify').checked
+ tracking_notify: $('#trackingNotify').checked,
+ tracking_push_mode: $('#trackingPushMode').value
}, quiet: { tracking: collectTrackingQuiet() } }) });
toast(r.msg || '已保存');
} catch (e) { toast('保存失败:' + e.message); }
diff --git a/templates/admin.html b/templates/admin.html
index 9ded37a..bba491a 100644
--- a/templates/admin.html
+++ b/templates/admin.html
@@ -47,13 +47,15 @@
系统维护
+
- 当前为模拟数据,用于功能演示与系统验证。接入真实数据后(行情/新闻/机构),修改 config.py 中的 IS_MOCK=False,并替换 seed_data.py 为真实数据源即可,其余分析/检索逻辑无需改动。
+ 当前为模拟数据,用于功能演示与系统验证。📅 更新行情会把模拟行情/新闻/指数从库内最新日期增量延续到最新交易日(服务每小时自动检查一次,无需手动)。
+ 接入真实数据后(行情/新闻/机构),修改 config.py 中的 IS_MOCK=False,并替换 seed_data.py 为真实数据源即可,其余分析/检索逻辑无需改动。
技术栈:Python Flask + SQLite + Chroma 向量库 + bge-large-zh 语义检索 + DeepSeek 大模型(RAG 增强研报)。
diff --git a/templates/automation.html b/templates/automation.html
index fd9432a..42adac5 100644
--- a/templates/automation.html
+++ b/templates/automation.html
@@ -12,30 +12,30 @@
开启后,落在静默时段内的扫描将不发送邮件(通知日志标记为“静默”,水位正常推进,不会在结束后补发积压)。
@@ -59,22 +59,32 @@
开启后,落在静默时段内发现的重大变化将不发送邮件(跟踪报告照常生成入库)。