diff --git a/fam-core/config/config.yaml.example b/fam-core/config/config.yaml.example new file mode 100644 index 0000000..3e834fd --- /dev/null +++ b/fam-core/config/config.yaml.example @@ -0,0 +1,36 @@ +# FAM-Core 配置文件 (NAS 端) +# 复制此文件为 config.yaml 并修改实际值 + +server: + host: "0.0.0.0" + port: 8000 + +database: + host: "127.0.0.1" + port: 3306 + user: "root" + password: "" + database: "sentinel_home_ai" + +scheduler: + scan_interval: 60 # 扫描间隔(秒) + video_dir: "/volume1/surveillance" # 视频录制目录 + video_extensions: [".mp4", ".mkv", ".avi"] + file_stable_seconds: 60 # 文件稳定判定秒数 + camera_name: "默认摄像头" + +dispatcher: + poll_interval: 30 # 轮询间隔(秒) + edge_url: "http://100.x.x.20:5000/api/edge/video/analyze" + webhook_url: "http://100.x.x.10:8000/api/core/callback/event" + max_retries: 3 + +video_server: + base_url: "http://100.x.x.10:8000/media" + token: "your-secret-token-here" # 视频访问 token + video_dir: "/volume1/surveillance" + +chat_handler: + ollama_url: "http://100.x.x.20:11434/api/generate" + model_name: "llava-phi3" + timeout: 120 diff --git a/fam-core/requirements.txt b/fam-core/requirements.txt new file mode 100644 index 0000000..1be6868 --- /dev/null +++ b/fam-core/requirements.txt @@ -0,0 +1,5 @@ +flask>=3.0.0 +gunicorn>=21.2.0 +mysql-connector-python>=8.3.0 +PyYAML>=6.0 +requests>=2.31.0 diff --git a/fam-core/src/fam_core/__init__.py b/fam-core/src/fam_core/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/app.py b/fam-core/src/fam_core/app.py new file mode 100644 index 0000000..f04f5d7 --- /dev/null +++ b/fam-core/src/fam_core/app.py @@ -0,0 +1,68 @@ +""" +FAM-Core 主应用 - Flask 单进程 + +承载: Task-Scheduler / Dispatcher / Event-Receiver / Chat-Handler / Member-Manager / Video-Server +""" +import os +import sys +from flask import Flask, jsonify + +# 确保包路径 +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + +from .config_loader import load_config +from .logger import setup_logger +from .scheduler.scheduler import TaskScheduler +from .dispatcher.dispatcher import Dispatcher +from .event_receiver.event_receiver import event_bp +from .chat_handler.chat_handler import chat_bp +from .member_manager.member_manager import member_bp +from .video_server.video_server import video_bp + +logger = setup_logger('fam-core.app') + +app = Flask(__name__) + +# 注册蓝图 +app.register_blueprint(event_bp) +app.register_blueprint(chat_bp) +app.register_blueprint(member_bp) +app.register_blueprint(video_bp) + +# 健康检查 +@app.route('/health', methods=['GET']) +def health(): + return jsonify({"status": "ok", "service": "fam-core"}), 200 + +# 初始化后台线程 +_scheduler = None +_dispatcher = None + +try: + _scheduler = TaskScheduler() + _scheduler.start() + logger.info("Task-Scheduler 已启动") +except Exception as e: + logger.error(f"Task-Scheduler 启动失败: {e}") + +try: + _dispatcher = Dispatcher() + _dispatcher.start() + logger.info("Dispatcher 已启动") +except Exception as e: + logger.error(f"Dispatcher 启动失败: {e}") + + +@app.route('/api/status', methods=['GET']) +def status(): + """系统状态""" + return jsonify({ + "scheduler_running": _scheduler._running if _scheduler else False, + "dispatcher_running": _dispatcher._running if _dispatcher else False, + }), 200 + + +if __name__ == '__main__': + cfg = load_config() + port = cfg.get('server', {}).get('port', 8000) + app.run(host='0.0.0.0', port=port, debug=False) diff --git a/fam-core/src/fam_core/chat_handler/__init__.py b/fam-core/src/fam_core/chat_handler/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/chat_handler/chat_handler.py b/fam-core/src/fam_core/chat_handler/chat_handler.py new file mode 100644 index 0000000..a82ddd4 --- /dev/null +++ b/fam-core/src/fam_core/chat_handler/chat_handler.py @@ -0,0 +1,173 @@ +""" +Chat-Handler - Flask 蓝图,接收用户问答 + +处理逻辑: +1. 根据 queried_person 和 queried_date 查询 event_details +2. 拼接上下文(每条明细一行) +3. 若明细条数 > 50,按小时聚合成摘要 +4. POST Oracle Ollama 纯文本模式,调问答 Prompt +5. 插入 chat_history +6. 返回回答 +""" +import requests +from flask import Blueprint, request, jsonify +from collections import defaultdict + +from ..logger import setup_logger +from ..config_loader import load_config +from .. import db_layer + +logger = setup_logger('fam-core.chat_handler') + +chat_bp = Blueprint('chat_handler', __name__) + +CHAT_SYSTEM_PROMPT = """你是家庭监控助手。根据以下今日监控数据,回答用户问题。 + +今日数据(按时间顺序,每条一行): +{context} + +已知家庭成员: {members} + +用户问题: {question} + +要求: +- 只基于上述数据回答,不要编造 +- 按时间顺序总结 +- 若有关注事件(跌倒、哭闹等),重点提示 +- 若当天没有该人员的数据,明确说"今天没有观察到{person}" +- 用自然语言回答,不要输出 JSON +""" + + +def _format_details(details): + """将 event_details 格式化为文本行""" + lines = [] + for d in details: + timestamp = d['frame_timestamp'].strftime('%H:%M') if hasattr(d['frame_timestamp'], 'strftime') else str(d['frame_timestamp']) + camera = d.get('camera_name', '') + person = d.get('person', '') + action = d.get('action', '') + clothing = d.get('clothing', '') + attention = ' [关注事件]' if d.get('is_attention_event') else '' + lines.append(f"[{timestamp} {camera}] {person} {action} ({clothing}){attention}") + return '\n'.join(lines) + + +def _aggregate_by_hour(details): + """当明细 > 50 条时,按小时聚合""" + hourly = defaultdict(list) + for d in details: + ts = d['frame_timestamp'] + hour_key = ts.strftime('%Y-%m-%d %H:00') if hasattr(ts, 'strftime') else str(ts) + hourly[hour_key].append(d) + + lines = [] + for hour, items in sorted(hourly.items()): + persons = set(i.get('person', '') for i in items) + actions = set(i.get('action', '') for i in items) + has_attention = any(i.get('is_attention_event') for i in items) + attention = ' [含关注事件]' if has_attention else '' + lines.append(f"[{hour}] {','.join(persons)}: {','.join(actions)}{attention}") + return '\n'.join(lines) + + +def _call_ollama(prompt: str) -> str: + """调用 Oracle Ollama 纯文本模式""" + cfg = load_config() + ollama_url = cfg.get('chat_handler', {}).get( + 'ollama_url', 'http://localhost:11434/api/generate' + ) + model_name = cfg.get('chat_handler', {}).get('model_name', 'llava-phi3') + timeout = cfg.get('chat_handler', {}).get('timeout', 120) + + resp = requests.post(ollama_url, json={ + "model": model_name, + "prompt": prompt, + "stream": False, + "options": {"temperature": 0.3} + }, timeout=timeout) + + if resp.status_code == 200: + return resp.json().get('response', '') + else: + logger.error(f"Ollama 调用失败: {resp.status_code} {resp.text}") + raise Exception(f"Ollama error: {resp.status_code}") + + +@chat_bp.route('/api/chat/ask', methods=['POST']) +def chat_ask(): + """接收用户问答""" + data = request.get_json(silent=True) + if not data: + return jsonify({"error": "Invalid JSON"}), 400 + + question = data.get('question', '') + queried_person = data.get('queried_person', '') + queried_date = data.get('queried_date', '') + + if not question or not queried_person or not queried_date: + return jsonify({"error": "缺少必填字段: question, queried_person, queried_date"}), 400 + + logger.info(f"Chat: person={queried_person}, date={queried_date}, question={question}") + + # 1. 查询 event_details + details = db_layer.query_event_details(queried_person, queried_date) + + # 2. 拼接上下文 + if len(details) == 0: + # 无数据 + answer = f"今天没有观察到{queried_person}。" + context_summary = "查询 event_details 0 条" + elif len(details) > 50: + context = _aggregate_by_hour(details) + context_summary = f"查询 event_details {len(details)} 条,按小时聚合为 {len(set(d['frame_timestamp'].strftime('%Y-%m-%d %H') for d in details))} 行" + else: + context = _format_details(details) + context_summary = f"查询 event_details {len(details)} 条,时间范围 {details[0]['frame_timestamp']} - {details[-1]['frame_timestamp']}" + + if len(details) > 0: + # 构建完整 Prompt + members = db_layer.get_known_members_context() + prompt = CHAT_SYSTEM_PROMPT.format( + context=context, + members=members or f"{queried_person}", + question=question, + person=queried_person + ) + + try: + answer = _call_ollama(prompt) + except Exception as e: + logger.error(f"Ollama 调用失败: {e}") + return jsonify({"error": f"AI 调用失败: {e}"}), 503 + + # 3. 写入 chat_history + chat_id = db_layer.insert_chat_history( + user_question=question, + ai_answer=answer, + context_summary=context_summary, + queried_date=queried_date, + queried_person=queried_person + ) + + return jsonify({ + "answer": answer, + "context_summary": context_summary, + "chat_id": chat_id + }), 200 + + +@chat_bp.route('/api/chat/history', methods=['GET']) +def chat_history(): + """获取对话历史""" + date = request.args.get('date') + person = request.args.get('person') + limit = int(request.args.get('limit', 20)) + + history = db_layer.get_chat_history(limit=limit, date_filter=date, person_filter=person) + # datetime 序列化 + for h in history: + for k, v in h.items(): + if hasattr(v, 'isoformat'): + h[k] = v.isoformat() + return jsonify({"history": history}), 200 diff --git a/fam-core/src/fam_core/config_loader.py b/fam-core/src/fam_core/config_loader.py new file mode 100644 index 0000000..914a0ff --- /dev/null +++ b/fam-core/src/fam_core/config_loader.py @@ -0,0 +1,32 @@ +""" +配置加载器 - 从 config.yaml 读取配置 +""" +import os +import re +import yaml + + +def _resolve_env_vars(value): + """递归解析字符串中的 ${ENV_VAR} 引用""" + if isinstance(value, str): + def replace_env(match): + env_name = match.group(1) + return os.environ.get(env_name, match.group(0)) + return re.sub(r'\$\{(\w+)\}', replace_env, value) + elif isinstance(value, dict): + return {k: _resolve_env_vars(v) for k, v in value.items()} + elif isinstance(value, list): + return [_resolve_env_vars(item) for item in value] + return value + + +def load_config(config_path=None): + """加载 YAML 配置文件,自动解析 ${ENV_VAR} 引用""" + if config_path is None: + config_path = os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), + 'config', 'config.yaml' + ) + with open(config_path, 'r', encoding='utf-8') as f: + raw = yaml.safe_load(f) + return _resolve_env_vars(raw) diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py new file mode 100644 index 0000000..31a5672 --- /dev/null +++ b/fam-core/src/fam_core/db_layer.py @@ -0,0 +1,467 @@ +""" +数据库访问层 - MariaDB 连接管理与 CRUD 操作 +""" +import json +import mysql.connector +from mysql.connector import pooling +from datetime import datetime +from typing import Optional, List, Dict, Any + +from .config_loader import load_config +from .logger import setup_logger + +logger = setup_logger('fam-core.db') + +_config = None +_pool = None + + +def get_config(): + global _config + if _config is None: + _config = load_config() + return _config + + +def get_pool(): + """获取数据库连接池(单例)""" + global _pool + if _pool is None: + cfg = get_config().get('database', {}) + _pool = pooling.MySQLConnectionPool( + host=cfg.get('host', '127.0.0.1'), + port=cfg.get('port', 3306), + user=cfg.get('user', 'root'), + password=cfg.get('password', ''), + database=cfg.get('database', 'sentinel_home_ai'), + charset='utf8mb4', + collation='utf8mb4_unicode_ci', + pool_name='fam_core_pool', + pool_size=5, + autocommit=False + ) + return _pool + + +def get_conn(): + """从连接池获取一个连接""" + return get_pool().get_connection() + + +# ============================================================ +# process_tasks 操作 +# ============================================================ + +def create_task(video_path: str, video_url: str) -> int: + """创建新任务""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + "INSERT INTO process_tasks (video_path, video_url, status) VALUES (%s, %s, 'PENDING')", + (video_path, video_url) + ) + conn.commit() + task_id = cursor.lastrowid + logger.info(f"[task_id={task_id}] task created: {video_path}") + return task_id + finally: + conn.close() + + +def get_pending_tasks(limit=10) -> List[Dict]: + """获取待处理任务""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + "SELECT * FROM process_tasks WHERE status = 'PENDING' ORDER BY created_at ASC LIMIT %s", + (limit,) + ) + return cursor.fetchall() + finally: + conn.close() + + +def get_tasks_by_status(status: str, limit=10) -> List[Dict]: + """按状态获取任务""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + "SELECT * FROM process_tasks WHERE status = %s ORDER BY created_at ASC LIMIT %s", + (status, limit) + ) + return cursor.fetchall() + finally: + conn.close() + + +def update_task_status(task_id: int, status: str, error_message: str = None, + failure_stage: str = None): + """更新任务状态""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + "UPDATE process_tasks SET status=%s, error_message=%s, failure_stage=%s WHERE task_id=%s", + (status, error_message, failure_stage, task_id) + ) + conn.commit() + finally: + conn.close() + + +def increment_retry(task_id: int, next_retry_at: datetime): + """递增重试次数""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + "UPDATE process_tasks SET retry_count=retry_count+1, next_retry_at=%s, status='PENDING' WHERE task_id=%s", + (next_retry_at, task_id) + ) + conn.commit() + finally: + conn.close() + + +def get_task(task_id: int) -> Optional[Dict]: + """获取单个任务""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute("SELECT * FROM process_tasks WHERE task_id = %s", (task_id,)) + return cursor.fetchone() + finally: + conn.close() + + +def get_video_url_exists(video_path: str) -> bool: + """检查视频是否已有对应任务(避免重复)""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + "SELECT COUNT(*) FROM process_tasks WHERE video_path = %s", + (video_path,) + ) + return cursor.fetchone()[0] > 0 + finally: + conn.close() + + +# ============================================================ +# monitor_events 操作 +# ============================================================ + +def insert_event(task_id: int, event_start_time: str, event_end_time: str, + camera_name: str, global_summary: str, entities_json: list, + compute_provider: list) -> int: + """插入事件聚合记录""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + """INSERT INTO monitor_events + (task_id, event_start_time, event_end_time, camera_name, + global_summary, entities_json, compute_provider) + VALUES (%s, %s, %s, %s, %s, %s, %s)""", + (task_id, event_start_time, event_end_time, camera_name, + global_summary, json.dumps(entities_json, ensure_ascii=False), + json.dumps(compute_provider, ensure_ascii=False)) + ) + conn.commit() + return cursor.lastrowid + finally: + conn.close() + + +# ============================================================ +# event_details 操作 +# ============================================================ + +def insert_event_detail(event_id: int, task_id: int, frame_index: int, + frame_timestamp: str, camera_name: str, + person: str, action: str, clothing: str, + is_attention_event: bool, source_providers: list): + """插入事件明细""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + """INSERT INTO event_details + (event_id, task_id, frame_index, frame_timestamp, camera_name, + person, action, clothing, is_attention_event, source_providers) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)""", + (event_id, task_id, frame_index, frame_timestamp, camera_name, + person, action, clothing, is_attention_event, + json.dumps(source_providers, ensure_ascii=False)) + ) + conn.commit() + finally: + conn.close() + + +def query_event_details(person: str, queried_date: str) -> List[Dict]: + """查询某人在某天的事件明细""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + """SELECT frame_timestamp, camera_name, person, action, clothing, + is_attention_event + FROM event_details + WHERE person = %s AND DATE(frame_timestamp) = %s + ORDER BY frame_timestamp ASC""", + (person, queried_date) + ) + return cursor.fetchall() + finally: + conn.close() + + +def get_recent_events(limit=20, offset=0, date_filter=None) -> List[Dict]: + """获取事件列表(分页 + 日期筛选)""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + if date_filter: + cursor.execute( + """SELECT me.event_id, me.task_id, me.event_start_time, me.event_end_time, + me.camera_name, me.global_summary, me.compute_provider, + me.created_at, + (SELECT COUNT(*) FROM event_details ed WHERE ed.event_id = me.event_id) AS detail_count + FROM monitor_events me + WHERE DATE(me.event_start_time) = %s + ORDER BY me.event_start_time DESC + LIMIT %s OFFSET %s""", + (date_filter, limit, offset) + ) + else: + cursor.execute( + """SELECT me.event_id, me.task_id, me.event_start_time, me.event_end_time, + me.camera_name, me.global_summary, me.compute_provider, + me.created_at, + (SELECT COUNT(*) FROM event_details ed WHERE ed.event_id = me.event_id) AS detail_count + FROM monitor_events me + ORDER BY me.event_start_time DESC + LIMIT %s OFFSET %s""", + (limit, offset) + ) + return cursor.fetchall() + finally: + conn.close() + + +# ============================================================ +# chat_history 操作 +# ============================================================ + +def insert_chat_history(user_question: str, ai_answer: str, + context_summary: str, queried_date: str, + queried_person: str) -> int: + """插入对话记录""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + """INSERT INTO chat_history + (user_question, ai_answer, context_summary, queried_date, queried_person) + VALUES (%s, %s, %s, %s, %s)""", + (user_question, ai_answer, context_summary, queried_date, queried_person) + ) + conn.commit() + return cursor.lastrowid + finally: + conn.close() + + +def get_chat_history(limit=20, date_filter=None, person_filter=None) -> List[Dict]: + """获取对话历史""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + conditions = [] + params = [] + if date_filter: + conditions.append("queried_date = %s") + params.append(date_filter) + if person_filter: + conditions.append("queried_person = %s") + params.append(person_filter) + where = f"WHERE {' AND '.join(conditions)}" if conditions else "" + params.extend([limit]) + cursor.execute( + f"SELECT * FROM chat_history {where} ORDER BY created_at DESC LIMIT %s", + params + ) + return cursor.fetchall() + finally: + conn.close() + + +# ============================================================ +# family_members 操作 +# ============================================================ + +def upsert_family_member(abstract_label: str, feature_description: str, + first_seen_at: str): + """upsert 家庭成员(abstract_label 唯一)""" + conn = get_conn() + try: + cursor = conn.cursor() + cursor.execute( + """INSERT INTO family_members (abstract_label, feature_description, first_seen_at) + VALUES (%s, %s, %s) + ON DUPLICATE KEY UPDATE abstract_label = abstract_label""", + (abstract_label, feature_description, first_seen_at) + ) + conn.commit() + finally: + conn.close() + + +def get_unnamed_members() -> List[Dict]: + """获取未命名成员列表""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + """SELECT fm.abstract_label, fm.feature_description, fm.first_seen_at, + (SELECT COUNT(*) FROM event_details ed WHERE ed.person = fm.abstract_label) AS event_count + FROM family_members fm + WHERE fm.real_name IS NULL AND fm.is_active = TRUE + ORDER BY fm.first_seen_at ASC""" + ) + return cursor.fetchall() + finally: + conn.close() + + +def get_all_members(include_named=True, include_unnamed=True) -> List[Dict]: + """获取所有成员""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + conditions = [] + if include_named and include_unnamed: + pass # 全部 + elif include_named: + conditions.append("real_name IS NOT NULL") + elif include_unnamed: + conditions.append("real_name IS NULL") + where = f"WHERE {' AND '.join(conditions)}" if conditions else "" + cursor.execute( + f"""SELECT * FROM family_members {where} + ORDER BY first_seen_at ASC""" + ) + return cursor.fetchall() + finally: + conn.close() + + +def name_member(abstract_label: str, real_name: str, named_by: str) -> Dict: + """命名成员 + 批量回溯更新历史记录""" + conn = get_conn() + try: + cursor = conn.cursor() + + # 1. 校验存在且未命名 + cursor.execute( + "SELECT member_id FROM family_members WHERE abstract_label = %s AND real_name IS NULL", + (abstract_label,) + ) + row = cursor.fetchone() + if not row: + return {"error": f"成员 {abstract_label} 不存在或已命名"} + + # 2. 更新 family_members + cursor.execute( + """UPDATE family_members + SET real_name = %s, named_at = NOW(), named_by = %s + WHERE abstract_label = %s""", + (real_name, named_by, abstract_label) + ) + + # 3. 批量回溯更新 event_details + cursor.execute( + "UPDATE event_details SET person = %s WHERE person = %s", + (real_name, abstract_label) + ) + updated_details_count = cursor.rowcount + + # 4. 批量回溯更新 monitor_events.entities_json + cursor.execute( + """UPDATE monitor_events + SET entities_json = JSON_REPLACE( + entities_json, '$[*].person', %s + ) + WHERE JSON_CONTAINS(entities_json->'$[*].person', JSON_QUOTE(%s))""", + (real_name, abstract_label) + ) + updated_events_count = cursor.rowcount + + conn.commit() + return { + "abstract_label": abstract_label, + "real_name": real_name, + "updated_event_details_count": updated_details_count, + "updated_monitor_events_count": updated_events_count + } + except Exception as e: + conn.rollback() + raise e + finally: + conn.close() + + +def get_known_members_context() -> str: + """获取已命名+未命名成员清单,用于注入 VLM Prompt""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + "SELECT abstract_label, real_name, feature_description FROM family_members WHERE is_active = TRUE" + ) + members = cursor.fetchall() + if not members: + return "" + parts = [] + for m in members: + name = m['real_name'] if m['real_name'] else m['abstract_label'] + feat = m['feature_description'] or '' + status = f"(real_name={m['real_name']})" if m['real_name'] else f"(abstract_label={m['abstract_label']}, 未命名)" + parts.append(f"{name}: {feat} {status}") + return "; ".join(parts) + finally: + conn.close() + + +# ============================================================ +# compute_provider 统计 +# ============================================================ + +def get_compute_provider_stats() -> List[Dict]: + """获取 compute_provider 分布统计""" + conn = get_conn() + try: + cursor = conn.cursor(dictionary=True) + cursor.execute( + """SELECT + JSON_UNQUOTE(JSON_EXTRACT(item, '$')) AS provider, + COUNT(*) AS count + FROM monitor_events, + JSON_TABLE(compute_provider, '$[*]' + COLUMNS(item VARCHAR(50) PATH '$' + )) AS jt + GROUP BY provider + ORDER BY count DESC""" + ) + return cursor.fetchall() + except Exception: + # MariaDB 旧版不支持 JSON_TABLE,降级方案 + cursor.execute("SELECT compute_provider, COUNT(*) AS count FROM monitor_events GROUP BY compute_provider") + return cursor.fetchall() + finally: + conn.close() diff --git a/fam-core/src/fam_core/dispatcher/__init__.py b/fam-core/src/fam_core/dispatcher/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py new file mode 100644 index 0000000..71f7a29 --- /dev/null +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -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) diff --git a/fam-core/src/fam_core/event_receiver/__init__.py b/fam-core/src/fam_core/event_receiver/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/event_receiver/event_receiver.py b/fam-core/src/fam_core/event_receiver/event_receiver.py new file mode 100644 index 0000000..74b8da5 --- /dev/null +++ b/fam-core/src/fam_core/event_receiver/event_receiver.py @@ -0,0 +1,110 @@ +""" +Event-Receiver - Flask 蓝图,接收 Edge 回调,写库 + +处理逻辑: +1. 成功回调: 插入 monitor_events 1 条 + 遍历 frame_details 逐条插入 event_details +2. 对未命名的 abstract_label 自动 upsert 到 family_members +3. 更新 process_tasks 状态为 SUCCESS +4. 失败回调: 更新任务状态为 FAILED,记录 failure_stage +""" +import re +import json +from flask import Blueprint, request, jsonify + +from ..logger import setup_logger +from .. import db_layer + +logger = setup_logger('fam-core.event_receiver') + +event_bp = Blueprint('event_receiver', __name__) + +# 匹配 "人物A" / "人物B" 等 abstract_label +_ABSTRACT_LABEL_PATTERN = re.compile(r'^人物[A-Z]$') + + +def _is_abstract_label(person: str) -> bool: + """判断是否为未命名的 abstract_label""" + return bool(_ABSTRACT_LABEL_PATTERN.match(person)) + + +def _upsert_abstract_members(frame_details: list): + """对未命名的 abstract_label 自动 upsert 到 family_members""" + seen = {} + for frame in frame_details: + person = frame.get('person', '') + if _is_abstract_label(person): + clothing = frame.get('clothing', '') + action = frame.get('action', '') + feature = f"{clothing}、{action}" if clothing and action else clothing or action + timestamp = frame.get('frame_timestamp', '') + if person not in seen: + seen[person] = (feature, timestamp) + + for label, (feature, ts) in seen.items(): + db_layer.upsert_family_member(label, feature, ts) + logger.info(f"upsert family_member: {label} (feature={feature})") + + +@event_bp.route('/api/core/callback/event', methods=['POST']) +def receive_event(): + """接收 Edge 回调""" + data = request.get_json(silent=True) + if not data: + return jsonify({"error": "Invalid JSON"}), 400 + + task_id = data.get('task_id') + status = data.get('status') + logger.info(f"[task_id={task_id}] 收到回调: status={status}") + + if status == 'success': + try: + # 1. 插入 monitor_events + event_id = db_layer.insert_event( + task_id=task_id, + event_start_time=data['event_start_time'], + event_end_time=data['event_end_time'], + camera_name=data.get('camera_name', ''), + global_summary=data.get('global_summary', ''), + entities_json=data.get('entities_json', []), + compute_provider=data.get('compute_provider', []) + ) + + # 2. 遍历 frame_details 逐条插入 + frame_details = data.get('frame_details', []) + for frame in frame_details: + db_layer.insert_event_detail( + event_id=event_id, + task_id=task_id, + frame_index=frame.get('frame_index', 0), + frame_timestamp=frame.get('frame_timestamp', ''), + camera_name=frame.get('camera_name', data.get('camera_name', '')), + person=frame.get('person', '未知'), + action=frame.get('action', ''), + clothing=frame.get('clothing', ''), + is_attention_event=frame.get('is_attention_event', False), + source_providers=frame.get('source_providers', []) + ) + + # 3. 对未命名的 abstract_label 自动 upsert + _upsert_abstract_members(frame_details) + + # 4. 更新任务状态 + db_layer.update_task_status(task_id, 'SUCCESS') + logger.info(f"[task_id={task_id}] 事件处理完成: event_id={event_id}, frame_details={len(frame_details)}条") + + return jsonify({"status": "ok", "event_id": event_id}), 200 + + except Exception as e: + logger.error(f"[task_id={task_id}] 处理回调失败: {e}", exc_info=True) + db_layer.update_task_status(task_id, 'FAILED', error_message=str(e), failure_stage='callback') + return jsonify({"error": str(e)}), 500 + + elif status == 'failed': + failure_stage = data.get('failure_stage', '') + error_message = data.get('error_message', '') + db_layer.update_task_status(task_id, 'FAILED', error_message=error_message, failure_stage=failure_stage) + logger.error(f"[task_id={task_id}] 任务失败: stage={failure_stage}, error={error_message}") + return jsonify({"status": "ok"}), 200 + + else: + return jsonify({"error": f"Unknown status: {status}"}), 400 diff --git a/fam-core/src/fam_core/logger.py b/fam-core/src/fam_core/logger.py new file mode 100644 index 0000000..f9c6f3d --- /dev/null +++ b/fam-core/src/fam_core/logger.py @@ -0,0 +1,31 @@ +""" +日志工具 - 统一格式,带 task_id 作为 trace_id +""" +import logging +import sys +from datetime import datetime + + +def setup_logger(name='fam-core', level=logging.INFO): + """配置并返回 logger""" + 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.setFormatter(formatter) + logger.addHandler(handler) + return logger + + +def log_task(logger, task_id, stage, message, level=logging.INFO, duration_ms=None): + """带 task_id 的结构化日志""" + parts = [f"[task_id={task_id}]", stage] + if duration_ms is not None: + parts.append(f"done in {duration_ms}ms") + parts.append(message) + logger.log(level, ' '.join(parts)) diff --git a/fam-core/src/fam_core/member_manager/__init__.py b/fam-core/src/fam_core/member_manager/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/member_manager/member_manager.py b/fam-core/src/fam_core/member_manager/member_manager.py new file mode 100644 index 0000000..53b977c --- /dev/null +++ b/fam-core/src/fam_core/member_manager/member_manager.py @@ -0,0 +1,79 @@ +""" +Member-Manager - Flask 蓝图,成员命名管理 + +1. GET /api/member/unnamed - 列出未命名人物 +2. POST /api/member/name - 命名人物 + 批量回溯更新 +3. GET /api/member/list - 列出所有成员 +""" +from flask import Blueprint, request, jsonify + +from ..logger import setup_logger +from .. import db_layer + +logger = setup_logger('fam-core.member_manager') + +member_bp = Blueprint('member_manager', __name__) + + +@member_bp.route('/api/member/unnamed', methods=['GET']) +def list_unnamed(): + """列出未命名人物""" + members = db_layer.get_unnamed_members() + # datetime 序列化 + result = [] + for m in members: + result.append({ + "abstract_label": m['abstract_label'], + "feature_description": m['feature_description'], + "first_seen_at": m['first_seen_at'].isoformat() if hasattr(m['first_seen_at'], 'isoformat') else str(m['first_seen_at']), + "event_count": m['event_count'] + }) + return jsonify({"unnamed_members": result}), 200 + + +@member_bp.route('/api/member/name', methods=['POST']) +def name_member(): + """命名人物 + 批量回溯更新历史记录""" + data = request.get_json(silent=True) + if not data: + return jsonify({"error": "Invalid JSON"}), 400 + + abstract_label = data.get('abstract_label') + real_name = data.get('real_name') + named_by = data.get('named_by', '管理员') + + if not abstract_label or not real_name: + return jsonify({"error": "缺少必填字段: abstract_label, real_name"}), 400 + + logger.info(f"命名: {abstract_label} -> {real_name}") + + try: + result = db_layer.name_member(abstract_label, real_name, named_by) + if 'error' in result: + return jsonify(result), 404 + return jsonify(result), 200 + except Exception as e: + logger.error(f"命名失败: {e}", exc_info=True) + return jsonify({"error": str(e)}), 500 + + +@member_bp.route('/api/member/list', methods=['GET']) +def list_members(): + """列出所有成员""" + include_named = request.args.get('include_named', 'true').lower() == 'true' + include_unnamed = request.args.get('include_unnamed', 'true').lower() == 'true' + + members = db_layer.get_all_members(include_named, include_unnamed) + result = [] + for m in members: + result.append({ + "member_id": m['member_id'], + "abstract_label": m['abstract_label'], + "real_name": m['real_name'], + "feature_description": m['feature_description'], + "first_seen_at": m['first_seen_at'].isoformat() if hasattr(m['first_seen_at'], 'isoformat') else str(m['first_seen_at']), + "named_at": m['named_at'].isoformat() if m.get('named_at') and hasattr(m['named_at'], 'isoformat') else None, + "named_by": m.get('named_by'), + "is_active": m['is_active'] + }) + return jsonify({"members": result}), 200 diff --git a/fam-core/src/fam_core/scheduler/__init__.py b/fam-core/src/fam_core/scheduler/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/scheduler/scheduler.py b/fam-core/src/fam_core/scheduler/scheduler.py new file mode 100644 index 0000000..6c91499 --- /dev/null +++ b/fam-core/src/fam_core/scheduler/scheduler.py @@ -0,0 +1,111 @@ +""" +Task-Scheduler - 60s 轮询视频目录,创建 PENDING 任务 + +判定视频完成(is_complete): +- 修改时间 > 60s(文件已停止写入) +- 文件大小稳定(连续两次检查大小一致) +""" +import os +import time +import threading +from datetime import datetime + +from ..logger import setup_logger, log_task +from ..config_loader import load_config +from .. import db_layer + +logger = setup_logger('fam-core.scheduler') + + +class TaskScheduler: + """视频目录扫描器,60s 轮询""" + + def __init__(self): + cfg = load_config() + self.scan_interval = cfg.get('scheduler', {}).get('scan_interval', 60) + self.video_dir = cfg.get('scheduler', {}).get('video_dir', '/volume1/surveillance') + self.media_base_url = cfg.get('video_server', {}).get('base_url', 'http://127.0.0.1:8000/media') + self.video_extensions = cfg.get('scheduler', {}).get('video_extensions', ['.mp4', '.mkv', '.avi']) + self.file_stable_seconds = cfg.get('scheduler', {}).get('file_stable_seconds', 60) + self.camera_name = cfg.get('scheduler', {}).get('camera_name', '默认摄像头') + + # 文件大小缓存,用于判断文件是否稳定 + self._file_sizes: dict = {} # path -> size + self._running = False + self._thread = None + + def _is_video(self, filename): + return any(filename.lower().endswith(ext) for ext in self.video_extensions) + + def _is_complete(self, filepath): + """判断视频是否已停止写入""" + try: + stat = os.stat(filepath) + now = time.time() + # 修改时间距当前 > stable_seconds + if now - stat.st_mtime < self.file_stable_seconds: + return False + # 文件大小稳定(与上次检查一致) + prev_size = self._file_sizes.get(filepath) + if prev_size is not None and prev_size == stat.st_size: + return True + self._file_sizes[filepath] = stat.st_size + return False + except OSError: + return False + + def _build_video_url(self, filepath): + """构建 Video-Server 下载 URL""" + filename = os.path.basename(filepath) + token = load_config().get('video_server', {}).get('token', '') + return f"{self.media_base_url}/{filename}?token={token}" + + def scan_once(self): + """执行一次扫描""" + if not os.path.isdir(self.video_dir): + logger.warning(f"视频目录不存在: {self.video_dir}") + return + + new_count = 0 + for root, dirs, files in os.walk(self.video_dir): + for filename in files: + if not self._is_video(filename): + continue + filepath = os.path.join(root, filename) + if not self._is_complete(filepath): + continue + # 检查是否已有任务 + if db_layer.get_video_url_exists(filepath): + continue + # 创建新任务 + video_url = self._build_video_url(filepath) + task_id = db_layer.create_task(filepath, video_url) + new_count += 1 + log_task(logger, task_id, 'scheduler', f'新任务: {filename}') + + if new_count > 0: + logger.info(f"本次扫描发现 {new_count} 个新视频") + + def _run(self): + """线程主循环""" + logger.info(f"Task-Scheduler 启动,扫描间隔 {self.scan_interval}s,目录: {self.video_dir}") + while self._running: + try: + self.scan_once() + except Exception as e: + logger.error(f"扫描异常: {e}", exc_info=True) + time.sleep(self.scan_interval) + + def start(self): + """启动调度线程""" + if self._running: + return + self._running = True + self._thread = threading.Thread(target=self._run, daemon=True, name='task-scheduler') + self._thread.start() + + def stop(self): + """停止调度线程""" + self._running = False + if self._thread: + self._thread.join(timeout=5) diff --git a/fam-core/src/fam_core/video_server/__init__.py b/fam-core/src/fam_core/video_server/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/video_server/video_server.py b/fam-core/src/fam_core/video_server/video_server.py new file mode 100644 index 0000000..fc41058 --- /dev/null +++ b/fam-core/src/fam_core/video_server/video_server.py @@ -0,0 +1,52 @@ +""" +Video-Server - Flask 蓝图,提供 mp4 静态下载 + +路由带 ?token=xxx 鉴权 +无 token 或 token 错误返回 403 +文件不存在返回 404 +""" +import os +from flask import Blueprint, request, send_from_directory, jsonify + +from ..logger import setup_logger +from ..config_loader import load_config + +logger = setup_logger('fam-core.video_server') + +video_bp = Blueprint('video_server', __name__) + +_config = None + + +def _get_config(): + global _config + if _config is None: + _config = load_config() + return _config + + +@video_bp.route('/media/', methods=['GET']) +def serve_video(filename): + """提供视频文件下载,带 token 鉴权""" + cfg = _get_config() + token = cfg.get('video_server', {}).get('token', '') + video_dir = cfg.get('video_server', {}).get('video_dir', '/volume1/surveillance') + + # 鉴权 + req_token = request.args.get('token', '') + if not token or req_token != token: + logger.warning(f"鉴权失败: {filename}, token={req_token}") + return jsonify({"error": "Forbidden"}), 403 + + # 检查文件 + filepath = os.path.join(video_dir, filename) + if not os.path.isfile(filepath): + logger.warning(f"文件不存在: {filepath}") + return jsonify({"error": "Not Found"}), 404 + + logger.info(f"提供视频: {filename}") + return send_from_directory( + os.path.dirname(filepath), + os.path.basename(filepath), + as_attachment=True + )