[阶段3] FAM-Core 瘦身为管理后台+Oracle-Sync 同步引擎,FAM-UI 改读同步镜像 - 删除 scheduler/dispatcher/poller/event_receiver/video_server,新增 oracle_sync 每30分钟拉增量写 sync_* 镜像表;member_manager/chat_handler 改走同步数据;UI 移除帧图改为事件时间线;DDL 新增 sync_videos/events/people/cursor

This commit is contained in:
ericwyuan
2026-08-21 10:38:26 +08:00
parent 9b1cc8f93b
commit 44026b65ea
22 changed files with 1359 additions and 2803 deletions

337
README.md
View File

@@ -13,13 +13,15 @@
### 1.1 首期范围(已基本完成)
- **FAM-Core**NAS 端单进程Task-Scheduler / Dispatcher / Event-Receiver / Chat-Handler / Member-Manager / Video-Server 六个子模块
- **FAM-Edge**Oracle 端单进程):接收视频上传 → FFmpeg 快速抽帧 → OpenCV 关键帧筛选 → 云端 VLM 视觉分析Gemini/NVIDIA NIM直出结构化 JSON → Edge 仅做格式化/校验 → 结果同步返回 全链路(**本地模型不参与视频分析**
- **FAM-UI**NAS 端Streamlit 直读 DB事件列表 + 成员命名页 + AI 对话页 + 对话历史
- **数据库六张表**process_tasks / monitor_events / event_details / chat_history / family_members / daily_summaries预留
- **任务状态机**PENDING → PROCESSING → SUCCESS/FAILED含退避重试与僵尸任务回收
- **AI 对话**:查 event_details 拼上下文 → 经 FAM-Edge 问答编排Gemini → NVIDIA → 本地 Ollama 兜底)生成回答 → 返回并写 chat_history
- **交互式成员命名**VLM 按特征提取"人物A/B/C"落库,用户命名后批量回溯更新历史,后续分析直接用真名
> **新架构 v22026-08-21 重构)**NAS 不再处理视频仅作管理后台视频分析全部上云Oracle
- **FAM-Core**NAS 端单进程):仅 **Oracle-Sync**(每 30 分钟拉增量镜像)+ **Chat-Handler** + **Member-Manager** 三个子模块CPU 占用极低
- **FAM-Edge**Oracle 端单进程rclone 实时同步 Google 硬盘视频 → 监听目录 → **整视频直传云端 VLM**Gemini 用 Files API / NVIDIA 用整视频 `video_url`,不切片不抽帧)→ 结构化 JSON 落本地 SQLite → 对外提供 `/api/oracle/sync` 增量拉取接口(**本地模型不参与视频分析**
- **FAM-UI**NAS 端Streamlit 读本地同步镜像sync_videos / sync_events / sync_people事件时间轴 + 人物管理 + AI 对话 + 对话历史 + 统计
- **数据库**Oracle 侧 SQLitevideos/events/people/sync_cursorNAS 侧 MariaDB 镜像sync_videos / sync_events / sync_people / sync_cursor+ chat_history
- **数据流向**Google 硬盘 ──rclone──► 甲骨文本地 ──整视频分析──► Oracle SQLite ──每 30 分钟 NAS 拉取──► NAS MariaDB 镜像 ──► FAM-UI
- **AI 对话**:查 sync_events 拼上下文 → 经 FAM-Edge 问答编排Gemini → NVIDIA → 本地 Ollama 兜底)生成回答 → 返回并写 chat_history
- **人物命名**:用户命名/合并某 label → 回推 Oracle `/api/oracle/people/correct`manual 优先)→ 下一周期同步回 NASOracle 独立 person_service 汇总全量人物 → LLM 合并为规范名 → 回灌视频提示
### 1.2 不在首期范围(推迟 v1.1+
@@ -70,59 +72,58 @@ Orchestrator 视觉阶段按 `fallback` 模式顺序降级Gemini → NVIDIA N
- Oracle 端 Ollama 端口 11434 不对外暴露,聊天请求经 FAM-Edge `/api/edge/chat` 代理转发
- Tailscale 两节点已安装在线,但 NAS tailscaled 为 userspace 模式且防火墙端口不通,暂走公网 IP
### 2.2 部署拓扑与数据流(推送模式
### 2.2 部署拓扑与数据流(新架构 v2Oracle 分析 + NAS 镜像
```
┌─────────────────────────────── NAS (192.168.50.64) ──────────────────────────────┐
Surveillance Station ──► /volume1/surveillance/Generic_ONVIF-001/
YYYYMMDDAM / YYYYMMDDPM 两级目录,~30min/350MB
│ │ │
│ ▼ │
│ FAM-Core (Flask :8000, gunicorn) │
│ ├─ Task-Scheduler: 60s 轮询视频目录,稳定文件建 PENDING 任务 │
│ ├─ Dispatcher: 30s 轮询multipart 上传视频 ──────────┐ │
│ │ (push_timeout=1800s同步等待响应) │ │
│ ├─ 僵尸回收: PROCESSING 超 push_timeout+120s 重置 PENDING │
│ ├─ Event-Receiver: /api/core/callback/event兼容保留
│ ├─ Chat-Handler ── /api/edge/chat 代理 ──────────────┐ │
│ ├─ Member-Manager / Video-Server(/media, token) │ │
│ ▼ │ │
│ MariaDB (sentinel_home_ai, 6 张表) │ │
│ │ │
│ FAM-UI (Streamlit :8501) 直读 DB │ │
└─────────────────────────────────────────────────────────┼─────────────────────────┘
│ HTTP (公网)
┌──────────── Google 硬盘 ────────────┐
oraclenas@...gserviceaccount.com
(Cloud Sync 落盘目录)
└──────────────┬───────────────────────┘
│ rclone 定时同步systemd timer
┌──────────────────────── Oracle Cloud (129.146.203.203) ──────────────────────────
│ FAM-Edge (Flask :5000, gunicorn --timeout 1800, 单 worker)
│ ├─ POST /api/edge/video/pushmultipart 视频,同步分析,结果随响应返回)
│ ├─ AI-Orchestrator: 健康检查 → 抽帧 → 选帧 → 压缩 → 云端VLM视觉直出结构化JSON → 格式化校验(无本地融合)
├─ POST /api/edge/chat/ask问答编排: Gemini→NVIDIA→本地Ollama 兜底)
┌──────────────────────── Oracle Cloud (129.146.203.203) ────────────────────────┐
│ FAM-Edge (Flask :5000)
│ ├─ Watch-Processor: 30s 轮询 /opt/fam-edge/gdrive_videos新视频串行处理
│ ├─ Video-Processor: 整视频直传云端 VLM不切片不抽帧
│ Gemini(Files API) → 失败 NVIDIA(整视频 video_url) → 再失败 FAILED
│ ├─ Person-Service: 汇总人物 → LLM 合并规范名 → 回灌视频提示 │
│ ├─ OracleDB (SQLite): videos / events / people / sync_cursor │
│ └─ API: /api/oracle/sync (增量拉取) · /api/oracle/people/correct (命名校正) · │
│ /api/edge/chat/ask (问答编排 Gemini→NVIDIA→Ollama) │
└───────────────────────────────┬───────────────────────────────────────────────┘
│ HTTP GET /api/oracle/sync?since=&token= (每 30 分钟)
┌─────────────────────────────── NAS (192.168.50.64) ────────────────────────────┐
│ FAM-Core (Flask :8000) │
│ ├─ Oracle-Sync: 唯一后台线程,拉增量写 MariaDB 镜像 + 维护 sync_cursor │
│ ├─ Chat-Handler: /api/chat/ask查 sync_events 拼上下文 → 走 Oracle 编排) │
│ ├─ Member-Manager: /api/member/name|merge回推 Oracle + 即时拉回) │
│ ▼ │
云端: Gemini / NVIDIA NIM (视觉直出结构化 + 问答) 本地: Ollama :11434 (qwen2.5:7b, 仅问答兜底)
└───────────────────────────────────────────────────────────────────────────────────┘
MariaDB (sentinel_home_ai): sync_videos / sync_events / sync_people /
│ sync_cursor / chat_history │
│ FAM-UI (Streamlit :8501) 读同步镜像 │
└─────────────────────────────────────────────────────────────────────────────────┘
```
### 2.3 主链路时序(推送模式)
**网络要点(新架构)**
- NAS → Oracle 仅一条出站 HTTPS/HTTP`GET /api/oracle/sync`(拉取)与 `POST /api/oracle/people/correct`(命名回推),均走 Oracle 公网 IP:5000token 鉴权
- Oracle Ollama :11434 不对外暴露,问答经 FAM-Edge `/api/edge/chat/ask` 代理
- Tailscale 两节点在线但 NAS 无法反向访问 Oracle故全部走 NAS 主动出站拉取模式
1. Scheduler 扫描到新视频(修改时间 > 60s 且大小稳定)→ 写 `process_tasks`PENDING
2. Dispatcher 领取 PENDING 任务 → 状态置 PROCESSING → 读本地视频文件multipart POST 到 Edge `/api/edge/video/push`payload 含 task_id / camera_name / event_start_time文件 mtime/ known_members_context
3. Edge 同步执行:
- a. 保存上传视频到临时目录(超时 60s
- b. FFmpeg 快速 seek`-ss <ts> -frames:v 1`)粗抽候选帧,帧数随视频时长自适应
- c. OpenCV MSE 帧差分析筛选关键帧 → 压缩(长边 ≤ 1024pxJPEG 质量 80
- d. 云端视觉模型按 `orchestrator.mode`fallback顺序降级Geminitimeout 30s→ NVIDIA NIMtimeout 20s首个**直出结构化 JSON** 成功的模型即采用,两云端全失败 → 任务 FAILED 走重试(绝不回退本地 Ollama本地模型不参与视频分析
- e. `format_cloud_result` 格式化校验(无模型调用):字段归一化、补 `source_providers=[provider]` / `compute_provider=[provider]`、缺失 `entities_json``frame_details` 推导、缺失 `global_summary` 时事实拼接 → 合法入库 schema
- f. `event_end_time` = event_start_time + 视频时长Edge 推算)
- g. `finally` 清理临时文件
4. Edge 把结果 JSON 直接作为 HTTP 响应返回(无 webhook
5. Dispatcher 收到响应 → 调用 `apply_success_event()``monitor_events`1 条聚合)+ `event_details`(每关键帧 1 条)+ upsert 未命名成员 → 任务置 SUCCESS失败则退避重试`min(60×(retry+1)×2, 600)`s超 3 次 FAILED
### 2.3 主链路时序(新架构 v2
1. Google 硬盘新视频 → rclone 定时同步到 Oracle `/opt/fam-edge/gdrive_videos`
2. Watch-Processor 轮询发现新文件 → 登记到 Oracle `videos`pending
3. Video-Processor 串行处理:整视频上传 Gemini Files API或 NVIDIA 整视频 `video_url`)→ 模型直出 `{global_summary, events[], people_mentioned[]}` → 写 Oracle `videos` + `events` + `people`
4. Person-Service 每 30 分钟汇总全量人物 → LLM 合并为规范名 → 更新 `people.canonical_name` → 生成 `known_members_context` 回灌后续视频提示
5. NAS Oracle-Sync 每 30 分钟 `GET /api/oracle/sync?since=<cursor>` → upsert 到本地 `sync_*` 镜像表 → 推进 `sync_cursor`
6. FAM-UI 读本地镜像展示;用户命名 → `POST /api/oracle/people/correct` 回推 Oracle下一周期同步生效
**容错设计**
- Dispatcher 僵尸回收PROCESSING 状态超过 `push_timeout + 120s` 自动重置 PENDING应对进程重启/Edge 重启导致 in-flight 请求丢失)
- fam-core 日志双写stdout + `fam-core/logs/fam-core.log`daemon 模式下 stdout 不可见)
- 所有日志带 `task_id` 作为 trace_id各阶段耗时打 INFO
- Oracle 单视频串行(`max_concurrent=1`)避免多视频抢占云端配额
- 视频分析失败(两云端均不可用)标记 `failed`,下一周期 cursor 仍包含它会被重试
- NAS 同步失败仅记日志下一周期30 分钟)自动重试,不阻塞 UI
- fam-core 日志双写stdout + `fam-core/logs/fam-core.log`
---
@@ -130,40 +131,42 @@ Orchestrator 视觉阶段按 `fallback` 模式顺序降级Gemini → NVIDIA N
### 3.1 FAM-CoreNAS 端)
> NAS 不再处理视频,仅作管理后台。唯一常驻后台线程是 Oracle-Sync。
| 模块 | 文件 | 职责 |
|------|------|------|
| Task-Scheduler | `scheduler/scheduler.py` | 60s 轮询视频目录(`os.walk` 递归,支持 AM/PM 子目录),`video_path` 去重,稳定文件建 PENDING 任务 |
| Dispatcher | `dispatcher/dispatcher.py` | 30s 轮询 PENDINGmultipart 上传视频至 Edge push 端点;收响应后经 `apply_success_event` 落库;僵尸 PROCESSING 回收;退避重试 |
| Event-Receiver | `event_receiver/event_receiver.py` | `/api/core/callback/event`webhook 兼容保留);核心逻辑抽为 `apply_success_event(data)` 供 Dispatcher 推送模式复用;未命名 abstract_label 自动 upsert `family_members` |
| Chat-Handler | `chat_handler/chat_handler.py` | `/api/chat/ask` 查 event_details 拼上下文 → 经 Edge `/api/edge/chat/ask` 问答编排Gemini→NVIDIA→本地 Ollama 兜底)→ 写 chat_history明细 > 50 条按小时聚合 |
| Member-Manager | `member_manager/member_manager.py` | `/api/member/unnamed` / `/api/member/name` / `/api/member/list`;命名后批量回溯 UPDATE event_detailsMariaDB 不支持 `$[*]` JSON 路径Python 层逐行更新) |
| Video-Server | `video_server/video_server.py` | `/media/<path>?token=xxx` 静态视频服务(推送模式下主链路不再使用,保留备用) |
| 公共层 | `db_layer.py` / `config_loader.py` / `logger.py` | PyMySQL 连接unix_socketdatetime 空串归一化 NULL + NOT NULL 列兜底;文件日志 |
| Oracle-Sync | `oracle_sync/oracle_sync.py` | 唯一后台线程:每 30 分钟 `GET /api/oracle/sync?since=<cursor>&token=` 拉增量 → upsert 到 `sync_videos`/`sync_events`/`sync_people` → 推进 `sync_cursor``push_name_correct()` 回推命名校正;`trigger_now()` 立即同步 |
| Chat-Handler | `chat_handler/chat_handler.py` | `/api/chat/ask``sync_events` 拼上下文 → 经 Oracle `/api/edge/chat/ask` 问答编排Gemini→NVIDIA→本地 Ollama 兜底)→ 写 chat_history |
| Member-Manager | `member_manager/member_manager.py` | `/api/member/unnamed` / `/api/member/list` / `/api/member/name` / `/api/member/merge`;命名/合并回推 Oracle 并即时拉回本地镜像 |
| 公共层 | `db_layer.py` / `config_loader.py` / `logger.py` | PyMySQL 连接unix_socket同步镜像 CRUD文件日志 |
> 已删除Task-Scheduler / Dispatcher / Poller / Event-Receiver / Video-Server(视频上传、切片、抽帧、关键帧落盘等职责全部迁移至 Oracle 端NAS CPU 占用大幅降低)。
### 3.2 FAM-EdgeOracle 端)
> 整视频分析,不切片、不抽帧、不依赖 OpenCV 人脸。
| 模块 | 文件 | 职责 |
|------|------|------|
| API-Gateway | `api_gateway/api_gateway.py` | `POST /api/edge/video/push`multipart 上传 + 同步分析 + 结果返回);`POST /api/edge/video/analyze`(旧拉取模式,兼容保留`POST /api/edge/chat`Ollama 代理`GET /health`;单并发控制(处理中返回 429 |
| Video-Preprocessor | `video_preprocessor/preprocessor.py` | `save_upload` 保存上传视频FFmpeg 快速 seek 粗抽候选帧帧数自适应OpenCV MSE 帧差筛选关键帧(首末帧必选);压缩;`compute_timestamps` 用 start+偏移算绝对时间戳;`video_duration` 供 event_end_time 推算 |
| AI-Orchestrator | `ai_orchestrator/orchestrator.py` | 模型健康检查 → 云端视觉适配器按 `orchestrator.mode`fallback 顺序降级)调度,**直出结构化 JSON** → `format_cloud_result` 格式化校验(无本地融合)→ JSON schema 校验;`run_qa` 实现问答编排Gemini→NVIDIA→本地 Ollama 兜底);`process_push_task` 为推送模式入口(不触发 webhook记录各模型实际执行耗时与成功状态到 `compute_provider` 数组 |
| Model-Adapters | `model_adapters/` | `BaseModelAdapter` 抽象基类(`__init__` / `health_check` / `analyze_frames` / `chat` / `get_timeout` / 熔断器实例);`build_adapter` 工厂函数按 `provider` 字段分发实例化视觉适配器gemini/nvidia直出结构化 JSON文本适配器ollama仅智能问答兜底 |
| └ OllamaAdapter | `model_adapters/ollama_adapter.py` | requests 直调本地 REST `/api/generate``num_predict` 可配;**role: text, usage: qa_fallback**(仅智能问答兜底,不参与视觉分析、不参与融合) |
| └ GeminiAdapter | `model_adapters/gemini_adapter.py` | requests 直调 Google REST `:generateContent`**多图单请求直出结构化 JSON****role: vision**`chat()` 参与问答 |
| └ NvidiaVisionAdapter | `model_adapters/nvidia_adapter.py` | **基于 openai SDK**NIM 兼容 OpenAI API 规范),`base_url=https://integrate.api.nvidia.com/v1``api_key``${NVIDIA_API_KEY}` 展开;**逐帧返回结构化单帧 JSON 并聚合为 frame_details**NIM 限 1 图/请求);`health_check``client.models.list()`**role: vision**`chat()` 参与问答 |
| Storage-Cleaner | `storage_cleaner/` | `finally` 删除临时视频与帧图片 |
| API-Gateway | `api_gateway/api_gateway.py` | `GET /api/oracle/sync`增量拉取since+token 校验);`POST /api/oracle/people/correct`(命名校正`POST /api/edge/chat/ask`(问答编排`GET /health` |
| Watch-Processor | `watch_processor.py` | 30s 轮询 rclone 同步落地目录,登记新视频,串行触发 Video-Processor |
| Video-Processor | `video_processor.py` | 按 `vision_order` 调适配器 `analyze_video`(整视频);首个成功即落库 Oracle `videos`+`events`+`people`;全失败标 `failed` |
| Person-Service | `person_service.py` | 汇总全量人物 → LLM 合并为规范名 → `set_canonical`;生成 `known_members_context` 回灌视频提示manual 命名优先不被覆盖 |
| OracleDB | `oracle_db.py` | SQLitevideos / events / people / sync_cursor`get_sync_delta(since)` 增量导出 |
| Model-Adapters | `model_adapters/` | `BaseModelAdapter.analyze_video(video_path, known_members_context, event_start_time)`GeminiFiles API 整视频)/ NVIDIA整视频 `video_url``num_frames=128`/ Ollama纯文本不参与视频 |
| QA-Orchestrator | `qa.py` | 遍历所有适配器 `chat()`Gemini→NVIDIA→Ollama 三级降级(仅问答 |
### 3.3 FAM-UINAS 端)
Streamlit 应用(`fam-ui/src/app.py`),侧边栏切换页面:
Streamlit 应用(`fam-ui/src/app.py`),侧边栏切换页面(均读本地同步镜像)
| 页面 | 功能 |
|------|------|
| 📊 事件列表 | 按日期筛选 + 分页20 条/页)展示 monitor_events |
| 👤 成员命名 | 列出未命名人物 + 特征描述,输入真名后调 `/api/member/name` 批量回溯 |
| 🕒 事件时间轴 | 视频会话列表(按处理后时间倒序)+ 选中会话的事件时间线(时间点 + 描述 + 人物/关注徽章,无帧图) |
| 💬 AI 对话 | 输入框 + 调 `/api/chat/ask`;按 queried_person 预设快捷提问 |
| 📜 对话历史 | chat_history 倒序展示 |
| 📈 统计图表 | compute_provider 占比bar_chart |
| 📝 对话历史 | chat_history 倒序展示 |
| 👤 人物管理 | 按规范名/标签聚合,命名/合并(回推 Oracle不再展示帧照片 |
| 📈 统计图表 | 模型来源占比 / 关注事件 / 同步状态 |
---
@@ -173,36 +176,46 @@ Streamlit 应用(`fam-ui/src/app.py`),侧边栏切换页面:
### 4.1 表清单
> 新架构 v2Oracle 侧用 SQLite`videos`/`events`/`people`/`sync_cursor`NAS 侧 MariaDB 仅保留 **同步镜像表 + 问答历史**。`process_tasks`/`monitor_events`/`event_details`/`family_members` 等旧表已不再写入(保留历史数据,未删除)。
**OracleSQLite`oracle_db.py`**
| 表 | 用途 | 关键字段 |
|----|------|---------|
| `process_tasks` | 视频处理任务 | task_id, video_path, video_url, status(PENDING/PROCESSING/SUCCESS/FAILED), retry_count, max_retries, next_retry_at, error_message, failure_stage(ENUM) |
| `monitor_events` | 事件聚合(每任务 1 条) | event_id, task_id, event_start_time, event_end_time, camera_name, global_summary, entities_json(JSON), compute_provider(JSON 数组) |
| `event_details` | 每关键帧一条明细 | detail_id, event_id, task_id, frame_index, frame_timestamp, person, action, clothing, is_attention_event, source_providers(JSON) |
| `family_members` | 交互式命名 | member_id, abstract_label(如"人物A"), real_name(NULL=未命名), feature_description, first_seen_at, named_at, named_by |
| `videos` | 视频会话(每视频 1 行) | id, filename(UNIQUE), camera_name, status, summary_json, events_json, people_json, compute_provider, event_start_time, updated_at |
| `events` | 视频内时间点事件 | id, video_id, ts, description, person_list_json, is_attention_event |
| `people` | 规范人物Oracle 维护) | id, label(UNIQUE), canonical_name, appearances, source(llm/manual) |
| `sync_cursor` | 同步游标 | key, value上次 server_time |
**NASMariaDB同步镜像`scripts/ddl.sql`**
| 表 | 用途 | 关键字段 |
|----|------|---------|
| `sync_videos` | 视频会话镜像(对齐 Oracle videos | id, filename, camera_name, status, summary_json, events_json, people_json, compute_provider, processed_at |
| `sync_events` | 事件镜像(对齐 Oracle events | id, video_id, ts, description, person_list_json, is_attention_event |
| `sync_people` | 人物镜像(对齐 Oracle people | id, label, canonical_name, appearances, source |
| `sync_cursor` | 同步游标 | key='last_since', value |
| `chat_history` | AI 问答记录 | chat_id, user_question, ai_answer, context_summary, queried_date, queried_person |
| `daily_summaries` | 每日摘要(预留) | target_date, summary_text |
### 4.2 表关系与命名回溯
### 4.2 表关系
```
process_tasks (1) ─── (N) monitor_events (1) ─── (N) event_details
family_members 独立表:
- event_details.person 存 abstract_label未命名或 real_name命名后
- 命名后: UPDATE event_details SET person = real_name WHERE person = abstract_label
- monitor_events.entities_json 由 Python 层解析逐行更新MariaDB 不支持 $[*] 路径)
chat_history 独立表
Oracle: videos (1) ─── (N) events people 独立label/canonical_name
NAS 镜像: sync_videos (1) ─── (N) sync_events sync_people 独立
chat_history 独立表(问答上下文摘要留存
```
### 4.3 compute_provider / source_providers
命名回溯:用户命名某 `label``POST /api/oracle/people/correct``canonical_name`manual 优先)→ 下一周期同步回 NAS `sync_people`Oracle `person_service` 用规范名回灌视频提示,后续事件 `person_list_json` 直接带真名。
- `monitor_events.compute_provider`JSON 数组,记录本次任务实际成功调用(**视觉分析**)的云端模型,如 `["gemini"]``["nvidia"]`;本地 Ollama 不参与视频分析,不会出现在该字段
- `event_details.source_providers`:该条明细被哪些模型识别到(可能少于 compute_provider
- 多模型交叉验证:多模型一致 → 可信度高;仅单一模型描述 → source_providers 仅含该模型;冲突 → 多数派为准
### 4.3 compute_provider
- `sync_videos.compute_provider`:字符串,记录该视频实际成功调用的视觉模型(`gemini` / `nvidia`);本地 Ollama 不参与视频分析,不会出现在该字段
- 问答链路Gemini→NVIDIA→Ollama 兜底)的 provider 体现在 `/api/edge/chat/ask` 响应的 `provider` 字段
### 4.4 兼容性注意
- MariaDB 10.11 严格模式:**空字符串不能插 DATETIME 列**1292 错误)。`db_layer._dt_or_none` 将空串归一化 NULL`event_end_time` NOT NULL 列按 end→start→NOW 兜底;`frame_timestamp` 空值兜底 NOW
- MariaDB 不支持 MySQL 的 `$[*]` JSON 通配路径`->` 操作符JSON 字段在 Python 层处理
- MariaDB 10.11 严格模式:**空字符串不能插 DATETIME 列**1292 错误)。同步表时间字段统一用 `VARCHAR(32)` 文本存储 Oracle 的 ISO 字符串,规避类型转换问题
- MariaDB 不支持 MySQL 的 `$[*]` JSON 通配路径,人物统计按 `person_list_json LIKE '%name%'` 字符串匹配在 Python 层完成
---
@@ -210,59 +223,70 @@ chat_history 独立表
### 5.1 FAM-EdgeOracle :5000
**POST /api/edge/video/push**(主链路,推送模式
**GET /api/oracle/sync**NAS 每 30 分钟拉增量,新架构主接口
- 请求:`multipart/form-data`,字段 `video`(文件) / `task_id` / `camera_name` / `event_start_time` / `known_members_context`
- 处理同步执行完整分析流水线可能耗时数分钟gunicorn timeout 1800
- 请求:`?since=<ISO 文本>&token=<ORACLE_SYNC_TOKEN>``since` 为空拉全量)
- 响应200
```json
{
"task_id": 289,
"status": "success",
"event_start_time": "2026-08-20 01:06:44",
"event_end_time": "2026-08-20 01:07:13",
"camera_name": "客厅",
"global_summary": "...",
"entities_json": [{"person": "汤圆", "action": "...", "clothing": "..."}],
"frame_details": [
{"frame_index": 1, "frame_timestamp": "...", "person": "...", "action": "...",
"clothing": "...", "is_attention_event": false, "source_providers": ["gemini"]}
"videos": [
{"id": 1, "filename": "2026-08-21_081500.mp4", "camera_name": "客厅",
"status": "done", "summary_json": "...", "events_json": "[...]",
"people_json": "[...]", "compute_provider": "gemini",
"event_start_time": "2026-08-21 08:15:00", "updated_at": "2026-08-21 08:40:12"}
],
"compute_provider": ["gemini"]
"events": [
{"id": 10, "video_id": 1, "ts": "00:01:23", "description": "汤圆在客厅玩耍",
"person_list_json": "[\"汤圆\"]", "is_attention_event": 0, "updated_at": "2026-08-21 08:40:12"}
],
"people": [
{"id": 1, "label": "人物A", "canonical_name": "汤圆", "source": "manual",
"appearances": 12, "updated_at": "2026-08-21 08:41:00"}
],
"server_time": "2026-08-21 08:41:30"
}
```
- 失败:`{"task_id": ..., "status": "failed", "failure_stage": "vlm_visual", "error_message": "..."}`
- 429已有任务处理中单并发503全部模型不健康
- 401token 校验失败
**POST /api/edge/video/analyze** — 旧拉取模式Edge 拉 video_url + webhook 回调),兼容保留,主链路不再使用
**POST /api/edge/chat/ask** — 智能问答编排FAM-Core Chat-Handler 调用):请求 `{"prompt"}` → 响应 `{"answer","provider"}`;内部按 Gemini → NVIDIA → 本地 Ollama 顺序,仅两云端都失败才用本地兜底
**POST /api/edge/chat** — Ollama 直连代理(兼容旧调用,保留
**GET /health** — 服务与模型健康状态(任务处理中可能无响应,单 worker 忙)
**POST /api/oracle/people/correct** — 命名校正回推:`{"label":"人物A","canonical_name":"汤圆","token":...}`manual 优先,不被 LLM 覆盖)→ `{"status":"ok"}`
**POST /api/edge/chat/ask** — 智能问答编排FAM-Core Chat-Handler 调用):请求 `{"prompt","max_tokens"}` → 响应 `{"answer","provider"}`;内部按 Gemini → NVIDIA → 本地 Ollama 顺序,仅两云端都失败才用本地兜底
**GET /health** — 服务状态(含已处理视频数
### 5.2 FAM-CoreNAS :8000
| 端点 | 方法 | 说明 |
|------|------|------|
| `/health` | GET | 服务健康 |
| `/api/status` | GET | scheduler/dispatcher 运行状态 |
| `/api/core/callback/event` | POST | Edge 回调webhook 兼容保留);推送模式下由 Dispatcher 内部调用 `apply_success_event` |
| `/api/chat/ask` | POST | 用户问答:`{"question","queried_person","queried_date"}``{"answer","context_summary","chat_id"}` |
| `/api/status` | GET | Oracle-Sync 同步状态running / last_sync_at / last_error / cursor / last_count |
| `/api/chat/ask` | POST | 用户问答:`{"question","queried_person","queried_date"}``{"answer","context_summary","chat_id"}`(上下文来自 sync_events |
| `/api/chat/history` | GET | 对话历史(`?date=``?person=&limit=` |
| `/api/member/unnamed` | GET | 未命名人物列表(含特征描述、出现次数 |
| `/api/member/name` | POST | 命名:`{"abstract_label","real_name","named_by"}` → 批量回溯 event_details/entities_json返回更新条数 |
| `/api/member/list` | GET | 全部成员 |
| `/media/<path>?token=xxx` | GET | 视频静态服务token 鉴权,推送模式下备用 |
| `/api/member/unnamed` | GET | 未命名人物列表(label / 出现次数 / 首见时间 |
| `/api/member/list` | GET | 全部人物label + canonical_name + 是否命名) |
| `/api/member/name` | POST | 命名:`{"label","canonical_name"}` → 回推 Oracle 并即时拉回本地镜像 |
| `/api/member/merge` | POST | 合并:`{"source","target"}` → 将 source 并入 target 身份(统一 canonical_name |
> 已删除:`/api/core/callback/event`、`/media/<path>`(视频处理职责已迁移至 Oracle
### 5.3 云端结构化输出 JSON Schema
云端 VLM 直接产出结构化 JSON,经两道处理入库
整视频直传云端 VLM,模型直接产出结构化 JSON`analyze_video` 返回)
1. **适配器内三层容错解析**`json_parser.parse_vlm_json`):直接 `json.loads` → 提取 markdown fence ` ```json ... ``` ` → 贪婪匹配最大 `{...}`;失败抛 `VLMOutputInvalidError`
2. **`format_cloud_result` 归一化/校验**(无模型调用):`frame_details` 必须为非空列表并做字段类型归一化;`source_providers` 缺失时补为 `[provider]``compute_provider` 置为本次成功 provider`entities_json` 缺失时由 `frame_details` 按人物去重推导;`global_summary` 缺失时格式化拼接生成
```json
{
"global_summary": "客厅监控摘要……",
"events": [
{"timestamp": "00:01:23", "description": "汤圆在客厅玩耍",
"people": ["汤圆"], "is_attention_event": false}
],
"people_mentioned": ["汤圆"]
}
```
任一步骤失败 → 任务 FAILED 走重试。`action` 由 AI 自由生成无枚举过滤,`is_attention_event` 由 AI 自行判断。
- Gemini 用 Files API 上传整视频后 `generateContent`NVIDIA 用整视频 `video_url` + `num_frames=128`(模型内部自行采样帧),均不切片、不抽帧、不依赖 OpenCV
- `person_service` 汇总全量 `people_mentioned` → LLM 合并为规范名 → 生成 `known_members_context` 回灌后续视频提示,使模型用真名指代
- `action` / 描述由 AI 自由生成无枚举过滤,`is_attention_event` 由 AI 自行判断
---
@@ -436,51 +460,62 @@ task_id=28930s 测试片段)全链路打通:推送 5.7MB → Edge 分析
### 8.3 配置文件要点
**fam-core/config/config.yaml**NAS生产值
**fam-core/config/config.yaml**NAS新架构 v2 —— 仅同步 + 问答
```yaml
scheduler:
video_dir: "/volume1/surveillance/Generic_ONVIF-001" # 生产目录YYYYMMDDAM/PM 两级子目录)
# 285 个历史视频由占位 FAILED 任务占用路径scheduler dedup 自动跳过forward-only 模式
dispatcher:
edge_url: "http://129.146.203.203:5000/api/edge/video/push"
push_timeout: 1800
server:
port: 8000
database: # MariaDBunix_socket 优先
unix_socket: "/run/mysqld/mysqld10.sock"
oracle_sync: # 唯一后台线程配置
base_url: "http://129.146.203.203:5000"
token: "${ORACLE_SYNC_TOKEN}" # 与 Oracle 端 sync_api.token 一致
interval_sec: 1800 # 每 30 分钟拉一次增量
timeout: 120
chat_handler:
qa_url: "http://129.146.203.203:5000/api/edge/chat/ask" # 问答统一走 Edge 编排Gemini→NVIDIA→Ollama
qa_url: "http://129.146.203.203:5000/api/edge/chat/ask" # 问答统一走 Oracle 编排
timeout: 120
```
**fam-edge/config/config.yaml**Oracle多模型池配置,新架构
**fam-edge/config/config.yaml**Oracle整视频分析 + 同步 + 人物服务
```yaml
# 编排调度模式: fallback(顺序降级, 默认) | ensemble(并行交叉验证)
orchestrator:
mode: "fallback"
overall_timeout: 600
# 多模型池配置(新框架:本地大模型不参与视频分析,仅智能问答兜底)
# 视频分析链路: 云端 VLM 直出结构化 JSON → Edge format_cloud_result 格式化/校验 → 直存 NAS DB无本地融合
# 智能问答链路: Gemini → NVIDIA → 本地 Ollama仅两云端都失败才启用本地兜底
server:
port: 5000
gdrive_sync: # rclone 同步落地目录监听
enabled: true
local_dir: "/opt/fam-edge/gdrive_videos"
watch_interval_sec: 30
camera_name: "客厅"
parse_start_from_filename: true
oracle_db:
path: "/opt/fam-edge/data/oracle.db"
sync_api:
token: "${ORACLE_SYNC_TOKEN}" # NAS 拉取鉴权(与 NAS oracle_sync.token 一致)
person_service:
schedule_interval_sec: 1800 # 每 30 分钟重新汇总人物
model: "gemini"
video_processing:
max_concurrent: 1 # 单视频串行,避免抢占云端配额
timeout: 900
vision_order: ["gemini", "nvidia"]
models:
# 1. Google Gemini视觉主 + 参与问答)
- provider: "gemini"
role: "vision" # 视觉分析 + 问答vision role 也参与 chat
enabled: true
model_name: "gemini-flash-latest" # v1beta 下 gemini-1.5-flash 会 404
role: "vision"
model_name: "gemini-flash-latest"
api_key: "${GEMINI_API_KEY}"
timeout: 30
circuit_breaker:
enabled: true
threshold: 3
cooldown: 600
# 2. NVIDIA NIM 托管 API视觉备 + 参与问答)
timeout: 600
- provider: "nvidia"
role: "vision" # 视觉分析 + 问答vision role 也参与 chat
enabled: true
model_name: "meta/llama-3.2-11b-vision-instruct" # 或 qwen/qwen2-vl-72b-instruct
role: "vision"
model_name: "nvidia/nemotron-nano-12b-v2-vl" # 整视频 video_url 输入(内部采样帧)
base_url: "https://integrate.api.nvidia.com/v1"
api_key: "${NVIDIA_API_KEY}"
timeout: 600
- provider: "ollama"
role: "text"
usage: "qa_fallback" # 仅智能问答兜底,不参与视频
model_name: "qwen2.5:7b"
``` api_key: "${NVIDIA_API_KEY}"
timeout: 20
circuit_breaker:
enabled: true

View File

@@ -1,10 +1,11 @@
# FAM-Core 配置文件 (NAS 端) - 实际部署配置
# 注: Tailscale 防火墙待修复,当前 edge_url 使用 Oracle 公网 IP
# 异步队列模式 + 分块断点续传:
# 小文件(<=50MB): 直接上传 /enqueue
# 大文件(>50MB): 分块(20MB/块)上传 /chunk → /assemble 合并入队
# → NAS Poller 定期从 /api/edge/results 拉取结果写库
# Ollama 未对外暴露chat_handler 通过 FAM-Edge 代理
# FAM-Core 配置文件 (NAS 端) - 新架构 v22026-08-21
#
# NAS 仅作管理后台,不再处理视频。唯一后台线程 Oracle-Sync 每 30 分钟
# 从甲骨文 FAM-Edge 拉取增量镜像到本地 MariaDBsync_videos/events/people
# 所有视频分析在 Oracle 完成。
#
# Tailscale 当前无法从 NAS 反向访问 Oracle故 base_url 用 Oracle 公网 IP。
# token 与 Oracle 端 sync_api.token 一致,均取自环境变量 ORACLE_SYNC_TOKEN。
server:
host: "0.0.0.0"
@@ -18,36 +19,16 @@ database:
database: "sentinel_home_ai"
unix_socket: "/run/mysqld/mysqld10.sock"
scheduler:
scan_interval: 60
# 正式目录 /volume1/surveillance/Generic_ONVIF-001
video_dir: "/volume1/surveillance/Generic_ONVIF-001"
video_extensions: [".mp4", ".mkv", ".avi"]
file_stable_seconds: 60
camera_name: "客厅"
dispatcher:
poll_interval: 30
edge_url: "http://129.146.203.203:5000/api/edge/video/enqueue"
max_retries: 5 # 文件级重试次数分块级重试另计每块3次
stale_timeout: 600 # PROCESSING 超时回收10分钟
poller:
poll_interval: 30
results_url: "http://129.146.203.203:5000/api/edge/results"
batch_size: 10
timeout: 30
video_server:
base_url: "http://127.0.0.1:8000/media"
token: "sentinel-media-2026"
video_dir: "/volume1/surveillance"
# 甲骨文同步(每 30 分钟拉增量镜像)
oracle_sync:
# FAM-Edge 对外同步接口地址(端口同其 server.port=5000
base_url: "http://129.146.203.203:5000"
# 与 Oracle 端 sync_api.token 一致(环境变量注入,避免明文入库)
token: "${ORACLE_SYNC_TOKEN}"
interval_sec: 1800 # 拉取间隔(秒),默认 30 分钟
timeout: 120 # 单次拉取超时(秒)
chat_handler:
# 智能问答统一走 FAM-Edge 编排端点Gemini → NVIDIA → 本地 Ollama 兜底)
qa_url: "http://129.146.203.203:5000/api/edge/chat/ask"
timeout: 120
storage:
# 关键帧落盘目录event_receiver 写入fam-ui 读取展示时间轴)
frame_image_dir: "/volume1/web/sentinel-home-ai/fam-ui/static/frames"

View File

@@ -1,12 +1,14 @@
"""
FAM-Core 主应用 - Flask 单进程
FAM-Core 主应用 - Flask 单进程(新架构 v2
承载: Task-Scheduler / Dispatcher / Poller / Event-Receiver / Chat-Handler / Member-Manager / Video-Server
承载: Oracle-Sync每 30 分钟拉取增量镜像)+ Member-Manager + Chat-Handler
NAS 不再处理视频:无 Scheduler / Dispatcher / Poller / Event-Receiver / Video-Server。
所有视频分析在 Oracle 完成NAS 仅作管理后台拉取展示CPU 占用大幅降低。
"""
import os
import sys
import time
import threading
from flask import Flask, jsonify
# 确保包路径
@@ -14,83 +16,49 @@ 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 .poller.poller import Poller
from .event_receiver.event_receiver import event_bp
from .oracle_sync import get_sync
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
_poller = None
# 初始化后台同步线程NAS 唯一常驻线程)
_sync = None
try:
_scheduler = TaskScheduler()
_scheduler.start()
logger.info("Task-Scheduler 已启动")
except Exception as e:
logger.error(f"Task-Scheduler 启动失败: {e}")
_sync = get_sync()
_sync.start()
logger.info("Oracle-Sync 已启动")
# 启动后立刻拉一次,前端无需等待首个周期即有数据
try:
_dispatcher = Dispatcher()
_dispatcher.start()
logger.info("Dispatcher 已启动")
_sync.trigger_now()
logger.info("启动首次同步完成")
except Exception as e:
logger.error(f"Dispatcher 启动失败: {e}")
try:
_poller = Poller()
_poller.start()
logger.info("Poller 已启动")
logger.warning(f"启动首次同步失败(后续周期会重试): {e}")
except Exception as e:
logger.error(f"Poller 启动失败: {e}")
logger.error(f"Oracle-Sync 启动失败: {e}")
@app.route('/api/status', methods=['GET'])
def status():
"""系统状态(检查线程实际存活)"""
"""系统状态"""
return jsonify({
"scheduler_running": _scheduler.is_alive() if _scheduler else False,
"dispatcher_running": _dispatcher.is_alive() if _dispatcher else False,
"poller_running": _poller.is_alive() if _poller else False,
"service": "fam-core",
"sync": _sync.status() if _sync else {"running": False, "error": "未初始化"},
}), 200
def _watchdog_run():
"""看门狗:每 60s 检查线程存活,崩溃自动重启"""
logger.info("Watchdog 线程启动,检查间隔 60s")
while True:
time.sleep(60)
for comp, name in [(_scheduler, 'Scheduler'), (_dispatcher, 'Dispatcher'), (_poller, 'Poller')]:
if comp and hasattr(comp, 'check_and_restart'):
try:
comp.check_and_restart()
except Exception as e:
logger.error(f"Watchdog 重启 {name} 失败: {e}", exc_info=True)
_watchdog_thread = threading.Thread(target=_watchdog_run, daemon=True, name='watchdog')
_watchdog_thread.start()
if __name__ == '__main__':
cfg = load_config()
port = cfg.get('server', {}).get('port', 8000)

View File

@@ -1,17 +1,19 @@
"""
Chat-Handler - Flask 蓝图,接收用户问答
Chat-Handler - Flask 蓝图,接收用户问答(新架构 v2
处理逻辑:
1. 根据 queried_person queried_date 查询 event_details
2. 拼接上下文(每条明细一行
3. 若明细条数 > 50按小时聚合成摘要
4. POST Oracle Ollama 纯文本模式,调问答 Prompt
5. 插入 chat_history
6. 返回回答
逻辑:
1. queried_person + queried_date 从 sync_events 拉取相关事件
person_list_json 含该名且 ts 落在日期内
2. 拼接上下文文本(每事件一行:时间 + 摄像头 + 描述 + 人物 + 是否关注)
3. 调 Oracle FAM-Edge 问答编排端点Gemini → NVIDIA → 本地 Ollama 兜底)
4. 插入 chat_history
5. 返回回答
注意: 上下文来自 Oracle 已分析好的事件摘要,不做本地视频处理。
"""
import json
import requests
from flask import Blueprint, request, jsonify
from collections import defaultdict
from ..logger import setup_logger
from ..config_loader import load_config
@@ -21,9 +23,9 @@ logger = setup_logger('fam-core.chat_handler')
chat_bp = Blueprint('chat_handler', __name__)
CHAT_SYSTEM_PROMPT = """你是家庭监控助手。根据以下今日监控数据,回答用户问题。
CHAT_SYSTEM_PROMPT = """你是家庭监控助手。根据以下监控数据,回答用户问题。
今日数据(按时间顺序,每条一行):
监控数据(按时间顺序,每条一行):
{context}
已知家庭成员: {members}
@@ -33,41 +35,27 @@ CHAT_SYSTEM_PROMPT = """你是家庭监控助手。根据以下今日监控数
要求:
- 只基于上述数据回答,不要编造
- 按时间顺序总结
- 若有关注事件(跌倒、哭闹等),重点提示
- 若有关注事件(跌倒、哭闹、陌生人等),重点提示
- 若当天没有该人员的数据,明确说"今天没有观察到{person}"
- 用自然语言回答,不要输出 JSON
"""
def _format_details(details):
"""event_details 格式化为文本"""
def _format_events(rows):
"""sync_events 查询行格式化为上下文文本"""
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}")
for r in rows:
ts = (r.get('ts') or '')[:16] # 'YYYY-MM-DD HH:MM'
camera = r.get('camera_name') or ''
desc = r.get('description') or ''
try:
persons = json.loads(r['person_list_json']) if isinstance(r['person_list_json'], str) else (r['person_list_json'] or [])
except (ValueError, TypeError):
persons = []
persons = [str(p) for p in persons]
attention = ' [关注事件]' if r.get('is_attention_event') else ''
person_str = ','.join(persons) if persons else '无人'
lines.append(f"[{ts} {camera}] {person_str}: {desc}{attention}")
return '\n'.join(lines)
@@ -80,7 +68,6 @@ def _call_edge_qa(prompt: str) -> str:
timeout = cfg.get('chat_handler', {}).get('timeout', 120)
resp = requests.post(qa_url, json={"prompt": prompt}, timeout=timeout)
if resp.status_code == 200:
data = resp.json()
answer = data.get('answer', '')
@@ -108,38 +95,27 @@ def chat_ask():
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)
rows = db_layer.query_sync_events_for_person_date(queried_person, queried_date)
# 2. 拼接上下文
if len(details) == 0:
# 无数据
if len(rows) == 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))}"
context_summary = "查询 sync_events 0 条"
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()
context = _format_events(rows)
context_summary = f"查询 sync_events {len(rows)}"
members = db_layer.get_sync_known_members_context()
prompt = CHAT_SYSTEM_PROMPT.format(
context=context,
members=members or f"{queried_person}",
members=members or queried_person,
question=question,
person=queried_person
)
try:
answer = _call_edge_qa(prompt)
except Exception as e:
logger.error(f"问答编排调用失败: {e}")
return jsonify({"error": f"AI 调用失败: {e}"}), 503
# 3. 写入 chat_history
chat_id = db_layer.insert_chat_history(
user_question=question,
ai_answer=answer,
@@ -163,7 +139,6 @@ def chat_history():
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'):

View File

@@ -1,6 +1,17 @@
"""
数据库访问层 - MariaDB 连接管理与 CRUD 操作
使用 PyMySQL (纯 Python, ~45KB) 连接 MariaDB 服务器
数据库访问层 - MariaDB 连接管理与同步镜像 CRUD
新架构 (2026-08-21 重构):
NAS 不再处理视频,仅作为管理后台。
Oracle (FAM-Edge) 处理整视频分析后存 SQLiteNAS 每 30 分钟拉增量,
镜像到本地三张表:
sync_videos : 视频会话(全局摘要 + 事件 JSON + 人物 JSON
sync_events : 视频拆出的时间点事件(描述 + 涉及人物 + 是否关注)
sync_people : 规范人物表label + canonical_nameOracle 维护)
sync_cursor : 同步游标(上次成功拉取到的 server_time
本层只服务同步镜像 + 问答历史,旧 process_tasks/event_details/monitor_events/
family_members 相关逻辑已全部移除(视频处理职责已迁移至 Oracle
"""
import json
import pymysql
@@ -41,283 +52,339 @@ def get_conn():
# ============================================================
# process_tasks 操作
# 同步镜像sync_videos
# ============================================================
def create_task(video_path: str, video_url: str) -> int:
"""创建新任务"""
def upsert_sync_videos(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cursor = conn.cursor()
cursor.execute(
"INSERT INTO process_tasks (video_path, video_url, status) VALUES (%s, %s, 'PENDING')",
(video_path, video_url)
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_videos
(id, drive_file_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
drive_file_id=VALUES(drive_file_id),
filename=VALUES(filename),
camera_name=VALUES(camera_name),
duration_sec=VALUES(duration_sec),
event_start_time=VALUES(event_start_time),
status=VALUES(status),
summary_json=VALUES(summary_json),
events_json=VALUES(events_json),
people_json=VALUES(people_json),
compute_provider=VALUES(compute_provider),
created_at=VALUES(created_at),
updated_at=VALUES(updated_at),
processed_at=VALUES(processed_at),
synced_at=NOW()""",
(r.get('id'), r.get('drive_file_id'), r.get('filename'),
r.get('camera_name'), r.get('duration_sec') or 0,
r.get('event_start_time'), r.get('status'),
r.get('summary_json'), r.get('events_json'), r.get('people_json'),
r.get('compute_provider'), r.get('created_at'),
r.get('updated_at'), r.get('processed_at'))
)
n += 1
conn.commit()
task_id = cursor.lastrowid
logger.info(f"[task_id={task_id}] task created: {video_path}")
return task_id
return n
finally:
conn.close()
def get_pending_tasks(limit=10) -> List[Dict]:
"""获取待处理任务"""
conn = get_conn()
try:
cursor = conn.cursor(pymysql.cursors.DictCursor)
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_sync_videos(limit=15, offset=0, date_filter=None) -> List[Dict]:
"""获取视频会话列表(已完成优先),支持日期筛选与分页。
def get_tasks_by_status(status: str, limit=10) -> List[Dict]:
"""按状态获取任务"""
conn = get_conn()
try:
cursor = conn.cursor(pymysql.cursors.DictCursor)
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):
"""更新任务状态"""
valid_stages = {'download', 'extract', 'vlm_visual', 'vlm_fusion', 'callback', 'process'}
if failure_stage and failure_stage not in valid_stages:
failure_stage = 'callback'
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 reclaim_stale_processing(timeout_seconds: int) -> List[int]:
"""回收僵尸 PROCESSING 任务updated_at 早于 timeout_seconds 前的任务重置为 PENDING
场景Dispatcher 推送过程中进程重启/Edge 重启导致 in-flight 请求丢失,
任务停留在 PROCESSING 无人处理。重置后由常规重试机制接管。
返回被回收的 task_id 列表。
排序按 COALESCE(processed_at, updated_at, created_at) 降序。
date_filter 形如 '2026-08-21',匹配 processed_at 前缀。
"""
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()
try:
cursor = conn.cursor(pymysql.cursors.DictCursor)
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 _dt_or_none(value):
"""datetime 字段归一化MariaDB 严格模式兼容)
- 空串/None/'None'/'null' → NULL
- ISO 86012026-08-20T01:06:44Z / 2026-08-20T01:06:44.123+00:00'2026-08-20 01:06:44'
- 已为标准格式则原样返回
"""
if value is None:
return None
s = str(value).strip()
if s in ('', 'None', 'null', 'NaN'):
return None
# 归一化 ISO 8601 -> 'YYYY-MM-DD HH:MM:SS'
s2 = s.replace('T', ' ').replace('Z', '').replace('z', '')
if '+' in s2[10:]:
s2 = s2[:s2.index('+')]
if '.' in s2:
s2 = s2[:s2.index('.')]
try:
dt = datetime.strptime(s2, '%Y-%m-%d %H:%M:%S')
return dt.strftime('%Y-%m-%d %H:%M:%S')
except Exception:
return None
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()
# event_start/end_time 为 NOT NULL 列: 空值兜底
# end 缺失 → 用 startstart 也缺失 → 用当前时间
dt_start = _dt_or_none(event_start_time)
dt_end = _dt_or_none(event_end_time)
if not dt_end:
dt_end = dt_start
if not dt_start:
from datetime import datetime as _dt
dt_start = dt_end = _dt.now().strftime('%Y-%m-%d %H:%M:%S')
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, dt_start, dt_end, 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()
# frame_timestamp 为 NOT NULL 列: 空值兜底为当前时间
dt_ts = _dt_or_none(frame_timestamp)
if not dt_ts:
from datetime import datetime as _dt
dt_ts = _dt.now().strftime('%Y-%m-%d %H:%M:%S')
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, dt_ts, 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(pymysql.cursors.DictCursor)
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(pymysql.cursors.DictCursor)
cur = conn.cursor(pymysql.cursors.DictCursor)
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
cur.execute(
"""SELECT id, filename, camera_name, event_start_time, status,
summary_json, events_json, people_json, compute_provider,
processed_at, updated_at,
(SELECT COUNT(*) FROM sync_events se WHERE se.video_id = sync_videos.id) AS event_count
FROM sync_videos
WHERE status='done' AND processed_at LIKE %s
ORDER BY COALESCE(processed_at, updated_at, created_at) DESC
LIMIT %s OFFSET %s""",
(date_filter, limit, offset)
)
(f'{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
cur.execute(
"""SELECT id, filename, camera_name, event_start_time, status,
summary_json, events_json, people_json, compute_provider,
processed_at, updated_at,
(SELECT COUNT(*) FROM sync_events se WHERE se.video_id = sync_videos.id) AS event_count
FROM sync_videos
WHERE status='done'
ORDER BY COALESCE(processed_at, updated_at, created_at) DESC
LIMIT %s OFFSET %s""",
(limit, offset)
)
return cursor.fetchall()
(limit, offset))
return cur.fetchall()
finally:
conn.close()
def get_sync_video(video_id: int) -> Optional[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute("SELECT * FROM sync_videos WHERE id=%s", (video_id,))
return cur.fetchone()
finally:
conn.close()
# ============================================================
# chat_history 操作
# 同步镜像sync_events
# ============================================================
def upsert_sync_events(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 events 增量。"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_events
(id, video_id, ts, description, person_list_json,
is_attention_event, updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
video_id=VALUES(video_id),
ts=VALUES(ts),
description=VALUES(description),
person_list_json=VALUES(person_list_json),
is_attention_event=VALUES(is_attention_event),
updated_at=VALUES(updated_at),
synced_at=NOW()""",
(r.get('id'), r.get('video_id'), r.get('ts'), r.get('description'),
r.get('person_list_json'), 1 if r.get('is_attention_event') else 0,
r.get('updated_at'))
)
n += 1
conn.commit()
return n
finally:
conn.close()
def get_sync_events_for_video(video_id: int) -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"""SELECT id, video_id, ts, description, person_list_json,
is_attention_event, updated_at
FROM sync_events WHERE video_id=%s ORDER BY ts ASC""",
(video_id,))
return cur.fetchall()
finally:
conn.close()
def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]:
"""问答上下文:某人在某天的事件。
说明: Oracle 事件 ts 为视频内相对时间点(如 00:01:23不是绝对日期
因此按所属视频的 processed_at 日期过滤,再按 person_list_json 命中人名。
person 可为真名或抽象标签Oracle 回灌上下文用真名,但历史标签也保留)。
"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"""SELECT e.ts, e.description, e.person_list_json, e.is_attention_event,
v.camera_name, v.filename, v.event_start_time, v.processed_at
FROM sync_events e
JOIN sync_videos v ON e.video_id = v.id
WHERE v.processed_at LIKE %s AND e.person_list_json LIKE %s
ORDER BY v.processed_at ASC, e.ts ASC""",
(f'{date_str}%', f'%{person}%'))
return cur.fetchall()
finally:
conn.close()
# ============================================================
# 同步镜像sync_people
# ============================================================
def upsert_sync_people(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 people 增量。"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_people
(id, label, canonical_name, first_seen, appearances,
source, updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
label=VALUES(label),
canonical_name=VALUES(canonical_name),
first_seen=VALUES(first_seen),
appearances=VALUES(appearances),
source=VALUES(source),
updated_at=VALUES(updated_at),
synced_at=NOW()""",
(r.get('id'), r.get('label'), r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('updated_at')))
n += 1
conn.commit()
return n
finally:
conn.close()
def get_sync_people() -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"SELECT id, label, canonical_name, first_seen, appearances, source, updated_at "
"FROM sync_people ORDER BY id ASC")
return cur.fetchall()
finally:
conn.close()
def get_sync_named_members() -> List[str]:
"""已命名成员的真名列表(供 UI 下拉 / 快捷选择)。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"SELECT DISTINCT canonical_name FROM sync_people "
"WHERE canonical_name IS NOT NULL AND canonical_name != '' "
"ORDER BY canonical_name ASC")
return [r['canonical_name'] for r in cur.fetchall()]
finally:
conn.close()
def get_sync_known_members_context() -> str:
"""获取人物清单文本,注入问答 Prompt让模型用真名指代。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"SELECT label, canonical_name FROM sync_people ORDER BY id ASC")
rows = cur.fetchall()
if not rows:
return ""
parts = []
for r in rows:
name = r['canonical_name'] or r['label']
if r['canonical_name'] and r['canonical_name'] != r['label']:
parts.append(f"- {name}(标识:{r['label']}")
else:
parts.append(f"- {name}")
return "\n".join(parts)
finally:
conn.close()
# ============================================================
# 同步游标
# ============================================================
def get_sync_cursor() -> str:
conn = get_conn()
try:
cur = conn.cursor()
cur.execute("SELECT value FROM sync_cursor WHERE key='last_since'")
row = cur.fetchone()
return row[0] if row else ''
finally:
conn.close()
def set_sync_cursor(value: str):
conn = get_conn()
try:
cur = conn.cursor()
cur.execute(
"""INSERT INTO sync_cursor (key, value) VALUES ('last_since', %s)
ON DUPLICATE KEY UPDATE value=VALUES(value)""",
(value,))
conn.commit()
finally:
conn.close()
# ============================================================
# 统计
# ============================================================
def get_sync_stats(date_str: str = None) -> Dict:
"""概览统计:视频数 / 事件数 / 关注事件数 / 出现人物数(按 date 可选过滤)。
人物数通过对 sync_events.person_list_json LIKE 统计MariaDB 不支持 JSON 数组展开)。
"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
# 视频/事件/关注数
if date_str:
cur.execute(
"""SELECT
COUNT(*) AS videos,
(SELECT COUNT(*) FROM sync_events se
JOIN sync_videos sv ON se.video_id=sv.id
WHERE sv.processed_at LIKE %s) AS events,
(SELECT COALESCE(SUM(se.is_attention_event),0) FROM sync_events se
JOIN sync_videos sv ON se.video_id=sv.id
WHERE sv.processed_at LIKE %s) AS attention
FROM sync_videos sv WHERE sv.processed_at LIKE %s""",
(f'{date_str}%', f'{date_str}%', f'{date_str}%'))
else:
cur.execute(
"""SELECT
(SELECT COUNT(*) FROM sync_videos WHERE status='done') AS videos,
(SELECT COUNT(*) FROM sync_events) AS events,
(SELECT COALESCE(SUM(is_attention_event),0) FROM sync_events) AS attention""")
stat = cur.fetchone() or {}
# 人物数distinct label 命中 sync_events
cur.execute("SELECT id, label, canonical_name FROM sync_people")
people = cur.fetchall()
# 基于 person_list_json 命中计数:逐 label 统计命中事件数
person_hits = 0
for p in people:
name = p['canonical_name'] or p['label']
cur.execute(
"SELECT COUNT(*) c FROM sync_events WHERE person_list_json LIKE %s",
(f'%{name}%',))
if cur.fetchone()['c'] > 0:
person_hits += 1
stat['people'] = person_hits
return stat
finally:
conn.close()
# ============================================================
# chat_history保留问答历史
# ============================================================
def insert_chat_history(user_question: str, ai_answer: str,
@@ -326,15 +393,15 @@ def insert_chat_history(user_question: str, ai_answer: str,
"""插入对话记录"""
conn = get_conn()
try:
cursor = conn.cursor()
cursor.execute(
cur = conn.cursor()
cur.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
return cur.lastrowid
finally:
conn.close()
@@ -343,7 +410,7 @@ def get_chat_history(limit=20, date_filter=None, person_filter=None) -> List[Dic
"""获取对话历史"""
conn = get_conn()
try:
cursor = conn.cursor(pymysql.cursors.DictCursor)
cur = conn.cursor(pymysql.cursors.DictCursor)
conditions = []
params = []
if date_filter:
@@ -353,327 +420,11 @@ def get_chat_history(limit=20, date_filter=None, person_filter=None) -> List[Dic
conditions.append("queried_person = %s")
params.append(person_filter)
where = f"WHERE {' AND '.join(conditions)}" if conditions else ""
params.extend([limit])
cursor.execute(
params.append(limit)
cur.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(pymysql.cursors.DictCursor)
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(pymysql.cursors.DictCursor)
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:
"""命名/重命名成员 + 批量回溯更新历史记录
- 标签未入库(如 AI 新输出的 人物1自动注册
- 已命名成员可重命名(旧真名一并回溯替换)
- person 为多人组合字符串('张三, 汤圆'),用 REPLACE 替换其中目标
"""
conn = get_conn()
try:
cursor = conn.cursor()
# 1. 查现有记录,拿到旧真名
cursor.execute(
"SELECT member_id, real_name FROM family_members WHERE abstract_label = %s",
(abstract_label,)
)
row = cursor.fetchone()
old_name = None
if row:
old_name = row[1]
else:
cursor.execute(
"""INSERT INTO family_members (abstract_label, feature_description, first_seen_at)
VALUES (%s, %s, NOW())""",
(abstract_label, f'由命名操作自动注册: {real_name}')
)
# 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
# 目标出现的三种形态: 独占整字段 / 多人组合内 / AI 直呼旧真名
updated_details_count = 0
for old in {abstract_label, old_name} - {None}:
if old == real_name:
continue
cursor.execute(
"UPDATE event_details SET person = %s WHERE person = %s",
(real_name, old)
)
updated_details_count += cursor.rowcount
cursor.execute(
"""UPDATE event_details
SET person = REPLACE(person, %s, %s)
WHERE person LIKE %s AND person <> %s""",
(old, real_name, f'%{old}%', real_name)
)
updated_details_count += cursor.rowcount
# 4. 批量回溯更新 monitor_events.entities_json
# MariaDB 10.11 不支持 MySQL 的 $[*] 通配符 JSON 路径,
# 改用 Python 层解析 + 逐行更新
import json as _json
updated_events_count = 0
targets = {abstract_label, old_name} - {None}
if targets and targets != {real_name}:
like_conds = ' OR '.join(['entities_json LIKE %s'] * len(targets))
like_args = [f'%{t}%' for t in targets]
cursor.execute(
f"SELECT event_id, entities_json FROM monitor_events WHERE {like_conds}",
tuple(like_args)
)
for eid, entities_raw in cursor.fetchall():
if not entities_raw:
continue
try:
entities = _json.loads(entities_raw) if isinstance(entities_raw, str) else entities_raw
except (ValueError, TypeError):
continue
changed = False
if isinstance(entities, list):
for ent in entities:
if isinstance(ent, dict) and ent.get('person') in targets:
ent['person'] = real_name
changed = True
if changed:
cursor.execute(
"UPDATE monitor_events SET entities_json = %s WHERE event_id = %s",
(_json.dumps(entities, ensure_ascii=False), eid)
)
updated_events_count += 1
conn.commit()
return {
"abstract_label": abstract_label,
"real_name": real_name,
"renamed_from": old_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 merge_member(source_key: str, target_key: str, named_by: str = '管理员') -> Dict:
"""合并人物: source 并入 target用户判断两帧是同一人时
source_key/target_key 可为 abstract_label 或 real_name。
- event_details.person: 独占/组合字符串内的 source 一律替换为 target 显示名
- monitor_events.entities_json: person 字段替换
- family_members: source 行置 is_active=0保留历史target 未入库则注册
"""
if source_key == target_key:
return {"error": "source 与 target 不能相同"}
conn = get_conn()
try:
cursor = conn.cursor(pymysql.cursors.DictCursor)
def resolve(key):
cursor.execute(
"""SELECT member_id, abstract_label, real_name FROM family_members
WHERE is_active = TRUE AND (abstract_label = %s OR real_name = %s)
ORDER BY real_name IS NULL LIMIT 1""",
(key, key))
return cursor.fetchone()
src = resolve(source_key)
tgt = resolve(target_key)
if not src:
return {"error": f"人物 {source_key} 不存在"}
if not tgt:
# target 是未入库的裸标签(如 人物1注册后作为目标
cursor.execute(
"""INSERT INTO family_members (abstract_label, feature_description, first_seen_at)
VALUES (%s, %s, NOW())""",
(target_key, f'合并操作自动注册: {source_key} 并入'))
tgt = {'abstract_label': target_key, 'real_name': None}
target_display = tgt['real_name'] or tgt['abstract_label']
# source 的所有称呼: 抽象标签 + 旧真名(多人组合里两种都可能出现)
source_names = {src['abstract_label']}
if src['real_name']:
source_names.add(src['real_name'])
updated_details = 0
for name in source_names:
if name == target_display:
continue
cursor.execute(
"UPDATE event_details SET person = %s WHERE person = %s",
(target_display, name))
updated_details += cursor.rowcount
cursor.execute(
"""UPDATE event_details
SET person = REPLACE(person, %s, %s)
WHERE person LIKE %s AND person <> %s""",
(name, target_display, f'%{name}%', target_display))
updated_details += cursor.rowcount
# entities_json 逐行替换
import json as _json
updated_events = 0
like_conds = ' OR '.join(['entities_json LIKE %s'] * len(source_names))
cursor.execute(
f"SELECT event_id, entities_json FROM monitor_events WHERE {like_conds}",
tuple(f'%{n}%' for n in source_names))
for row in cursor.fetchall():
raw = row['entities_json']
if not raw:
continue
try:
entities = _json.loads(raw) if isinstance(raw, str) else raw
except (ValueError, TypeError):
continue
changed = False
if isinstance(entities, list):
for ent in entities:
if isinstance(ent, dict) and ent.get('person') in source_names:
ent['person'] = target_display
changed = True
if changed:
cursor.execute(
"UPDATE monitor_events SET entities_json = %s WHERE event_id = %s",
(_json.dumps(entities, ensure_ascii=False), row['event_id']))
updated_events += 1
# source 行停用
cursor.execute(
"UPDATE family_members SET is_active = 0, updated_at = NOW() WHERE member_id = %s",
(src['member_id'],))
conn.commit()
return {
"source": source_key,
"target": target_display,
"merged_names": sorted(source_names),
"updated_event_details_count": updated_details,
"updated_monitor_events_count": updated_events
}
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(pymysql.cursors.DictCursor)
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(pymysql.cursors.DictCursor)
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()
return cur.fetchall()
finally:
conn.close()

View File

@@ -1,4 +0,0 @@
"""Dispatcher 包"""
from .dispatcher import Dispatcher
__all__ = ["Dispatcher"]

View File

@@ -1,517 +0,0 @@
"""
Dispatcher - 30s 轮询 PENDING 任务,上传视频至 Edge 异步队列
流程(异步队列模式 + NAS 预压缩 + 分块断点续传):
1. 读取任务对应的本地视频文件
2. 大文件 (>20MB): NAS 端 FFmpeg 预压缩 (480p/CRF28, ~36x 压缩比)
(根因: NAS→Oracle 跨境上行带宽仅 ~0.5-0.9MB/s360MB 原始上传需 10-20min 且频繁超时)
3a. 压缩后小文件 (<=20MB): 直接 multipart 上传至 /enqueue
3b. 仍超阈值: 分块上传 (5MB/块) 至 /chunk支持断点续传最后调 /assemble 合并入队
4. Edge 保存视频 + 入 SQLite 队列,返回 202
5. Dispatcher 标记任务为 PROCESSING已派发等待 Poller 拉取结果)
6. Poller 线程定期从 Edge /api/edge/results 拉取结果,写库后标记 SUCCESS
退避重试: min(30 * (retry_count + 1), 300) 秒
分块级重试: 每块最多重试 3 次
看门狗: 线程崩溃后自动重启
"""
import os
import io
import re
import time
import math
import shutil
import subprocess
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')
CHUNK_SIZE = 5 * 1024 * 1024 # 5MB per chunk (reliable at ~1Mbps upload)
CHUNK_THRESHOLD = 20 * 1024 * 1024 # files > 20MB use chunked upload
COMPRESS_THRESHOLD = 20 * 1024 * 1024 # files > 20MB get pre-compressed before upload
COMPRESS_DIR = '/tmp/fam_compressed'
COMPRESS_CACHE_TTL = 24 * 3600 # 压缩缓存保留 24h供上传失败重试复用
MAX_CHUNK_RETRIES = 3
# Synology 系统 ffmpeg 被裁剪(无 h264 编解码CodecPack 的 ffmpeg41 带 libx264
FFMPEG_CANDIDATES = [
'/var/packages/CodecPack/target/bin/ffmpeg41',
'/usr/local/bin/ffmpeg',
]
def _safe_remove(path):
"""忽略不存在/清理失败的删除"""
try:
if os.path.isfile(path):
os.remove(path)
except OSError:
pass
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/enqueue')
self.max_retries = cfg.get('dispatcher', {}).get('max_retries', 3)
self.camera_name = cfg.get('scheduler', {}).get('camera_name', '默认摄像头')
self.stale_timeout = cfg.get('dispatcher', {}).get('stale_timeout', 600)
self.compress_timeout = cfg.get('dispatcher', {}).get('compress_timeout', 3600)
self.ffmpeg = self._find_ffmpeg(cfg)
# 推导 Edge base URL
self.edge_base = self.edge_url.rsplit('/api/edge/video/enqueue', 1)[0]
self.chunk_url = f"{self.edge_base}/api/edge/video/chunk"
self.chunks_query_url = f"{self.edge_base}/api/edge/video/chunks"
self.assemble_url = f"{self.edge_base}/api/edge/video/assemble"
self._running = False
self._thread = None
@staticmethod
def _find_ffmpeg(cfg):
"""查找可用 ffmpeg配置优先其次 CodecPack带 libx264最后 PATH"""
configured = cfg.get('dispatcher', {}).get('ffmpeg_path')
candidates = ([configured] if configured else []) + FFMPEG_CANDIDATES
for path in candidates:
if os.path.isfile(path) and os.access(path, os.X_OK):
return path
return shutil.which('ffmpeg')
def _calculate_backoff(self, retry_count):
"""退避策略: min(30 * (retry_count + 1), 300)"""
return min(30 * (retry_count + 1), 300)
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 _probe_duration(self, video_path):
"""用 ffmpeg 解析视频时长NAS 无独立 ffprobe失败返回 0"""
if not self.ffmpeg:
return 0
try:
r = subprocess.run(
[self.ffmpeg, '-i', video_path],
capture_output=True, timeout=30)
m = re.search(r'Duration:\s*(\d+):(\d+):(\d+(?:\.\d+)?)',
r.stderr.decode('utf-8', 'ignore'))
if m:
return int(m.group(1)) * 3600 + int(m.group(2)) * 60 + float(m.group(3))
except (subprocess.TimeoutExpired, OSError):
pass
return 0
def _build_payload(self, task):
"""构建推送元数据,注入已知成员清单与事件时间"""
video_path = task['video_path']
payload = {
"task_id": str(task['task_id']),
"camera_name": self.camera_name,
"known_members_context": db_layer.get_known_members_context(),
"event_start_time": "",
"event_end_time": "",
}
try:
mtime = os.path.getmtime(video_path)
# mtime 是录制结束时刻,开始时间 = 结束时间 - 视频时长
duration = self._probe_duration(video_path)
start_dt = datetime.fromtimestamp(mtime - duration)
payload["event_start_time"] = start_dt.strftime('%Y-%m-%d %H:%M:%S')
except OSError:
pass
return payload
def _compress_video(self, task_id, video_path):
"""NAS 端预压缩: 480p/CRF28/veryfast静态监控场景实测 ~36x 压缩比
成功返回压缩文件路径;失败返回 None回退原始文件分块上传
压缩产物缓存在 /tmp/fam_compressed/task_{id}/,重试时源文件未变则复用。
"""
if not self.ffmpeg:
logger.warning(f"[task_id={task_id}] ffmpeg 不可用,跳过预压缩")
return None
self._cleanup_compress_cache()
out_dir = os.path.join(COMPRESS_DIR, f'task_{task_id}')
out_path = os.path.join(out_dir, os.path.basename(video_path))
# 缓存复用:源文件未变且压缩产物有效
try:
if (os.path.isfile(out_path)
and os.path.getsize(out_path) > 0
and os.path.getmtime(out_path) >= os.path.getmtime(video_path)):
logger.info(f"[task_id={task_id}] 复用压缩缓存: {out_path}")
return out_path
except OSError:
pass
os.makedirs(out_dir, exist_ok=True)
src_size = os.path.getsize(video_path)
start = time.time()
# 写临时文件(带 PID 防多进程冲突),成功后原子 rename —
# 服务被 kill 时 ffmpeg 成为孤儿继续写 tmp缓存目录中
# 只会出现完整产物,杜绝半成品被复用(曾导致上传损坏视频)
tmp_path = f"{out_path}.{os.getpid()}.tmp"
cmd = [
self.ffmpeg, '-y',
'-i', video_path,
# 第二级 scale 向下取偶h264 要求偶数尺寸force_divisible_by 需 ffmpeg>=4.3
'-vf', ('scale=854:480:force_original_aspect_ratio=decrease,'
'scale=trunc(iw/2)*2:trunc(ih/2)*2'),
'-c:v', 'libx264', '-preset', 'veryfast', '-crf', '28',
'-an',
'-f', 'mp4', # .tmp 扩展名无法推断 muxer必须显式指定
tmp_path,
]
try:
result = subprocess.run(
cmd, capture_output=True, text=True,
timeout=self.compress_timeout,
)
except subprocess.TimeoutExpired:
_safe_remove(tmp_path)
logger.error(f"[task_id={task_id}] 压缩超时 ({self.compress_timeout}s),回退原始上传")
return None
except OSError as e:
logger.error(f"[task_id={task_id}] 启动 ffmpeg 失败: {e}")
return None
if result.returncode != 0 or not os.path.isfile(tmp_path) or os.path.getsize(tmp_path) == 0:
_safe_remove(tmp_path)
stderr_tail = (result.stderr or '')[-300:]
logger.error(f"[task_id={task_id}] 压缩失败 (rc={result.returncode}): {stderr_tail}")
return None
os.replace(tmp_path, out_path)
dst_size = os.path.getsize(out_path)
elapsed = time.time() - start
logger.info(f"[task_id={task_id}] 预压缩完成: {src_size/1048576:.1f}MB → "
f"{dst_size/1048576:.1f}MB ({src_size/max(dst_size,1):.1f}x),耗时 {elapsed:.0f}s")
return out_path
@staticmethod
def _cleanup_compress_cache():
"""清理超过 TTL 的压缩缓存目录 + 孤儿 .tmp 残留tmpfs 空间有限)"""
try:
if not os.path.isdir(COMPRESS_DIR):
return
cutoff = time.time() - COMPRESS_CACHE_TTL
for entry in os.listdir(COMPRESS_DIR):
path = os.path.join(COMPRESS_DIR, entry)
try:
if os.path.isdir(path):
if os.path.getmtime(path) < cutoff:
shutil.rmtree(path, ignore_errors=True)
continue
# 清理超过 1h 的 .tmp 残留(孤儿 ffmpeg 产物)
for fn in os.listdir(path):
if fn.endswith('.tmp') and \
os.path.getmtime(os.path.join(path, fn)) < time.time() - 3600:
_safe_remove(os.path.join(path, fn))
except OSError:
continue
except OSError:
pass
def _dispatch_one(self, task):
"""上传视频至 Edge 异步队列(大文件先预压缩,自动选择直接/分块模式)"""
task_id = task['task_id']
video_path = task['video_path']
if not video_path or not os.path.isfile(video_path):
db_layer.update_task_status(
task_id, 'FAILED',
error_message=f"视频文件不存在: {video_path}",
failure_stage='callback'
)
logger.error(f"[task_id={task_id}] 视频文件不存在,标记 FAILED: {video_path}")
return
file_size = os.path.getsize(video_path)
size_mb = file_size / (1024 * 1024)
db_layer.update_task_status(task_id, 'PROCESSING')
# 跨境上行带宽受限(实测 ~0.5-0.9MB/s大文件先预压缩再上传
upload_path = video_path
if file_size > COMPRESS_THRESHOLD:
compressed = self._compress_video(task_id, video_path)
if compressed:
upload_path = compressed
file_size = os.path.getsize(upload_path)
size_mb = file_size / (1024 * 1024)
# 压缩耗时较长,重置 stale 计时基准reclaim 按 updated_at 判断)
db_layer.update_task_status(task_id, 'PROCESSING')
else:
logger.warning(f"[task_id={task_id}] 压缩失败,回退原始文件上传 ({size_mb:.1f}MB)")
payload = self._build_payload(task)
if file_size > CHUNK_THRESHOLD:
logger.info(f"[task_id={task_id}] 大文件分块上传: {size_mb:.1f}MB, "
f"{math.ceil(file_size / CHUNK_SIZE)}")
self._dispatch_chunked(task, payload, upload_path, file_size)
else:
log_task(logger, task_id, 'dispatcher',
f'直接上传: {self.edge_url} ({size_mb:.1f}MB)')
self._dispatch_direct(task, payload, upload_path)
def _dispatch_direct(self, task, payload, video_path):
"""小文件直接上传至 /enqueue
注意 timeout 第一参数: urllib3 发送 multipart body 期间 socket
timeout 取的是 connect timeout 值(实测传 60s 则 60s 整超时),
而非 read timeout——跨境 1.4MB/s 下 17MB 需 ~12s必须给足。
"""
task_id = task['task_id']
try:
with open(video_path, 'rb') as fh:
resp = requests.post(
self.edge_url,
data=payload,
files={'video': (os.path.basename(video_path), fh, 'video/mp4')},
timeout=(120, 300)
)
except requests.RequestException as e:
logger.error(f"[task_id={task_id}] 上传失败: {e}")
self._schedule_retry(task)
return
if resp.status_code == 202:
try:
data = resp.json()
queue_id = data.get('queue_id', '?')
logger.info(f"[task_id={task_id}] 已入 Edge 队列 (queue_id={queue_id}),等待 Poller 拉取结果")
except ValueError:
logger.info(f"[task_id={task_id}] 已入 Edge 队列,等待 Poller 拉取结果")
return
if resp.status_code == 429:
logger.warning(f"[task_id={task_id}] Edge 队列满 (429),回到 PENDING 稍后重试")
db_layer.update_task_status(task_id, 'PENDING')
return
logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}")
self._schedule_retry(task)
def _query_uploaded_chunks(self, task_id, expected_total=None):
"""查询 Edge 端已上传分块列表
返回 (uploaded_set, edge_total_chunks)。
如果 expected_total 与 edge_total 不匹配chunk_size 变更),
返回空集让 Edge 自动清理旧分块。
"""
try:
resp = requests.get(
self.chunks_query_url,
params={'task_id': task_id},
timeout=(30, 15)
)
if resp.status_code == 200:
data = resp.json()
uploaded = set(data.get('uploaded_chunks', []))
edge_total = data.get('total_chunks', 0)
if expected_total and edge_total and edge_total != expected_total:
logger.warning(f"[task_id={task_id}] Edge total_chunks={edge_total} "
f"≠ expected={expected_total}chunk_size 已变更),从头上传")
return set(), edge_total
return uploaded, edge_total
logger.warning(f"[task_id={task_id}] 查询已上传分块返回 {resp.status_code},将全量重传")
except requests.RequestException as e:
logger.warning(f"[task_id={task_id}] 查询已上传分块失败(将全量重传): {e}")
return set(), 0
def _dispatch_chunked(self, task, payload, video_path, file_size):
"""大文件分块上传 + 断点续传
1. 查询 Edge 端已上传分块(断点续传)
2. 上传缺失分块(每块最多重试 3 次,失败后查询 Edge 确认是否实际收到)
3. 全部分块上传后调用 /assemble 合并入队
"""
task_id = task['task_id']
total_chunks = math.ceil(file_size / CHUNK_SIZE)
filename = os.path.basename(video_path)
# 1. 查询已上传分块(断点续传)
uploaded_set, edge_total = self._query_uploaded_chunks(task_id, expected_total=total_chunks)
if uploaded_set:
logger.info(f"[task_id={task_id}] 断点续传: 已有 {len(uploaded_set)}/{total_chunks}")
# 2. 上传缺失分块
try:
with open(video_path, 'rb') as fh:
for idx in range(total_chunks):
if idx in uploaded_set:
continue
chunk_data = fh.read(CHUNK_SIZE)
if not chunk_data:
break
success = False
for attempt in range(MAX_CHUNK_RETRIES):
try:
cresp = requests.post(
self.chunk_url,
data={
'task_id': str(task_id),
'chunk_index': str(idx),
'total_chunks': str(total_chunks),
'filename': filename,
},
files={'chunk': (f'chunk_{idx}', io.BytesIO(chunk_data))},
timeout=(60, 180)
)
if cresp.status_code == 200:
success = True
break
logger.warning(f"[task_id={task_id}] 分块 {idx} 返回 {cresp.status_code}(尝试 {attempt+1}/{MAX_CHUNK_RETRIES}")
except requests.RequestException as e:
logger.warning(f"[task_id={task_id}] 分块 {idx} 上传失败(尝试 {attempt+1}/{MAX_CHUNK_RETRIES}: {e}")
if attempt < MAX_CHUNK_RETRIES - 1:
time.sleep(5 * (attempt + 1))
if not success:
uploaded_now, _ = self._query_uploaded_chunks(task_id)
if idx in uploaded_now:
logger.info(f"[task_id={task_id}] 分块 {idx} 虽超时但 Edge 已收到,继续下一块")
uploaded_set.add(idx)
continue
logger.error(f"[task_id={task_id}] 分块 {idx} 确认未收到,安排文件级重试")
self._schedule_retry(task)
return
uploaded_set.add(idx)
if (idx + 1) % 5 == 0 or idx == total_chunks - 1:
logger.info(f"[task_id={task_id}] 分块进度: {idx + 1}/{total_chunks}")
except IOError as e:
logger.error(f"[task_id={task_id}] 读取视频文件失败: {e}")
self._schedule_retry(task)
return
# 3. 合并 + 入队
try:
aresp = requests.post(
self.assemble_url,
data={
'task_id': str(task_id),
'camera_name': payload.get('camera_name', ''),
'event_start_time': payload.get('event_start_time', ''),
'known_members_context': payload.get('known_members_context', ''),
},
timeout=(10, 60)
)
except requests.RequestException as e:
logger.error(f"[task_id={task_id}] 合并请求失败: {e}")
self._schedule_retry(task)
return
if aresp.status_code == 202:
try:
data = aresp.json()
queue_id = data.get('queue_id', '?')
asm_size = data.get('size_mb', '?')
logger.info(f"[task_id={task_id}] 分块合并入队成功 (queue_id={queue_id}, {asm_size}MB),等待 Poller 拉取结果")
except ValueError:
logger.info(f"[task_id={task_id}] 分块合并入队成功,等待 Poller 拉取结果")
return
logger.error(f"[task_id={task_id}] 合并端点返回 {aresp.status_code}: {aresp.text[:200]}")
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):
"""执行一次轮询"""
# 回收僵尸任务PROCESSING 超过 stale_timeout 说明 Edge 丢失了任务
try:
stale_ids = db_layer.reclaim_stale_processing(self.stale_timeout)
for tid in stale_ids:
logger.warning(f"[task_id={tid}] PROCESSING 超时 {self.stale_timeout}s回收为 PENDING 重试")
except Exception as e:
logger.error(f"僵尸任务回收失败: {e}", exc_info=True)
tasks = db_layer.get_pending_tasks(limit=1)
for task in tasks:
if task['retry_count'] >= task['max_retries']:
db_layer.update_task_status(
task['task_id'], 'FAILED',
error_message=f"超过最大重试次数 {task['max_retries']}",
failure_stage='callback'
)
logger.warning(f"[task_id={task['task_id']}] retry_count={task['retry_count']} >= max_retries={task['max_retries']},标记 FAILED")
continue
if self._should_retry(task):
try:
self._dispatch_one(task)
except Exception as e:
logger.error(f"[task_id={task['task_id']}] dispatch 异常: {e}", exc_info=True)
def _run(self):
"""线程主循环"""
logger.info(f"Dispatcher 启动 (enqueue + 分块模式),轮询间隔 {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 is_alive(self):
"""线程是否存活"""
return self._thread is not None and self._thread.is_alive()
def check_and_restart(self):
"""看门狗:线程崩溃后自动重启"""
if self._running and not self.is_alive():
logger.warning("Dispatcher 线程已死亡,正在重启...")
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)

View File

@@ -1,4 +0,0 @@
"""Event-Receiver 包"""
from .event_receiver import event_bp
__all__ = ["event_bp"]

View File

@@ -1,172 +0,0 @@
"""
Event-Receiver - Flask 蓝图,接收 Edge 回调,写库
处理逻辑:
1. 成功回调: 插入 monitor_events 1 条 + 遍历 frame_details 逐条插入 event_details
2. frame_details 携带的关键帧 base64 落盘到 fam-ui 静态目录(供时间轴展示)
3. 对未命名的 abstract_label 自动 upsert 到 family_members
4. 更新 process_tasks 状态为 SUCCESS
5. 失败回调: 更新任务状态为 FAILED记录 failure_stage
"""
import os
import re
import json
import base64
from flask import Blueprint, request, jsonify
from ..logger import setup_logger
from ..config_loader import load_config
from .. import db_layer
logger = setup_logger('fam-core.event_receiver')
event_bp = Blueprint('event_receiver', __name__)
# 关键帧落盘目录fam-ui 读取展示fam-core 与 fam-ui 同机部署)
_cfg = load_config()
FRAME_IMAGE_DIR = _cfg.get('storage', {}).get(
'frame_image_dir',
'/volume1/web/sentinel-home-ai/fam-ui/static/frames')
# 匹配 "人物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 _save_frame_images(event_id: int, frame_details: list) -> int:
"""把 frame_details 中的 base64 关键帧落盘,返回成功张数
同时写 meta.jsonframe_index -> face_countUI 据此挑选有人像的帧做头像。
"""
saved = 0
face_counts = {}
for frame in frame_details:
img_b64 = frame.pop('frame_image', None)
faces = frame.pop('face_count', None)
if not img_b64:
continue
idx = frame.get('frame_index', 0)
try:
out_dir = os.path.join(FRAME_IMAGE_DIR, f'event_{event_id}')
os.makedirs(out_dir, exist_ok=True)
out_path = os.path.join(out_dir, f'frame_{idx}.jpg')
with open(out_path, 'wb') as f:
f.write(base64.b64decode(img_b64))
if faces is not None:
face_counts[str(idx)] = int(faces)
saved += 1
except Exception as e:
logger.warning(f"[event_id={event_id}] 关键帧落盘失败 frame_{idx}: {e}")
if face_counts:
try:
import json as _json
with open(os.path.join(out_dir, 'meta.json'), 'w') as f:
_json.dump(face_counts, f)
except Exception as e:
logger.warning(f"[event_id={event_id}] meta.json 写入失败: {e}")
if saved:
logger.info(f"[event_id={event_id}] 关键帧落盘 {saved} 张 -> {FRAME_IMAGE_DIR}")
return saved
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})")
def apply_success_event(task_id, data: dict) -> int:
"""将成功结果写库,返回 event_id
供两条路径复用:
- webhook 回调路由 (拉取模式)
- Dispatcher 收到推送模式同步响应后直接落库
"""
# 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. 关键帧图片落盘先落盘再入库pop 掉 base64 后 insert避免大字段进 DB
frame_details = data.get('frame_details', [])
try:
_save_frame_images(event_id, frame_details)
except Exception as e:
logger.warning(f"[event_id={event_id}] 关键帧落盘异常(不影响入库): {e}")
# 3. 遍历 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 event_id
@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:
event_id = apply_success_event(task_id, data)
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

View File

@@ -1,14 +1,20 @@
"""
Member-Manager - Flask 蓝图,成员命名管理
Member-Manager - Flask 蓝图,人物命名管理(新架构 v2
1. GET /api/member/unnamed - 列出未命名人物
2. POST /api/member/name - 命名人物 + 批量回溯更新
3. GET /api/member/list - 列出所有成员
Oracle 是人物规范的唯一真源。NAS 命名操作:
1. POST /api/member/name 命名/重命名某 label -> 回推 Oracle + 立即拉回
2. POST /api/member/merge 将两个 label 合并为同一身份(统一 canonical_name
3. GET /api/member/list 列出所有人物label + canonical_name
4. GET /api/member/unnamed 列出未命名人物canonical_name 为空)
命名流程: 调 Oracle /api/oracle/people/correct 设置 manual 规范名 ->
立即 trigger_now() 拉回最新 people 镜像 -> 前端刷新即见结果。
"""
from flask import Blueprint, request, jsonify
from ..logger import setup_logger
from .. import db_layer
from ..oracle_sync import get_sync
logger = setup_logger('fam-core.member_manager')
@@ -17,86 +23,122 @@ member_bp = Blueprint('member_manager', __name__)
@member_bp.route('/api/member/unnamed', methods=['GET'])
def list_unnamed():
"""列出未命名人物"""
members = db_layer.get_unnamed_members()
# datetime 序列化
"""列出未命名人物canonical_name 为空)"""
members = db_layer.get_sync_people()
result = []
for m in members:
canonical = m.get('canonical_name')
if not canonical:
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']
"label": m['label'],
"appearances": m.get('appearances', 0),
"first_seen": m.get('first_seen'),
})
return jsonify({"unnamed_members": result}), 200
@member_bp.route('/api/member/list', methods=['GET'])
def list_members():
"""列出所有人物(按 canonical_name 或 label 展示)"""
members = db_layer.get_sync_people()
result = []
for m in members:
canonical = m.get('canonical_name')
result.append({
"label": m['label'],
"canonical_name": canonical,
"display_name": canonical or m['label'],
"is_named": bool(canonical),
"appearances": m.get('appearances', 0),
"source": m.get('source'),
"first_seen": m.get('first_seen'),
})
return jsonify({"members": result}), 200
@member_bp.route('/api/member/name', methods=['POST'])
def name_member():
"""命名人物 + 批量回溯更新历史记录"""
"""命名人物(回推 Oracle + 立即拉回本地镜像)
请求: {"label": "人物A", "canonical_name": "张三"}
"""
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', '管理员')
label = (data.get('label') or '').strip()
canonical_name = (data.get('canonical_name') or '').strip()
if not label or not canonical_name:
return jsonify({"error": "缺少必填字段: label, canonical_name"}), 400
if not abstract_label or not real_name:
return jsonify({"error": "缺少必填字段: abstract_label, real_name"}), 400
logger.info(f"命名: {abstract_label} -> {real_name}")
logger.info(f"命名: {label} -> {canonical_name}(回推 Oracle")
ok, err = get_sync().push_name_correct(label, canonical_name)
if not ok:
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
# 立即拉回最新 people 镜像,前端无需等待下一个 30 分钟周期
try:
result = db_layer.name_member(abstract_label, real_name, named_by)
if 'error' in result:
return jsonify(result), 404
return jsonify(result), 200
get_sync().trigger_now()
except Exception as e:
logger.error(f"命名失败: {e}", exc_info=True)
return jsonify({"error": str(e)}), 500
logger.warning(f"命名后即时拉回失败(下一个周期会自动同步): {e}")
members = db_layer.get_sync_people()
return jsonify({
"status": "ok",
"label": label,
"canonical_name": canonical_name,
"members": [{
"label": m['label'],
"canonical_name": m.get('canonical_name'),
"display_name": m.get('canonical_name') or m['label'],
} for m in members]
}), 200
@member_bp.route('/api/member/merge', methods=['POST'])
def merge_member():
"""合并人物(用户判断两个标签是同一人时,source 并入 target"""
"""合并人物:将 source 并入 target 身份(统一 canonical_name
若 target 已命名 -> 用其 canonical_name否则用 target label 作为规范名。
请求: {"source": "人物B", "target": "张三""人物A"}
"""
data = request.get_json(silent=True)
if not data:
return jsonify({"error": "Invalid JSON"}), 400
source_key = data.get('source')
target_key = data.get('target')
if not source_key or not target_key:
source = (data.get('source') or '').strip()
target = (data.get('target') or '').strip()
if not source or not target:
return jsonify({"error": "缺少必填字段: source, target"}), 400
if source == target:
return jsonify({"error": "source 与 target 不能相同"}), 400
# 解析 target 的规范名
members = {m['label']: m for m in db_layer.get_sync_people()}
target_row = members.get(target)
if target_row and target_row.get('canonical_name'):
canonical = target_row['canonical_name']
else:
canonical = target # target 未命名 -> 以 label 作为规范名
logger.info(f"合并: {source} -> {canonical}(回推 Oracle")
ok, err = get_sync().push_name_correct(source, canonical)
if not ok:
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
logger.info(f"合并人物: {source_key} -> {target_key}")
try:
result = db_layer.merge_member(source_key, target_key)
if 'error' in result:
return jsonify(result), 404
return jsonify(result), 200
get_sync().trigger_now()
except Exception as e:
logger.error(f"合并失败: {e}", exc_info=True)
return jsonify({"error": str(e)}), 500
logger.warning(f"合并后即时拉回失败(下一个周期会自动同步): {e}")
@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
members = db_layer.get_sync_people()
return jsonify({
"status": "ok",
"source": source,
"target_canonical": canonical,
"members": [{
"label": m['label'],
"canonical_name": m.get('canonical_name'),
"display_name": m.get('canonical_name') or m['label'],
} for m in members]
}), 200

View File

@@ -0,0 +1,170 @@
"""
Oracle-Sync - NAS 端唯一后台线程
职责:
1. 每 interval_sec默认 1800s = 30 分钟)从甲骨文 FAM-Edge 拉取增量:
GET {base_url}/api/oracle/sync?since=<cursor>&token=<token>
返回 {videos, events, people, server_time},写入本地 MariaDB 镜像表
(sync_videos / sync_events / sync_people),并推进 sync_cursor。
2. 接收命名校正回推: POST {base_url}/api/oracle/people/correct
{label, canonical_name, token} —— 手动命名manual 优先,不被 LLM 覆盖)。
数据流向(新架构 v2:
Google 硬盘 --rclone--> 甲骨文本地 --> 整视频分析 --> Oracle SQLite
--> [本线程每 30 分钟拉增量] --> NAS MariaDB 镜像 --> fam-ui 读取展示
NAS 不再处理任何视频CPU 占用显著降低。
"""
import time
import threading
import requests
from datetime import datetime
from .logger import setup_logger
from .config_loader import load_config
from . import db_layer
logger = setup_logger('fam-core.oracle_sync')
_SYNC_INSTANCE = None
def get_sync():
"""模块级单例app.py 启动时创建并 start其余模块经此获取"""
global _SYNC_INSTANCE
if _SYNC_INSTANCE is None:
_SYNC_INSTANCE = OracleSync()
return _SYNC_INSTANCE
class OracleSync:
def __init__(self):
cfg = load_config().get('oracle_sync', {})
self.base_url = cfg.get('base_url', 'http://129.146.203.203:5000').rstrip('/')
self.token = cfg.get('token', '')
self.interval_sec = int(cfg.get('interval_sec', 1800))
self.timeout = int(cfg.get('timeout', 120))
self._running = False
self._thread = None
self._last_sync_at = None
self._last_error = None
self._last_count = None
# ------------------------------------------------------------------
def _pull_once(self) -> bool:
"""执行一次增量拉取。返回是否成功。"""
since = db_layer.get_sync_cursor() or ''
params = {'since': since, 'token': self.token}
try:
resp = requests.get(
f"{self.base_url}/api/oracle/sync",
params=params, timeout=(10, self.timeout))
except requests.RequestException as e:
self._last_error = f"请求失败: {e}"
logger.error(f"拉取同步失败: {e}")
return False
if resp.status_code == 401:
self._last_error = "token 校验失败"
logger.error("同步 token 校验失败 (401),请检查 oracle_sync.token 配置")
return False
if resp.status_code != 200:
self._last_error = f"HTTP {resp.status_code}"
logger.error(f"同步返回异常: {resp.status_code} {resp.text[:200]}")
return False
try:
data = resp.json()
except ValueError:
self._last_error = "非 JSON 响应"
logger.error("同步返回非 JSON 响应")
return False
videos = data.get('videos', []) or []
events = data.get('events', []) or []
people = data.get('people', []) or []
server_time = data.get('server_time', '') or ''
n_videos = db_layer.upsert_sync_videos(videos)
n_events = db_layer.upsert_sync_events(events)
n_people = db_layer.upsert_sync_people(people)
if server_time:
db_layer.set_sync_cursor(server_time)
self._last_sync_at = datetime.now()
self._last_error = None
self._last_count = (n_videos, n_events, n_people)
logger.info(
f"同步完成: videos+{n_videos} events+{n_events} people+{n_people} "
f"since={since!r} -> server_time={server_time}")
return True
# ------------------------------------------------------------------
def push_name_correct(self, label: str, canonical_name: str):
"""回推命名校正到 Oracle手动命名优先级最高不被 LLM 覆盖)。
返回 (success: bool, error: str)
"""
label = (label or '').strip()
canonical_name = (canonical_name or '').strip()
if not label or not canonical_name:
return False, "缺少 label / canonical_name"
try:
resp = requests.post(
f"{self.base_url}/api/oracle/people/correct",
json={"label": label, "canonical_name": canonical_name,
"token": self.token},
timeout=(10, 30))
except requests.RequestException as e:
logger.error(f"命名校正回推失败: {e}")
return False, str(e)
if resp.status_code == 200:
return True, ""
msg = f"HTTP {resp.status_code}: {resp.text[:200]}"
logger.error(f"命名校正回推失败: {msg}")
return False, msg
# ------------------------------------------------------------------
def _run(self):
logger.info(f"OracleSync 线程启动,间隔 {self.interval_sec}s目标 {self.base_url}")
while self._running:
try:
self._pull_once()
except Exception as e:
self._last_error = str(e)
logger.error(f"同步异常: {e}", exc_info=True)
# 分段休眠,便于 stop 快速唤醒
for _ in range(self.interval_sec):
if not self._running:
break
time.sleep(1)
def start(self):
if self._running:
return
self._running = True
self._thread = threading.Thread(target=self._run, daemon=True, name='oracle-sync')
self._thread.start()
def is_alive(self):
return self._thread is not None and self._thread.is_alive()
def stop(self):
self._running = False
if self._thread:
self._thread.join(timeout=5)
def trigger_now(self) -> bool:
"""立即触发一次同步(命名后即时拉回 / 手动)。"""
return self._pull_once()
def status(self) -> dict:
return {
"running": self.is_alive(),
"last_sync_at": self._last_sync_at.isoformat() if self._last_sync_at else None,
"last_error": self._last_error,
"last_count": self._last_count,
"cursor": db_layer.get_sync_cursor(),
"interval_sec": self.interval_sec,
}

View File

@@ -1,139 +0,0 @@
"""
Poller - 定期从 Edge 拉取已完成的任务结果,写入 MariaDB
流程:
1. 每 N 秒请求 Edge /api/edge/results?limit=10
2. 遍历结果列表,对每个 nas_task_id:
- success: 调用 apply_success_event 写入 monitor_events + event_details标记 SUCCESS
- failed: 更新任务状态为 FAILED记录 failure_stage 和 error_message
3. Edge 端自动标记已拉取的结果为 delivered
"""
import time
import threading
import requests
from ..logger import setup_logger, log_task
from ..config_loader import load_config
from .. import db_layer
from ..event_receiver.event_receiver import apply_success_event
logger = setup_logger('fam-core.poller')
class Poller:
"""结果拉取器,定期从 Edge 拉取处理结果"""
def __init__(self):
cfg = load_config()
poller_cfg = cfg.get('poller', {})
self.poll_interval = poller_cfg.get('poll_interval', 30)
self.results_url = poller_cfg.get('results_url', 'http://localhost:5000/api/edge/results')
self.batch_size = poller_cfg.get('batch_size', 10)
self.timeout = poller_cfg.get('timeout', 30)
self._running = False
self._thread = None
def _handle_result(self, item: dict):
"""处理单个结果"""
nas_task_id = item.get('nas_task_id')
result = item.get('result')
if not nas_task_id:
logger.warning(f"结果缺少 nas_task_id跳过: {item}")
return
if result is None:
logger.error(f"[task_id={nas_task_id}] Edge 返回空结果,标记 FAILED")
db_layer.update_task_status(
nas_task_id, 'FAILED',
error_message='Edge returned empty result',
failure_stage='callback'
)
return
status = result.get('status')
if status == 'success':
try:
event_id = apply_success_event(nas_task_id, result)
log_task(logger, nas_task_id, 'poller', f'结果落库成功: event_id={event_id}')
except Exception as e:
logger.error(f"[task_id={nas_task_id}] 结果落库失败: {e}", exc_info=True)
db_layer.update_task_status(
nas_task_id, 'FAILED', error_message=str(e), failure_stage='callback')
elif status == 'failed':
error_message = result.get('error_message', 'unknown')
failure_stage = result.get('failure_stage', '')
logger.error(f"[task_id={nas_task_id}] Edge 处理失败: stage={failure_stage}, error={error_message}")
db_layer.update_task_status(
nas_task_id, 'FAILED', error_message=error_message, failure_stage=failure_stage)
else:
logger.warning(f"[task_id={nas_task_id}] 未知状态: {status}")
def _poll_once(self):
"""执行一次拉取"""
try:
resp = requests.get(
self.results_url,
params={'limit': self.batch_size},
timeout=(10, 15)
)
except requests.RequestException as e:
logger.error(f"拉取结果失败: {e}")
return
if resp.status_code != 200:
logger.warning(f"Edge 返回 {resp.status_code}")
return
try:
data = resp.json()
except ValueError:
logger.error("Edge 返回非 JSON 响应")
return
results = data.get('results', [])
if not results:
return
logger.info(f"拉取到 {len(results)} 条结果")
for item in results:
try:
self._handle_result(item)
except Exception as e:
task_id = item.get('nas_task_id', '?')
logger.error(f"[task_id={task_id}] 处理结果异常: {e}", exc_info=True)
def _run(self):
"""线程主循环"""
logger.info(f"Poller 启动,轮询间隔 {self.poll_interval}s目标: {self.results_url}")
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='poller')
self._thread.start()
def is_alive(self):
"""线程是否存活"""
return self._thread is not None and self._thread.is_alive()
def check_and_restart(self):
"""看门狗:线程崩溃后自动重启"""
if self._running and not self.is_alive():
logger.warning("Poller 线程已死亡,正在重启...")
self._thread = threading.Thread(target=self._run, daemon=True, name='poller')
self._thread.start()
def stop(self):
"""停止拉取线程"""
self._running = False
if self._thread:
self._thread.join(timeout=5)

View File

@@ -1,4 +0,0 @@
"""Task-Scheduler 包"""
from .scheduler import TaskScheduler
__all__ = ["TaskScheduler"]

View File

@@ -1,122 +0,0 @@
"""
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 is_alive(self):
"""线程是否存活"""
return self._thread is not None and self._thread.is_alive()
def check_and_restart(self):
"""看门狗:线程崩溃后自动重启"""
if self._running and not self.is_alive():
logger.warning("Scheduler 线程已死亡,正在重启...")
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)

View File

@@ -1,4 +0,0 @@
"""Video-Server 包"""
from .video_server import video_bp
__all__ = ["video_bp"]

View File

@@ -1,52 +0,0 @@
"""
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/<path:filename>', 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
)

View File

@@ -1,164 +0,0 @@
"""为历史事件补抽关键帧图(视频仍存于 NAS 时)
Edge 注入关键帧 base64 上线前落库的事件没有帧图,本脚本从原始视频
按 frame_timestamp - event_start_time 偏移重新抽帧,补齐到 UI 静态目录。
用法:
venv/bin/python tools/backfill_frames.py [--dry-run] [--event-id N]
特性:
- 幂等: 单帧文件已存在即跳过,帧数齐全的事件整条跳过
- NAS ffmpeg41 无 image2 muxer必须用 -f singlejpeg 输出 jpg
- 抽帧尺寸与 NAS 预压缩一致480p 等比缩放),-q:v 5 约 60-100KB/张
"""
import os
import re
import sys
import argparse
import subprocess
from datetime import datetime
sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), '..', 'src'))
import pymysql
from fam_core.config_loader import load_config
FFMPEG = '/var/packages/CodecPack/target/bin/ffmpeg41'
def video_duration(path: str) -> float:
"""解析 ffmpeg header 里的 Duration无 ffprobe 环境)"""
try:
r = subprocess.run([FFMPEG, '-i', path], capture_output=True, text=True, timeout=60)
m = re.search(r'Duration:\s*(\d+):(\d+):(\d+)', r.stderr)
if m:
h, mi, s = (int(x) for x in m.groups())
return h * 3600 + mi * 60 + s
except Exception:
pass
return 0.0
def extract_frame(video: str, offset: float, out_path: str) -> bool:
cmd = [
FFMPEG, '-y',
'-ss', f'{offset:.1f}',
'-i', video,
'-frames:v', '1',
'-vf', 'scale=854:480:force_original_aspect_ratio=decrease,scale=trunc(iw/2)*2:trunc(ih/2)*2',
'-q:v', '5',
'-f', 'singlejpeg',
out_path,
]
try:
r = subprocess.run(cmd, capture_output=True, timeout=120)
return r.returncode == 0 and os.path.getsize(out_path) > 1024
except Exception:
return False
def main():
ap = argparse.ArgumentParser()
ap.add_argument('--dry-run', action='store_true', help='只统计不落盘')
ap.add_argument('--event-id', type=int, default=None, help='只处理指定事件')
args = ap.parse_args()
cfg = load_config()
frame_dir = cfg.get('storage', {}).get(
'frame_image_dir', '/volume1/web/sentinel-home-ai/fam-ui/static/frames')
db = cfg['database']
conn = pymysql.connect(
host=db.get('host', '127.0.0.1'),
port=db.get('port', 3306),
user=db.get('user', 'root'),
password=db.get('password', ''),
database=db.get('database', 'sentinel_home_ai'),
unix_socket=db.get('unix_socket'),
charset='utf8mb4',
cursorclass=pymysql.cursors.DictCursor,
)
try:
with conn.cursor() as cur:
sql = """
SELECT me.event_id, me.task_id, me.event_start_time, pt.video_path,
(SELECT COUNT(*) FROM event_details ed
WHERE ed.event_id = me.event_id) AS detail_count
FROM monitor_events me
LEFT JOIN process_tasks pt ON pt.task_id = me.task_id
"""
params = ()
if args.event_id:
sql += ' WHERE me.event_id = %s'
params = (args.event_id,)
sql += ' ORDER BY me.event_id'
cur.execute(sql, params)
events = cur.fetchall()
stat = {'skip_complete': 0, 'skip_no_video': 0, 'extracted': 0, 'failed': 0, 'events': 0}
for ev in events:
eid = ev['event_id']
edir = os.path.join(frame_dir, f'event_{eid}')
existing = {f for f in os.listdir(edir) if f.endswith('.jpg')} if os.path.isdir(edir) else set()
with conn.cursor() as cur:
cur.execute(
"""SELECT frame_index, frame_timestamp FROM event_details
WHERE event_id = %s ORDER BY frame_index""",
(eid,))
frames = cur.fetchall()
missing = [f for f in frames if f'frame_{f["frame_index"]}.jpg' not in existing]
if not missing:
stat['skip_complete'] += 1
continue
video = ev.get('video_path') or ''
source = '原始视频'
if not video or not os.path.isfile(video):
# 原始视频被监控保留策略清理时,回退到 dispatcher 压缩缓存480p 副本)
cached = os.path.join(
f'/tmp/fam_compressed/task_{ev["task_id"]}', os.path.basename(video)) if video else ''
if cached and os.path.isfile(cached):
video = cached
source = '压缩缓存'
if not video or not os.path.isfile(video):
print(f'[event {eid}] 原始视频与压缩缓存均不存在,跳过 {len(missing)} 帧: {video}')
stat['skip_no_video'] += 1
continue
dur = video_duration(video)
start = ev['event_start_time']
if isinstance(start, str):
start = datetime.strptime(start[:19], '%Y-%m-%d %H:%M:%S')
os.makedirs(edir, exist_ok=True)
stat['events'] += 1
print(f'[event {eid}] 补 {len(missing)}/{len(frames)} 帧 · {source} (视频 {os.path.basename(video)}, {dur:.0f}s)')
for f in missing:
out_path = os.path.join(edir, f'frame_{f["frame_index"]}.jpg')
ts = f['frame_timestamp']
if isinstance(ts, str):
ts = datetime.strptime(ts[:19], '%Y-%m-%d %H:%M:%S')
offset = (ts - start).total_seconds()
offset = max(1.0, min(offset, max(1.0, dur - 2)))
if args.dry_run:
print(f' dry-run frame_{f["frame_index"]} @ {offset:.0f}s')
continue
ok = extract_frame(video, offset, out_path)
stat['extracted' if ok else 'failed'] += 1
if not ok:
print(f' 失败 frame_{f["frame_index"]} @ {offset:.0f}s')
if os.path.exists(out_path):
os.remove(out_path)
print(f"\n完成: 事件 {stat['events']} 个已补 | 抽帧成功 {stat['extracted']} 失败 {stat['failed']} "
f"| 齐全跳过 {stat['skip_complete']} | 视频缺失跳过 {stat['skip_no_video']}")
finally:
conn.close()
if __name__ == '__main__':
main()

View File

@@ -1,118 +0,0 @@
#!/usr/bin/env python3
"""存量关键帧批量补红框(计算在 Edge/OracleNAS 只编排与存图)
流程: 遍历 frames/event_*/frame_*.jpg -> 分批(8张)上传 Edge /api/edge/mark_frames
-> 用标记后的图覆盖原文件 -> 写 meta.json (frame_index -> face_count)
幂等: 已有 meta.json 的事件跳过;--force 强制重跑
备份: 首次覆盖前原文件备份到 frames_orig/event_*/
"""
import argparse
import base64
import json
import os
import shutil
import sys
import requests
sys.path.insert(0, os.path.join(os.path.dirname(__file__), '..', 'src'))
from fam_core.config_loader import load_config # noqa: E402
FRAME_DIR = load_config().get('storage', {}).get(
'frame_image_dir', '/volume1/web/sentinel-home-ai/fam-ui/static/frames')
MARK_URL = load_config().get('poller', {}).get(
'results_url', 'http://129.146.203.203:5000/api/edge/results'
).rsplit('/', 1)[0] + '/mark_frames'
BATCH = 8
def mark_batch(images):
"""images: [(key, jpeg_bytes)] -> {key: (marked_bytes, faces)}"""
payload = {
'images': [
{'key': k, 'data': base64.b64encode(b).decode('ascii')}
for k, b in images
]
}
resp = requests.post(MARK_URL, json=payload, timeout=(30, 120))
resp.raise_for_status()
out = {}
for r in resp.json().get('results', []):
if r.get('data'):
out[r['key']] = (base64.b64decode(r['data']), r.get('faces', 0))
else:
out[r['key']] = (None, 0)
return out
def main():
ap = argparse.ArgumentParser()
ap.add_argument('--force', action='store_true', help='忽略已有 meta.json 重跑')
ap.add_argument('--event', type=int, help='只处理指定 event_id')
args = ap.parse_args()
orig_root = os.path.join(os.path.dirname(FRAME_DIR.rstrip('/')), 'frames_orig')
stat = {'events': 0, 'frames': 0, 'faces': 0, 'skip': 0, 'fail': 0}
for name in sorted(os.listdir(FRAME_DIR)):
if not name.startswith('event_'):
continue
eid = int(name.split('_')[1])
if args.event and eid != args.event:
continue
edir = os.path.join(FRAME_DIR, name)
meta_path = os.path.join(edir, 'meta.json')
frames = sorted(
f for f in os.listdir(edir)
if f.startswith('frame_') and f.endswith('.jpg'))
if not frames:
continue
if os.path.exists(meta_path) and not args.force:
stat['skip'] += 1
continue
# 备份原件(一次)
bak_dir = os.path.join(orig_root, name)
if not os.path.isdir(bak_dir):
os.makedirs(bak_dir, exist_ok=True)
for f in frames:
shutil.copy2(os.path.join(edir, f), os.path.join(bak_dir, f))
face_counts = {}
for i in range(0, len(frames), BATCH):
batch = []
for f in frames[i:i + BATCH]:
with open(os.path.join(edir, f), 'rb') as fh:
batch.append((f, fh.read()))
try:
marked = mark_batch(batch)
except Exception as e:
print(f'[event {eid}] 批次失败 (跳过 {len(batch)} 帧): {e}')
stat['fail'] += len(batch)
continue
for f, _ in batch:
data, faces = marked.get(f, (None, 0))
if data is None:
stat['fail'] += 1
continue
with open(os.path.join(edir, f), 'wb') as fh:
fh.write(data)
idx = f[len('frame_'):-len('.jpg')]
face_counts[idx] = faces
stat['frames'] += 1
stat['faces'] += faces
if face_counts:
with open(meta_path, 'w') as fh:
json.dump(face_counts, fh)
stat['events'] += 1
print(f'[event {eid}] 标记 {len(face_counts)}/{len(frames)} 帧, '
f'人脸合计 {sum(face_counts.values())}')
print(f"\n完成: 事件 {stat['events']} | 帧标记 {stat['frames']} "
f"(含人脸帧人脸数 {stat['faces']}) | 已标跳过 {stat['skip']} | 失败 {stat['fail']}")
print(f'备份目录: {orig_root}')
if __name__ == '__main__':
main()

View File

@@ -1,5 +1,5 @@
# FAM-UI 配置文件 (NAS 端) - 实际部署配置
# Tailscale: NAS=100.70.234.39
# FAM-UI 配置文件 (NAS 端) - 新架构 v22026-08-21
# NAS 仅作管理后台,前端读本地 MariaDB 同步镜像,不再读取关键帧图片。
core_url: "http://127.0.0.1:8000"
@@ -10,7 +10,3 @@ database:
password: "iLoveJava5!"
database: "sentinel_home_ai"
unix_socket: "/run/mysqld/mysqld10.sock"
storage:
# 关键帧目录fam-core event_receiver 落盘UI 读取展示时间轴)
frame_image_dir: "/volume1/web/sentinel-home-ai/fam-ui/static/frames"

File diff suppressed because it is too large Load Diff

View File

@@ -121,6 +121,67 @@ CREATE TABLE IF NOT EXISTS family_members (
INDEX idx_is_active (is_active)
) ENGINE=InnoDB COMMENT='家庭成员表(交互式命名)';
-- ============================================================
-- 7. 甲骨文同步镜像表(新架构 v22026-08-21
-- NAS 每 30 分钟从甲骨文 FAM-Edge 拉增量,镜像到本地,仅作展示
-- 字段对齐 Oracle 端 SQLite 库oracle_db.py
-- ============================================================
-- 7.1 视频会话表Oracle videos 镜像)
CREATE TABLE IF NOT EXISTS sync_videos (
id INT PRIMARY KEY COMMENT 'Oracle videos.id',
drive_file_id VARCHAR(255) COMMENT 'Google 硬盘文件 ID',
filename VARCHAR(500) NOT NULL UNIQUE COMMENT '视频文件名(唯一)',
camera_name VARCHAR(50) COMMENT '摄像头名称/位置',
duration_sec DOUBLE DEFAULT 0 COMMENT '视频时长(秒)',
event_start_time VARCHAR(32) COMMENT '视频开始时间(文本)',
status VARCHAR(20) DEFAULT 'pending' COMMENT 'pending/done/failed',
summary_json LONGTEXT COMMENT '全局摘要文本',
events_json LONGTEXT COMMENT '事件列表 JSON 数组(冗余,便于查询)',
people_json LONGTEXT COMMENT '人物列表 JSON 数组',
compute_provider VARCHAR(255) COMMENT '模型来源,如 gemini / nvidia',
created_at VARCHAR(32),
updated_at VARCHAR(32),
processed_at VARCHAR(32),
synced_at DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '最近一次同步写入时间',
INDEX idx_processed (processed_at),
INDEX idx_status (status)
) ENGINE=InnoDB COMMENT='甲骨文视频会话镜像表';
-- 7.2 事件明细表Oracle events 镜像)
CREATE TABLE IF NOT EXISTS sync_events (
id INT PRIMARY KEY COMMENT 'Oracle events.id',
video_id INT NOT NULL COMMENT '关联 sync_videos.id',
ts VARCHAR(32) COMMENT '事件时间点(文本)',
description TEXT COMMENT '事件描述',
person_list_json LONGTEXT COMMENT '涉及人物 JSON 数组(字符串或标签)',
is_attention_event TINYINT(1) DEFAULT 0 COMMENT 'AI 判断是否为关注事件',
updated_at VARCHAR(32),
synced_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_video (video_id),
INDEX idx_ts (ts)
) ENGINE=InnoDB COMMENT='甲骨文事件镜像表';
-- 7.3 人物规范表Oracle people 镜像)
CREATE TABLE IF NOT EXISTS sync_people (
id INT PRIMARY KEY COMMENT 'Oracle people.id',
label VARCHAR(100) NOT NULL UNIQUE COMMENT '抽象标识,如"人物A"',
canonical_name VARCHAR(100) COMMENT '规范名(用户命名或 LLM 合并NULL 表示未命名',
first_seen VARCHAR(32) COMMENT '首次出现时间',
appearances INT DEFAULT 0 COMMENT '出现次数',
source VARCHAR(20) DEFAULT 'llm' COMMENT 'llm / manualmanual 优先不被覆盖)',
updated_at VARCHAR(32),
synced_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_label (label),
INDEX idx_canonical (canonical_name)
) ENGINE=InnoDB COMMENT='甲骨文人物规范镜像表';
-- 7.4 同步游标表
CREATE TABLE IF NOT EXISTS sync_cursor (
key VARCHAR(50) PRIMARY KEY,
value VARCHAR(64) COMMENT '上次成功拉取到的 server_timeISO 文本)'
) ENGINE=InnoDB COMMENT='同步游标表';
-- ============================================================
-- 验证
-- ============================================================