fix: Edge端三处修复 — API key丢失致全模型失败 + 失败结果黑洞 + 重派发死锁
根因: Edge服务手动重启未source .env,GEMINI/NVIDIA key丢失,
所有视觉模型调用失败("All models failed"),且失败结果永不回传NAS。
1. config_loader: 启动时自动加载.env(export KEY=VALUE),已存在环境变量不覆盖
2. queue_manager.enqueue: FAILED(或已交付SUCCESS)行重置为PENDING复用重跑,
修复 INSERT OR IGNORE + UNIQUE(nas_task_id) 导致的静默忽略死锁
3. get_undelivered_results: 包含FAILED状态,/results构造failed负载,
NAS Poller可感知失败避免僵尸循环空转重传
4. ddl.sql: failure_stage ENUM 增加 'process'(与db_layer valid_stages对齐),
修复 pymysql DataError 1265
验证: task 293(压缩后22MB) NVIDIA降级分析成功→event_id=13;
task 295 全链路成功→event_id=14;NAS端ENUM已同步ALTER
This commit is contained in:
@@ -358,6 +358,12 @@ def get_results():
|
|||||||
result_json = json.loads(r['result_json']) if r['result_json'] else None
|
result_json = json.loads(r['result_json']) if r['result_json'] else None
|
||||||
except json.JSONDecodeError:
|
except json.JSONDecodeError:
|
||||||
result_json = None
|
result_json = None
|
||||||
|
if r['status'] == 'FAILED':
|
||||||
|
result_json = {
|
||||||
|
"status": "failed",
|
||||||
|
"error_message": r['error_message'] or 'unknown',
|
||||||
|
"failure_stage": r['failure_stage'] or '',
|
||||||
|
}
|
||||||
payload.append({
|
payload.append({
|
||||||
"nas_task_id": r['nas_task_id'],
|
"nas_task_id": r['nas_task_id'],
|
||||||
"result": result_json,
|
"result": result_json,
|
||||||
|
|||||||
@@ -6,6 +6,37 @@ import re
|
|||||||
import yaml
|
import yaml
|
||||||
|
|
||||||
|
|
||||||
|
def _load_env_file():
|
||||||
|
"""加载项目根目录 .env(支持 export KEY=VALUE 格式)
|
||||||
|
|
||||||
|
根因: 服务手动启动时未 source .env 导致 GEMINI_API_KEY/NVIDIA_API_KEY
|
||||||
|
丢失,所有视觉模型调用失败。此处兜底加载,已存在的环境变量不覆盖。
|
||||||
|
"""
|
||||||
|
path = os.environ.get('FAM_ENV_FILE') or os.path.join(
|
||||||
|
os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))),
|
||||||
|
'.env'
|
||||||
|
)
|
||||||
|
if not os.path.isfile(path):
|
||||||
|
return
|
||||||
|
with open(path, 'r', encoding='utf-8') as f:
|
||||||
|
for line in f:
|
||||||
|
line = line.strip()
|
||||||
|
if not line or line.startswith('#'):
|
||||||
|
continue
|
||||||
|
if line.startswith('export '):
|
||||||
|
line = line[7:].strip()
|
||||||
|
if '=' not in line:
|
||||||
|
continue
|
||||||
|
key, _, value = line.partition('=')
|
||||||
|
key = key.strip()
|
||||||
|
value = value.strip().strip('"').strip("'")
|
||||||
|
if key and key not in os.environ:
|
||||||
|
os.environ[key] = value
|
||||||
|
|
||||||
|
|
||||||
|
_load_env_file()
|
||||||
|
|
||||||
|
|
||||||
def _resolve_env_vars(value):
|
def _resolve_env_vars(value):
|
||||||
if isinstance(value, str):
|
if isinstance(value, str):
|
||||||
def replace_env(match):
|
def replace_env(match):
|
||||||
|
|||||||
@@ -73,9 +73,24 @@ def enqueue(nas_task_id: int, video_filename: str, video_path: str,
|
|||||||
conn.commit()
|
conn.commit()
|
||||||
if cur.rowcount == 0:
|
if cur.rowcount == 0:
|
||||||
row = conn.execute(
|
row = conn.execute(
|
||||||
"SELECT id FROM task_queue WHERE nas_task_id=?", (nas_task_id,)
|
"SELECT id, status, delivered FROM task_queue WHERE nas_task_id=?",
|
||||||
|
(nas_task_id,)
|
||||||
).fetchone()
|
).fetchone()
|
||||||
return row['id'] if row else 0
|
if row is None:
|
||||||
|
return 0
|
||||||
|
# NAS 重新派发: FAILED(或已交付的 SUCCESS)重置为 PENDING,用新上传的视频重跑。
|
||||||
|
# SUCCESS 且未交付的行不动,避免丢失待 Poller 拉取的结果。
|
||||||
|
if row['status'] == 'FAILED' or (row['status'] == 'SUCCESS' and row['delivered']):
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE task_queue SET status='PENDING', result_json=NULL, error_message=NULL, "
|
||||||
|
"failure_stage=NULL, retry_count=0, delivered=0, video_filename=?, video_path=?, "
|
||||||
|
"camera_name=?, event_start_time=?, known_members_context=?, "
|
||||||
|
"updated_at=datetime('now','localtime') WHERE id=?",
|
||||||
|
(video_filename, video_path, camera_name, event_start_time,
|
||||||
|
known_members_context, row['id'])
|
||||||
|
)
|
||||||
|
conn.commit()
|
||||||
|
return row['id']
|
||||||
return cur.lastrowid
|
return cur.lastrowid
|
||||||
finally:
|
finally:
|
||||||
conn.close()
|
conn.close()
|
||||||
@@ -132,8 +147,10 @@ def mark_failed(task_id: int, error_message: str, failure_stage: str = ''):
|
|||||||
def get_undelivered_results(limit: int = 10) -> List[Dict]:
|
def get_undelivered_results(limit: int = 10) -> List[Dict]:
|
||||||
conn = _get_conn()
|
conn = _get_conn()
|
||||||
try:
|
try:
|
||||||
|
# FAILED 也需交付: 否则 NAS 永远收不到失败结果,任务卡 PROCESSING
|
||||||
|
# 直至僵尸回收后无意义地重传 22MB 视频
|
||||||
rows = conn.execute(
|
rows = conn.execute(
|
||||||
"SELECT * FROM task_queue WHERE status='SUCCESS' AND delivered=0 "
|
"SELECT * FROM task_queue WHERE status IN ('SUCCESS','FAILED') AND delivered=0 "
|
||||||
"ORDER BY id LIMIT ?", (limit,)
|
"ORDER BY id LIMIT ?", (limit,)
|
||||||
).fetchall()
|
).fetchall()
|
||||||
return [dict(r) for r in rows]
|
return [dict(r) for r in rows]
|
||||||
@@ -166,7 +183,8 @@ def get_queue_stats() -> Dict:
|
|||||||
).fetchone()
|
).fetchone()
|
||||||
stats[status] = row['cnt']
|
stats[status] = row['cnt']
|
||||||
row = conn.execute(
|
row = conn.execute(
|
||||||
"SELECT COUNT(*) as cnt FROM task_queue WHERE status='SUCCESS' AND delivered=0"
|
"SELECT COUNT(*) as cnt FROM task_queue "
|
||||||
|
"WHERE status IN ('SUCCESS','FAILED') AND delivered=0"
|
||||||
).fetchone()
|
).fetchone()
|
||||||
stats['UNDELIVERED'] = row['cnt']
|
stats['UNDELIVERED'] = row['cnt']
|
||||||
return stats
|
return stats
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ CREATE TABLE IF NOT EXISTS process_tasks (
|
|||||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
|
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
|
||||||
error_message TEXT NULL,
|
error_message TEXT NULL,
|
||||||
failure_stage ENUM('download','extract','vlm_visual','vlm_fusion','callback') NULL COMMENT '失败阶段',
|
failure_stage ENUM('download','extract','vlm_visual','vlm_fusion','callback','process') NULL COMMENT '失败阶段(process=Edge 消费者处理异常)',
|
||||||
heartbeat_at DATETIME NULL COMMENT '心跳时间戳(首期预留,不写入)',
|
heartbeat_at DATETIME NULL COMMENT '心跳时间戳(首期预留,不写入)',
|
||||||
INDEX idx_status (status),
|
INDEX idx_status (status),
|
||||||
INDEX idx_created (created_at),
|
INDEX idx_created (created_at),
|
||||||
|
|||||||
Reference in New Issue
Block a user