diff --git a/apps/api/routers/settings.py b/apps/api/routers/settings.py index 9bfa29d..c721295 100644 --- a/apps/api/routers/settings.py +++ b/apps/api/routers/settings.py @@ -57,6 +57,19 @@ class DailyReportConfigUpdate(BaseModel): include_graph_insights: bool | None = None +class WorkerScheduleConfigUpdate(BaseModel): + """Worker 调度配置更新请求(cron 表达式 + 闲时处理器开关) + + 所有字段可选,仅传入字段会被更新(None 被过滤)。worker 轮询线程 + 检测到 updated_at 变化后热重载 APScheduler job,30s 内生效。 + """ + + topic_dispatch_cron: str | None = None + cs_feed_dispatch_cron: str | None = None + weekly_graph_cron: str | None = None + idle_processor_enabled: bool | None = None + + # ---------- 辅助函数 ---------- @@ -310,6 +323,51 @@ def update_daily_report_config(body: DailyReportConfigUpdate): return {"message": "每日报告配置已更新", "config": config} +# ---------- Worker 调度配置 ---------- + + +def _worker_schedule_to_out(cfg) -> dict: + """序列化 WorkerScheduleConfig 为前端可读的 dict(含 iso 时间)""" + return { + "id": cfg.id, + "topic_dispatch_cron": cfg.topic_dispatch_cron, + "cs_feed_dispatch_cron": cfg.cs_feed_dispatch_cron, + "weekly_graph_cron": cfg.weekly_graph_cron, + "idle_processor_enabled": cfg.idle_processor_enabled, + "last_applied_at": iso_dt(cfg.last_applied_at), + "updated_at": iso_dt(cfg.updated_at), + } + + +@router.get("/settings/worker-schedule") +def get_worker_schedule_config(): + """获取 worker 调度配置 + + 返回 cron 表达式 + 闲时处理器开关 + last_applied_at(worker 热重载时间)。 + 前端据此显示"已生效"状态(last_applied_at >= updated_at 即已同步)。 + """ + from packages.storage.repositories import WorkerScheduleConfigRepository + + with session_scope() as session: + cfg = WorkerScheduleConfigRepository(session).get_config() + return _worker_schedule_to_out(cfg) + + +@router.put("/settings/worker-schedule") +def update_worker_schedule_config(body: WorkerScheduleConfigUpdate): + """更新 worker 调度配置 + + 写入 DB 后,worker 轮询线程在 30s 内检测到 updated_at 变化并热重载 + APScheduler job。无需重启 worker 容器。 + """ + from packages.storage.repositories import WorkerScheduleConfigRepository + + update_data = {k: v for k, v in body.model_dump().items() if v is not None} + with session_scope() as session: + cfg = WorkerScheduleConfigRepository(session).update_config(**update_data) + return {"message": "worker 调度配置已更新", "config": _worker_schedule_to_out(cfg)} + + # ---------- SMTP 配置预设 ---------- diff --git a/apps/worker/main.py b/apps/worker/main.py index 15fb196..8f217f2 100644 --- a/apps/worker/main.py +++ b/apps/worker/main.py @@ -9,6 +9,7 @@ import contextlib import logging import signal +import threading import time from datetime import UTC, datetime from pathlib import Path @@ -31,7 +32,11 @@ from packages.config import get_settings from packages.logging_setup import setup_logging from packages.storage.db import session_scope -from packages.storage.repositories import TopicRepository +from packages.storage.repositories import ( + DailyReportConfigRepository, + TopicRepository, + WorkerScheduleConfigRepository, +) setup_logging() logger = logging.getLogger(__name__) @@ -42,6 +47,10 @@ # 心跳健康判定:最近一次心跳距现在超过此秒数视为不健康(status 页展示用) _HEARTBEAT_STALE_SECONDS = 1200 # 20 分钟 +# 配置轮询间隔(秒):worker 每 N 秒查 DB 配置 updated_at,检测变化即热重载。 +# 30s 是"网页端及时控制"与"DB 查询频率"的平衡点(用户改完最多 30s 生效)。 +_CONFIG_POLL_INTERVAL = 30 + def _write_heartbeat(error: str | None = None) -> None: """写入心跳文件供 status 页查询(记录最近一次错误,不再掩盖故障)。 @@ -106,6 +115,179 @@ def _retry_with_backoff(fn, *args, max_retries: int = 3, base_delay: float = 5.0 cs_orchestrator = CSFeedOrchestrator() +# 模块级 scheduler 单例 + 访问器(镜像 get_idle_processor 模式)。 +# BlockingScheduler.start() 会阻塞主线程,scheduler 对象本应活在 run_worker() +# 栈帧;提升到模块级后,reload_worker_schedule() / 轮询线程才能拿到它热重载 job。 +_scheduler: BlockingScheduler | None = None +# 内存中缓存上次应用的配置时间戳,用于轮询时检测 DB 是否有新变化 +_last_worker_cfg_ts: datetime | None = None +_last_daily_cfg_ts: datetime | None = None +# 闲时处理器当前运行状态镜像(避免无谓的 stop/start 循环) +_idle_running: bool = False + + +def get_scheduler() -> BlockingScheduler | None: + """获取运行中的 scheduler 单例(run_worker 启动前为 None)""" + return _scheduler + + +# 公共 job 配置:单实例 + 5 分钟 misfire 容忍 + 合并错过的触发 +# 模块级常量,_register_all_jobs / reload_worker_schedule 共用 +_JOB_KWARGS = { + "max_instances": 1, + "misfire_grace_time": 300, + "coalesce": True, + "replace_existing": True, +} + + +def _read_schedule_config() -> dict: + """从 DB 读所有调度 cron + idle 开关,DB 不可用时回退到安全默认值。 + + 返回 dict: + topic_dispatch_cron, cs_feed_dispatch_cron, weekly_graph_cron, + daily_brief_cron, idle_processor_enabled, + worker_cfg_ts, daily_cfg_ts # 用于变更检测 + """ + result = { + "topic_dispatch_cron": "0 * * * *", + "cs_feed_dispatch_cron": "5 * * * *", + "weekly_graph_cron": getattr(settings, "weekly_cron", "0 22 * * 0"), + "daily_brief_cron": "0 4 * * *", + "idle_processor_enabled": True, + "worker_cfg_ts": None, + "daily_cfg_ts": None, + } + try: + with session_scope() as session: + wcfg = WorkerScheduleConfigRepository(session).get_config() + result["topic_dispatch_cron"] = ( + wcfg.topic_dispatch_cron or result["topic_dispatch_cron"] + ) + result["cs_feed_dispatch_cron"] = ( + wcfg.cs_feed_dispatch_cron or result["cs_feed_dispatch_cron"] + ) + result["weekly_graph_cron"] = wcfg.weekly_graph_cron or result["weekly_graph_cron"] + result["idle_processor_enabled"] = wcfg.idle_processor_enabled + result["worker_cfg_ts"] = wcfg.updated_at + dcfg = DailyReportConfigRepository(session).get_config() + result["daily_brief_cron"] = dcfg.cron_expression or result["daily_brief_cron"] + result["daily_cfg_ts"] = dcfg.updated_at + except Exception as e: + logger.warning("读取调度配置失败,使用安全默认值:%s", e) + return result + + +def _register_all_jobs(scheduler: BlockingScheduler, cfg: dict) -> None: + """(重新)注册全部 4 个 APScheduler job。 + + 所有 job 都用 replace_existing=True,所以重复调用安全 —— 等价于 reschedule。 + 这是热重载的核心:每次轮询检测到配置变化后调用此函数即可。 + """ + scheduler.add_job( + topic_dispatch_job, + trigger=CronTrigger.from_crontab(cfg["topic_dispatch_cron"]), + id="topic_dispatch", + **_JOB_KWARGS, + ) + scheduler.add_job( + cs_feed_dispatch_job, + trigger=CronTrigger.from_crontab(cfg["cs_feed_dispatch_cron"]), + id="cs_feed_dispatch", + **_JOB_KWARGS, + ) + scheduler.add_job( + brief_job, + trigger=CronTrigger.from_crontab(cfg["daily_brief_cron"]), + id="daily_brief", + **_JOB_KWARGS, + ) + scheduler.add_job( + weekly_graph_job, + trigger=CronTrigger.from_crontab(cfg["weekly_graph_cron"]), + id="weekly_graph", + **_JOB_KWARGS, + ) + + +def _sync_idle_processor(enabled: bool) -> None: + """按配置开关同步 idle_processor 运行状态(避免无谓 stop/start)""" + global _idle_running + if enabled and not _idle_running: + logger.info("🤖 启动闲时自动处理器") + start_idle_processor() + _idle_running = True + elif not enabled and _idle_running: + logger.info("🛑 停止闲时自动处理器(配置已关闭)") + stop_idle_processor() + _idle_running = False + + +def reload_worker_schedule() -> bool: + """热重载:重读 DB 配置 → reschedule 所有 job → 同步 idle_processor → 写 last_applied_at。 + + 返回 True 表示成功应用新配置(含无变化时),False 表示失败(DB 不可用等)。 + 供轮询线程和未来 HTTP 端点调用。 + """ + global _scheduler, _last_worker_cfg_ts, _last_daily_cfg_ts + if _scheduler is None: + logger.warning("reload_worker_schedule: scheduler 尚未启动,跳过") + return False + try: + cfg = _read_schedule_config() + _register_all_jobs(_scheduler, cfg) + _sync_idle_processor(cfg["idle_processor_enabled"]) + # 更新内存时间戳缓存(下次轮询与此对比) + _last_worker_cfg_ts = cfg["worker_cfg_ts"] + _last_daily_cfg_ts = cfg["daily_cfg_ts"] + # 写回 last_applied_at 供前端显示"已生效" + try: + with session_scope() as session: + WorkerScheduleConfigRepository(session).update_last_applied_at(datetime.now(UTC)) + except Exception as e: + logger.warning("写回 last_applied_at 失败(不影响调度):%s", e) + logger.info( + "🔄 worker 调度已热重载:topic=%s, cs=%s, brief=%s, weekly=%s, idle=%s", + cfg["topic_dispatch_cron"], + cfg["cs_feed_dispatch_cron"], + cfg["daily_brief_cron"], + cfg["weekly_graph_cron"], + cfg["idle_processor_enabled"], + ) + return True + except Exception as e: + logger.exception("reload_worker_schedule 失败:%s", e) + return False + + +def _config_poll_loop(stop: Event) -> None: + """配置轮询线程主循环:每 _CONFIG_POLL_INTERVAL 秒查 DB,检测到变化即热重载。 + + 变更检测基于 WorkerScheduleConfig.updated_at + DailyReportConfig.updated_at + 与内存缓存对比 —— 任何配置写入都会 bump 对应表的 updated_at。 + """ + global _last_worker_cfg_ts, _last_daily_cfg_ts + logger.info("📡 配置轮询线程已启动(间隔 %ds)", _CONFIG_POLL_INTERVAL) + while not stop.wait(_CONFIG_POLL_INTERVAL): + try: + with session_scope() as session: + wcfg = WorkerScheduleConfigRepository(session).get_config() + dcfg = DailyReportConfigRepository(session).get_config() + w_ts = wcfg.updated_at + d_ts = dcfg.updated_at + changed = ( + _last_worker_cfg_ts is None + or _last_daily_cfg_ts is None + or w_ts != _last_worker_cfg_ts + or d_ts != _last_daily_cfg_ts + ) + if changed: + logger.info("📥 检测到调度配置变化(worker_ts=%s, daily_ts=%s)", w_ts, d_ts) + reload_worker_schedule() + except Exception as e: + logger.warning("配置轮询查询失败(下次重试):%s", e) + logger.info("📡 配置轮询线程已退出") + def _should_run(freq: str, time_utc: int, hour: int, weekday: int) -> bool: """判断当前 UTC 小时是否匹配主题的调度规则""" @@ -231,96 +413,52 @@ def cs_feed_dispatch_job(): def run_worker() -> None: """ - Worker 主函数 - UTC 时间智能调度 + Worker 主函数 - UTC 时间智能调度 + 网页端实时热重载 - 调度时间表(UTC): + 调度时间表(UTC,默认值,均可在网页端 Worker / 调度 tab 修改): ┌─────────────────────────────────────────────────────────┐ - │ 任务 │ 时间 (UTC) │ 北京时间 │ + │ 任务 │ 默认 cron (UTC) │ 北京时间 │ ├─────────────────────────────────────────────────────────┤ - │ 主题论文抓取 │ 02:00 每小时 │ 10:00 每小时 │ - │ 论文处理缓冲 │ 02:00-04:00 │ 10:00-12:00 │ - │ 每日简报生成 │ 04:00 │ 12:00 │ - │ 简报邮件发送 │ 04:30 │ 12:30 (午饭时间) │ - │ 每周图谱维护 │ 22:00 周日 │ 周一 06:00 │ - │ 闲时自动处理 │ 全天检测 │ 全天检测 │ + │ 主题论文抓取 │ 0 * * * * │ 每小时整点 │ + │ CS分类订阅 │ 5 * * * * │ 每小时 :05 │ + │ 每日简报生成 │ 0 4 * * * │ 12:00 │ + │ 每周图谱维护 │ 0 22 * * 0 │ 周一 06:00 │ + │ 闲时自动处理 │ 全天检测 │ 全天检测 │ └─────────────────────────────────────────────────────────┘ + + 热重载:网页端改配置 → 写 DB → 轮询线程 30s 内检测 updated_at 变化 + → reload_worker_schedule() 重排 job + 同步 idle_processor → 写 last_applied_at """ + global _scheduler, _idle_running + from apscheduler.executors.pool import ThreadPoolExecutor as APSThreadPoolExecutor + # High 3e:显式配置 max_instances / misfire_grace_time / coalesce,避免 # 重复触发与 misfire 丢失;用 apscheduler 的 ThreadPoolExecutor(max_workers=3) # 替代默认单线程池,允许 topic_dispatch / cs_feed / brief 适度并发 - from apscheduler.executors.pool import ThreadPoolExecutor as APSThreadPoolExecutor - scheduler = BlockingScheduler(timezone="UTC") scheduler.add_executor(APSThreadPoolExecutor(max_workers=3)) + _scheduler = scheduler # 暴露给模块级访问器,供 reload 使用 - # 公共 job 配置:单实例 + 5 分钟 misfire 容忍 + 合并错过的触发 - _job_kwargs = { - "max_instances": 1, - "misfire_grace_time": 300, - "coalesce": True, - "replace_existing": True, - } - - settings = get_settings() - - # 每整点检查主题调度(UTC 时间)—— 整点第 0 分钟 - scheduler.add_job( - topic_dispatch_job, - trigger=CronTrigger(minute=0), - id="topic_dispatch", - **_job_kwargs, - ) - logger.info("✅ 已添加:主题分发任务(每小时整点,UTC)") - - # CS 分类订阅调度 —— 错开 5 分钟,避免与 topic_dispatch 同分钟抢线程 - scheduler.add_job( - cs_feed_dispatch_job, - trigger=CronTrigger(minute=5), - id="cs_feed_dispatch", - **_job_kwargs, - ) - logger.info("✅ 已添加:CS分类订阅调度任务(每小时 :05,UTC)") - - # 每日简报(从数据库读取 cron 表达式) - from packages.storage.db import session_scope - from packages.storage.repositories import DailyReportConfigRepository - - try: - with session_scope() as session: - config = DailyReportConfigRepository(session).get_config() - daily_cron = config.cron_expression or "0 4 * * *" - except Exception as e: - logger.warning(f"从数据库读取 cron 失败:{e},使用默认值") - daily_cron = "0 4 * * *" - - daily_trigger = CronTrigger.from_crontab(daily_cron) - scheduler.add_job( - brief_job, - trigger=daily_trigger, - id="daily_brief", - **_job_kwargs, - ) + # 初始注册:从 DB 读配置(失败回退默认值),replace_existing=True 保证可重复 + cfg = _read_schedule_config() + _register_all_jobs(scheduler, cfg) + _last_worker_cfg_ts = cfg["worker_cfg_ts"] + _last_daily_cfg_ts = cfg["daily_cfg_ts"] logger.info( - "✅ 已添加:每日简报任务(cron: %s)", - daily_cron, + "✅ 初始调度已注册:topic=%s, cs=%s, brief=%s, weekly=%s", + cfg["topic_dispatch_cron"], + cfg["cs_feed_dispatch_cron"], + cfg["daily_brief_cron"], + cfg["weekly_graph_cron"], ) - # 每周图谱维护(UTC 周日 22 点 = 北京时间周一 6 点) - weekly_trigger = CronTrigger.from_crontab(getattr(settings, "weekly_cron", "0 22 * * 0")) - scheduler.add_job( - weekly_graph_job, - trigger=weekly_trigger, - id="weekly_graph", - **_job_kwargs, - ) - logger.info("✅ 已添加:每周图谱维护任务(UTC 周日 22:00)") - # 优雅关闭(High 3f:等待进行中任务跑完,避免已下载 PDF 未 set_pdf_path # 的中间态丢失;wait=True + 60s 超时兜底) def _graceful_stop(*_: object) -> None: logger.info("收到终止信号,正在关闭...") stop_event.set() stop_idle_processor() # 停止闲时处理器 + _idle_running = False scheduler.shutdown(wait=True) logger.info("Worker 已关闭") @@ -330,18 +468,25 @@ def _graceful_stop(*_: object) -> None: # 写入初始心跳 _write_heartbeat() - # 启动闲时处理器 - logger.info("🤖 启动闲时自动处理器...") - start_idle_processor() + # 启动闲时处理器(按配置开关) + _sync_idle_processor(cfg["idle_processor_enabled"]) + + # 启动配置轮询线程(daemon,随主进程退出) + poll_thread = threading.Thread( + target=_config_poll_loop, args=(stop_event,), daemon=True, name="cfg-poll" + ) + poll_thread.start() # 启动调度器 - logger.info("🚀 Worker 启动完成 - UTC 智能调度 + 闲时处理") + logger.info("🚀 Worker 启动完成 - UTC 智能调度 + 闲时处理 + 配置热重载") logger.info("=" * 60) - logger.info("调度时间表(UTC → 北京时间):") - logger.info(" • 主题抓取:每小时整点 → 每小时整点") - logger.info(" • 每日简报:04:00 → 12:00") - logger.info(" • 每周图谱:周日 22:00 → 周一 06:00") - logger.info(" • 闲时处理:全天自动检测 → 全天自动检测") + logger.info("调度时间表(UTC → 北京时间,可在网页端 Worker / 调度 tab 修改):") + logger.info(" • 主题抓取:%s", cfg["topic_dispatch_cron"]) + logger.info(" • CS订阅: %s", cfg["cs_feed_dispatch_cron"]) + logger.info(" • 每日简报:%s", cfg["daily_brief_cron"]) + logger.info(" • 每周图谱:%s", cfg["weekly_graph_cron"]) + logger.info(" • 闲时处理:%s", "开启" if cfg["idle_processor_enabled"] else "关闭") + logger.info(" • 配置热重载:每 %ds 轮询 DB", _CONFIG_POLL_INTERVAL) logger.info("=" * 60) scheduler.start() diff --git a/frontend/src/components/settings/WorkerSettingsTab.tsx b/frontend/src/components/settings/WorkerSettingsTab.tsx new file mode 100644 index 0000000..5c6971c --- /dev/null +++ b/frontend/src/components/settings/WorkerSettingsTab.tsx @@ -0,0 +1,240 @@ +import { useState, useCallback, useEffect } from "react"; +import { Clock, RefreshCw } from "lucide-react"; +import { useToast } from "@/contexts/ToastContext"; +import { Spinner } from "@/components/ui/Spinner"; +import { workerScheduleApi } from "@/services/api"; +import { getErrorMessage } from "@/lib/errorHandler"; +import { cn } from "@/lib/utils"; +import type { WorkerScheduleConfig } from "@/types"; + +/** + * Worker 调度配置 Tab + * + * 网页端修改 cron 表达式 / 闲时处理器开关后,worker 轮询线程在 30s 内 + * 检测到 updated_at 变化并热重载 APScheduler job,无需重启容器。 + * + * 同步状态:last_applied_at(worker 写回)vs updated_at(前端写入) + * - last_applied_at >= updated_at → 已生效 + * - 否则 → 等待 worker 同步中(最多 30s) + */ +type CronField = keyof Pick; + +const CRON_FIELDS: { field: CronField; label: string; desc: string; placeholder: string; hint: string }[] = [ + { + field: "topic_dispatch_cron", + label: "主题论文抓取", + desc: "按主题订阅关键词从 ArXiv 抓取论文", + placeholder: "0 * * * *", + hint: "默认 0 * * * *(每小时整点)", + }, + { + field: "cs_feed_dispatch_cron", + label: "CS 分类订阅", + desc: "按 CS 分类目录同步并抓取订阅论文", + placeholder: "5 * * * *", + hint: "默认 5 * * * *(每小时 :05,错开避免抢线程)", + }, + { + field: "weekly_graph_cron", + label: "每周图谱维护", + desc: "同步论文引用关系 + 图谱增量维护", + placeholder: "0 22 * * 0", + hint: "默认 0 22 * * 0(UTC 周日 22:00 = 北京周一 06:00)", + }, +]; + +export function WorkerSettingsTab() { + const { toast } = useToast(); + const [config, setConfig] = useState(null); + const [localConfig, setLocalConfig] = useState(null); + const [loading, setLoading] = useState(true); + const [submitting, setSubmitting] = useState(false); + // 手动刷新同步状态用(点击刷新按钮重新拉 last_applied_at) + const [refreshing, setRefreshing] = useState(false); + + const loadConfig = useCallback(async () => { + try { + const data = await workerScheduleApi.getConfig(); + setConfig(data); + setLocalConfig(data); + } catch { + toast("error", "加载 Worker 调度配置失败"); + } + }, [toast]); + + useEffect(() => { + loadConfig().finally(() => setLoading(false)); + }, [loadConfig]); + + const handleUpdateConfig = async (updates: Partial) => { + setSubmitting(true); + try { + const data = await workerScheduleApi.updateConfig(updates as Record); + if (data.config) { + setConfig(data.config); + setLocalConfig(data.config); + toast("success", "已保存,worker 将在 30s 内自动同步"); + } + } catch (err) { + toast("error", getErrorMessage(err)); + await loadConfig(); + } finally { + setSubmitting(false); + } + }; + + const handleInputChange = (field: CronField, value: string) => { + setLocalConfig((prev) => (prev ? { ...prev, [field]: value } : prev)); + }; + + const handleInputBlur = (field: CronField) => { + if (!localConfig || !config) return; + if (localConfig[field] !== config[field]) { + handleUpdateConfig({ [field]: localConfig[field] }); + } + }; + + const handleRefresh = async () => { + setRefreshing(true); + try { + const data = await workerScheduleApi.getConfig(); + setConfig(data); + // 不覆盖 localConfig(避免用户正在编辑的内容被覆盖) + setLocalConfig((prev) => (prev ? { ...data, ...prev } : data)); + } catch (err) { + toast("error", getErrorMessage(err)); + } finally { + setRefreshing(false); + } + }; + + // 同步状态判定:last_applied_at >= updated_at 即已生效 + const computeSyncStatus = (): { synced: boolean; label: string } => { + if (!config) return { synced: false, label: "未知" }; + if (!config.last_applied_at) { + return { synced: false, label: "worker 尚未应用过配置" }; + } + const applied = new Date(config.last_applied_at).getTime(); + const updated = new Date(config.updated_at).getTime(); + if (applied >= updated) { + const ageSec = Math.floor((Date.now() - applied) / 1000); + return { synced: true, label: `已生效(${ageSec}s 前 worker 已同步)` }; + } + return { synced: false, label: "等待 worker 同步中(最多 30s)" }; + }; + + if (loading) return
; + + return ( +
+
+

Worker / 调度

+

+ 管理 worker 定时任务调度。修改后 worker 在 30s 内自动热重载,无需重启。 +

+
+ + {/* 同步状态指示器 */} +
+
+
+ +
+
+

同步状态

+

{computeSyncStatus().label}

+
+
+ +
+ + {/* cron 配置卡片 */} +
+ {CRON_FIELDS.map((item) => ( +
+
+
+

{item.label}

+

{item.desc}

+
+
+
+ + handleInputChange(item.field, e.target.value)} + onBlur={() => handleInputBlur(item.field)} + disabled={submitting} + className="w-full rounded border border-border bg-surface px-2.5 py-1.5 font-mono text-xs text-ink placeholder:text-ink-placeholder disabled:opacity-60" + /> +

+ {item.hint} +
+ 格式:分 时 日 月 周(UTC,北京时间 = UTC + 8) +

+
+
+ ))} +
+ + {/* 闲时处理器开关 */} +
+
+
+
+ +
+
+

闲时自动处理器

+

+ worker 空闲时自动嵌入并粗读未处理论文 +

+
+
+ +
+

+ 关闭后 worker 不会在空闲时段自动处理论文,但已调度的 cron 任务仍正常执行 +

+
+ + {/* 说明区 */} +
+

+ 说明: + 所有 cron 基于 UTC 时间,北京时间 = UTC + 8。修改保存后 worker 轮询线程在 30s 内 + 检测到变化并重排 APScheduler job(replace_existing=True), + 无需重启容器。每日简报的 cron 在「邮箱与报告」tab 单独配置。 +

+
+
+ ); +} diff --git a/frontend/src/pages/Settings.tsx b/frontend/src/pages/Settings.tsx index 1dffeb6..b0f9a41 100644 --- a/frontend/src/pages/Settings.tsx +++ b/frontend/src/pages/Settings.tsx @@ -2,18 +2,20 @@ * Claude 风格的设置页面 - 左侧导航 + 右侧内容 */ import { useState } from "react"; -import { Cpu, Mail, GitBranch, Settings, ChevronRight } from "lucide-react"; +import { Cpu, Mail, GitBranch, Settings, ChevronRight, Clock } from "lucide-react"; import { cn } from "@/lib/utils"; import { LLMSettingsTab } from "@/components/settings/LLMSettingsTab"; import { EmailSettingsTab } from "@/components/settings/EmailSettingsTab"; import { PipelineSettingsTab } from "@/components/settings/PipelineSettingsTab"; import { OpsSettingsTab } from "@/components/settings/OpsSettingsTab"; +import { WorkerSettingsTab } from "@/components/settings/WorkerSettingsTab"; -type SettingsTab = "llm" | "email" | "pipeline" | "ops"; +type SettingsTab = "llm" | "email" | "pipeline" | "ops" | "worker"; const NAV_ITEMS: { key: SettingsTab; label: string; icon: typeof Cpu }[] = [ { key: "llm", label: "LLM 配置", icon: Cpu }, { key: "email", label: "邮箱与报告", icon: Mail }, + { key: "worker", label: "Worker / 调度", icon: Clock }, { key: "pipeline", label: "Pipeline", icon: GitBranch }, { key: "ops", label: "运维", icon: Settings }, ]; @@ -60,6 +62,7 @@ export default function SettingsPage() { {activeTab === "email" && } {activeTab === "pipeline" && } {activeTab === "ops" && } + {activeTab === "worker" && } diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index b01c042..b9e9ecf 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -54,6 +54,7 @@ import type { EmailConfig, EmailConfigForm, DailyReportConfig, + WorkerScheduleConfig, TaskStatus, ActiveTaskInfo, LoginResponse, @@ -75,6 +76,7 @@ export type { EmailConfig, EmailConfigForm, DailyReportConfig, + WorkerScheduleConfig, TaskStatus, ActiveTaskInfo, LoginResponse, @@ -683,6 +685,15 @@ export const dailyReportApi = { post<{ html: string }>(`/jobs/daily-report/generate-only?use_cache=${useCache}`), }; +/* ========== Worker 调度配置 ========== */ +// 网页端修改 cron / 闲时处理器开关,worker 轮询线程 30s 内热重载, +// 无需重启容器。getConfig 返回 last_applied_at 供 UI 显示"已生效"。 +export const workerScheduleApi = { + getConfig: () => get("/settings/worker-schedule"), + updateConfig: (data: Record) => + put<{ config: WorkerScheduleConfig }>("/settings/worker-schedule", data), +}; + /* ========== 后台任务 ========== */ export const tasksApi = { active: () => get<{ tasks: ActiveTaskInfo[] }>("/tasks/active"), diff --git a/frontend/src/types/index.ts b/frontend/src/types/index.ts index ca51151..20484e1 100644 --- a/frontend/src/types/index.ts +++ b/frontend/src/types/index.ts @@ -1015,6 +1015,19 @@ export interface DailyReportConfig { include_graph_insights: boolean; } +/* ========== Worker 调度配置 ========== */ +// Worker 的 cron 调度 + 闲时处理器开关。网页端修改后 worker 轮询线程 +// 在 30s 内检测 updated_at 变化并热重载,无需重启容器。 +// last_applied_at 由 worker 写回,前端据此显示"已生效"状态。 +export interface WorkerScheduleConfig { + topic_dispatch_cron: string; // 主题分发 cron(UTC) + cs_feed_dispatch_cron: string; // CS分类订阅 cron(UTC) + weekly_graph_cron: string; // 每周图谱维护 cron(UTC) + idle_processor_enabled: boolean; // 闲时自动处理器开关 + last_applied_at: string | null; // worker 最后应用配置的时间(ISO) + updated_at: string; // 配置最后更新时间(ISO) +} + /* ========== 后台任务 ========== */ export interface TaskStatus { task_id: string; diff --git a/infra/migrations/versions/a7b8c9d0e1f2_add_worker_schedule_config.py b/infra/migrations/versions/a7b8c9d0e1f2_add_worker_schedule_config.py new file mode 100644 index 0000000..d763fb2 --- /dev/null +++ b/infra/migrations/versions/a7b8c9d0e1f2_add_worker_schedule_config.py @@ -0,0 +1,67 @@ +"""add worker_schedule_configs table + +Revision ID: a7b8c9d0e1f2 +Revises: f6a7b8c9d0e1 +Create Date: 2026-07-20 12:00:00.000000 + +目的:把 worker 调度 cron 从硬编码/env 提取到 DB 单例表,支持网页端实时控制。 +单例表(单行),lazy 创建默认行。worker 启动读一次,运行时轮询 updated_at 热重载。 +last_applied_at 由 worker 写回,供前端显示"已生效"状态。 +""" +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision = "a7b8c9d0e1f2" +down_revision = "f6a7b8c9d0e1" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + bind = op.get_bind() + # PG 支持 CREATE TABLE IF NOT EXISTS;SQLite 用 try/except 兜底 + if bind.dialect.name == "postgresql": + op.execute( + """ + CREATE TABLE IF NOT EXISTS worker_schedule_configs ( + id VARCHAR(36) NOT NULL, + topic_dispatch_cron VARCHAR(64) NOT NULL, + cs_feed_dispatch_cron VARCHAR(64) NOT NULL, + weekly_graph_cron VARCHAR(64) NOT NULL, + idle_processor_enabled BOOLEAN NOT NULL, + last_applied_at DATETIME, + created_at DATETIME NOT NULL, + updated_at DATETIME NOT NULL, + PRIMARY KEY (id) + ) + """ + ) + else: + try: + op.create_table( + "worker_schedule_configs", + sa.Column("id", sa.String(length=36), nullable=False), + sa.Column("topic_dispatch_cron", sa.String(length=64), nullable=False), + sa.Column("cs_feed_dispatch_cron", sa.String(length=64), nullable=False), + sa.Column("weekly_graph_cron", sa.String(length=64), nullable=False), + sa.Column("idle_processor_enabled", sa.Boolean(), nullable=False), + sa.Column("last_applied_at", sa.DateTime(), nullable=True), + sa.Column("created_at", sa.DateTime(), nullable=False), + sa.Column("updated_at", sa.DateTime(), nullable=False), + sa.PrimaryKeyConstraint("id"), + ) + except Exception: + pass + + +def downgrade() -> None: + bind = op.get_bind() + if bind.dialect.name == "postgresql": + op.execute("DROP TABLE IF EXISTS worker_schedule_configs") + else: + try: + op.drop_table("worker_schedule_configs") + except Exception: + pass diff --git a/packages/ai/brief_service.py b/packages/ai/brief_service.py index 0c64bd7..464bb6b 100644 --- a/packages/ai/brief_service.py +++ b/packages/ai/brief_service.py @@ -597,7 +597,6 @@ def publish(self, recipient: str | None = None) -> dict: # 如果没有指定收件人,从数据库读取配置 if not recipient: - from packages.storage.db import session_scope from packages.storage.repositories import DailyReportConfigRepository with session_scope() as session: diff --git a/packages/storage/models.py b/packages/storage/models.py index 85a42ae..3bd1d80 100644 --- a/packages/storage/models.py +++ b/packages/storage/models.py @@ -449,6 +449,38 @@ class EmailConfig(Base): ) +class WorkerScheduleConfig(Base): + """Worker 调度配置 - cron 表达式 + 闲时处理器开关 + + 单例表(单行)。worker 启动时读一次注册 APScheduler job, + 运行时通过轮询线程(30s)检测 updated_at 变化后热重载。 + last_applied_at 由 worker 写回,供前端显示"已生效"状态。 + """ + + __tablename__ = "worker_schedule_configs" + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: str(uuid4())) + topic_dispatch_cron: Mapped[str] = mapped_column( + String(64), nullable=False, default="0 * * * *", doc="主题分发 cron(UTC)" + ) + cs_feed_dispatch_cron: Mapped[str] = mapped_column( + String(64), nullable=False, default="5 * * * *", doc="CS分类订阅 cron(UTC)" + ) + weekly_graph_cron: Mapped[str] = mapped_column( + String(64), nullable=False, default="0 22 * * 0", doc="每周图谱维护 cron(UTC)" + ) + idle_processor_enabled: Mapped[bool] = mapped_column( + Boolean, nullable=False, default=True, doc="闲时自动处理器开关" + ) + last_applied_at: Mapped[datetime] = mapped_column( + DateTime, nullable=True, doc="worker 最后一次应用配置的时间(worker 写、前端读)" + ) + created_at: Mapped[datetime] = mapped_column(DateTime, default=_utcnow, nullable=False) + updated_at: Mapped[datetime] = mapped_column( + DateTime, default=_utcnow, onupdate=_utcnow, nullable=False + ) + + class DailyReportConfig(Base): """每日报告配置 - 自动精读和邮件发送设置""" diff --git a/packages/storage/repositories/__init__.py b/packages/storage/repositories/__init__.py index 107c2f3..c0cc1b9 100644 --- a/packages/storage/repositories/__init__.py +++ b/packages/storage/repositories/__init__.py @@ -27,6 +27,7 @@ from packages.storage.repositories.prompt_trace import PromptTraceRepository from packages.storage.repositories.tag import TagRepository from packages.storage.repositories.topic import TopicRepository +from packages.storage.repositories.worker_schedule import WorkerScheduleConfigRepository __all__ = [ "BaseQuery", @@ -48,4 +49,5 @@ "AgentPendingActionRepository", "CSFeedRepository", "BatchJobRepository", + "WorkerScheduleConfigRepository", ] diff --git a/packages/storage/repositories/worker_schedule.py b/packages/storage/repositories/worker_schedule.py new file mode 100644 index 0000000..c2e6e46 --- /dev/null +++ b/packages/storage/repositories/worker_schedule.py @@ -0,0 +1,49 @@ +""" +Worker 调度配置数据仓储 +@author Color2333 +""" + +from __future__ import annotations + +from typing import TYPE_CHECKING + +from sqlalchemy import select + +if TYPE_CHECKING: + from datetime import datetime + + from sqlalchemy.orm import Session + +from packages.storage.models import WorkerScheduleConfig + + +class WorkerScheduleConfigRepository: + """Worker 调度配置仓储(单例)""" + + def __init__(self, session: Session): + self.session = session + + def get_config(self) -> WorkerScheduleConfig: + """获取 worker 调度配置(单例,不存在则创建默认行)""" + config = self.session.execute(select(WorkerScheduleConfig)).scalar_one_or_none() + + if not config: + config = WorkerScheduleConfig() + self.session.add(config) + self.session.flush() + + return config + + def update_config(self, **kwargs) -> WorkerScheduleConfig: + """更新 worker 调度配置""" + config = self.get_config() + for key, value in kwargs.items(): + if hasattr(config, key): + setattr(config, key, value) + return config + + def update_last_applied_at(self, ts: datetime) -> WorkerScheduleConfig: + """worker 热重载成功后写回最后应用时间(前端读以显示"已生效")""" + config = self.get_config() + config.last_applied_at = ts + return config diff --git a/scripts/auto_deploy.sh b/scripts/auto_deploy.sh index 39839d6..761aecd 100755 --- a/scripts/auto_deploy.sh +++ b/scripts/auto_deploy.sh @@ -97,6 +97,12 @@ if ! docker compose up -d 2>&1 | tee -a "$LOG_FILE" | tail -10; then exit 1 fi +# 6.1 清理悬空镜像(防止磁盘满导致 postgres 崩溃进 recovery mode) +# 历史教训:docker images 积累到 15GB 占满磁盘 → postgres 崩 → 每日简报任务连不上 DB 失败 +log "清理悬空 docker 镜像..." +docker image prune -f 2>&1 | tail -2 | while read -r line; do log "image prune: $line"; done || true +docker builder prune -f 2>&1 | tail -2 | while read -r line; do log "builder prune: $line"; done || true + # 7. 健康检查 log "等待健康检查 (40s)..." sleep 40 diff --git a/tests/test_repositories.py b/tests/test_repositories.py index 53d9181..92b9b91 100644 --- a/tests/test_repositories.py +++ b/tests/test_repositories.py @@ -490,3 +490,73 @@ def test_compensate_runs_independently_of_unread(self, db_session, monkeypatch): stuck = ip._get_stuck_skimmed_papers(limit=5) assert len(stuck) == 1, "应捞到 1 篇 stuck skimmed 论文" assert stuck[0][0] == paper.id + + +class TestWorkerScheduleConfigRepository: + """Worker 调度配置仓储 —— 单例 + 热重载写回时间戳""" + + def test_get_config_creates_default_singleton(self, db_session): + """首次 get_config 惰性创建默认单例行(cron 默认值 + idle 开)""" + from packages.storage.repositories import WorkerScheduleConfigRepository + + repo = WorkerScheduleConfigRepository(db_session) + cfg = repo.get_config() + db_session.flush() + + assert cfg.id is not None + assert cfg.topic_dispatch_cron == "0 * * * *" + assert cfg.cs_feed_dispatch_cron == "5 * * * *" + assert cfg.weekly_graph_cron == "0 22 * * 0" + assert cfg.idle_processor_enabled is True + assert cfg.last_applied_at is None # worker 尚未应用 + + def test_get_config_returns_same_singleton(self, db_session): + """二次 get_config 返回同一行(不创建新行)""" + from packages.storage.repositories import WorkerScheduleConfigRepository + + repo = WorkerScheduleConfigRepository(db_session) + first = repo.get_config() + db_session.flush() + first_id = first.id + second = repo.get_config() + assert second.id == first_id, "get_config 应返回同一单例" + + def test_update_config_partial_fields(self, db_session): + """update_config 仅更新传入字段,未传字段保留原值""" + from packages.storage.repositories import WorkerScheduleConfigRepository + + repo = WorkerScheduleConfigRepository(db_session) + repo.get_config() + db_session.flush() + + # 只改一个 cron + idle 开关 + updated = repo.update_config( + topic_dispatch_cron="*/10 * * * *", + idle_processor_enabled=False, + ) + db_session.flush() + + assert updated.topic_dispatch_cron == "*/10 * * * *" + assert updated.idle_processor_enabled is False + # 未传字段保留默认 + assert updated.cs_feed_dispatch_cron == "5 * * * *" + assert updated.weekly_graph_cron == "0 22 * * 0" + + def test_update_last_applied_at_sets_timestamp(self, db_session): + """update_last_applied_at 写入时间戳(worker 热重载后回写)""" + from packages.storage.repositories import WorkerScheduleConfigRepository + + repo = WorkerScheduleConfigRepository(db_session) + repo.get_config() + db_session.flush() + + ts = datetime.now(UTC) + repo.update_last_applied_at(ts) + db_session.flush() + + cfg = repo.get_config() + assert cfg.last_applied_at is not None + # SQLite DateTime 不存 tz,比较到秒级(剥离 tzinfo 与微秒) + stored = cfg.last_applied_at.replace(tzinfo=None, microsecond=0) + expected = ts.replace(tzinfo=None, microsecond=0) + assert stored == expected