diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index f5b8718..89b4fa2 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -41,7 +41,10 @@ video_processing: max_concurrent: 1 # 消费者线程数(串行处理,避免云端并发超额) timeout: 900 # 兜底单视频分析超时 timeout_multiplier: 2 # 模型消费超时倍数:在 models[i].timeout 原值上 ×2(大视频上传+分析耗时) - max_retries: 2 # 单视频失败最大重试次数(failed 且 retry_count 超限不再消费) + max_retries: 10 # 单视频失败最大重试次数(配额/过载等瞬时故障给足重试机会) + retry_interval_sec: 3600 # 失败重试最小间隔:距上次失败 ≥1h 才重新入队,等配额恢复 + file_validate: true # 登记入队前用 OpenCV 校验文件可解码;失败标记 invalid 不入队 + stable_window_sec: 60 # 文件 mtime 稳定窗口:写入中(rclone 同步未完成)的文件跳过本轮 # 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed vision_order: ["gemini", "nvidia"] diff --git a/fam-edge/src/fam_edge/oracle_db.py b/fam-edge/src/fam_edge/oracle_db.py index 5defc4b..c169b8e 100644 --- a/fam-edge/src/fam_edge/oracle_db.py +++ b/fam-edge/src/fam_edge/oracle_db.py @@ -51,6 +51,10 @@ class OracleDB: event_start_time TEXT, status TEXT DEFAULT 'pending', retry_count INTEGER DEFAULT 0, + file_valid INTEGER DEFAULT 1, + file_error TEXT, + media_meta_json TEXT, + last_fail_at TEXT, summary_json TEXT, events_json TEXT, people_json TEXT, @@ -97,10 +101,17 @@ class OracleDB: CREATE INDEX IF NOT EXISTS idx_events_video ON events(video_id); CREATE INDEX IF NOT EXISTS idx_model_calls_created ON model_calls(created_at); """) - # 兼容旧库:补 retry_count 列(生产-消费队列重试上限用) + # 兼容旧库:补 retry_count / file_valid / media 等列(生产-消费队列用) cols = [r[1] for r in c.execute("PRAGMA table_info(videos)").fetchall()] - if 'retry_count' not in cols: - c.execute("ALTER TABLE videos ADD COLUMN retry_count INTEGER DEFAULT 0") + for col, ddl in [ + ('retry_count', "ALTER TABLE videos ADD COLUMN retry_count INTEGER DEFAULT 0"), + ('file_valid', "ALTER TABLE videos ADD COLUMN file_valid INTEGER DEFAULT 1"), + ('file_error', "ALTER TABLE videos ADD COLUMN file_error TEXT"), + ('media_meta_json', "ALTER TABLE videos ADD COLUMN media_meta_json TEXT"), + ('last_fail_at', "ALTER TABLE videos ADD COLUMN last_fail_at TEXT"), + ]: + if col not in cols: + c.execute(ddl) self._conn.commit() # ------------------------------------------------------------------ @@ -128,6 +139,18 @@ class OracleDB: duration_sec, 1 if success else 0, error or '', now)) self._conn.commit() + def set_video_file_status(self, video_id: int, valid: bool, + error: str = '', media_meta: dict = None): + """登记/更新文件校验结果:valid / file_error / media_meta_json""" + now = _now_iso() + self._conn.execute( + "UPDATE videos SET file_valid=?, file_error=?, media_meta_json=?, " + "updated_at=? WHERE id=?", + (1 if valid else 0, error or '', + json.dumps(media_meta, ensure_ascii=False) if media_meta else None, + now, video_id)) + self._conn.commit() + def ensure_video(self, filename: str, local_path: str, camera_name: str = '', event_start_time: str = '', duration_sec: float = 0.0, drive_file_id: str = '') -> int: @@ -183,8 +206,8 @@ class OracleDB: now = _now_iso() self._conn.execute( "UPDATE videos SET status='failed', retry_count=retry_count+1, " - "summary_json=?, updated_at=? WHERE id=?", - (error, now, video_id)) + "summary_json=?, updated_at=?, last_fail_at=? WHERE id=?", + (error, now, now, video_id)) self._conn.commit() def get_all_videos(self) -> List[sqlite3.Row]: diff --git a/fam-edge/src/fam_edge/video_processor.py b/fam-edge/src/fam_edge/video_processor.py index ef3b4f3..58ad6c2 100644 --- a/fam-edge/src/fam_edge/video_processor.py +++ b/fam-edge/src/fam_edge/video_processor.py @@ -86,6 +86,46 @@ def _parse_event_ts(ts: str, start_dt): return ts, 0.0 +def validate_video(path: str) -> tuple: + """校验视频文件是否为正常可解码视频。 + + 返回 (ok: bool, error: str, meta: dict|None) + - meta: {fps, frames, duration_sec, width, height} + 用 OpenCV 打开并读取至少 1 帧(不校验会导致空/半成品文件浪费云端配额)。 + """ + meta = None + try: + if not path or not os.path.isfile(path): + return False, "file_missing", None + if os.path.getsize(path) == 0: + return False, "file_empty", None + try: + import cv2 + except ImportError: + return True, "", None # 无 cv2 时跳过深度校验(仅大小检查) + cap = cv2.VideoCapture(path) + try: + if not cap.isOpened(): + return False, "cannot_open", None + ok, frame = cap.read() + if not ok or frame is None: + return False, "no_decodable_frame", None + fps = float(cap.get(cv2.CAP_PROP_FPS) or 0) + frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT) or 0) + meta = { + "fps": round(fps, 2), + "frames": frames, + "duration_sec": round(frames / max(fps, 0.01), 1), + "width": int(cap.get(cv2.CAP_PROP_FRAME_WIDTH) or 0), + "height": int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT) or 0), + } + finally: + cap.release() + return True, "", meta + except Exception as e: + return False, f"validate_exc: {e}", None + + class VideoProcessor: def __init__(self, db: oracle_db.OracleDB): self.config = load_config() @@ -93,6 +133,8 @@ class VideoProcessor: self.vision_order = self.config.get('video_processing', {}).get( 'vision_order', ['gemini', 'nvidia']) self.vision_timeout = self.config.get('video_processing', {}).get('timeout', 900) + self.file_validate = bool(self.config.get('video_processing', {}).get( + 'file_validate', True)) self.parse_start = self.config.get('gdrive_sync', {}).get( 'parse_start_from_filename', True) adapters = build_adapters(self.config.get('models', [])) @@ -122,6 +164,17 @@ class VideoProcessor: self.db.mark_video_failed(video_id, "file_missing") return False + # 处理前二次确认文件有效性(防止登记后文件被破坏/截断;校验结果落库) + if self.file_validate: + ok, verr, vmeta = validate_video(local_path) + if not ok: + logger.error(f"[video_id={video_id}] 文件校验失败({verr}),标记 failed: {local_path}") + self.db.set_video_file_status(video_id, False, verr) + self.db.mark_video_failed(video_id, f"invalid_file:{verr}") + return False + if vmeta: + self.db.set_video_file_status(video_id, True, '', vmeta) + camera_name = self.db.get_video_by_filename(filename)['camera_name'] or '' event_start = '' if self.parse_start: diff --git a/fam-edge/src/fam_edge/video_queue.py b/fam-edge/src/fam_edge/video_queue.py index 9f60424..88458ee 100644 --- a/fam-edge/src/fam_edge/video_queue.py +++ b/fam-edge/src/fam_edge/video_queue.py @@ -21,7 +21,7 @@ from typing import List, Optional from .logger import setup_logger from .config_loader import load_config from . import oracle_db -from .video_processor import VideoProcessor +from .video_processor import VideoProcessor, validate_video logger = setup_logger('fam-edge.video_queue') @@ -44,7 +44,12 @@ class VideoQueue: # 模型超时放大倍数:消费时在模型原配置 timeout 上 ×N self.timeout_multiplier = float(vp.get('timeout_multiplier', 2.0)) # 单视频失败最大重试次数(failed 且 retry_count>=max_retries 不再消费) - self.max_retries = int(vp.get('max_retries', 2)) + self.max_retries = int(vp.get('max_retries', 10)) + # 失败重试最小间隔(秒):配额类瞬时故障(如 429)需等其恢复再重试,避免短时间重复打爆 + self.retry_interval_sec = int(vp.get('retry_interval_sec', 3600)) + # 文件校验:入队前用 OpenCV 确认可解码;mtime 稳定窗口防 rclone 写入中的半成品 + self.file_validate = bool(vp.get('file_validate', True)) + self.stable_window_sec = int(vp.get('stable_window_sec', 60)) self._queue: "queue.Queue[int]" = queue.Queue() self._queued: set = set() # 已在队中的 video_id(防重复入队) @@ -75,6 +80,26 @@ class VideoQueue: fn = os.path.basename(path) row = self.db.get_video_by_filename(fn) if row is None: + # 新文件:先过 mtime 稳定窗口 + 可解码校验,通过才登记入队;失败标记 invalid + if self.file_validate: + if time.time() - os.path.getmtime(path) < self.stable_window_sec: + logger.info(f"文件仍在写入(mtime 未稳定),跳过本轮: {fn}") + continue + ok, verr, vmeta = validate_video(path) + if not ok: + vid = self.db.ensure_video(fn, path, camera_name=self.camera_name) + self.db.set_video_file_status(vid, False, verr) + self.db._conn.execute( + "UPDATE videos SET status='invalid' WHERE id=?", (vid,)) + self.db._conn.commit() + logger.warning(f"文件校验失败,标记 invalid 不入队: {fn} ({verr})") + continue + if vmeta: + vid = self.db.ensure_video(fn, path, camera_name=self.camera_name) + self.db.set_video_file_status(vid, True, '', vmeta) + logger.info(f"登记新视频并入队: {fn} (id={vid}, meta={vmeta})") + self._enqueue(vid) + continue vid = self.db.ensure_video(fn, path, camera_name=self.camera_name) logger.info(f"登记新视频并入队: {fn} (id={vid})") self._enqueue(vid) @@ -91,10 +116,26 @@ class VideoQueue: logger.info(f"启动恢复入队 {len(rows)} 个待处理视频") def _retry_allowed(self, row) -> bool: - if row['status'] == 'pending': + """是否允许消费/重试该视频:pending 直接放行;failed 需未超重试上限且距上次失败够久;invalid 永不放行""" + status = row['status'] + if status == 'pending': return True - # failed:未超重试上限才允许再次消费 - return int(row['retry_count'] or 0) < self.max_retries + if status == 'invalid': + return False + # failed:未超重试上限,且距上次失败 >= retry_interval_sec(等配额类故障恢复) + if int(row['retry_count'] or 0) >= self.max_retries: + return False + last_fail = row['last_fail_at'] or row['updated_at'] or '' + if last_fail: + try: + from datetime import datetime + lt = datetime.strptime(last_fail[:19], '%Y-%m-%d %H:%M:%S') + elapsed = (datetime.now() - lt).total_seconds() + if elapsed < self.retry_interval_sec: + return False + except ValueError: + pass + return True def _enqueue(self, video_id: int): with self._lock: