[2.1-2.5] FAM-Core 六个子模块 - Scheduler/Dispatcher/Event-Receiver/Chat-Handler/Member-Manager/Video-Server + DB层 + 配置
This commit is contained in:
0
fam-core/src/fam_core/dispatcher/__init__.py
Normal file
0
fam-core/src/fam_core/dispatcher/__init__.py
Normal file
127
fam-core/src/fam_core/dispatcher/dispatcher.py
Normal file
127
fam-core/src/fam_core/dispatcher/dispatcher.py
Normal file
@@ -0,0 +1,127 @@
|
||||
"""
|
||||
Dispatcher - 30s 轮询 PENDING 任务,下发至 Edge
|
||||
|
||||
退避重试: min(60 * (retry_count + 1) * 2, 600) 秒
|
||||
payload 注入 family_members 表的已命名+未命名成员清单
|
||||
"""
|
||||
import time
|
||||
import threading
|
||||
import requests
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from ..logger import setup_logger, log_task
|
||||
from ..config_loader import load_config
|
||||
from .. import db_layer
|
||||
|
||||
logger = setup_logger('fam-core.dispatcher')
|
||||
|
||||
|
||||
class Dispatcher:
|
||||
"""任务下发器,30s 轮询"""
|
||||
|
||||
def __init__(self):
|
||||
cfg = load_config()
|
||||
self.poll_interval = cfg.get('dispatcher', {}).get('poll_interval', 30)
|
||||
self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/analyze')
|
||||
self.max_retries = cfg.get('dispatcher', {}).get('max_retries', 3)
|
||||
self._running = False
|
||||
self._thread = None
|
||||
|
||||
def _calculate_backoff(self, retry_count):
|
||||
"""退避策略: min(60 * (retry_count + 1) * 2, 600)"""
|
||||
return min(60 * (retry_count + 1) * 2, 600)
|
||||
|
||||
def _should_retry(self, task):
|
||||
"""检查任务是否可以重试"""
|
||||
if task['retry_count'] >= task['max_retries']:
|
||||
return False
|
||||
if task['next_retry_at']:
|
||||
now = datetime.now()
|
||||
if now < task['next_retry_at']:
|
||||
return False
|
||||
return True
|
||||
|
||||
def _build_payload(self, task):
|
||||
"""构建下发 payload,注入已知成员清单"""
|
||||
known_members = db_layer.get_known_members_context()
|
||||
return {
|
||||
"task_id": task['task_id'],
|
||||
"video_url": task['video_url'],
|
||||
"webhook_url": load_config().get('dispatcher', {}).get(
|
||||
'webhook_url', 'http://localhost:8000/api/core/callback/event'
|
||||
),
|
||||
"known_members_context": known_members
|
||||
}
|
||||
|
||||
def _dispatch_one(self, task):
|
||||
"""下发单个任务"""
|
||||
task_id = task['task_id']
|
||||
payload = self._build_payload(task)
|
||||
|
||||
try:
|
||||
log_task(logger, task_id, 'dispatcher', f'下发至 Edge: {self.edge_url}')
|
||||
resp = requests.post(self.edge_url, json=payload, timeout=30)
|
||||
|
||||
if resp.status_code == 202:
|
||||
db_layer.update_task_status(task_id, 'PROCESSING')
|
||||
log_task(logger, task_id, 'dispatcher', 'Edge 接受任务,状态切换为 PROCESSING')
|
||||
elif resp.status_code == 429:
|
||||
logger.warning(f"[task_id={task_id}] Edge 队列已满 (429),稍后重试")
|
||||
elif resp.status_code == 503:
|
||||
logger.warning(f"[task_id={task_id}] Edge Ollama 不可用 (503),退避重试")
|
||||
self._schedule_retry(task)
|
||||
else:
|
||||
logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}")
|
||||
self._schedule_retry(task)
|
||||
|
||||
except requests.RequestException as e:
|
||||
logger.error(f"[task_id={task_id}] 下发失败: {e}")
|
||||
self._schedule_retry(task)
|
||||
|
||||
def _schedule_retry(self, task):
|
||||
"""调度重试"""
|
||||
task_id = task['task_id']
|
||||
if task['retry_count'] >= self.max_retries:
|
||||
db_layer.update_task_status(
|
||||
task_id, 'FAILED',
|
||||
error_message=f"超过最大重试次数 {self.max_retries}",
|
||||
failure_stage='callback'
|
||||
)
|
||||
logger.error(f"[task_id={task_id}] 超过最大重试次数,标记为 FAILED")
|
||||
return
|
||||
|
||||
backoff = self._calculate_backoff(task['retry_count'])
|
||||
next_retry = datetime.now() + timedelta(seconds=backoff)
|
||||
db_layer.increment_retry(task_id, next_retry)
|
||||
logger.info(f"[task_id={task_id}] 安排重试 #{task['retry_count']+1},{backoff}s 后执行 (at {next_retry})")
|
||||
|
||||
def _poll_once(self):
|
||||
"""执行一次轮询"""
|
||||
tasks = db_layer.get_pending_tasks(limit=10)
|
||||
for task in tasks:
|
||||
if self._should_retry(task):
|
||||
self._dispatch_one(task)
|
||||
|
||||
def _run(self):
|
||||
"""线程主循环"""
|
||||
logger.info(f"Dispatcher 启动,轮询间隔 {self.poll_interval}s")
|
||||
while self._running:
|
||||
try:
|
||||
self._poll_once()
|
||||
except Exception as e:
|
||||
logger.error(f"轮询异常: {e}", exc_info=True)
|
||||
time.sleep(self.poll_interval)
|
||||
|
||||
def start(self):
|
||||
"""启动下发线程"""
|
||||
if self._running:
|
||||
return
|
||||
self._running = True
|
||||
self._thread = threading.Thread(target=self._run, daemon=True, name='dispatcher')
|
||||
self._thread.start()
|
||||
|
||||
def stop(self):
|
||||
"""停止下发线程"""
|
||||
self._running = False
|
||||
if self._thread:
|
||||
self._thread.join(timeout=5)
|
||||
Reference in New Issue
Block a user