fix(scheduler): _log_task 重复 INSERT 导致 zombie running 行

_log_task 在任务开始(running)和结束(success/failed)时分别
INSERT,导致 running 行永不更新,被 task monitor 标记为失败。

修复:新增 log_id 参数,传入时 UPDATE 已有行。所有 _run_* 方法
save running log_id 并在结束时传入。早期 return 处也标记失败。
This commit is contained in:
Yuzhiran Dev
2026-06-01 10:01:50 +08:00
parent 9b9ee4964b
commit 468c9066cb
+55 -32
View File
@@ -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))