From 68841af662edee1301cec34d7b623eb24a7bccd3 Mon Sep 17 00:00:00 2001 From: shy3130 <415333856@qq.com> Date: Tue, 30 Jun 2026 17:28:18 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=9A=E6=97=B6=E5=A4=8D=E7=9B=98?= =?UTF-8?q?=E9=A3=9E=E4=B9=A6=E6=8E=A8=E9=80=81=20+=20SSE=20=E5=AE=9E?= =?UTF-8?q?=E6=97=B6=E8=BF=9B=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 定时复盘支持推送到飞书(卡片消息, 完整报告), 渠道多选(飞书可选, 微信开发中) - 推送开关独立常驻, 与定时/实时行情解耦; 手动/定时生成都可推送 - 定时复盘改流式生成, 通过 SSE(review_progress)实时推给前端, 开着页面可见边生成边显示 - 修复定时任务协程未 await 的 bug(lambda 包裹 async → 改传函数对象 + args) - LLM 断流自动重试(最多2次), 后端归档 + 飞书推送, 异常兜底通知前端 - 复盘定时下限改为 15:00, 默认 15:10 - 版本号 0.1.66 → 0.1.67 --- backend/app/api/intraday.py | 11 ++ backend/app/api/market_recap.py | 7 ++ backend/app/api/screener.py | 26 ++++- backend/app/api/settings.py | 20 +++- backend/app/jobs/daily_pipeline.py | 138 +++++++++++++++++++++++- backend/app/services/depth_service.py | 13 ++- backend/app/services/preferences.py | 61 +++++++++-- backend/app/services/quote_service.py | 122 ++++++--------------- backend/app/services/strategy_cache.py | 21 ++-- backend/app/services/webhook_adapter.py | 100 +++++++++++++---- backend/app/strategy/monitor.py | 121 ++++++++++++++++++++- backend/pyproject.toml | 2 +- frontend/package.json | 2 +- frontend/src/components/AlertToast.tsx | 4 +- frontend/src/components/Layout.tsx | 3 +- frontend/src/lib/api.ts | 6 ++ frontend/src/lib/queryKeys.ts | 2 +- frontend/src/lib/reviewStore.ts | 60 ++++++++++- frontend/src/lib/useQuoteStream.ts | 18 +++- frontend/src/lib/useSharedQueries.ts | 14 ++- frontend/src/pages/Review.tsx | 84 ++++++++++++++- frontend/src/pages/Screener.tsx | 6 +- frontend/src/pages/Watchlist.tsx | 25 ++++- 23 files changed, 707 insertions(+), 159 deletions(-) diff --git a/backend/app/api/intraday.py b/backend/app/api/intraday.py index 6271b8b..64b5d8a 100644 --- a/backend/app/api/intraday.py +++ b/backend/app/api/intraday.py @@ -136,6 +136,9 @@ async def quote_stream(request: Request): "depth": asyncio.ensure_future( asyncio.to_thread(qs.wait_for_depth_update, timeout=5.0) if qs else asyncio.sleep(5) ), + "review": asyncio.ensure_future( + asyncio.to_thread(qs.wait_for_review, timeout=5.0) if qs else asyncio.sleep(5) + ), } done, pending = await asyncio.wait( @@ -160,6 +163,14 @@ async def quote_stream(request: Request): }, ensure_ascii=False), } + # 推送复盘进度 (定时复盘流式生成时) — 前端 reviewStore 直接消费 + # 事件已是 recap_market_stream 产出的 JSON 字符串, 逐条转发 + for evt_json in qs.pop_review_events(): + yield { + "event": "review_progress", + "data": evt_json, + } + # 推送行情更新 (行情信号触发) if tasks["quote"] in done: try: diff --git a/backend/app/api/market_recap.py b/backend/app/api/market_recap.py index 8f7a337..7215d23 100644 --- a/backend/app/api/market_recap.py +++ b/backend/app/api/market_recap.py @@ -98,6 +98,13 @@ def save_report(request: Request, req: SaveReportRequest): "emotion_score": req.emotion_score, "emotion_label": req.emotion_label, }) + # 推送到飞书(可选): 与定时复盘共用同一开关 review_push_enabled 与 _maybe_push_review。 + # 内部 try/except 静默降级, 不影响归档返回值。 + from app.jobs.daily_pipeline import _maybe_push_review + _maybe_push_review(req.content, { + "as_of": req.as_of, + "emotion_label": req.emotion_label, + }) return {"ok": True, "report": report} diff --git a/backend/app/api/screener.py b/backend/app/api/screener.py index a64a837..d35faef 100644 --- a/backend/app/api/screener.py +++ b/backend/app/api/screener.py @@ -278,11 +278,35 @@ def get_cached( request: Request, ext_columns: Optional[str] = Query(None, description="逗号分隔: config_id.field_name"), ): - """读取策略结果缓存。返回 None 表示无缓存。""" + """读取策略结果缓存, 并叠加监控引擎本轮实时算出的结果。 + + - 盘后缓存 (strategy_cache.json): 非监控策略 / 页面秒加载用, run_all 写入。 + - 监控引擎内存结果 (latest_strategy_results): 实时行情每轮对「加入监控的策略」算出, + 不落盘 (避免与 read_cache 的 mtime 校验冲突), 在此直接叠加覆盖盘后结果。 + 被监控的策略拿到新鲜数据, 非监控策略仍用盘后缓存。 + """ data_dir = request.app.state.repo.store.data_dir cached = strategy_cache.read_cache(data_dir) if cached is None: + cached = {"as_of": None, "results": {}, "updated_at": None} + + # 叠加监控引擎内存里的实时结果 (若有), 用新鲜数据覆盖同策略的盘后结果 + monitor_engine = getattr(request.app.state, "monitor_engine", None) + if monitor_engine is not None: + realtime_results = monitor_engine.latest_strategy_results() + if realtime_results: + results = dict(cached.get("results") or {}) + results.update(realtime_results) + cached = dict(cached) + cached["results"] = results + # 有实时数据时, 以最新时间戳为准 + import time as _time + cached["updated_at"] = int(_time.time() * 1000) + + # 无任何数据 (盘后缓存空 + 无实时结果) → 返回空标记, 前端据此提示 + if not cached.get("results") and cached.get("as_of") is None: return {"as_of": None, "results": {}, "updated_at": None} + ext_values = _load_ext_value_maps(request.app.state.repo, ext_columns) return _cache_payload_with_ext(cached, ext_values) diff --git a/backend/app/api/settings.py b/backend/app/api/settings.py index 2fe3f43..2ec0fcd 100644 --- a/backend/app/api/settings.py +++ b/backend/app/api/settings.py @@ -342,6 +342,7 @@ def get_preferences() -> dict: "depth_polling_interval": preferences.get_depth_polling_interval(), "depth_finalize_time": preferences.get_depth_finalize_time(), "review_schedule": preferences.get_review_schedule(), + "review_push_channels": preferences.get_review_push_channels(), } @@ -1014,7 +1015,7 @@ def update_review_schedule(req: ReviewScheduleIn, request: Request) -> dict: - enabled=True: 注册/更新 job(工作日定时生成复盘报告) - enabled=False: 移除 job(停止定时复盘) - 校验: 开启时若 AI Key 未配置则拒绝(复盘依赖 AI), 提示用户先配置。 - - 时间下限 15:30(盘后数据就绪), 由 preferences 层强制。 + - 时间下限 15:00(A股收盘), 由 preferences 层强制。 """ from app.services import preferences @@ -1045,3 +1046,20 @@ def update_review_schedule(req: ReviewScheduleIn, request: Request) -> dict: return sched + +class ReviewPushIn(BaseModel): + channels: list[str] # 多选: ['feishu'] 等; 空数组=不推送。微信等开发中 + + +@router.put("/preferences/review-push") +def update_review_push(req: ReviewPushIn) -> dict: + """复盘推送渠道(多选) — 选定把复盘报告(手动生成 / 定时生成归档后)推送到哪些外部工具。 + + 纯偏好, 与定时复盘 / 实时行情完全独立, 常驻可单独设置。空数组=不推送。 + 实际推送由归档端点(POST /api/market-recap/reports)与定时任务(_run_scheduled_review) + 在归档后读取本列表逐个推送。白名单外的渠道会被过滤掉。 + """ + from app.services import preferences + saved = preferences.set_review_push_channels(req.channels) + return {"review_push_channels": saved} + diff --git a/backend/app/jobs/daily_pipeline.py b/backend/app/jobs/daily_pipeline.py index bf02072..6bb7099 100644 --- a/backend/app/jobs/daily_pipeline.py +++ b/backend/app/jobs/daily_pipeline.py @@ -558,13 +558,17 @@ REVIEW_JOB_ID = "scheduled_review" async def _run_scheduled_review(repo) -> None: - """定时复盘 job: 调用非流式复盘生成 → 落盘归档(与手动生成同格式)。 + """定时复盘 job: 流式生成复盘 → 实时推 SSE(开着页面可见) → 落盘归档 → 推飞书。 - 静默执行, 不推送 SSE/系统通知 —— 用户下次打开复盘页即可看到新报告。 + 与手动「生成复盘」体验一致: 流式事件经 quote_service.push_review_event → + /api/intraday/stream 的 review_progress 事件 → 前端 reviewStore, 用户开着复盘页 + 即可看到报告边生成边显示, 切走再回来也能看到生成中/已生成。 + LLM 偶发断流(peer closed connection)时自动重试最多 2 次。 任何异常都吞掉只记日志, 绝不影响调度器主循环。 """ + import json + try: - from app.services.market_recap import recap_market_once from app.services import market_recap_reports from app import secrets_store as ss @@ -577,9 +581,14 @@ async def _run_scheduled_review(repo) -> None: quote_service = getattr(app_state, "quote_service", None) if app_state else None depth_service = getattr(app_state, "depth_service", None) if app_state else None - content, meta = await recap_market_once(repo, quote_service, depth_service) + content, meta = await _stream_review_with_retry(repo, quote_service, depth_service) if not content: logger.warning("scheduled review produced no content (meta=%s)", meta) + # 通知前端进入 error 态(若有页面在听) + if quote_service: + quote_service.push_review_event(json.dumps( + {"type": "error", "message": "复盘生成失败,请稍后手动重试"}, + ensure_ascii=False)) return # 落盘: 与手动生成完全相同的归档格式 @@ -592,8 +601,122 @@ async def _run_scheduled_review(repo) -> None: "emotion_label": meta.get("emotion_label", ""), }) logger.info("scheduled review saved: as_of=%s", meta.get("as_of")) + + # 通知前端: 生成完成且已归档(archived=true 让前端只刷新列表, 不重复归档) + if quote_service: + quote_service.push_review_event(json.dumps( + {"type": "done", "archived": True}, ensure_ascii=False)) + + # 推送到飞书(可选): 运行时读取配置, 用户改设置下次触发即生效。 + # 失败静默降级, 不影响已归档的报告。 + _maybe_push_review(content, meta) except Exception as e: # noqa: BLE001 logger.exception("scheduled review failed: %s", e) + # 兜底: 异常时通知前端停止「生成中」状态, 避免页面卡在 streaming + try: + app_state = _get_app_state() + qs = getattr(app_state, "quote_service", None) if app_state else None + if qs: + import json as _json + qs.push_review_event(_json.dumps( + {"type": "error", "message": "复盘生成异常,请稍后手动重试"}, + ensure_ascii=False)) + except Exception: # noqa: BLE001 + pass + + +async def _stream_review_with_retry(repo, quote_service, depth_service) -> tuple[str, dict]: + """流式生成复盘, 每个事件推 SSE + 累积内容。LLM 断流时最多重试 2 次。 + + 返回 (content, meta)。重试时推一个 retry 事件让前端清空已累积内容重新开始。 + 成功(收到 done/无 error)或耗尽重试后返回。 + """ + import asyncio + import json + from app.services.market_recap import recap_market_stream + + max_attempts = 3 # 初次 + 2 次重试 + last_meta: dict = {} + content_parts: list[str] = [] + + for attempt in range(1, max_attempts + 1): + content_parts = [] # 每次重试重新累积 + failed = False + try: + async for evt_json in recap_market_stream(repo, quote_service, depth_service): + evt = json.loads(evt_json) + t = evt.get("type") + + # 推给前端(让开着页面的用户实时看到, 与手动一致) + if quote_service: + quote_service.push_review_event(evt_json) + + if t == "meta": + last_meta = evt + elif t == "delta" and evt.get("content"): + content_parts.append(evt["content"]) + elif t == "error": + failed = True + logger.warning("scheduled review stream error (attempt %d/%d): %s", + attempt, max_attempts, evt.get("message")) + break # 触发重试 + elif t == "done": + # 正常完成 + return "".join(content_parts), last_meta + # 流自然结束(无 done 事件)且有内容, 视为成功 + if content_parts and not failed: + return "".join(content_parts), last_meta + except Exception as e: # noqa: BLE001 + # LLM 断流等异常(httpx.RemoteProtocolError)落到这里 + failed = True + logger.warning("scheduled review stream exception (attempt %d/%d): %s", + attempt, max_attempts, e) + + # 失败: 决定是否重试 + if attempt < max_attempts: + logger.info("scheduled review retrying in 3s (attempt %d → %d)", attempt, attempt + 1) + # 通知前端: 即将重试, 清空已累积内容重新开始 + if quote_service: + quote_service.push_review_event(json.dumps( + {"type": "retry", "attempt": attempt + 1}, ensure_ascii=False)) + await asyncio.sleep(3) + + # 耗尽重试, 返回已累积内容(可能为空)和最后 meta + return "".join(content_parts), last_meta + + +def _maybe_push_review(content: str, meta: dict) -> None: + """复盘报告归档后, 按 review_push_channels 选定的外部工具逐个推送完整报告。 + + 定时生成与手动生成共用本函数 (手动归档端点 POST /api/market-recap/reports 也会调用)。 + channels 为空则不推送; 'feishu' 复用监控中心的全局飞书 Webhook 通道。 + 推送失败静默降级 (Webhook 是辅助通道), 不影响已归档的报告。 + """ + try: + from app.services import preferences, webhook_adapter + + channels = preferences.get_review_push_channels() + if not channels: + return + + emotion = f"{meta.get('emotion_label') or ''}".strip() + as_of = meta.get("as_of") or "" + subtitle = as_of + (f" · 情绪 {emotion}" if emotion else "") + + for ch in channels: + if ch == "feishu": + url = preferences.get_feishu_webhook_url() + if not url: + logger.info("review push(feishu) skipped: webhook not configured") + continue + secret = preferences.get_feishu_webhook_secret() + ok = webhook_adapter.send_feishu_card( + url, "TickFlow · 每日复盘", subtitle, content, secret + ) + logger.info("review push(feishu) %s", "sent" if ok else "failed") + # 未来更多渠道在此追加分支 + except Exception as e: # noqa: BLE001 + logger.warning("review push error: %s", e) def _register_review_job(scheduler, repo, hour: int, minute: int) -> None: @@ -601,9 +724,14 @@ def _register_review_job(scheduler, repo, hour: int, minute: int) -> None: 供 start_scheduler(启动时) 和 settings API(改时间时) 共用。 用 replace_existing=True, 重复注册只更新 trigger。 + + 注意: _run_scheduled_review 是协程函数, 必须把函数对象本身(配合 args)传给 + add_job, 而非用 lambda 包裹 —— 否则 APScheduler 会把 lambda 当同步函数在线程池 + 执行, 仅得到一个未 await 的协程对象, 复盘实际不会运行。 """ scheduler.add_job( - lambda: _run_scheduled_review(repo), + _run_scheduled_review, + args=[repo], trigger=CronTrigger(day_of_week="mon-fri", hour=hour, minute=minute, timezone="Asia/Shanghai"), diff --git a/backend/app/services/depth_service.py b/backend/app/services/depth_service.py index 4978a92..1e5c423 100644 --- a/backend/app/services/depth_service.py +++ b/backend/app/services/depth_service.py @@ -316,7 +316,18 @@ class DepthService: "status": e.get("status"), "fetched_at": e.get("fetched_ts"), }) - df = pl.DataFrame(rows) + # 显式 schema: sealed_up/sealed_down 是 bool 与 None 混合, 不指定 schema + # polars 会按首行推断类型, 后续遇到不一致 (bool vs null) 报 + # "could not append value: false of type: bool to the builder"。 + df = pl.DataFrame(rows, schema={ + "symbol": pl.Utf8, + "sealed_up": pl.Boolean, + "sealed_down": pl.Boolean, + "ask1_vol": pl.Int64, + "bid1_vol": pl.Int64, + "status": pl.Utf8, + "fetched_at": pl.Float64, + }) ds = today.isoformat() out = self._repo.store.data_dir / "depth5" / f"date={ds}" / "part.parquet" out.parent.mkdir(parents=True, exist_ok=True) diff --git a/backend/app/services/preferences.py b/backend/app/services/preferences.py index 9f06a2e..f4e31db 100644 --- a/backend/app/services/preferences.py +++ b/backend/app/services/preferences.py @@ -265,34 +265,73 @@ def set_depth_finalize_time(hour: int, minute: int) -> dict: return {"hour": h, "minute": m} -def get_review_schedule() -> dict: - """定时复盘调度 {"enabled": False, "hour": 16, "minute": 30}。默认关闭。 +# 复盘推送可选渠道白名单 (微信等暂未实现, 不在白名单内, 前端仅作占位) +# 多选: 不推送 = 空数组, 而非 'none' +REVIEW_PUSH_CHANNELS = {"feishu"} - 复盘依赖盘后数据(日K/enriched, 盘后管道默认 15:30 跑完), - 故默认时间设为 16:30(数据就绪后), 强制下限 15:30。 + +def get_review_schedule() -> dict: + """定时复盘调度 {"enabled": False, "hour": 15, "minute": 10}。默认关闭。 + + A股 15:00 收盘, 默认时间设为 15:10(收盘后即时复盘), 强制下限 15:00。 """ - d = load().get("review_schedule", {"enabled": False, "hour": 16, "minute": 30}) + d = load().get("review_schedule", {"enabled": False, "hour": 15, "minute": 10}) return { "enabled": bool(d.get("enabled", False)), - "hour": d.get("hour", 16), - "minute": d.get("minute", 30), + "hour": d.get("hour", 15), + "minute": d.get("minute", 10), } def set_review_schedule(enabled: bool, hour: int, minute: int) -> dict: - """保存定时复盘调度。强制时间下限 15:30(盘后数据就绪)。 + """保存定时复盘调度。强制时间下限 15:00(A股收盘)。 enabled=False 时时间仍保存(下次开启可沿用), 但调度器不会注册 job。 """ h = max(0, min(23, hour)) m = max(0, min(59, minute)) - # 下限 15:30: 盘后管道(默认 15:30)完成后才有完整数据复盘 - if h * 60 + m < 15 * 60 + 30: - h, m = 15, 30 + # 下限 15:00: A股 15:00 收盘, 收盘后才有当日完整数据复盘 + if h * 60 + m < 15 * 60: + h, m = 15, 0 save({"review_schedule": {"enabled": bool(enabled), "hour": h, "minute": m}}) return {"enabled": bool(enabled), "hour": h, "minute": m} +def get_review_push_channels() -> list[str]: + """复盘推送渠道(多选) — 选定的外部工具列表, 复盘归档后逐个推送。 + + 与 review_schedule / 实时行情完全独立, 常驻可单独设置。 + 空列表 = 不推送; ['feishu'] = 推送到飞书(复用监控中心全局 feishu_webhook_url/secret)。 + + 向后兼容: + - 老多版本单选 review_push_channel=='feishu' → ['feishu'] + - 更老布尔 review_push_enabled==True → ['feishu'] + """ + d = load() + raw = d.get("review_push_channels") + if isinstance(raw, list): + return [c for c in raw if c in REVIEW_PUSH_CHANNELS] + # 兼容老单选字符串 + if d.get("review_push_channel") == "feishu": + return ["feishu"] + # 兼容更老布尔开关 + if d.get("review_push_enabled") is True: + return ["feishu"] + return [] + + +def set_review_push_channels(channels: list[str]) -> list[str]: + """保存复盘推送渠道(多选)。过滤白名单外的值、去重、保序。空列表 = 不推送。""" + seen: set[str] = set() + cleaned: list[str] = [] + for c in channels or []: + if c in REVIEW_PUSH_CHANNELS and c not in seen: + seen.add(c) + cleaned.append(c) + save({"review_push_channels": cleaned}) + return cleaned + + # ===== 实时监控 ===== diff --git a/backend/app/services/quote_service.py b/backend/app/services/quote_service.py index 9a8000d..8fca1be 100644 --- a/backend/app/services/quote_service.py +++ b/backend/app/services/quote_service.py @@ -60,6 +60,10 @@ class QuoteService: self._depth_update_event = threading.Event() # SSE 通知: depth 五档修正后 set (刷新连板梯队) self._pending_alerts: list[dict] = [] # 待推送的告警 self._max_pending_alerts: int = 1000 # 背压上限: 超出丢弃最旧 + # 复盘进度 SSE 通道: 定时复盘流式生成时, 把 meta/delta/done 事件推给开着页面的前端 + self._review_event = threading.Event() # SSE 通知: 有复盘进度事件时 set + self._pending_review: list[str] = [] # 待推送的复盘事件(JSON 字符串) + self._max_pending_review: int = 200 # 背压上限: 超出丢弃最旧 self._strategy_monitor = None # 延迟注入 self._app_state = None # 延迟注入 (FastAPI app.state) @@ -190,6 +194,34 @@ class QuoteService: self._pending_alerts = [] return alerts + # ================================================================ + # 复盘进度 SSE 通道 — 定时复盘流式生成时, 把事件实时推给前端 + # ================================================================ + def push_review_event(self, event_json: str) -> None: + """追加一条复盘进度事件(JSON 字符串), 并唤醒 SSE generator。 + + 事件格式与 recap_market_stream 的产出一致(meta/delta/error/done), + 前端 reviewStore 直接消费。背压: 超过上限丢弃最旧(复盘流几百条 delta, 200 够用)。 + """ + with self._lock: + self._pending_review.append(event_json) + if len(self._pending_review) > self._max_pending_review: + overflow = len(self._pending_review) - self._max_pending_review + self._pending_review = self._pending_review[overflow:] + self._review_event.set() + + def wait_for_review(self, timeout: float = 30.0) -> bool: + """阻塞等待复盘进度事件 (供 SSE 线程使用)。""" + self._review_event.clear() + return self._review_event.wait(timeout=timeout) + + def pop_review_events(self) -> list[str]: + """取走所有待推送的复盘事件 (线程安全)。""" + with self._lock: + events = self._pending_review + self._pending_review = [] + return events + # ================================================================ # 档位感知间隔限制 # ================================================================ @@ -681,9 +713,9 @@ class QuoteService: "logic": ev.get("logic") or "and", }) - # Free 自选实时只刷新少量标的, 不写全市场策略缓存。 - if self._enabled and self._app_state and self.realtime_mode() == "full_market": - self._refresh_strategy_cache(enriched_today, enriched_date) + # 策略页实时回显: 不写文件 (实时行情每轮更新 enriched, 写文件会被 read_cache + # 的 mtime 校验判过期, 反复读不到)。监控引擎本轮已算出的结果存在内存 + # (latest_strategy_results), 由 /api/screener/cached 端点直接叠加读取。 # 推入待推送队列 + 通知 SSE (含背压保护) if all_alerts: @@ -791,90 +823,6 @@ class QuoteService: except Exception as e: # noqa: BLE001 logger.debug("系统通知发送异常 (不影响告警主流程): %s", e) - def _refresh_strategy_cache(self, enriched_today: pl.DataFrame, enriched_date: date | None) -> None: - """利用已计算好的 enriched 数据,运行策略池并写入缓存。""" - import math - from dataclasses import asdict - from app.services import strategy_cache - from app.services.screener import PRESET_STRATEGIES, ScreenerService - from app.strategy import config as strategy_config - - try: - if enriched_date is None: - return - as_of = enriched_date - data_dir = self._repo.store.data_dir - svc = ScreenerService(self._repo) - engine = getattr(self._app_state, "strategy_engine", None) - - # 确定要运行的策略: 策略监控池中的策略 - monitor_ids = self._get_monitor_pool_ids() - if not monitor_ids: - return - - # 一次加载所有 override - all_overrides = strategy_config.list_overrides(data_dir) - - # 历史策略: 只在需要时加载 - shared_history = None - history_strats = [] - if engine: - id_set = set(monitor_ids) - history_strats = [ - (sid, s) for sid, s in engine._strategies.items() - if s.filter_history_fn and sid in id_set - ] - if history_strats: - max_lb = max(s.lookback_days for _, s in history_strats) - shared_history = svc._load_enriched_history(as_of, max(1, max_lb)) - - results: dict[str, dict] = {} - for sid in monitor_ids: - try: - overrides = all_overrides.get(sid, {}) - bf = overrides.get("basic_filter") if overrides else None - dl = overrides.get("display_limit") if overrides else None - if dl is None and overrides and "display_limit" in overrides: - dl = 0 - - if sid in PRESET_STRATEGIES: - r = svc.run_preset(sid, as_of=as_of, precomputed=enriched_today, basic_filter=bf, display_limit=dl) - elif engine: - r = engine.run( - sid, as_of, overrides=overrides or None, - precomputed=enriched_today, precomputed_history=shared_history, - ) - if dl is not None and dl > 0: - r.rows = r.rows[:dl] - r.total = min(r.total, dl) - else: - continue - - # sanitize NaN/Inf - rows = [] - for row_dict in asdict(r).get("rows", []): - for k, v in list(row_dict.items()): - if isinstance(v, float) and not math.isfinite(v): - row_dict[k] = None - rows.append(row_dict) - results[sid] = {"total": r.total, "as_of": str(as_of), "rows": rows} - except Exception: # noqa: BLE001 - continue - - if results: - strategy_cache.write_cache(data_dir, str(as_of), results) - - except Exception as e: # noqa: BLE001 - logger.warning("策略缓存刷新失败: %s", e) - - def _get_monitor_pool_ids(self) -> list[str]: - """获取策略监控池中的策略 ID 列表。""" - from app.services import preferences - ids = preferences.get_strategy_monitor_ids() - if not ids: - return [] - return [sid for sid in ids if sid] - @staticmethod def _get_strategy_monitor(): """获取 StrategyMonitorService — 不再使用, 改用 _app_state 注入。""" diff --git a/backend/app/services/strategy_cache.py b/backend/app/services/strategy_cache.py index 7dc700f..201134d 100644 --- a/backend/app/services/strategy_cache.py +++ b/backend/app/services/strategy_cache.py @@ -54,7 +54,14 @@ def _get_enriched_mtime(data_dir: Path, as_of: str) -> float | None: def read_cache(data_dir: Path) -> dict | None: - """读取策略缓存文件。返回 None 表示无缓存、读取失败或 enriched 数据已更新导致缓存过期。""" + """读取策略缓存文件。返回 None 表示无缓存或读取失败。 + + 说明: 原先有 enriched mtime 过期校验 (数据文件变化 → 判过期返回 None), + 但在有实时行情的系统里, enriched parquet 每轮被刷新 → mtime 必然变化 → + 缓存被永久判死, 策略页读不到数据。且判过期后不触发重算, 只能让用户手动重跑, + 保护价值有限。故移除: 盘后缓存总能读出, 实时新鲜度由 /api/screener/cached + 端点叠加监控引擎的内存实时结果 (latest_strategy_results) 来保证。 + """ path = _cache_path(data_dir) if not path.exists(): return None @@ -67,15 +74,6 @@ def read_cache(data_dir: Path) -> dict | None: logger.warning("读取策略缓存失败: %s", e) return None - # 校验 enriched mtime: 数据文件变化 → 缓存过期 - as_of = cached.get("as_of") - stored_mtime = cached.get("enriched_mtime") - if as_of and stored_mtime: - current_mtime = _get_enriched_mtime(data_dir, as_of) - if current_mtime is not None and current_mtime != stored_mtime: - logger.info("策略缓存过期: enriched 数据已更新 (as_of=%s)", as_of) - return None - return cached @@ -130,7 +128,8 @@ def write_cache( # 从 ever_rows 提取 symbol 列表 (用于快速计数) today_ever_matched = {sid: sorted(maps.keys()) for sid, maps in today_ever_rows.items()} - # 记录 enriched parquet 文件的 mtime,用于后续校验缓存是否过期 + # enriched_mtime: 盘后缓存写入时记录 (向后兼容旧字段)。read_cache 已不再用它 + # 做过期校验, 实时新鲜度改由 /cached 端点叠加监控引擎内存结果保证。 enriched_mtime = _get_enriched_mtime(data_dir, as_of) payload = { diff --git a/backend/app/services/webhook_adapter.py b/backend/app/services/webhook_adapter.py index 28c356f..afb76f5 100644 --- a/backend/app/services/webhook_adapter.py +++ b/backend/app/services/webhook_adapter.py @@ -25,6 +25,9 @@ logger = logging.getLogger(__name__) # 单次推送最长字符 (飞书单条文本消息上限 30KB, 这里保守截断避免刷屏) _MAX_LEN = 500 +# 卡片消息正文最长字符 (飞书 interactive 卡片上限 30KB, 保守留余量给标题/结构) +_CARD_MAX_LEN = 28000 + # 飞书自定义机器人 Webhook 前缀 (用于 URL 合法性校验) FEISHU_HOOK_PREFIX = "https://open.feishu.cn/open-apis/bot/v2/hook/" @@ -54,30 +57,20 @@ def _gen_sign(timestamp: str, secret: str) -> str: return base64.b64encode(hmac_code).decode("utf-8") -def send_feishu(webhook_url: str, title: str, body: str, secret: str = "") -> bool: - """推送一条文本消息到飞书群机器人。 +def _truncate_card(text: str) -> str: + """截断卡片正文 (留余量给标题与卡片结构)。""" + text = (text or "").strip() + return text[:_CARD_MAX_LEN] + ("…" if len(text) > _CARD_MAX_LEN else "") - Args: - webhook_url: 飞书自定义机器人 Webhook 地址 - title: 消息标题 (与正文拼接为一条文本) - body: 消息正文 - secret: 签名密钥 (机器人启用了「签名校验」时必填; 留空则不带签名) - Returns: - True=成功送达, False=失败或 URL 非法。 - 失败静默, 不抛异常 (Webhook 是辅助通道, 不能阻断告警主流程)。 +def _post_feishu(webhook_url: str, payload: dict, secret: str) -> bool: + """发送一次飞书 webhook 请求并判定成败 (供 text / card 共用)。 + + 成功响应: HTTP 200 且业务 code=0 (或非 JSON 的 200)。失败静默返回 False。 """ - if not is_valid_feishu_url(webhook_url): - return False - - text = _truncate(f"{title}\n{body}".strip()) - if not text: - return False - try: import httpx - payload: dict = {"msg_type": "text", "content": {"text": text}} # 启用签名校验时, 请求体须带 timestamp + sign (秒级时间戳) if secret: timestamp = str(int(time.time())) @@ -104,3 +97,74 @@ def send_feishu(webhook_url: str, title: str, body: str, secret: str = "") -> bo except Exception as e: # noqa: BLE001 logger.debug("飞书 Webhook 推送失败: %s", e) return False + + +def send_feishu(webhook_url: str, title: str, body: str, secret: str = "") -> bool: + """推送一条文本消息到飞书群机器人。 + + Args: + webhook_url: 飞书自定义机器人 Webhook 地址 + title: 消息标题 (与正文拼接为一条文本) + body: 消息正文 + secret: 签名密钥 (机器人启用了「签名校验」时必填; 留空则不带签名) + + Returns: + True=成功送达, False=失败或 URL 非法。 + 失败静默, 不抛异常 (Webhook 是辅助通道, 不能阻断告警主流程)。 + """ + if not is_valid_feishu_url(webhook_url): + return False + + text = _truncate(f"{title}\n{body}".strip()) + if not text: + return False + + payload: dict = {"msg_type": "text", "content": {"text": text}} + return _post_feishu(webhook_url, payload, secret) + + +def send_feishu_card(webhook_url: str, title: str, subtitle: str, body_md: str, secret: str = "") -> bool: + """推送一条 interactive 卡片消息到飞书群机器人 —— 用 lark_md 渲染完整 markdown 报告。 + + 飞书「自定义机器人」webhook 不支持文件附件, 但 interactive 卡片的 lark_md 元素 + 可渲染 markdown, 能承载完整复盘报告(通常 2-5KB, 远小于卡片 30KB 上限)。 + + Args: + webhook_url: 飞书自定义机器人 Webhook 地址 + title: 卡片标题 (显示在蓝色 header) + subtitle: 副标题 (加粗显示, 如日期/情绪标签; 留空则省略) + body_md: 卡片正文 markdown (报告全文) + secret: 签名密钥 (启用签名校验时必填) + + Returns: + True=成功送达, False=失败或 URL 非法。 + 失败静默, 不抛异常 (与 send_feishu 一致, 不阻断告警主流程)。 + """ + if not is_valid_feishu_url(webhook_url): + return False + + body = _truncate_card(body_md) + elements: list[dict] = [] + if subtitle.strip(): + elements.append({ + "tag": "div", + "text": {"tag": "lark_md", "content": f"**{subtitle.strip()}**"}, + }) + elements.append({"tag": "hr"}) + elements.append({ + "tag": "div", + "text": {"tag": "lark_md", "content": body}, + }) + + payload: dict = { + "msg_type": "interactive", + "card": { + "config": {"wide_screen_mode": True}, + "header": { + "title": {"tag": "plain_text", "content": title}, + "template": "blue", + }, + "elements": elements, + }, + } + return _post_feishu(webhook_url, payload, secret) diff --git a/backend/app/strategy/monitor.py b/backend/app/strategy/monitor.py index 8215815..23325df 100644 --- a/backend/app/strategy/monitor.py +++ b/backend/app/strategy/monitor.py @@ -24,6 +24,44 @@ from app.strategy import config as _strategy_config logger = logging.getLogger(__name__) +# 信号 / 字段中文名映射 — 与前端 lib/signals.ts 对齐, 用于告警 message / 推送文案。 +# signal_* 为内置原子信号, 其余为技术指标/行情字段。 +_SIGNAL_CN: dict[str, str] = { + # 内置信号 + "signal_ma_golden_5_20": "MA5上穿MA20", "signal_ma_dead_5_20": "MA5下穿MA20", + "signal_ma_golden_20_60": "MA20上穿MA60", "signal_macd_golden": "MACD金叉", + "signal_macd_dead": "MACD死叉", "signal_ma20_breakout": "突破MA20", + "signal_ma20_breakdown": "跌破MA20", "signal_n_day_high": "60日新高", + "signal_n_day_low": "60日新低", "signal_boll_breakout_upper": "突破布林上轨", + "signal_boll_breakdown_lower": "跌破布林下轨", "signal_volume_surge": "放量", + "signal_limit_up": "涨停", "signal_limit_down": "跌停", + "signal_limit_down_recovery": "跌停翘板", "signal_broken_limit_up": "炸板", + # 行情字段 + "close": "收盘价", "open": "开盘价", "high": "最高价", "low": "最低价", + "change_pct": "涨跌幅", "change_amount": "涨跌额", "amplitude": "振幅", + "turnover_rate": "换手率", "volume": "成交量", "amount": "成交额", + # 均线 + "ma5": "MA5", "ma10": "MA10", "ma20": "MA20", "ma30": "MA30", "ma60": "MA60", + "ema5": "EMA5", "ema10": "EMA10", "ema20": "EMA20", + # MACD / BOLL / KDJ / RSI + "macd_dif": "MACD-DIF", "macd_dea": "MACD-DEA", "macd_hist": "MACD柱", + "boll_upper": "布林上轨", "boll_lower": "布林下轨", + "kdj_k": "KDJ-K", "kdj_d": "KDJ-D", "kdj_j": "KDJ-J", + "rsi_6": "RSI6", "rsi_14": "RSI14", "rsi_24": "RSI24", + # 量能 / 动量 / 波动 + "vol_ratio_5d": "5日量比", "vol_ratio_20d": "20日量比", + "vol_ma5": "5日均量", "vol_ma10": "10日均量", + "high_60d": "60日最高", "low_60d": "60日最低", + "momentum_5d": "5日动量", "momentum_20d": "20日动量", "momentum_60d": "60日动量", + "atr_14": "ATR14", "annual_vol_20d": "20日年化波动", + "consecutive_limit_ups": "连板数", "consecutive_limit_downs": "跌停连板", +} + + +def _signal_cn_name(name: str) -> str: + """返回信号/字段的中文名, 找不到原样返回 (与前端 cnSignal 对齐)。""" + return _SIGNAL_CN.get(name, name) + @dataclass class StrategyAlert: @@ -276,6 +314,9 @@ class MonitorRuleEngine: self._strategy_pools: dict[str, set[str]] = {} # 数据目录 (用于加载策略 overrides) self._data_dir = None + # 本轮 evaluate() 产出的策略选股结果: strategy_id → {rows, total, as_of} + # 供策略页实时回显复用 (/api/screener/cached 端点直接读取此内存结果), 避免重跑 + self._latest_strategy_results: dict[str, dict] = {} def set_strategy_engine(self, engine) -> None: """注入 StrategyEngine, type=strategy 规则据此跑选股。""" @@ -325,6 +366,14 @@ class MonitorRuleEngine: def rule_count(self) -> int: return len(self._rules) + def latest_strategy_results(self) -> dict[str, dict]: + """返回本轮 evaluate() 产出的策略选股结果 (strategy_id → {rows, total, as_of})。 + + 供策略页实时回显复用: /api/screener/cached 端点直接读取此内存结果, + 避免对被监控的策略重跑第二遍。无 type=strategy 规则时返回空 dict。 + """ + return self._latest_strategy_results + # ── 评估 ─────────────────────────────────────────── def evaluate(self, df: pl.DataFrame) -> list[dict]: """行情更新后评估所有规则。 @@ -339,6 +388,8 @@ class MonitorRuleEngine: now = time.time() events: list[dict] = [] + # 每轮重置: 只保留本次 evaluate 产出的策略结果 + self._latest_strategy_results = {} for rule_id, rule in self._rules.items(): try: @@ -395,7 +446,11 @@ class MonitorRuleEngine: message = name # name 字段即批量消息 else: resolved_name = name if name else self._name_map.get(sym) - message = rule.get("message", "") or self._default_message(rule, ev_type=ev_type, sym=sym, name=resolved_name, pct=pct) + message = rule.get("message", "") or self._default_message( + rule, ev_type=ev_type, sym=sym, name=resolved_name, + pct=pct, price=price, + conditions=list(rule.get("conditions", [])) if rule.get("type") != "strategy" else None, + ) ev = { "ts": int(now * 1000), @@ -486,6 +541,22 @@ class MonitorRuleEngine: logger.warning("策略 %s 选股执行失败: %s", sid, e) return [] + # 记录本轮完整选股结果 (供策略页实时回显: /cached 端点直接读取, 不落盘)。 + # 与下面的 diff 事件无关 — 无论是否产生 new_entry/dropped, 结果都该可用于回显。 + try: + import math + self._latest_strategy_results[sid] = { + "total": result.total, + "as_of": str(_dt.date.today()), + "rows": [ + {k: (None if isinstance(v, float) and not math.isfinite(v) else v) + for k, v in row.items()} + for row in result.rows + ], + } + except Exception: # noqa: BLE001 + pass + current_pool: set[str] = {r["symbol"] for r in result.rows} prev_pool = self._strategy_pools.get(sid) @@ -582,8 +653,13 @@ class MonitorRuleEngine: return results def _default_message(self, rule: dict, ev_type: str = "", sym: str = "", - name: str = "", pct: Any = None) -> str: - """生成默认 message。策略类型按变更方向生成。""" + name: str = "", pct: Any = None, price: Any = None, + conditions: list[dict] | None = None) -> str: + """生成默认 message。 + + - strategy: 按变更方向生成 (进入/移出 + 涨跌幅) + - signal/price/market: 条件摘要 + 现价 + 涨跌幅 (避免笼统的「信号触发」) + """ rtype = rule.get("type", "signal") if rtype == "strategy": # 从 StrategyEngine 取策略名; 失败则退化为 rule_name 里截取的部分 @@ -613,5 +689,40 @@ class MonitorRuleEngine: return f"策略「{sname}」移出 {name}{pct_text}" return f"策略「{sname}」变更" - name_map = {"signal": "信号触发", "price": "价格触发", "market": "市场异动"} - return name_map.get(rtype, "监控触发") + # signal / price / market: 条件摘要 + 现价 + 涨跌幅 + # 条件摘要: 把 conditions (truth/比较) 拼成可读串, 如 "MA20金叉 且 量比>2" + cond_text = self._format_conditions_text(rule, conditions) + price_text = f"现价 {price}" if price is not None else "" + pct_text = "" + if pct is not None: + sign = "+" if pct >= 0 else "" + pct_text = f"{sign}{pct * 100:.1f}%" + tail = " · ".join(s for s in (price_text, pct_text) if s) + if cond_text and tail: + return f"{cond_text} · {tail}" + return cond_text or tail or "监控触发" + + @staticmethod + def _format_conditions_text(rule: dict, conditions: list[dict] | None) -> str: + """把 rule.conditions 拼成可读文本 (用于 message / 推送)。 + + op=truth: 直接用信号中文名 (如 "MA20金叉") + op=比较: 字段中文名 + 操作符 + 值 (如 "涨跌幅≥5") + logic: and → "且", or → "或" + """ + conds = conditions if conditions is not None else list(rule.get("conditions", [])) + if not conds: + return "" + logic_word = "且" if rule.get("logic", "and") == "and" else "或" + parts: list[str] = [] + for c in conds: + field = c.get("field", "") + op = c.get("op", "truth") + value = c.get("value") + label = _signal_cn_name(field) or field + if op == "truth": + parts.append(label) + else: + op_map = {"gte": "≥", "lte": "≤", "gt": ">", "lt": "<", "eq": "="} + parts.append(f"{label}{op_map.get(op, op)}{value}") + return f" {logic_word} ".join(parts) diff --git a/backend/pyproject.toml b/backend/pyproject.toml index 9766e2c..f9c95c6 100644 --- a/backend/pyproject.toml +++ b/backend/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "tickflow-stock-panel-backend" -version = "0.1.66" +version = "0.1.67" description = "A 股选股 + 监控 + 回测面板 — TickFlow 适配" readme = "../README.md" requires-python = ">=3.11" diff --git a/frontend/package.json b/frontend/package.json index 263f1a5..a7c66ff 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -1,7 +1,7 @@ { "name": "tickflow-stock-panel-frontend", "private": true, - "version": "0.1.66", + "version": "0.1.67", "type": "module", "scripts": { "dev": "vite", diff --git a/frontend/src/components/AlertToast.tsx b/frontend/src/components/AlertToast.tsx index cf2f898..bb9fc1a 100644 --- a/frontend/src/components/AlertToast.tsx +++ b/frontend/src/components/AlertToast.tsx @@ -161,10 +161,10 @@ export function AlertToastContainer() { {ev.price != null && {fmtPrice(ev.price)}} ) : ( -
+
+ {/* message 已含「条件摘要 · 现价 · 涨跌幅」(后端生成), 直接展示避免重复 */} {ev.message && {ev.message}} - {ev.price != null && {fmtPrice(ev.price)}}
)} diff --git a/frontend/src/components/Layout.tsx b/frontend/src/components/Layout.tsx index a557511..619330c 100644 --- a/frontend/src/components/Layout.tsx +++ b/frontend/src/components/Layout.tsx @@ -261,7 +261,8 @@ export function Layout() { const { data: settingsState } = useSettings() const { data: versionData } = useVersion() const { data: prefs } = usePreferences() - const { data: quoteStatus } = useQuoteStatus() + // poll=true: 全局唯一开启条件轮询 (非交易时段 60s 兜底, 交易时段靠 SSE) + const { data: quoteStatus } = useQuoteStatus({ poll: true }) const { data: analysisMenus } = useQuery({ queryKey: QK.analysisMenus, queryFn: api.analysisMenus, diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index 04229d5..6db5241 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -684,6 +684,7 @@ export interface Preferences { depth_polling_interval: number depth_finalize_time: { hour: number; minute: number } review_schedule: { enabled: boolean; hour: number; minute: number } + review_push_channels: string[] sse_refresh_pages: Record strategy_monitor_enabled: boolean strategy_monitor_ids: string[] @@ -867,6 +868,11 @@ export const api = { method: 'PUT', body: JSON.stringify({ enabled, hour, minute }), }), + updateReviewPush: (channels: string[]) => + request<{ review_push_channels: string[] }>('/api/settings/preferences/review-push', { + method: 'PUT', + body: JSON.stringify({ channels }), + }), updateDepthPollingInterval: (interval: number) => request<{ depth_polling_interval: number }>('/api/settings/preferences/depth-polling-interval', { method: 'PUT', diff --git a/frontend/src/lib/queryKeys.ts b/frontend/src/lib/queryKeys.ts index 8f821ec..61cbb69 100644 --- a/frontend/src/lib/queryKeys.ts +++ b/frontend/src/lib/queryKeys.ts @@ -84,5 +84,5 @@ export const SSE_INVALIDATE_PREFIXES = [ 'index-quotes', 'overview-market', 'limit-ladder', - 'screener-cached', + 'screener', ] as const diff --git a/frontend/src/lib/reviewStore.ts b/frontend/src/lib/reviewStore.ts index 78473ef..4d82fb7 100644 --- a/frontend/src/lib/reviewStore.ts +++ b/frontend/src/lib/reviewStore.ts @@ -35,6 +35,10 @@ const INITIAL: ReviewState = { phase: 'idle', content: '', error: '', meta: null let state: ReviewState = { ...INITIAL } let abortCtrl: AbortController | null = null +// 当前生成来源: 'manual'(手动点生成) | 'sse'(定时任务 SSE 推送) | null(空闲) +// 用于区分两条流, 避免互相丢弃事件或重复归档。 +let generatingSource: 'manual' | 'sse' | null = null + // ===== 订阅机制 ===== type Listener = () => void const listeners = new Set() @@ -76,6 +80,7 @@ export async function startReviewGeneration( // 已在生成中,不重复启动 if (isReviewGenerating()) return + generatingSource = 'manual' state = { phase: 'loading', content: '', error: '', meta: null, focus } notify() @@ -109,7 +114,7 @@ export async function startReviewGeneration( if (buf && !failed) { state = { ...state, phase: 'done' } notify() - // 自动归档 + // 自动归档(仅手动流: 定时流由后端归档, SSE done 不走这里) if (buf && !failed) { onDone?.(buf, doneMeta) } @@ -121,6 +126,7 @@ export async function startReviewGeneration( } } finally { abortCtrl = null + generatingSource = null } } @@ -160,3 +166,55 @@ export function resetReview(): void { state = { ...INITIAL } notify() } + +/** + * 喂入一条来自 SSE 的复盘事件(定时生成时后端推来的)。 + * + * 用途: 定时复盘在后端流式生成, 通过 /api/intraday/stream 的 review_progress 事件 + * 把 meta/delta/done 等实时推给前端, 前端调本函数把事件写进 store —— + * 这样开着复盘页的用户能看到「边生成边显示」, 和手动点生成完全一致。 + * + * 事件格式与 recap_market_stream 产出一致: + * {type:'meta'|'delta'|'error'|'done'|'retry', ...} + * + * 与手动生成的并发: + * - 若手动正在生成(isReviewGenerating), 忽略 SSE 事件(手动流优先, 避免冲突)。 + * - done 带 archived=true(定时场景后端已归档): 不重复调归档接口, 仅切到 done 态。 + * - retry: 后端 LLM 断流重试, 清空已累积内容重新开始。 + */ +export function feedReviewEvent(evt: any): void { + if (!evt || typeof evt !== 'object') return + const t = evt.type + + // 并发控制: 手动流进行中时, SSE 事件一律忽略(手动流优先, 避免两条流抢同一个 store) + // 但若当前是 SSE 流自己在跑(generatingSource==='sse'), 则正常处理后续事件 + if (generatingSource === 'manual') return + + if (t === 'meta') { + // 定时流的第一个事件: 标记来源为 sse, 进入 streaming 态, 重置 content + generatingSource = 'sse' + state = { phase: 'streaming', content: '', error: '', meta: evt, focus: '' } + notify() + } else if (t === 'delta' && evt.content) { + // 只有 sse 流进行中时才累积(防止 meta 丢失时的孤立 delta) + if (generatingSource !== 'sse') return + state = { ...state, content: state.content + evt.content, phase: 'streaming' } + notify() + } else if (t === 'retry') { + if (generatingSource !== 'sse') return + // 后端重试: 清空已累积内容, 等待新一轮 meta/delta + state = { ...state, content: '', phase: 'streaming' } + notify() + } else if (t === 'error') { + if (generatingSource !== 'sse') return + state = { ...state, error: evt.message ?? '复盘生成失败', phase: 'error' } + notify() + generatingSource = null + } else if (t === 'done') { + if (generatingSource !== 'sse') return + // 定时场景 done 带 archived=true: 后端已归档, 前端只切 done 态, 不调归档接口。 + state = { ...state, phase: 'done' } + notify() + generatingSource = null + } +} diff --git a/frontend/src/lib/useQuoteStream.ts b/frontend/src/lib/useQuoteStream.ts index ae15e15..483eced 100644 --- a/frontend/src/lib/useQuoteStream.ts +++ b/frontend/src/lib/useQuoteStream.ts @@ -1,9 +1,10 @@ import { useEffect, useRef, useCallback } from 'react' import { useQueryClient } from '@tanstack/react-query' -import { SSE_INVALIDATE_PREFIXES } from './queryKeys' +import { SSE_INVALIDATE_PREFIXES, QK } from './queryKeys' import { getQueryConfig } from './useQueryConfig' import { toast } from '@/components/Toast' import { pushAlertToasts } from '@/components/AlertToast' +import { feedReviewEvent } from './reviewStore' import type { StrategyAlertEvent } from './api' /** @@ -110,6 +111,21 @@ export function useQuoteStream( } }) + // 定时复盘流式进度: 后端到点生成时把 meta/delta/done 推来, 喂进 reviewStore + // 开着复盘页可看到「边生成边显示」, 切走再回来也能看到生成中/已生成 + es.addEventListener('review_progress', (e: MessageEvent) => { + try { + const evt = JSON.parse(e.data) + feedReviewEvent(evt) + // done(后端已归档) → 刷新历史列表, 让新报告出现并可查看 + if (evt.type === 'done') { + qc.invalidateQueries({ queryKey: QK.reviewReports }) + } + } catch { + // 忽略解析错误 + } + }) + es.onerror = () => { es.close() esRef.current = null diff --git a/frontend/src/lib/useSharedQueries.ts b/frontend/src/lib/useSharedQueries.ts index 0209976..a75845e 100644 --- a/frontend/src/lib/useSharedQueries.ts +++ b/frontend/src/lib/useSharedQueries.ts @@ -34,12 +34,22 @@ export function usePreferences() { }) } -/** 行情状态 — SSE quotes_updated 自动刷新 */ -export function useQuoteStatus(opts?: { enabled?: boolean }) { +/** 行情状态 — SSE quotes_updated 自动刷新。 + + * poll=true 时启用条件轮询兜底: 仅在非交易时段每 60s 轮询一次, + * 用于在交易时段边界 (11:30午休 / 12:55开盘 / 15:05收盘) 同步 is_trading_hours。 + * 交易时段不轮询 (SSE 已驱动刷新), 非交易时段无 SSE 推送, 需要兜底。 + * 只应在全局唯一挂载处 (Layout) 传 poll=true, 避免多页面重复轮询; + * 其他调用方共享同一 queryKey 缓存, 无需自行轮询。 + */ +export function useQuoteStatus(opts?: { enabled?: boolean; poll?: boolean }) { return useQuery({ queryKey: QK.quoteStatus, queryFn: api.quoteStatus, enabled: opts?.enabled ?? true, + refetchInterval: opts?.poll + ? (query) => (query.state.data?.is_trading_hours ? false : 60_000) + : false, }) } diff --git a/frontend/src/pages/Review.tsx b/frontend/src/pages/Review.tsx index 8f22b5f..dbe2f46 100644 --- a/frontend/src/pages/Review.tsx +++ b/frontend/src/pages/Review.tsx @@ -12,7 +12,7 @@ import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query' import { motion, AnimatePresence } from 'framer-motion' import { BookOpenCheck, RefreshCw, Sparkles, Trash2, History, ChevronRight, AlertTriangle, - Database, Wand2, Copy, Download, Clock, X, + Database, Wand2, Copy, Download, Clock, X, Check, } from 'lucide-react' import { api, type OverviewMarket, type AiReviewReport } from '@/lib/api' @@ -102,7 +102,11 @@ export function Review() { // ===== 定时复盘 ===== const [showSchedule, setShowSchedule] = useState(false) const prefs = usePreferences() - const reviewSched = prefs.data?.review_schedule ?? { enabled: false, hour: 16, minute: 30 } + const reviewSched = prefs.data?.review_schedule ?? { enabled: false, hour: 15, minute: 10 } + const feishuConfigured = !!(prefs.data?.feishu_webhook_url) + // 推送渠道是独立的顶层偏好(多选), 与定时 / 实时行情无关, 常驻可单独设置 + // []=不推送, ['feishu']=飞书(微信开发中, 仅占位) + const reviewPushChannels = prefs.data?.review_push_channels ?? [] // 弹窗内的本地草稿: 开关和时间都在本地改, 点「保存」才真正提交(避免开关一拨就关弹窗) const [draft, setDraft] = useState(reviewSched) const openSchedule = useCallback(() => { @@ -119,6 +123,21 @@ export function Review() { }, onError: () => { /* request() 已 toast */ }, }) + // 推送渠道(多选): 独立常驻, 即时生效(勾选渠道即开关, 改了立刻提交) + const pushMut = useMutation({ + mutationFn: (channels: string[]) => api.updateReviewPush(channels), + onSuccess: (_data, vars) => { + qc.invalidateQueries({ queryKey: QK.preferences }) + toast(vars.length === 0 ? '已关闭复盘推送' : '已更新复盘推送渠道', 'success') + }, + onError: () => { /* request() 已 toast */ }, + }) + const togglePushChannel = useCallback((ch: string) => { + const next = reviewPushChannels.includes(ch) + ? reviewPushChannels.filter(c => c !== ch) + : [...reviewPushChannels, ch] + pushMut.mutate(next) + }, [reviewPushChannels, pushMut]) // 自动滚动到报告底部(streaming 时) useEffect(() => { @@ -127,6 +146,15 @@ export function Review() { } }, [content, phase]) + // 当进入生成中(streaming)时, 清掉「查看历史」状态, 让主区域显示流内容。 + // 手动 generate 已自带 setViewing(null), 这里主要补定时 SSE 流的场景: + // 用户若正看着历史报告, 定时触发生成时也要切回主区域显示流式内容。 + useEffect(() => { + if (phase === 'streaming' && viewing) { + setViewing(null) + } + }, [phase, viewing]) + // 自动归档(生成完成后台静默保存)—— 通过回调注入 store,避免 store 直接依赖 qc/marketQuery const onGenerationDone = useCallback(async (fullContent: string, doneMeta: { as_of?: string; summary?: string; emotion_score?: number; emotion_label?: string } | null) => { const reportAsOf = doneMeta?.as_of ?? marketQuery.data?.as_of ?? asOf ?? new Date().toISOString().slice(0, 10) @@ -347,8 +375,8 @@ export function Review() {

- 开启后,每个交易日到点自动生成大盘复盘报告并归档,静默执行(不弹通知)。 - 下次打开本页即可在历史列表看到新报告。 + 开启后,每个交易日到点自动生成大盘复盘报告并归档,静默执行。 + 下次打开本页即可在历史列表看到新报告;也可选推送到飞书。

{/* 开关(只改本地草稿, 不提交) */} @@ -383,10 +411,56 @@ export function Review() { onChange={e => setDraft(d => ({ ...d, minute: Math.max(0, Math.min(59, Number(e.target.value))) }))} className="w-12 px-1.5 py-1 rounded-btn bg-base border border-border text-xs font-mono text-foreground text-center focus:outline-none focus:border-accent/50" /> - 不早于 15:30 · 工作日执行 + 不早于 15:00 · 工作日执行 )} + {/* 推送渠道(多选, 独立常驻, 与定时无关, 即时生效) */} +
+
+ 生成后推送完整报告 + {reviewPushChannels.length === 0 ? '未开启' : `${reviewPushChannels.length} 个渠道`} +
+
+ {/* 飞书(可用, 多选) */} + + {/* 微信(开发中, 占位不可选) */} +
+ + 微信 + 公众号/企业微信 + 开发中 +
+
+

+ 手动或定时生成的复盘都会以卡片消息推送完整报告。复用「设置 → 实时监控」的飞书 Webhook。 + {reviewPushChannels.includes('feishu') && !feishuConfigured && ( + setShowSchedule(false)}> + 前往配置 → + + )} +

+
+ {!draft.enabled && (

当前: 已关闭。开启后将按设定时间自动复盘。 diff --git a/frontend/src/pages/Screener.tsx b/frontend/src/pages/Screener.tsx index 41a04e9..2970546 100644 --- a/frontend/src/pages/Screener.tsx +++ b/frontend/src/pages/Screener.tsx @@ -424,7 +424,7 @@ export function Screener() { }, }) - // 开发用:重载策略文件 + // 重新运行策略:重载策略文件 + 重跑全部策略,刷新符合条件的个股 const reloadStrategies = useMutation({ mutationFn: api.strategyReload, onSuccess: () => { @@ -496,11 +496,11 @@ export function Screener() { subtitle="基于本地 enriched 表 · 毫秒级 SQL" right={

- {/* 开发用:重载策略 */} + {/* 重新运行策略:重载策略文件并重跑全部策略,更新命中个股 */}