mirror of
https://ghfast.top/https://github.com/aeroxw/tick-stock-panel.git
synced 2026-09-12 15:34:16 +08:00
feat: 定时复盘飞书推送 + SSE 实时进度
- 定时复盘支持推送到飞书(卡片消息, 完整报告), 渠道多选(飞书可选, 微信开发中) - 推送开关独立常驻, 与定时/实时行情解耦; 手动/定时生成都可推送 - 定时复盘改流式生成, 通过 SSE(review_progress)实时推给前端, 开着页面可见边生成边显示 - 修复定时任务协程未 await 的 bug(lambda 包裹 async → 改传函数对象 + args) - LLM 断流自动重试(最多2次), 后端归档 + 飞书推送, 异常兜底通知前端 - 复盘定时下限改为 15:00, 默认 15:10 - 版本号 0.1.66 → 0.1.67
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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}
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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}
|
||||
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
# ===== 实时监控 =====
|
||||
|
||||
|
||||
@@ -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 注入。"""
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "tickflow-stock-panel-frontend",
|
||||
"private": true,
|
||||
"version": "0.1.66",
|
||||
"version": "0.1.67",
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"dev": "vite",
|
||||
|
||||
@@ -161,10 +161,10 @@ export function AlertToastContainer() {
|
||||
{ev.price != null && <span className="text-[10px] font-mono text-muted shrink-0">{fmtPrice(ev.price)}</span>}
|
||||
</div>
|
||||
) : (
|
||||
<div className="mt-1 flex items-center gap-2 pl-0.5">
|
||||
<div className="mt-1 flex items-center gap-1.5 pl-0.5">
|
||||
<Bell className={cn('h-3 w-3 shrink-0', sev.replace('bg-', 'text-'))} />
|
||||
{/* message 已含「条件摘要 · 现价 · 涨跌幅」(后端生成), 直接展示避免重复 */}
|
||||
{ev.message && <span className="text-[11px] text-foreground/70 truncate flex-1">{ev.message}</span>}
|
||||
{ev.price != null && <span className="text-[10px] font-mono text-muted shrink-0">{fmtPrice(ev.price)}</span>}
|
||||
</div>
|
||||
)}
|
||||
</motion.div>
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<string, boolean>
|
||||
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',
|
||||
|
||||
@@ -84,5 +84,5 @@ export const SSE_INVALIDATE_PREFIXES = [
|
||||
'index-quotes',
|
||||
'overview-market',
|
||||
'limit-ladder',
|
||||
'screener-cached',
|
||||
'screener',
|
||||
] as const
|
||||
|
||||
@@ -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<Listener>()
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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() {
|
||||
</div>
|
||||
|
||||
<p className="mb-4 text-[11px] leading-relaxed text-muted">
|
||||
开启后,每个交易日到点自动生成大盘复盘报告并归档,静默执行(不弹通知)。
|
||||
下次打开本页即可在历史列表看到新报告。
|
||||
开启后,每个交易日到点自动生成大盘复盘报告并归档,静默执行。
|
||||
下次打开本页即可在历史列表看到新报告;也可选推送到飞书。
|
||||
</p>
|
||||
|
||||
{/* 开关(只改本地草稿, 不提交) */}
|
||||
@@ -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"
|
||||
/>
|
||||
<span className="text-[10px] text-muted/70">不早于 15:30 · 工作日执行</span>
|
||||
<span className="text-[10px] text-muted/70">不早于 15:00 · 工作日执行</span>
|
||||
</div>
|
||||
)}
|
||||
|
||||
{/* 推送渠道(多选, 独立常驻, 与定时无关, 即时生效) */}
|
||||
<div className="mt-3 rounded-btn bg-elevated/40 px-3 py-2.5">
|
||||
<div className="flex items-center justify-between">
|
||||
<span className="text-xs text-foreground">生成后推送完整报告</span>
|
||||
<span className="text-[10px] text-muted/70">{reviewPushChannels.length === 0 ? '未开启' : `${reviewPushChannels.length} 个渠道`}</span>
|
||||
</div>
|
||||
<div className="mt-2 space-y-1.5">
|
||||
{/* 飞书(可用, 多选) */}
|
||||
<button
|
||||
type="button"
|
||||
disabled={pushMut.isPending}
|
||||
onClick={() => togglePushChannel('feishu')}
|
||||
className={cn(
|
||||
'flex w-full items-center gap-2 rounded-btn border px-2.5 py-1.5 text-left transition-colors disabled:opacity-50',
|
||||
reviewPushChannels.includes('feishu')
|
||||
? 'border-accent/40 bg-accent/10'
|
||||
: 'border-border/60 bg-base/40 hover:bg-base/60',
|
||||
)}
|
||||
>
|
||||
<span className={cn('flex h-3 w-3 shrink-0 items-center justify-center rounded border', reviewPushChannels.includes('feishu') ? 'border-accent bg-accent text-white' : 'border-border')}>
|
||||
{reviewPushChannels.includes('feishu') && <Check className="h-2.5 w-2.5" />}
|
||||
</span>
|
||||
<span className="text-[11px] text-foreground">飞书</span>
|
||||
<span className="text-[9px] text-muted">群机器人</span>
|
||||
<span className={cn('ml-auto text-[9px]', feishuConfigured ? 'text-emerald-500' : 'text-warning')}>
|
||||
{feishuConfigured ? '已配置' : '未配置'}
|
||||
</span>
|
||||
</button>
|
||||
{/* 微信(开发中, 占位不可选) */}
|
||||
<div className="flex items-center gap-2 rounded-btn border border-border/40 bg-base/20 px-2.5 py-1.5 opacity-60">
|
||||
<span className="flex h-3 w-3 shrink-0 items-center justify-center rounded border border-border" />
|
||||
<span className="text-[11px] text-secondary">微信</span>
|
||||
<span className="text-[9px] text-muted">公众号/企业微信</span>
|
||||
<span className="ml-auto rounded bg-muted/10 px-1 py-px text-[9px] text-muted">开发中</span>
|
||||
</div>
|
||||
</div>
|
||||
<p className="mt-1.5 text-[10px] leading-relaxed text-muted/70">
|
||||
手动或定时生成的复盘都会以卡片消息推送完整报告。复用「设置 → 实时监控」的飞书 Webhook。
|
||||
{reviewPushChannels.includes('feishu') && !feishuConfigured && (
|
||||
<Link to="/settings?tab=monitoring" className="ml-1 text-accent hover:underline" onClick={() => setShowSchedule(false)}>
|
||||
前往配置 →
|
||||
</Link>
|
||||
)}
|
||||
</p>
|
||||
</div>
|
||||
|
||||
{!draft.enabled && (
|
||||
<p className="mt-3 text-[10px] text-muted/70">
|
||||
当前: 已关闭。开启后将按设定时间自动复盘。
|
||||
|
||||
@@ -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={
|
||||
<div className="flex items-center gap-2">
|
||||
{/* 开发用:重载策略 */}
|
||||
{/* 重新运行策略:重载策略文件并重跑全部策略,更新命中个股 */}
|
||||
<button
|
||||
onClick={() => reloadStrategies.mutate()}
|
||||
disabled={reloadStrategies.isPending}
|
||||
title="重载策略文件(开发用)"
|
||||
title="重新加载策略并运行全部策略,刷新当前符合条件的个股"
|
||||
className="inline-flex items-center gap-1.5 h-7 px-2.5 rounded-btn
|
||||
border border-border bg-surface text-xs font-medium text-muted
|
||||
hover:text-accent hover:border-accent/50 transition-colors cursor-pointer
|
||||
|
||||
@@ -768,11 +768,34 @@ export function Watchlist() {
|
||||
[visibleColumns]
|
||||
)
|
||||
|
||||
// 被过滤掉的个股数 (筛选/板块过滤导致的隐藏)
|
||||
const hiddenCount = Math.max(0, allSymbols.length - sortedRows.length)
|
||||
|
||||
return (
|
||||
<div className="flex flex-col h-full">
|
||||
<PageHeader
|
||||
title="自选股"
|
||||
subtitle={`${sortedRows.length}/${allSymbols.length} 只`}
|
||||
titleExtra={
|
||||
<span className="inline-flex items-center gap-1.5">
|
||||
{/* 计数胶囊: 显示数/总数, mono 字体突出数字 */}
|
||||
<span className="inline-flex items-baseline gap-0.5 px-2 py-0.5 rounded-md bg-elevated/70 text-[11px]">
|
||||
<span className="font-mono font-semibold text-secondary tabular-nums">{sortedRows.length}</span>
|
||||
<span className="text-muted/50">/</span>
|
||||
<span className="font-mono text-muted tabular-nums">{allSymbols.length}</span>
|
||||
<span className="text-muted/60 ml-0.5">只</span>
|
||||
</span>
|
||||
{/* 过滤提示: 仅在有隐藏时出现, 柔和橙色融入整体 */}
|
||||
{hiddenCount > 0 && (
|
||||
<span
|
||||
className="inline-flex items-center gap-1 px-1.5 py-0.5 rounded-md text-[10px] font-medium bg-warning/12 text-warning/90 border border-warning/25 whitespace-nowrap"
|
||||
title={`当前有 ${hiddenCount} 只被筛选条件隐藏,清除筛选可查看全部`}
|
||||
>
|
||||
<Filter className="h-2.5 w-2.5" />
|
||||
已过滤 {hiddenCount}
|
||||
</span>
|
||||
)}
|
||||
</span>
|
||||
}
|
||||
right={
|
||||
<div className="flex items-center gap-2">
|
||||
{/* 筛选 / 搜索 */}
|
||||
|
||||
Reference in New Issue
Block a user