diff --git a/platform/backend/app/core/scheduler.py b/platform/backend/app/core/scheduler.py index 99446b2..9aaf988 100644 --- a/platform/backend/app/core/scheduler.py +++ b/platform/backend/app/core/scheduler.py @@ -61,7 +61,9 @@ def _log_to_file(module_id: str, status: str, message: str = None, error_trace: def _log_task(module_id: str, status: str, message: str = None, error_trace: str = None, result_data: dict = None, started_at: datetime = None, finished_at: datetime = None, - triggered_by: str = "scheduler", next_run_time: datetime = None): + triggered_by: str = "scheduler", next_run_time: datetime = None, + log_id: int = None): + """写入/更新 TaskLog。如果 log_id 不为 None 则 UPDATE 已有行,否则 INSERT。返回 log_id。""" _log_to_file(module_id, status, message, error_trace) try: from ..database import SessionLocal @@ -71,6 +73,21 @@ def _log_task(module_id: str, status: str, message: str = None, duration = None if started_at and finished_at: duration = int((finished_at - started_at).total_seconds()) + if log_id: + log = db.query(TaskLog).filter(TaskLog.id == log_id).first() + if log: + log.status = status + log.message = message + if error_trace: + log.error_trace = error_trace + if result_data is not None: + log.result_data = result_data + if finished_at: + log.finished_at = finished_at + if duration is not None: + log.duration = duration + db.commit() + return log_id log = TaskLog( module_id=module_id, task_name=MODULES.get(module_id, {}).get("name", module_id), @@ -86,6 +103,8 @@ def _log_task(module_id: str, status: str, message: str = None, ) db.add(log) db.commit() + db.refresh(log) + return log.id finally: db.close() except Exception: @@ -182,7 +201,7 @@ class TaskScheduler: def _run_fetch_trends(self): """定时刷新热点趋势(百度/微博/知乎实时热搜 + LLM补充)""" started = datetime.now(timezone.utc) - _log_task("scheduled_fetch_trends", "running", started_at=started) + log_id = _log_task("scheduled_fetch_trends", "running", started_at=started) try: logger.info("[Scheduled] Fetching hot trends...") import subprocess @@ -194,17 +213,17 @@ class TaskScheduler: for line in result.stdout.strip().split("\n"): if line.strip(): logger.info("[Trends] %s", line.strip()) - _log_task("scheduled_fetch_trends", "success", + _log_task("scheduled_fetch_trends", "success", log_id=log_id, message="趋势刷新成功", result_data={"output_lines": len(result.stdout.splitlines())}, started_at=started, finished_at=datetime.now(timezone.utc)) else: - _log_task("scheduled_fetch_trends", "failed", + _log_task("scheduled_fetch_trends", "failed", log_id=log_id, message=f"返回码 {result.returncode}", error_trace=result.stderr[-500:], started_at=started, finished_at=datetime.now(timezone.utc)) except Exception as e: - _log_task("scheduled_fetch_trends", "failed", + _log_task("scheduled_fetch_trends", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -213,7 +232,7 @@ class TaskScheduler: def _run_refresh_search_cache(self): """定时刷新搜索缓存(通过 opencode webfetch)""" started = datetime.now(timezone.utc) - _log_task("scheduled_refresh_search_cache", "running", started_at=started) + log_id = _log_task("scheduled_refresh_search_cache", "running", started_at=started) try: logger.info("[Scheduled] Refreshing search cache via opencode...") import subprocess @@ -228,24 +247,24 @@ class TaskScheduler: if line.strip(): logger.warning("[SearchCache] %s", line.strip()) if result.returncode == 0: - _log_task("scheduled_refresh_search_cache", "success", + _log_task("scheduled_refresh_search_cache", "success", log_id=log_id, message="搜索缓存刷新成功", result_data={"output_lines": len(result.stdout.splitlines())}, started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Search cache refreshed") else: - _log_task("scheduled_refresh_search_cache", "failed", + _log_task("scheduled_refresh_search_cache", "failed", log_id=log_id, message="部分失败", error_trace=result.stderr[-500:], started_at=started, finished_at=datetime.now(timezone.utc)) logger.warning("[Scheduled] Search cache refresh may have partial failures") except subprocess.TimeoutExpired: - _log_task("scheduled_refresh_search_cache", "failed", + _log_task("scheduled_refresh_search_cache", "failed", log_id=log_id, message="超时", started_at=started, finished_at=datetime.now(timezone.utc)) logger.warning("[Scheduled] Search cache refresh timed out") except Exception as e: - _log_task("scheduled_refresh_search_cache", "failed", + _log_task("scheduled_refresh_search_cache", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -253,7 +272,7 @@ class TaskScheduler: def _run_generate(self): started = datetime.now(timezone.utc) - _log_task("scheduled_generate", "running", started_at=started) + log_id = _log_task("scheduled_generate", "running", started_at=started) try: logger.info("[Scheduled] Starting content generation...") result = run_creator_blocking() @@ -265,12 +284,12 @@ class TaskScheduler: review_result = run_optimizer_blocking([created_id]) if review_result.get("ok"): logger.info("[Scheduled] Review completed for %s", created_id) - _log_task("scheduled_generate", "success", + _log_task("scheduled_generate", "success", log_id=log_id, message=f"创作完成" + (f", 选题 {created_id}" if created_id else ""), result_data={"topic_id": created_id, "review_ok": review_result.get("ok") if review_result else None}, started_at=started, finished_at=datetime.now(timezone.utc)) except Exception as e: - _log_task("scheduled_generate", "failed", + _log_task("scheduled_generate", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -278,17 +297,17 @@ class TaskScheduler: def _run_optimize(self): started = datetime.now(timezone.utc) - _log_task("scheduled_optimize", "running", started_at=started) + log_id = _log_task("scheduled_optimize", "running", started_at=started) try: logger.info("[Scheduled] Starting compliance review...") result = run_optimizer_blocking() - _log_task("scheduled_optimize", "success", + _log_task("scheduled_optimize", "success", log_id=log_id, message="合规审查完成", result_data={"processed": result.get("processed", 0), "passed": result.get("passed", 0)}, started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Review completed: %s", result) except Exception as e: - _log_task("scheduled_optimize", "failed", + _log_task("scheduled_optimize", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -296,18 +315,18 @@ class TaskScheduler: def _run_collect(self): started = datetime.now(timezone.utc) - _log_task("scheduled_collect", "running", started_at=started) + log_id = _log_task("scheduled_collect", "running", started_at=started) try: logger.info("[Scheduled] Starting topic collection...") result = run_collector_blocking() topics_count = result.get("topics_count", 0) - _log_task("scheduled_collect", "success", + _log_task("scheduled_collect", "success", log_id=log_id, message=f"采集完成,找到 {topics_count} 个选题", result_data={"topics_count": topics_count, "output": str(result.get("output", ""))[:200]}, started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Collection completed: %s", result.get("output", "")[-200:]) except Exception as e: - _log_task("scheduled_collect", "failed", + _log_task("scheduled_collect", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -316,7 +335,7 @@ class TaskScheduler: def _run_optimize_sources(self, triggered_by="scheduler"): """AI自动优化采集类别与信息源:对比市场热点和当前配置,给出调整建议""" started = datetime.now(timezone.utc) - _log_task("scheduled_optimize_sources", "running", started_at=started, triggered_by=triggered_by) + log_id = _log_task("scheduled_optimize_sources", "running", started_at=started, triggered_by=triggered_by) try: logger.info("[Scheduled] Starting source optimization with AI...") from .nvidia_client import call_llm @@ -331,6 +350,8 @@ class TaskScheduler: except Exception: logger.warning("[Scheduled] DB not ready for source optimization") db.close() + _log_task("scheduled_optimize_sources", "failed", log_id=log_id, + message="数据库未就绪", started_at=started, finished_at=datetime.now(timezone.utc)) return cat_names = [c.name for c in cats] @@ -359,7 +380,7 @@ class TaskScheduler: db.add(SystemConfig(key="collector_ai_advice", value=json.dumps(result, ensure_ascii=False), description="AI每日采集优化建议")) db.commit() logger.info("[Scheduled] Source AI optimization completed: %s", result.get("summary", "")) - _log_task("scheduled_optimize_sources", "success", + _log_task("scheduled_optimize_sources", "success", log_id=log_id, message=result.get("summary", "优化完成"), result_data={"categories_assessed": len(result.get("category_assessment", [])), "sources_assessed": len(result.get("source_assessment", [])), @@ -368,7 +389,7 @@ class TaskScheduler: started_at=started, finished_at=datetime.now(timezone.utc), triggered_by=triggered_by) db.close() except Exception as e: - _log_task("scheduled_optimize_sources", "failed", + _log_task("scheduled_optimize_sources", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc), triggered_by=triggered_by) @@ -377,7 +398,7 @@ class TaskScheduler: def _run_metrics_sync(self, triggered_by="scheduler"): """定时从各平台公开API获取发布文章的效果数据(当前仅支持知乎)""" started = datetime.now(timezone.utc) - _log_task("scheduled_metrics_sync", "running", started_at=started, triggered_by=triggered_by) + log_id = _log_task("scheduled_metrics_sync", "running", started_at=started, triggered_by=triggered_by) try: logger.info("[Scheduled] Starting metrics sync (zhihu auto-fetch)...") from ..database import SessionLocal @@ -392,6 +413,8 @@ class TaskScheduler: except Exception: logger.warning("[Scheduled] DB not ready for metrics sync") db.close() + _log_task("scheduled_metrics_sync", "failed", log_id=log_id, + message="数据库未就绪", started_at=started, finished_at=datetime.now(timezone.utc)) return ua = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36" @@ -436,7 +459,7 @@ class TaskScheduler: if count: db.commit() logger.info("[Scheduled] Metrics sync completed: synced %d zhihu articles", count) - _log_task("scheduled_metrics_sync", "success", + _log_task("scheduled_metrics_sync", "success", log_id=log_id, message=f"同步完成,{count} 篇知乎文章", result_data={"articles_synced": count}, started_at=started, finished_at=datetime.now(timezone.utc), triggered_by=triggered_by) @@ -471,13 +494,13 @@ class TaskScheduler: except Exception as e_fb: logger.warning("[Scheduled] Metrics feedback generation failed: %s", e_fb) else: - _log_task("scheduled_metrics_sync", "success", + _log_task("scheduled_metrics_sync", "success", log_id=log_id, message="无已发布的知乎文章", started_at=started, finished_at=datetime.now(timezone.utc), triggered_by=triggered_by) logger.info("[Scheduled] Metrics sync: no zhihu articles to sync") db.close() except Exception as e: - _log_task("scheduled_metrics_sync", "failed", + _log_task("scheduled_metrics_sync", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc), triggered_by=triggered_by) @@ -486,7 +509,7 @@ class TaskScheduler: def _run_reset_search_usage(self): """每日凌晨重置搜索 API 提供商用量计数""" started = datetime.now(timezone.utc) - _log_task("scheduled_reset_search_usage", "running", started_at=started) + log_id = _log_task("scheduled_reset_search_usage", "running", started_at=started) try: from ..database import SessionLocal from ..models import SearchProvider @@ -494,7 +517,7 @@ class TaskScheduler: try: total = db.query(SearchProvider).update({SearchProvider.usage_today: 0, SearchProvider.last_used_at: None}) db.commit() - _log_task("scheduled_reset_search_usage", "success", + _log_task("scheduled_reset_search_usage", "success", log_id=log_id, message=f"已重置 {total} 个提供商用量", result_data={"reset_count": total}, started_at=started, finished_at=datetime.now(timezone.utc)) @@ -502,7 +525,7 @@ class TaskScheduler: finally: db.close() except Exception as e: - _log_task("scheduled_reset_search_usage", "failed", + _log_task("scheduled_reset_search_usage", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc)) @@ -511,7 +534,7 @@ class TaskScheduler: def _run_task_monitor(self): """每小时检查卡死/中断的任务,标记为失败""" started = datetime.now(timezone.utc) - _log_task("scheduled_task_monitor", "running", started_at=started) + log_id = _log_task("scheduled_task_monitor", "running", started_at=started) stuck_tasklog_timeout = 7200 # 超过2小时视为卡死 stuck_contenttask_timeout = 10800 # 超过3小时视为卡死 try: @@ -558,7 +581,7 @@ class TaskScheduler: db.commit() logger.info("[TaskMonitor] 已标记 %d 个卡死任务为失败", marked) - _log_task("scheduled_task_monitor", "success", + _log_task("scheduled_task_monitor", "success", log_id=log_id, message=f"检查完成,标记 {marked} 个卡死任务", result_data={"marked_failed": marked}, started_at=started, finished_at=datetime.now(timezone.utc)) @@ -566,7 +589,7 @@ class TaskScheduler: db.close() except Exception as e: import traceback - _log_task("scheduled_task_monitor", "failed", + _log_task("scheduled_task_monitor", "failed", log_id=log_id, message=str(e), error_trace=traceback.format_exc(), started_at=started, finished_at=datetime.now(timezone.utc))