From 3a195d69b9404e003e861fee7cd56092f987fbe4 Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Thu, 20 Aug 2026 01:42:59 +0800 Subject: [PATCH] =?UTF-8?q?fix(dispatcher):=20=E5=83=B5=E5=B0=B8PROCESSING?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E5=9B=9E=E6=94=B6=20+=20fam-core=E6=96=87?= =?UTF-8?q?=E4=BB=B6=E6=97=A5=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - db_layer 新增 reclaim_stale_processing: PROCESSING 超过 push_timeout+120s 重置 PENDING - dispatcher 轮询前先回收僵尸任务(进程重启/Edge重启导致 in-flight 请求丢失的场景) - logger 增加 fam-core/logs/fam-core.log 文件输出(daemon 模式 stdout 不可见) --- fam-core/src/fam_core/db_layer.py | 28 +++++++++++++++++++ .../src/fam_core/dispatcher/dispatcher.py | 10 +++++++ fam-core/src/fam_core/logger.py | 17 +++++++++-- 3 files changed, 53 insertions(+), 2 deletions(-) diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py index c565ebc..7a633d1 100644 --- a/fam-core/src/fam_core/db_layer.py +++ b/fam-core/src/fam_core/db_layer.py @@ -118,6 +118,34 @@ def increment_retry(task_id: int, next_retry_at: datetime): conn.close() +def reclaim_stale_processing(timeout_seconds: int) -> List[int]: + """回收僵尸 PROCESSING 任务:updated_at 早于 timeout_seconds 前的任务重置为 PENDING + + 场景:Dispatcher 推送过程中进程重启/Edge 重启导致 in-flight 请求丢失, + 任务停留在 PROCESSING 无人处理。重置后由常规重试机制接管。 + 返回被回收的 task_id 列表。 + """ + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + "SELECT task_id FROM process_tasks " + "WHERE status='PROCESSING' AND updated_at < NOW() - INTERVAL %s SECOND", + (timeout_seconds,) + ) + task_ids = [row[0] for row in cursor.fetchall()] + if task_ids: + placeholders = ','.join(['%s'] * len(task_ids)) + cursor.execute( + f"UPDATE process_tasks SET status='PENDING' WHERE task_id IN ({placeholders})", + task_ids + ) + conn.commit() + return task_ids + finally: + conn.close() + + def get_task(task_id: int) -> Optional[Dict]: """获取单个任务""" conn = get_conn() diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py index cf8a3a5..e1c1af1 100644 --- a/fam-core/src/fam_core/dispatcher/dispatcher.py +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -149,6 +149,16 @@ class Dispatcher: def _poll_once(self): """执行一次轮询""" + # 先回收僵尸任务:PROCESSING 超过 push_timeout+缓冲 说明推送进程已丢失 + # (如 fam-core 重启、Edge 重启掐断 in-flight 连接),重置回 PENDING 走重试 + stale_timeout = self.push_timeout + 120 + try: + stale_ids = db_layer.reclaim_stale_processing(stale_timeout) + for tid in stale_ids: + logger.warning(f"[task_id={tid}] PROCESSING 超时 {stale_timeout}s,回收为 PENDING 重试") + except Exception as e: + logger.error(f"僵尸任务回收失败: {e}", exc_info=True) + tasks = db_layer.get_pending_tasks(limit=10) for task in tasks: if self._should_retry(task): diff --git a/fam-core/src/fam_core/logger.py b/fam-core/src/fam_core/logger.py index f9c6f3d..9db3a88 100644 --- a/fam-core/src/fam_core/logger.py +++ b/fam-core/src/fam_core/logger.py @@ -2,23 +2,36 @@ 日志工具 - 统一格式,带 task_id 作为 trace_id """ import logging +import os import sys from datetime import datetime +# 日志目录:fam-core/logs/(相对 src 的上一级),失败则退化为仅 stdout +_LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), 'logs') + def setup_logger(name='fam-core', level=logging.INFO): - """配置并返回 logger""" + """配置并返回 logger(stdout + 文件双写)""" logger = logging.getLogger(name) if logger.handlers: return logger logger.setLevel(level) - handler = logging.StreamHandler(sys.stdout) formatter = logging.Formatter( '%(asctime)s [%(name)s] %(levelname)s %(message)s', datefmt='%Y-%m-%d %H:%M:%S' ) + handler = logging.StreamHandler(sys.stdout) handler.setFormatter(formatter) logger.addHandler(handler) + # 文件输出(daemon 模式下 stdout 不可见,文件是唯一可追溯日志) + try: + os.makedirs(_LOG_DIR, exist_ok=True) + file_handler = logging.FileHandler( + os.path.join(_LOG_DIR, 'fam-core.log'), encoding='utf-8') + file_handler.setFormatter(formatter) + logger.addHandler(file_handler) + except OSError: + pass # 目录不可写时退化为仅 stdout return logger