diff --git a/PROGRESS.md b/PROGRESS.md index 38ba586..c5376c8 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -541,3 +541,74 @@ AI 分析瓶颈: 帧5耗时 229s (疑似 ARM CPU 热降频), 其余帧 50-65s - `/api/ui/videos`(无 cookie)→ 401 → 前端据此自动跳 `/login` 登录链路已通:未登录访问任意页面 → 首个 401 → 自动跳 auth-hub 登录 → 回调种 `fam_session` cookie → 回 `/timeline` 正常加载。 + +## 统一登录从 NAS 迁到甲骨文 fam-edge(2026-09-12) + +**触发**:`smart-camera.zichuan.xyz/login` 报 502。排查链路:静态页 `/`、`/timeline` 200(Caddy 读本机磁盘), +只有 `/login` 和 `/api/*` 502;甲骨文本机 `curl 127.0.0.1:8000/health` 0.6 秒空响应(curl exit 52, +Caddy 日志 `msg:"EOF"`),frps 正常、同隧道的 gitea 也正常 → NAS 上 fam-core 进程没了。 +局域网直扫 NAS:22/2222/3000/5000/5001/3306 全 OPEN,**只有 8000 closed**,NAS 本身没事。 +根因是登录入口挂在 NAS 上,而 NAS 上的 fam-core 没有任何守护(DSM 无 systemd),挂了不会自启。 + +**决策**:登录是入口,不该依赖家里的机器。整体迁到甲骨文的 fam-edge(前端静态文件、auth-hub 本来就在这台), +NAS fam-core 退化成纯数据接口。 + +**改动**: +- 新增 `fam-edge/src/fam_edge/auth.py`:`/login`、`/api/auth/callback`、`/api/logout`、 + `/api/auth/check`、`/api/auth/verify`(给 Caddy forward_auth)。跟旧实现三处关键差异: + 1. 换 token / 拉 JWKS 走 `AUTH_HUB_INTERNAL_BASE`(本机 :5300),不再跨公网 TLS——旧链路上 + `PyJWKClient` 用 urllib + 系统 CA(群晖易 CERTIFICATE_VERIFY_FAILED)、两机时钟偏差会让 `iat` + 显得来自未来,这两个坑一起消失;但 `iss` 校验和浏览器跳转仍用公网 issuer + 2. 会话改无状态 HS256 签名 cookie,服务重启不掉线(旧实现进程内 token 表) + 3. **回调失败渲染错误页,不再 302 回 `/login`**——旧实现失败即跳 `/login`,而 auth-hub 只要还有 + 会话就立刻再签一个 code 跳回来,两边对跳成死循环,浏览器只报「重定向次数过多」, + 既看不到登录页也看不到原因(这正是 9/1 那次「跳不到登录页」的成因) +- 删除 `fam-core/src/fam_core/auth.py` + `tests/test_auth.py`,`app.py` 去掉 `init_auth`; + fam-core 不再有任何鉴权,改由甲骨文 Caddy `forward_auth` 前置拦截 +- `fam-edge/tests/test_auth.py` 16 个用例(含「失败分支绝不 302」的回归测试、 + 「服务端走内网但 iss 按公网校验」、「cookie 无状态」);fam-edge 157 / fam-core 15 全绿 +- 写测试时逮到自己写的一个 bug:会话 cookie 也套了 60 秒 leeway(本进程自签自验根本不需要), + 会让每个会话白白多活 60 秒,已改成只有 id_token 用 leeway + +**遗留风险(已记入 README/DEPLOY)**:frps 在甲骨文绑的是 `*:8000` 且 iptables 明确放行, +fam-core 去掉鉴权后,绕过 Caddy 直连 `129.146.26.249:8000` 就是无门禁的全量数据接口, +必须靠 `iptables -I INPUT 1 -p tcp --dport 8000 ! -i lo -j DROP` 兜底。 +另:登录进程并到 fam-edge 后,fam-edge 正在跑视频分析时登录响应可能变慢(单 worker 4 线程)。 + +## 全量迁云:NAS 只剩推送进程(2026-09-13) + +**触发**:前一天刚把登录迁到甲骨文,隔天 NAS 上的 fam-core 又挂了导致数据接口 502。 +用户一句话点破:「NAS 上只有一个无状态的通知甲骨文的服务,剩下的全部在甲骨文呀」。 +盘点后确认这个判断成立——甲骨文的 SQLite 才是权威数据源(videos 3113 / events 16159 / +people 60 / model_calls 9876),NAS 的 MariaDB 全是它的镜像,前端读的数据本来就产自 +甲骨文,绕了一圈回家又绕回来。 + +**改动**: +- `fam-core/db_layer.py` 从 725 行重写成 377 行:MySQL 镜像查询 → 直读 fam-edge 的 + SQLite。5 个 `upsert_sync_*`(约 300 行去重逻辑,8/29 和 9/3 两次 1062 事故的发源地) + 连同 `oracle_sync.py` 整个删除。SQL 方言:`JSON_CONTAINS` → `json_each`(前置 + `json_valid`,历史脏数据不会把查询搞崩)、`LEFT()` → `substr()`、`%s` → `?`。 + 函数名 `get_sync_*` 一并改掉——已经没有 sync 这回事了,留着名字会误导。 +- 新增 `fam-core/edge_client.py`:写操作(改名/删除)、帧图头像、服务状态都打给同机 + fam-edge,全走 127.0.0.1,不出公网。 +- 新增 `fam-notifier/`:把 `motion_notifier.py` 从 fam-core 拆出来,游标从 MariaDB + 换成本地 JSON 文件。NAS 上从此没有 Flask、没有数据库、没有监听端口,只有一个 + 单向推送进程。 +- fam-core 移到甲骨文 `/opt/fam-core`(systemd,gunicorn -w 2,**只绑 127.0.0.1:5401**—— + 5400 被 chat-relay 占了)。Caddy 的 `/api/*` 从"frp 隧道回源 NAS"改成同机反代, + forward_auth 闸门不变。 +- 前端跟着删:侧边栏同步面板、统计页同步状态、服务状态页的"NAS 同步"卡片和 + "立即同步"按钮(背后的镜像层已不存在)。"NAS 同步"换成"NAS 运动推送",读 + fam-edge activity 新增的 `motion` 段(心跳年龄 + 最近事件)。 +- 顺带修掉一个隐蔽 bug:镜像表为保外键稳定用的是 NAS 本地自增 id,而帧图接口要的是 + 甲骨文的 id,两边在 9/3 那次 id 重排后就对不上了。现在只有一套 id。 + +**测试**:fam-core 21(新增 12 个 db_layer 用例:脏 JSON 不崩、人物精确匹配不误伤 +"人物B"、日期过滤、统计口径、chat_history 懒建表)、fam-notifier 6、fam-edge 157,全绿。 + +**验证**:甲骨文 `/api/ui/stats` 返回 videos 2965 / events 16159 / attention 13 / +people 59;外网 `/` `/timeline` 200、`/login` 302、`/api/*` 未登录 401——**全程 NAS +上的 fam-core 是停着的**,这就是迁云的验收标准。 + +**待办**:NAS 侧部署 fam-notifier(只能用户手动,密码登录);chat_history 一次性迁移 +(`fam-core/scripts/import_chat_history.py`,幂等);frpc.toml 里的 8000 映射可删。 diff --git a/README.md b/README.md index 954838e..d8d291f 100644 --- a/README.md +++ b/README.md @@ -63,9 +63,9 @@ Orchestrator 视觉阶段按 `fallback` 模式顺序降级:Gemini → NVIDIA N | 节点 | 角色 | 硬件 | IP | 服务与端口 | |------|------|------|-----|-----------| -| Synology NAS | FAM-Core + 数据库(采集 + 聚合 + 镜像层) | DS220+ (Geminilake), DSM 7 | 家庭局域网 192.168.50.64 | FAM-Core :8000(API 聚合 + 轮询 SS + OracleSync), MariaDB :3306, Surveillance Station :5000;FAM-UI 已迁至云服务器(见下) | -| Oracle Cloud(云服务器) | FAM-Edge + AI-Gateway + Caddy(前端) | Ampere A1 4C23G ARM64(无 GPU), Ubuntu 20.04 | 公网 129.146.26.249 / Tailscale(已安装未启用,备用) | FAM-Edge :5000(systemd 守护,视频分析 + 转发问答), **AI-Gateway :5100(systemd 守护,对外监听,问答模型降级链)**, Caddy :80(托管 FAM-UI SPA + 反代 /api/* → NAS FAM-Core), Ollama :11434(仅本地,被 AI-Gateway 调用) | -| 家庭网络 | 用户入口 | 普通终端 | 公网 / 家庭网络 | 浏览器访问 `http://129.146.26.249/`(Caddy 托管前端,/api 经 frp 隧道回源 NAS FAM-Core :8000) | +| Synology NAS | 摄像头录像 + 运动事件推送(**仅此而已**,2026-09-13 起) | DS220+ (Geminilake), DSM 7 | 家庭局域网 192.168.50.64 | Surveillance Station :5000(录像机本体,搬不走), **fam-notifier**(无端口的推送进程,轮询 SS → 推甲骨文)。MariaDB 与 FAM-Core 已随镜像层一起下线 | +| Oracle Cloud(云服务器) | **除录像外的全部**:视频分析 + 数据 + 接口 + 登录 + 前端 | Ampere A1 4C23G ARM64(无 GPU), Ubuntu 20.04 | 公网 129.146.26.249 / Tailscale(已安装未启用,备用) | FAM-Edge :5000(systemd,视频分析 + 统一登录 + SQLite 权威库), **FAM-Core :5401(systemd,UI/对话/成员接口,仅监听 127.0.0.1)**, AI-Gateway :5100(systemd,问答模型降级链), Caddy :80/:443(前端 + 反代 + forward_auth 鉴权), Ollama :11434(仅本地) | +| 家庭网络 | 用户入口 | 普通终端 | 公网 / 家庭网络 | 浏览器访问 `https://smart-camera.zichuan.xyz/`(全部由甲骨文提供,**不依赖 NAS 在线**) | **网络要点(推送模式)**: - 服务间通信只有两条出站:**NAS → Oracle 公网 IP:5000**(① `GET /api/oracle/sync` 拉增量 + `POST /api/oracle/people/correct` 命名回推 ② `POST /api/ss/motion` 推送运动侦测事件)。Edge 不需要反向访问 NAS @@ -98,19 +98,22 @@ Orchestrator 视觉阶段按 `fallback` 模式顺序降级:Gemini → NVIDIA N │ ├─ /v1/chat/completions:NVIDIA → Gemini(多 Key 轮换)→ 本地 Ollama 降级链 │ │ ├─ Bearer token 鉴权(fail-closed),对外 0.0.0.0:5100,非本系统专属 │ │ └─ 独立 .env(/opt/ai-gateway/.env,NVIDIA/Gemini key 与 FAM-Edge 各自一份) │ -│ Caddy :80:托管 FAM-UI SPA + 反代 /api/* → NAS FAM-Core(frp 隧道 :8000) │ -└───────────────────────────────┬───────────────────────────────────────────────┘ - │ HTTP GET /api/oracle/sync?since=&token= (每 30 分钟) - ▼ +│ FAM-Core (Flask :5401,仅 127.0.0.1) │ +│ ├─ UI-API: /api/ui/*(直读下面那个 SQLite,无镜像、无后台线程) │ +│ ├─ Chat-Handler: /api/chat/ask(查 events 拼上下文 → 转 fam-edge 编排) │ +│ ├─ Member-Manager: /api/member/name|merge(转 fam-edge,它负责合并人物) │ +│ └─ chat_history: 本服务唯一写的表,落在同一个 SQLite 里 │ +│ SQLite /opt/fam-edge/data/oracle.db(WAL):videos / events / people / │ +│ model_calls / person_identity_map / ss_motion_events / chat_history │ +│ Caddy :80/:443:前端静态 + /login 与 /api/auth/* → fam-edge │ +│ + /api/* 先 forward_auth 再反代 fam-core(唯一鉴权闸门) │ +└───────────────────────────────▲───────────────────────────────────────────────┘ + │ POST /api/ss/motion(单向推送 + 5 分钟心跳) + │ ┌─────────────────────────────── 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 + 即时拉回) │ -│ ▼ │ -│ MariaDB (sentinel_home_ai): sync_videos / sync_events / sync_people / │ -│ sync_cursor / chat_history │ -│ (FAM-UI 已迁至云服务器 Caddy :80,见上) │ +│ Surveillance Station :5000(录像机本体,摄像头插在这台上) │ +│ fam-notifier:轮询 SS EventCenter → 推运动事件给甲骨文。无端口、无数据库, │ +│ 游标存本地 JSON 文件。挂了重启即可,不影响网站。 │ └─────────────────────────────────────────────────────────────────────────────────┘ ``` @@ -315,15 +318,23 @@ chat_history 独立表(问答上下文摘要留存) **GET /api/oracle/avatar** — 人物头像:`?label=&w=` → jpeg(从该人物候选事件 `person_appearances` bbox 裁剪) **GET /health** — 服务状态(含已处理视频数) -### 5.2 FAM-Core(NAS :8000) +**登录相关端点**(2026-09-12 从 NAS FAM-Core 迁入,见 `fam-edge/src/fam_edge/auth.py`) + +| 端点 | 方法 | 说明 | +|------|------|------| +| `/login` | GET | 登录入口:生成 PKCE 参数后 302 跳 auth-hub 公网 `/authorize` | +| `/api/auth/callback` | GET | auth-hub 回跳:用 `AUTH_HUB_INTERNAL_BASE`(本机 :5300)换 token + 拉 JWKS 验签,成功后种 HttpOnly cookie `fam_session`(HS256 无状态签名,2 小时);**任何一步失败都渲染错误页,绝不 302 回 `/login`**——旧实现失败即跳 `/login`,而 auth-hub 有会话时会立刻再签发 code,两边对跳成死循环 | +| `/api/logout` | POST | 退出登录(清 cookie,不影响 auth-hub 上的 SSO 会话) | +| `/api/auth/check` | GET | 登录态检查:`{"authed": true\|false, "username": "..."}` | +| `/api/auth/verify` | ANY | 给甲骨文 Caddy 的 `forward_auth` 用:已登录 204,未登录 401 JSON | + +### 5.2 FAM-Core(甲骨文 127.0.0.1:5401) + +> 2026-09-13 起 FAM-Core 与 FAM-Edge 同机,直读它的 SQLite(不再有 MariaDB 镜像层),且**只监听 127.0.0.1**——唯一客户端是同机 Caddy。登录端点在 FAM-Edge(见 §5.1),鉴权由 Caddy 的 `forward_auth` 前置完成,所以本服务自身不做任何校验。 | 端点 | 方法 | 说明 | |------|------|------| | `/health` | GET | 服务健康(**免登录**) | -| `/login` | GET | 登录入口:302 跳转 auth-hub `/authorize`(Authorization Code + PKCE,**免登录**) | -| `/api/auth/callback` | GET | auth-hub 登录回跳:拿 `code` 换 token、验 `id_token` 签名后种 HttpOnly cookie `fam_session`(2 小时,**免登录**) | -| `/api/logout` | POST | 退出登录(清本地 cookie,不影响 auth-hub 上的登录态) | -| `/api/auth/check` | GET | 登录态检查:`{"authed": true\|false}`(**免登录**) | | `/api/status` | GET | Oracle-Sync + MotionNotifier 状态(sync running/cursor;motion poll_enabled/pushed_total/heartbeat) | | `/api/chat/ask` | POST | 用户问答:`{"question","queried_person","queried_date"}` → `{"answer","context_summary","chat_id"}`(上下文来自 sync_events) | | `/api/chat/history` | GET | 对话历史(`?date=` 或 `?person=&limit=`) | @@ -540,17 +551,27 @@ task_id=289(30s 测试片段)全链路打通:推送 5.7MB → Edge 分析 | 组件 | 节点 | 路径 | 启动 | |------|------|------|------| -| FAM-Core | NAS | `/volume1/web/sentinel-home-ai/fam-core/` | `bash start_core.sh`(gunicorn -w 1 :8000,source .env 注入 DSM_*/ORACLE_SYNC_TOKEN/AUTH_HUB_*) | +| FAM-Core | Oracle | `/opt/fam-core/` | **systemd `fam-core.service`**(gunicorn -w 2 --threads 4,绑 127.0.0.1:5401);部署代码后 `sudo systemctl restart fam-core` | +| fam-notifier | NAS | `/volume1/web/sentinel-home-ai/fam-notifier/` | `setsid bash scripts/start_notifier.sh &`(DSM 无 systemd,挂了不自启);NAS 上唯一运行的本项目进程 | | FAM-UI | 云服务器(甲骨文,129.146.26.249) | `/var/www/fam-ui/` | Vue3 构建产物,由 Caddy :80 静态托管(无需独立进程);本地改代码后 `npm run build`,`dist/` rsync/tar 到云服务器 | | FAM-Edge | Oracle | `/opt/fam-edge/` | **systemd `fam-edge.service` 守护**(Restart=always);部署代码后 `sudo systemctl restart fam-edge`(勿手动 setsid,会端口冲突) | | **AI-Gateway** | Oracle | `/opt/ai-gateway/` | **systemd `ai-gateway.service` 守护**(Restart=always),独立 venv(Python 3.8);gunicorn 绑定 `0.0.0.0:5100`(对外直接开放,非仅本机);部署代码后 `sudo systemctl restart ai-gateway` | | Ollama | Oracle | systemd 托管 | 环境变量 `OLLAMA_KEEP_ALIVE=-1`;**被 AI-Gateway 调用,不再被 FAM-Edge 调用** | > **外网访问(frp 内网穿透)**:NAS 跑 `frpc`(`/etc/frp/frpc.toml`,S99frpc.sh 守护),映射到云服务器 `129.146.26.249`(frps :7000): -> - `3000` → NAS Gitea、`8500` → NAS WordPress(8088)、**`8000` → NAS FAM-Core(本系统)** -> - 前端入口 `http://129.146.26.249/`(Caddy 托管 FAM-UI,/api 经 frp 隧道反代回 NAS FAM-Core :8000),**需登录**(见 §5.2 `/login`) +> - `3000` → NAS Gitea、`8500` → NAS WordPress(8088) +> - **`8000` → NAS FAM-Core 这条已于 2026-09-13 迁云后废弃**(本系统不再有任何"回源 NAS"的流量,可从 frpc.toml 删除) +> - 前端入口 `https://smart-camera.zichuan.xyz/`:静态页、接口、登录全部由甲骨文提供,**NAS 离线也能正常访问**(只是不再有新的运动事件推进来),**需登录**(见 §5.1 `/login`) > -> **登录校验(2026-08-22 新增,2026-08-23 改为 fail-closed,2026-08-31 接入 auth-hub 统一登录)**:FAM-Core 不再自己保存/校验密码,全权委托给独立部署的 [auth-hub](http://129.146.26.249:5300)(OAuth2 Authorization Code + PKCE / OIDC,同一 Oracle 主机 :5300,与本系统完全独立部署)——`/login` 302 跳到 auth-hub `/authorize`,登录后带 `code` 跳回本服务 `/api/auth/callback`,服务端换 token、验完 `id_token` 签名(JWKS)后种本地会话 cookie。全站拦截仍在 `auth.py`:页面未登录 302 `/login`,`/api/*` 未登录 401;接入参数 `AUTH_HUB_ISSUER`/`AUTH_HUB_CLIENT_ID`/`AUTH_HUB_CLIENT_SECRET`/`AUTH_HUB_REDIRECT_URI`(NAS `.env` 配置,**无硬编码默认值**,任一没配置直接拒绝所有登录,fail closed);`AUTH_HUB_REDIRECT_URI` 必须与 auth-hub 用 `manage_clients create` 登记的 redirect_uri 逐字符一致(auth-hub 只做精确匹配,不做前缀/子串匹配)。白名单免登录:`/login` `/api/auth/callback` `/api/logout` `/api/auth/check` `/health` `/assets/*`(`/api/ss/webhook` 已随端点一起移除,2026-08-25)。本地会话仍是进程内 token + HttpOnly cookie(**2 小时**),过期或重启 fam-core 需重新走一遍 auth-hub 登录;auth-hub 侧注册是任何人都能自助注册(管理员审批后才能登录),**未对 fam-core 侧再加用户名白名单**——任何在 auth-hub 上审批通过的账号登录后都能访问本系统,这是有意识的取舍(自用场景,见项目记忆)。 +> **登录校验(2026-08-22 新增;2026-08-31 接入 auth-hub 统一登录;2026-09-12 整体迁到甲骨文 FAM-Edge)**:登录流程和会话校验现在都在甲骨文 FAM-Edge 的 `auth.py`(§5.1),NAS FAM-Core 自身不再做任何鉴权。迁移的直接原因:NAS 或 frp 隧道一挂,`smart-camera.zichuan.xyz/login` 直接 502,**连登录页都打不开**——登录是入口,不该依赖家里那台机器。 +> +> - **链路**:`/login`(Caddy → 本机 fam-edge :5000)→ 302 到 auth-hub 公网 `/authorize`(Authorization Code + PKCE)→ 用户在 [auth-hub](https://auth.zichuan.xyz) 登录 → 回跳 `/api/auth/callback`(仍是 fam-edge)→ 走**本机** `http://127.0.0.1:5300` 换 token、拉 JWKS 验 `id_token` 签名 → 种 cookie `fam_session` +> - **为什么服务端调用走本机**:旧实现是 NAS 跨公网访问 `https://auth.zichuan.xyz`,那条链路上 `jwt.PyJWKClient` 用 urllib + 系统 CA(不像 requests 自带 certifi),群晖上容易 `CERTIFICATE_VERIFY_FAILED`;两台机器的时钟偏差还会让 `iat` 看起来来自未来。同机直连把这两个坑一起消掉。但 `iss` 校验和浏览器跳转仍用公网 `AUTH_HUB_ISSUER` +> - **拦截**:甲骨文 Caddy 对 `/api/*` 做 `forward_auth` → fam-edge `/api/auth/verify`,200/204 才反代回 NAS;`/login`、`/api/auth/*`、`/api/logout` 直接由 fam-edge 处理(Caddy 按路径具体程度排序,这几条稳定排在 `/api/*` 前面) +> - **NAS :8000 必须靠防火墙兜底**:frps 把它暴露在公网(历史入口 `http://129.146.26.249:8000`),绕过 Caddy 直连就没有任何鉴权。甲骨文 iptables 只允许本机访问 :8000,见 `docs/DEPLOY.md` §2.3 +> - **参数**(`/opt/fam-edge/.env`,**无硬编码默认值**,缺任一项拒绝所有登录 fail closed):`AUTH_HUB_ISSUER`(公网,浏览器跳转 + `iss` 校验)/ `AUTH_HUB_INTERNAL_BASE`(本机,换 token + JWKS)/ `AUTH_HUB_CLIENT_ID` / `AUTH_HUB_CLIENT_SECRET` / `AUTH_HUB_REDIRECT_URI`(必须跟 auth-hub 登记的逐字符一致,只做精确匹配)/ `FAM_SESSION_SECRET`(会话签名密钥,≥32 字节) +> - **会话是无状态签名 cookie**(HS256,2 小时),fam-edge 重启不掉线(旧实现是进程内 token 表,fam-core 一重启全员下线);要强制全员下线就换掉 `FAM_SESSION_SECRET` 再重启 +> - auth-hub 侧任何审批通过的账号登录后都能访问本系统,**未再加用户名白名单**——自用场景的有意取舍(见项目记忆) > Oracle 部署方式:本地 git 提交 push Gitea → tar 管道到 `/opt/fam-edge`(`--strip-components=1` 解临时目录再 cp,避免动 data/venv/gdrive_videos)。**AI-Gateway 是独立 git 仓库**(http://192.168.50.64:3000/ericwyuan/ai-gateway,见 §10.4),同样 tar 管道部署到 `/opt/ai-gateway`,互不影响。 diff --git a/docs/DEPLOY.md b/docs/DEPLOY.md index fc6f894..030e4b8 100644 --- a/docs/DEPLOY.md +++ b/docs/DEPLOY.md @@ -1,54 +1,58 @@ -# 部署指南(v3 运动事件驱动架构,2026-08-22 同步) +# 部署指南(2026-09-13 全量迁云后同步) -## 1. NAS 端部署 (FAM-Core + MariaDB;FAM-UI 见 §1.3,部署在云服务器) +> **一句话架构**:除了录像本身,全部跑在甲骨文。NAS 只剩 Surveillance Station +> 和一个往甲骨文推运动事件的进程(fam-notifier)——NAS 离线不影响网站访问。 -### 1.1 MariaDB +## 1. NAS 端部署(只有 fam-notifier) + +2026-09-13 迁云后 NAS 上**不再有** MariaDB、FAM-Core、FAM-UI:镜像层整个删除 +(数据权威源本来就在甲骨文),接口和前端也都搬走了。旧部署的清理见 §1.2。 + +### 1.1 fam-notifier(轮询 Surveillance Station → 推甲骨文) ```bash -# 通过 Synology 套件中心安装 MariaDB(10.11),unix_socket=/run/mysqld/mysqld10.sock - -# 执行 DDL -python3 scripts/init_db.py --host 127.0.0.1 --port 3306 --user root --password <密码> -``` - -### 1.2 FAM-Core -```bash -cd /volume1/web/sentinel-home-ai/fam-core +cd /volume1/web/sentinel-home-ai/fam-notifier python3 -m venv venv source venv/bin/activate -pip install -r requirements.txt +pip install -r requirements.txt # 只有 requests + PyYAML -# 配置 cp config/config.yaml.example config/config.yaml -# 编辑 config.yaml:数据库密码、oracle_sync.base_url/token、motion_notifier(DSM 凭据走 .env) +# 编辑 config.yaml:camera_ids、oracle_base_url;DSM 凭据与 token 走仓库根 .env +# 仓库根 .env 需提供 DSM_ACCOUNT / DSM_PASSWORD / ORACLE_SYNC_TOKEN -# 仓库根 .env 提供 DSM_ACCOUNT/DSM_PASSWORD/ORACLE_SYNC_TOKEN/AUTH_HUB_*(start_core.sh 会 source) +# 启动(DSM 没有 systemd,用 setsid 脱离 SSH 会话;挂了不会自启,需要人工重启) +cd /volume1/web/sentinel-home-ai/fam-notifier +setsid bash scripts/start_notifier.sh >> logs/start.log 2>&1 & -# 登录改接 auth-hub 统一登录(2026-08-31),先在 auth-hub 侧注册本站点 client: -# ssh 到 auth-hub 所在 Oracle 主机,cd /opt/auth-hub -# python -m auth_hub.manage_clients create "FAM-Core" "http://129.146.26.249/api/auth/callback" -# 输出的 client_id/client_secret 明文只显示这一次,立即抄进 .env -# 仓库根 .env 补齐四项(缺任一项 fam-core 直接拒绝所有登录,fail closed): -# AUTH_HUB_ISSUER=http://129.146.26.249:5300 -# AUTH_HUB_CLIENT_ID=<上一步输出> -# AUTH_HUB_CLIENT_SECRET=<上一步输出> -# AUTH_HUB_REDIRECT_URI=http://129.146.26.249/api/auth/callback # 必须跟注册时的 redirect_uri 逐字符一致 - -# 启动(含 Oracle-Sync + MotionNotifier 轮询) -cd /volume1/web/sentinel-home-ai && bash start_core.sh +# 确认在跑(进程没有监听端口,只能看进程和日志) +ps aux | grep "[f]am_notifier" +tail -f logs/fam-notifier.log # 应每 60s 轮询、每 5min 心跳 ``` -### 1.3 FAM-UI(Vue3 SPA,托管于云服务器 Caddy :80) +游标存在 `fam-notifier/data/cursor.json`(迁云前存在 MariaDB)。删掉也不致命: +窗口回看会把最近的事件补推一遍,甲骨文按 event_id 幂等落库。 + +### 1.2 清理旧部署(迁云一次性动作) +```bash +# 1) 停掉 NAS 上的 fam-core(它已经不该再跑了;跑着也没用,前端不再连它) +ps aux | grep "[f]am-core/venv/bin/gunicorn" # 字符类写法:pkill 会杀掉 SSH 自己 +kill <上面的 pid> + +# 2) frpc 里的 8000 映射可以删了(/etc/frp/frpc.toml 的 fam-core 段), +# 甲骨文侧对应的 iptables DROP 规则也就没有存在意义了 +# 3) MariaDB 的 sentinel_home_ai 库确认 chat_history 已迁移(§2.4)后再考虑删 +``` + +### 1.3 FAM-UI(Vue3 SPA,构建后托管在甲骨文 Caddy) ```bash cd /Users/ericwyuan/Desktop/Work/sentinel-home-ai/fam-ui # 本地开发机 npm install npm run build # 产物 fam-ui/dist/ -# 部署 dist 到云服务器(Caddy 静态目录,/api 经 frp 隧道反代回 NAS): -tar czf - fam-ui/dist | ssh -i ~/.ssh/oracle_new ubuntu@129.146.26.249 \ - 'sudo tar xzf - --strip-components=2 -C /var/www/fam-ui && sudo chown -R ubuntu:ubuntu /var/www/fam-ui' -# 浏览器访问 http://129.146.26.249/(Caddy 托管,/api/* 反代 NAS fam-core) +tar czf - -C fam-ui dist | ssh -i ~/.ssh/oracle_new ubuntu@129.146.26.249 \ + 'sudo tar xzf - --strip-components=1 -C /var/www/fam-ui && sudo chown -R ubuntu:ubuntu /var/www/fam-ui' +# 浏览器访问 https://smart-camera.zichuan.xyz/ ``` -## 2. Oracle 端部署 (FAM-Edge + Ollama + FFmpeg) +## 2. Oracle 端部署 (FAM-Edge + FAM-Core + Ollama + FFmpeg) ### 2.1 系统依赖 ```bash @@ -74,6 +78,89 @@ sudo systemctl enable fam-edge sudo systemctl restart fam-edge # 部署代码后必须用 systemctl 重启,勿手动 setsid ``` +### 2.3 统一登录(2026-09-12 从 NAS 迁入,fam-edge + Caddy + 防火墙三件套) + +登录入口不再依赖 NAS。三处缺一不可: + +```bash +# (1) /opt/fam-edge/.env 追加六项(缺任一项 fam-edge 拒绝所有登录,fail closed) +# client_secret 用 rotate-secret 现拿,明文只显示一次: +# cd /opt/auth-hub && venv/bin/python -m auth_hub.manage_clients rotate-secret +# 会话密钥自己生成:openssl rand -hex 32 +export AUTH_HUB_ISSUER=https://auth.zichuan.xyz # 公网:浏览器跳转 + id_token 的 iss 校验 +export AUTH_HUB_INTERNAL_BASE=http://127.0.0.1:5300 # 本机:换 token + 拉 JWKS,不走公网 TLS +export AUTH_HUB_CLIENT_ID= +export AUTH_HUB_CLIENT_SECRET= +export AUTH_HUB_REDIRECT_URI=https://smart-camera.zichuan.xyz/api/auth/callback # 与 auth-hub 登记的逐字符一致 +export FAM_SESSION_SECRET= # 会话 cookie 签名密钥,换掉即全员下线 + +# (2) venv 补依赖(Python 3.8,pip 会自动选到兼容版本)后重启 +/opt/fam-edge/venv/bin/pip install "PyJWT>=2.8.0" "cryptography>=42.0.0" +sudo systemctl restart fam-edge + +# (3) Caddy:登录端点走本机 fam-edge,其余 /api/* 先 forward_auth 再回源 NAS +# 改完先 caddy validate 再 reload,见下方 Caddyfile 片段 +sudo caddy validate --adapter caddyfile --config /etc/caddy/Caddyfile && sudo systemctl reload caddy +``` + +```caddyfile +smart-camera.zichuan.xyz { + handle /login { reverse_proxy 127.0.0.1:5000 } + handle /api/auth/* { reverse_proxy 127.0.0.1:5000 } + handle /api/logout { reverse_proxy 127.0.0.1:5000 } + handle /api/* { + forward_auth 127.0.0.1:5000 { uri /api/auth/verify } # 删了等于数据接口全裸 + reverse_proxy 127.0.0.1:8000 + } + handle { root * /var/www/fam-ui; encode gzip; try_files {path} /index.html; file_server } +} +``` + +```bash +# (4) 防火墙:frps 把 NAS fam-core 的 :8000 暴露在公网,而 fam-core 自身已无鉴权, +# 必须只放行本机(Caddy)访问,否则绕过 Caddy 就能读到全部数据 +sudo iptables -I INPUT 1 -p tcp --dport 8000 ! -i lo -j DROP +sudo netfilter-persistent save # 否则重启后规则丢失 +``` + +### 2.4 FAM-Core(UI/对话/成员接口,2026-09-13 从 NAS 迁入) + +```bash +# 代码(本地开发机 → 甲骨文) +tar czf - --exclude=venv --exclude=__pycache__ --exclude=logs --exclude='config/config.yaml' fam-core | \ + ssh -i ~/.ssh/oracle_new ubuntu@129.146.26.249 'tar xzf - --strip-components=1 -C /opt/fam-core' + +# 甲骨文上:venv + 配置 +python3 -m venv /opt/fam-core/venv +/opt/fam-core/venv/bin/pip install -r /opt/fam-core/requirements.txt +cp /opt/fam-core/config/config.yaml.example /opt/fam-core/config/config.yaml +# .env 只需要 ORACLE_SYNC_TOKEN(与 fam-edge 的一致),config_loader 会自动加载 +grep "^export ORACLE_SYNC_TOKEN=" /opt/fam-edge/.env > /opt/fam-core/.env && chmod 600 /opt/fam-core/.env + +sudo systemctl restart fam-core # 单元见 /etc/systemd/system/fam-core.service +curl -s http://127.0.0.1:5401/api/status # {"db":{"ok":true},...} +``` + +**端口 5401 不是笔误**:5400 被这台机器上的 chat-relay 占了。fam-core 只绑 +127.0.0.1,外部进不来,唯一客户端是同机 Caddy(先 forward_auth 再反代)。 + +**chat_history 迁移(一次性)**:这是 NAS MariaDB 里唯一不是镜像的表。 +在 NAS 上导出,再从本地导入甲骨文的 SQLite: +```bash +# NAS 上导出(用 fam-core 旧 venv 里的 PyMySQL) +ssh -p 2222 ericwyuan@192.168.50.64 '/volume1/web/sentinel-home-ai/fam-core/venv/bin/python -c " +import pymysql, json, yaml +cfg = yaml.safe_load(open(\"/volume1/web/sentinel-home-ai/fam-core/config/config.yaml\"))[\"database\"] +c = pymysql.connect(host=\"127.0.0.1\", user=cfg[\"user\"], password=cfg[\"password\"], database=cfg[\"database\"], cursorclass=pymysql.cursors.DictCursor) +cur = c.cursor(); cur.execute(\"SELECT * FROM chat_history ORDER BY chat_id\") +print(json.dumps(cur.fetchall(), ensure_ascii=False, default=str)) +"' > /tmp/chat_history.json + +# 本地 → 甲骨文导入(幂等:按 chat_id 跳过已存在的) +cat /tmp/chat_history.json | ssh -i ~/.ssh/oracle_new ubuntu@129.146.26.249 \ + '/opt/fam-core/venv/bin/python /opt/fam-core/scripts/import_chat_history.py' +``` + ## 3. 代码同步(tar 管道,scp 在 NAS 被禁用) ```bash @@ -91,10 +178,12 @@ tar czf - --exclude=venv --exclude=__pycache__ --exclude=data --exclude=gdrive_v | 项目 | 命令 | 预期 | |------|------|------| -| NAS MariaDB | `mysql -u root -p -e "SHOW DATABASES"` | 包含 sentinel_home_ai | -| NAS FAM-Core | `curl http://localhost:8000/health` | `{"status":"ok"}` | -| NAS 运动监测 | `curl http://localhost:8000/api/ss/status` | `poll_enabled: true, running: true` | -| 云服务器 FAM-UI | 浏览器访问 `http://129.146.26.249/`(Caddy 托管) | 未登录被 302 到 auth-hub 登录页;登录(审批过的 auth-hub 账号)后跳回展示 Vue3 SPA(时间轴/人物管理/统计) | +| 云服务器 FAM-UI | 浏览器访问 `https://smart-camera.zichuan.xyz/` | 未登录时 `/api/*` 返回 401,前端自动跳 `/login` → auth-hub 登录页;登录后跳回展示 Vue3 SPA | +| 登录不依赖 NAS | `curl -sI https://smart-camera.zichuan.xyz/login`(此时哪怕 fam-core 是停的) | 302 到 `auth.zichuan.xyz/authorize?...`,**不再是 502** | +| 鉴权闸门 | `curl -s -o /dev/null -w '%{http_code}' https://smart-camera.zichuan.xyz/api/ui/videos` | 401(Caddy forward_auth 拦下,没到 NAS) | +| NAS fam-notifier | NAS 上 `ps aux \| grep "[f]am_notifier"` | 有进程;`fam-notifier/logs/fam-notifier.log` 每 60s 一条轮询、每 5min 一条心跳 | +| 甲骨文 FAM-Core | `curl -s http://127.0.0.1:5401/api/status` | `{"db":{"ok":true,...}}`;`/api/ui/stats` 返回真实计数 | +| NAS 离线也能用 | 关掉 NAS 后访问 `https://smart-camera.zichuan.xyz/timeline` | 页面与数据照常(只是不再有新运动事件);**这是本次迁云的验收标准** | | Oracle FAM-Edge | `curl http://localhost:5000/health` | `{"status":"ok","queue_alive":true}` | | Oracle 运动事件 | `sqlite3 /opt/fam-edge/data/oracle.db "SELECT COUNT(*) FROM ss_motion_events"` | >0(NAS 推送) | | Oracle 运动片段 | `ls /opt/fam-edge/motion_clips/` | 存在 motion_*.mp4(素材分割产物) | diff --git a/fam-core/config/config.yaml.example b/fam-core/config/config.yaml.example index aeefc04..5549e69 100644 --- a/fam-core/config/config.yaml.example +++ b/fam-core/config/config.yaml.example @@ -1,32 +1,28 @@ -# FAM-Core 配置文件 (NAS 端) - 新架构 v2 +# FAM-Core 配置文件(2026-09-13 起跑在甲骨文,与 fam-edge 同机) # 复制此文件为 config.yaml 并修改实际值 # -# NAS 仅作管理后台,不再处理视频。唯一后台线程 Oracle-Sync 每 30 分钟 -# 从甲骨文 FAM-Edge 拉取增量镜像到本地 MariaDB(sync_videos/events/people)。 -# 所有视频分析在 Oracle 完成。 +# 本服务是纯查询/转发层:直读 fam-edge 的 SQLite,写操作转给 fam-edge。 +# 没有后台线程,没有自己的数据库,重启不影响任何数据。 server: - host: "0.0.0.0" - port: 8000 + # 只监听回环:唯一的客户端是同机 Caddy(它做 forward_auth 鉴权后才反代进来)。 + # 迁云前这里是 0.0.0.0:8000 并经 frp 暴露到公网,绕过 Caddy 就能拿到全部数据。 + host: "127.0.0.1" + port: 5401 database: - host: "127.0.0.1" - port: 3306 - user: "root" - password: "" - database: "sentinel_home_ai" - unix_socket: "/run/mysqld/mysqld10.sock" + # fam-edge 的库,本服务只读它写的表(唯一写的是 chat_history)。 + # fam-edge 侧已开 WAL,多进程读写安全。 + path: "/opt/fam-edge/data/oracle.db" -# 甲骨文同步(每 30 分钟拉增量镜像) -oracle_sync: - # FAM-Edge 对外同步接口地址(端口同其 server.port=5000) - base_url: "http://:5000" - # 与 Oracle 端 sync_api.token 一致 - token: "" - interval_sec: 1800 # 拉取间隔(秒),默认 30 分钟 - timeout: 120 # 单次拉取超时(秒) +# 同机 fam-edge:写操作(改名/删除)、帧图头像、服务状态都打给它 +edge: + base_url: "http://127.0.0.1:5000" + token: "${ORACLE_SYNC_TOKEN}" # 与 fam-edge 的 sync_api.token 一致 + timeout: 60 chat_handler: - # 智能问答统一走 FAM-Edge 编排端点(Gemini → NVIDIA → 本地 Ollama 兜底) - qa_url: "http://:5000/api/edge/chat/ask" + # 智能问答走 fam-edge 编排端点(NVIDIA → Gemini → 本地 Ollama 兜底) + qa_url: "http://127.0.0.1:5000/api/edge/chat/ask" + qa_stream_url: "http://127.0.0.1:5000/api/edge/chat/ask/stream" timeout: 120 diff --git a/fam-core/requirements.txt b/fam-core/requirements.txt index 78ee41a..e0aac59 100644 --- a/fam-core/requirements.txt +++ b/fam-core/requirements.txt @@ -1,7 +1,4 @@ flask>=3.0.0 gunicorn>=21.2.0 -PyMySQL>=1.1.0 PyYAML>=6.0 requests>=2.31.0 -PyJWT>=2.8.0 -cryptography>=42.0.0 diff --git a/fam-core/scripts/import_chat_history.py b/fam-core/scripts/import_chat_history.py new file mode 100644 index 0000000..d3d9fdf --- /dev/null +++ b/fam-core/scripts/import_chat_history.py @@ -0,0 +1,57 @@ +"""把 NAS MariaDB 导出的 chat_history 导入甲骨文 SQLite(2026-09-13 迁云一次性脚本)。 + +chat_history 是 NAS 那套库里唯一"不是镜像"的表——其余 sync_* 都能从甲骨文重新 +读出来,只有问答历史是本地产生的,迁云时必须搬过来。 + +用法(JSON 从 stdin 进来,导出命令见 docs/DEPLOY.md §2.4): + cat chat_history.json | /opt/fam-core/venv/bin/python scripts/import_chat_history.py + +幂等:按 chat_id 跳过已存在的行,重复执行不会产生重复记录。 +""" +import json +import os +import sys + +sys.path.insert(0, os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), 'src')) + +from fam_core import db_layer # noqa: E402 + +_COLS = ('chat_id', 'user_question', 'ai_answer', 'context_summary', + 'queried_date', 'queried_person', 'created_at') + + +def main(): + try: + rows = json.load(sys.stdin) + except ValueError as e: + print(f"stdin 不是合法 JSON: {e}", file=sys.stderr) + return 1 + if not isinstance(rows, list): + print("期望一个 JSON 数组", file=sys.stderr) + return 1 + + conn = db_layer.get_conn() + try: + db_layer._ensure_chat_schema(conn) + inserted = skipped = 0 + for r in rows: + cid = r.get('chat_id') + if cid is not None and conn.execute( + "SELECT 1 FROM chat_history WHERE chat_id=?", (cid,)).fetchone(): + skipped += 1 + continue + conn.execute( + "INSERT INTO chat_history ({}) VALUES ({})".format( + ','.join(_COLS), ','.join('?' * len(_COLS))), + tuple(None if r.get(c) is None else str(r.get(c)) for c in _COLS)) + inserted += 1 + conn.commit() + finally: + conn.close() + print(f"导入 {inserted} 条,跳过 {skipped} 条(chat_id 已存在)") + return 0 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/fam-core/src/fam_core/app.py b/fam-core/src/fam_core/app.py index 178b1f7..577017d 100644 --- a/fam-core/src/fam_core/app.py +++ b/fam-core/src/fam_core/app.py @@ -1,10 +1,20 @@ """ -FAM-Core 主应用 - Flask 单进程(新架构 v2) +FAM-Core 主应用 - Flask(2026-09-13 从 NAS 迁到甲骨文) -承载: Oracle-Sync(每 30 分钟拉取增量镜像)+ Member-Manager + Chat-Handler +承载: UI-API + Member-Manager + Chat-Handler + Img-Proxy —— 全是查询与转发, +**没有任何后台线程**,进程随时可重启,不持有任何状态。 -NAS 不再处理视频:无 Scheduler / Dispatcher / Poller / Event-Receiver / Video-Server。 -所有视频分析在 Oracle 完成,NAS 仅作管理后台拉取展示,CPU 占用大幅降低。 +迁云前它跑在 NAS 上,还扛着两个后台线程,现在都不在这里了: + - Oracle-Sync(每 30 分钟把甲骨文数据拉一份镜像进 MariaDB)—— 整个删除。 + 本服务现在与 fam-edge 同机,直接读它的 SQLite(见 db_layer.py 开头)。 + - MotionNotifier(轮询 Surveillance Station 推运动事件)—— 留在 NAS, + 拆成独立的 fam-notifier(摄像头插在 NAS 上,这部分搬不走)。 + +两件跟安全有关的事: + - 本服务只监听 127.0.0.1(config.yaml 的 server.host),唯一的客户端是同机 + Caddy。不像迁云前那样经 frp 把 :8000 暴露到公网。 + - 登录校验也不在这里:Caddy 用 forward_auth 打 fam-edge 的 /api/auth/verify, + 通过了才反代进来。 """ import os import sys @@ -16,87 +26,52 @@ sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from .config_loader import load_config from .logger import setup_logger -from .oracle_sync import get_sync -from .motion_notifier.motion_notifier import get_motion_notifier from .chat_handler.chat_handler import chat_bp from .member_manager.member_manager import member_bp from .img_proxy import img_bp -from .motion_bp import motion_bp from .ui_api import ui_bp -from .auth import auth_bp, init_auth +from . import db_layer, edge_client logger = setup_logger('fam-core.app') app = Flask(__name__) -# 注册蓝图:全部 /api/* 路由(前端已迁云服务器 Caddy:80,NAS 不再托管 SPA 静态文件) app.register_blueprint(chat_bp) app.register_blueprint(member_bp) app.register_blueprint(img_bp) -app.register_blueprint(motion_bp) app.register_blueprint(ui_bp) -app.register_blueprint(auth_bp) -# 登录校验全局拦截(页面未登录 302 /login;/api/* 未登录 401) -init_auth(app) -# 健康检查 @app.route('/health', methods=['GET']) def health(): return jsonify({"status": "ok", "service": "fam-core"}), 200 -# 初始化后台同步线程(NAS 唯一常驻线程) -_sync = None -try: - _sync = get_sync() - _sync.start() - logger.info("Oracle-Sync 已启动") - # 启动后立刻拉一次,前端无需等待首个周期即有数据 - try: - _sync.trigger_now() - logger.info("启动首次同步完成") - except Exception as e: - logger.warning(f"启动首次同步失败(后续周期会重试): {e}") -except Exception as e: - logger.error(f"Oracle-Sync 启动失败: {e}") - -# 初始化运动监测通知服务(NAS 轮询 SS EventCenter 主路径) -_notifier = None -try: - _notifier = get_motion_notifier() - _notifier.start() - logger.info("MotionNotifier 已初始化") -except Exception as e: - logger.error(f"MotionNotifier 初始化失败: {e}") - - @app.route('/api/status', methods=['GET']) def status(): - """系统状态""" + """服务自检:库能不能读、fam-edge 在不在。 + + 迁云前这里报告的是镜像同步状态(游标/周期/上次增量条数),镜像没了之后 + 那些字段不再存在,前端侧边栏的同步面板也一并去掉了。 + """ + db_ok, db_err = True, None + try: + conn = db_layer.get_conn() + try: + conn.execute("SELECT 1 FROM videos LIMIT 1").fetchone() + finally: + conn.close() + except Exception as e: + db_ok, db_err = False, str(e) return jsonify({ "service": "fam-core", - "sync": _sync.status() if _sync else {"running": False, "error": "未初始化"}, - "motion": _notifier.status() if _notifier else {"running": False, "error": "未初始化"}, + "db": {"ok": db_ok, "error": db_err}, + "edge_base_url": edge_client.base_url(), }), 200 -@app.route('/api/sync/trigger', methods=['POST']) -def sync_trigger(): - """手动立即触发一次甲骨文增量同步(服务状态页"立即同步"按钮)。 - - 正常情况下后台线程每 30 分钟自动拉一次;这个接口给用户想立刻看到最新数据 - 时用,跟后台线程共用同一把拉取锁(oracle_sync._pull_lock),不会并发重复拉。 - """ - if not _sync: - return jsonify({"error": "同步服务未初始化"}), 503 - ok = _sync.trigger_now() - if not ok: - return jsonify({"status": "failed", "error": _sync.status().get("last_error")}), 502 - return jsonify({"status": "ok", **_sync.status()}), 200 - - if __name__ == '__main__': cfg = load_config() - port = cfg.get('server', {}).get('port', 8000) - app.run(host='0.0.0.0', port=port, debug=False) + server = cfg.get('server', {}) + app.run(host=server.get('host', '127.0.0.1'), + port=server.get('port', 5401), debug=False) diff --git a/fam-core/src/fam_core/chat_handler/chat_handler.py b/fam-core/src/fam_core/chat_handler/chat_handler.py index e56b016..eb6535d 100644 --- a/fam-core/src/fam_core/chat_handler/chat_handler.py +++ b/fam-core/src/fam_core/chat_handler/chat_handler.py @@ -47,7 +47,7 @@ def _call_edge_qa(prompt: str) -> str: """调用 FAM-Edge 问答编排端点(Gemini → NVIDIA → 本地 Ollama 兜底)""" cfg = load_config() qa_url = cfg.get('chat_handler', {}).get( - 'qa_url', 'http://129.146.26.249:5000/api/edge/chat/ask' + 'qa_url', 'http://127.0.0.1:5000/api/edge/chat/ask' ) timeout = cfg.get('chat_handler', {}).get('timeout', 120) @@ -79,7 +79,7 @@ def chat_ask(): logger.info(f"Chat: person={queried_person}, date={queried_date}, question={question}") - rows = db_layer.query_sync_events_for_person_date(queried_person, queried_date) + rows = db_layer.query_events_for_person_date(queried_person, queried_date) if len(rows) == 0: answer = f"今天没有观察到{queried_person}。" @@ -87,7 +87,7 @@ def chat_ask(): else: context = _format_events(rows) context_summary = f"查询 sync_events {len(rows)} 条" - members = db_layer.get_sync_known_members_context() + members = db_layer.get_known_members_context() prompt = build_chat_prompt( context=context, members=members or queried_person, @@ -134,7 +134,7 @@ def chat_ask_stream(): return jsonify({"error": "缺少必填字段: question, queried_person, queried_date"}), 400 logger.info(f"Chat(stream): person={queried_person}, date={queried_date}, question={question}") - rows = db_layer.query_sync_events_for_person_date(queried_person, queried_date) + rows = db_layer.query_events_for_person_date(queried_person, queried_date) def sse(obj): return f"data: {json.dumps(obj, ensure_ascii=False)}\n\n" @@ -157,14 +157,14 @@ def chat_ask_stream(): yield sse({"type": "context", "count": len(rows), "summary": context_summary, "preview": context[:800]}) - members = db_layer.get_sync_known_members_context() + members = db_layer.get_known_members_context() prompt = build_chat_prompt( context=context, members=members or queried_person, question=question, queried_person=queried_person) cfg = load_config() stream_url = cfg.get('chat_handler', {}).get( - 'qa_stream_url', 'http://129.146.26.249:5000/api/edge/chat/ask/stream') + 'qa_stream_url', 'http://127.0.0.1:5000/api/edge/chat/ask/stream') timeout = cfg.get('chat_handler', {}).get('timeout', 120) full_answer = [] diff --git a/fam-core/src/fam_core/config_loader.py b/fam-core/src/fam_core/config_loader.py index a884eba..d029ca2 100644 --- a/fam-core/src/fam_core/config_loader.py +++ b/fam-core/src/fam_core/config_loader.py @@ -1,11 +1,42 @@ """ -配置加载器 - 从 config.yaml 读取配置 +配置加载器 - 从 config.yaml 读取配置,支持 ${ENV_VAR} 解析 """ import os import re import yaml +def _load_env_file(): + """加载部署目录下的 .env(支持 export KEY=VALUE 格式)。 + + 迁云前是 start_core.sh 负责 source .env 再起 gunicorn;现在由 systemd 拉起, + 没有那一步,所以在这里兜底加载(跟 fam-edge 的做法一致)。 + 已存在的环境变量不覆盖。 + """ + path = os.environ.get('FAM_ENV_FILE') or os.path.join( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), + '.env') + if not os.path.isfile(path): + return + with open(path, 'r', encoding='utf-8') as f: + for line in f: + line = line.strip() + if not line or line.startswith('#'): + continue + if line.startswith('export '): + line = line[7:].strip() + if '=' not in line: + continue + key, _, value = line.partition('=') + key = key.strip() + value = value.strip().strip('"').strip("'") + if key and key not in os.environ: + os.environ[key] = value + + +_load_env_file() + + def _resolve_env_vars(value): """递归解析字符串中的 ${ENV_VAR} 引用""" if isinstance(value, str): diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py index e4168b7..7011db2 100644 --- a/fam-core/src/fam_core/db_layer.py +++ b/fam-core/src/fam_core/db_layer.py @@ -1,300 +1,176 @@ """ -数据库访问层 - MariaDB 连接管理与同步镜像 CRUD +DB-Layer - 直读 fam-edge 的 SQLite 库(2026-09-13 迁云重写) -新架构 (2026-08-21 重构): - NAS 不再处理视频,仅作为管理后台。 - Oracle (FAM-Edge) 处理整视频分析后存 SQLite,NAS 每 30 分钟拉增量, - 镜像到本地三张表: - sync_videos : 视频会话(全局摘要 + 事件 JSON + 人物 JSON) - sync_events : 视频拆出的时间点事件(描述 + 涉及人物 + 是否关注) - sync_people : 规范人物表(label + canonical_name,Oracle 维护) - sync_cursor : 同步游标(上次成功拉取到的 server_time) +**这个文件以前是什么样**:fam-core 跑在 NAS 上,每 30 分钟把甲骨文的数据拉一份 +镜像进 MariaDB(`sync_videos`/`sync_events`/…),前端读镜像。725 行里有一半是 +镜像 upsert 的去重逻辑——8/29 和 9/3 两次线上事故(1062 主键冲突、游标卡死) +都出在那一半。 - 本层只服务同步镜像 + 问答历史,旧 process_tasks/event_details/monitor_events/ - family_members 相关逻辑已全部移除(视频处理职责已迁移至 Oracle)。 +**现在**:fam-core 跟 fam-edge 同机,直接打开它的 SQLite 库读,镜像层整个不存在了。 +连带消失的还有一个隐蔽 bug:镜像表为了保外键稳定用的是 NAS 本地自增 id,而帧图 +接口要的是甲骨文的 id,两边在 9/3 那次 id 重排后就对不上了。现在只有一套 id。 + +要点: +- 库文件是 fam-edge 的(`database.path`),**本模块只读它写的表**,唯一写的表是 + `chat_history`(问答历史,迁云时从 NAS MariaDB 搬过来的,fam-edge 不碰)。 +- fam-edge 那边开了 WAL,读不会阻塞它的写;这边同样设 busy_timeout 兜底。 +- 视频/人物的写操作(改名、删除)不在这里做,走 `edge_client` 打给 fam-edge—— + 它除了改库还要合并人物、删磁盘素材。 """ -import json -import pymysql -from datetime import datetime -from typing import Optional, List, Dict, Any +import sqlite3 +from datetime import datetime, timedelta, timezone +from typing import Dict, List, Optional from .config_loader import load_config from .logger import setup_logger logger = setup_logger('fam-core.db') -_config = None +_db_path = None +_chat_schema_ready = False -def get_config(): - global _config - if _config is None: - _config = load_config() - return _config +def _now() -> str: + return datetime.now(timezone(timedelta(hours=8))).strftime('%Y-%m-%d %H:%M:%S') -def get_conn(): - """获取数据库连接(单 worker gunicorn,无需连接池)""" - cfg = get_config().get('database', {}) - kwargs = dict( - host=cfg.get('host', '127.0.0.1'), - port=cfg.get('port', 3306), - user=cfg.get('user', 'root'), - password=cfg.get('password', ''), - database=cfg.get('database', 'sentinel_home_ai'), - charset='utf8mb4', - autocommit=False - ) - unix_socket = cfg.get('unix_socket') - if unix_socket: - kwargs['unix_socket'] = unix_socket - return pymysql.connect(**kwargs) +def get_config() -> Dict: + return load_config() + + +def _path() -> str: + global _db_path + if _db_path is None: + _db_path = load_config().get('database', {}).get( + 'path', '/opt/fam-edge/data/oracle.db') + return _db_path + + +def get_conn() -> sqlite3.Connection: + """每次调用开一个连接(请求级,跟改写前的 MySQL 用法一致)。 + + busy_timeout:fam-edge 的写事务提交时会短暂持锁,这里等而不是立刻报 + `database is locked`(9/3 那次 Oracle 过载时刷过一片这个错)。 + """ + conn = sqlite3.connect(_path(), timeout=10) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA busy_timeout=10000") + return conn + + +def _rows(cur) -> List[Dict]: + return [dict(r) for r in cur.fetchall()] + + +def _row(cur) -> Optional[Dict]: + r = cur.fetchone() + return dict(r) if r else None + + +# 人物命中判断:person_list_json 是 JSON 数组文本,用 json_each 展开精确匹配。 +# 必须先 json_valid——历史数据里有非 JSON 的脏值,直接 json_each 会整条查询报错。 +_PERSON_HIT = ("{col} IS NOT NULL AND json_valid({col}) " + "AND EXISTS(SELECT 1 FROM json_each({col}) WHERE json_each.value = ?)") + +# 日期维度:录制时间优先(文件名解析出来的),回退分析时间 +_DATE_EXPR = "COALESCE(NULLIF({a}.event_start_time,''), {a}.processed_at, {a}.updated_at, {a}.created_at)" + +# 只展示"有内容"的会话:运动片段,或含事件的视频。 +# 整段素材分割 0 段的空会话(历史素材无运动事件)不展示,避免淹没时间轴。 +_CONTENT_FILTER = ("substr(v.filename, 1, 7) = 'motion_' " + "OR EXISTS(SELECT 1 FROM events se WHERE se.video_id = v.id)") # ============================================================ -# 同步镜像:sync_videos +# videos / events # ============================================================ - -def upsert_sync_videos(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。 - - v2 修复 (2026-09-03):不再以 Oracle 的 videos.id 作为 NAS 主键。Oracle 库重建/ - 重排后 id 会复用/跳变(实测 2303 段跳到 3408/4600+ 段),旧写法按 id 插入时与 - 已存在的同 filename 行在 UNIQUE filename 上二次冲突报 1062,整批中止、游标不推进。 - 现改为:以 filename 为业务唯一键去重,命中则就地 UPDATE(保留 NAS 原 id,防止 - sync_events/model_calls/identity_map 的 video_id 外键失效),Oracle id 仅落 - oracle_id 列溯源;未命中则插入(优先用 Oracle id 作主键以对齐子表引用,主键冲突时 - 回退本地自增,避免 1062)。 - """ - if not rows: - return 0 - conn = get_conn() - n = 0 - try: - cur = conn.cursor() - for r in rows: - fn = r.get('filename') - oracle_id = r.get('id') - cur.execute("SELECT id FROM sync_videos WHERE filename=%s", (fn,)) - row = cur.fetchone() - if row: - # 命中业务键:就地更新,保留原 NAS id(外键不失效) - cur.execute( - """UPDATE sync_videos SET - oracle_id=%s, drive_file_id=%s, camera_name=%s, - duration_sec=%s, event_start_time=%s, status=%s, - summary_json=%s, events_json=%s, people_json=%s, - compute_provider=%s, created_at=%s, updated_at=%s, - processed_at=%s, synced_at=NOW() - WHERE id=%s""", - (oracle_id, r.get('drive_file_id'), 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'), row[0])) - else: - # 未命中:新文件。优先用 Oracle id 作主键(对齐子表 video_id 引用) - try: - cur.execute( - """INSERT INTO sync_videos - (id, oracle_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())""", - (oracle_id, oracle_id, fn, 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'))) - except Exception: - # 主键 id 被复用(极端情况):回退本地自增,避免 1062 - cur.execute( - """INSERT INTO sync_videos - (oracle_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, NOW())""", - (oracle_id, fn, 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() - return n - finally: - conn.close() - - -def get_sync_videos(limit=15, offset=0, date_filter=None) -> List[Dict]: - """获取视频会话列表(已完成优先),支持日期筛选与分页。 - - 排序/日期维度按视频实际录制时间(event_start_time,文件名解析), - 为空回退 processed_at/updated_at/created_at。 - date_filter 形如 '2026-08-21'。 - 只返回"有内容"的会话:运动片段(filename 前缀 motion_)或含事件的视频; - 整段素材分割 0 段的空会话(历史素材无运动事件)不展示,避免淹没时间轴。 - """ +def get_videos(limit=15, offset=0, date_filter=None) -> List[Dict]: + """视频会话列表(事件时间轴左侧),按录制时间倒序,支持日期筛选与分页。""" + d = _DATE_EXPR.format(a='v') + sql = f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status, + v.summary_json, v.events_json, v.people_json, v.compute_provider, + v.processed_at, v.updated_at, + (SELECT COUNT(*) FROM events se WHERE se.video_id = v.id) AS event_count + FROM videos v + WHERE v.status='done' AND ({_CONTENT_FILTER}) {{extra}} + ORDER BY {d} DESC + LIMIT ? OFFSET ?""" conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) - # 录制时间优先,回退分析时间 - date_expr = "COALESCE(NULLIF(event_start_time,''), processed_at, updated_at, created_at)" - # 注意:pymysql execute 用 % 做参数占位符,SQL 字面量不能含 %,故用 LEFT 判断 motion_ 前缀 - content_filter = ("LEFT(v.filename, 7) = 'motion_' " - "OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id = v.id)") if date_filter: - cur.execute( - f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status, - v.summary_json, v.events_json, v.people_json, v.compute_provider, - v.processed_at, v.updated_at, - (SELECT COUNT(*) FROM sync_events se WHERE se.video_id = v.id) AS event_count - FROM sync_videos v - WHERE v.status='done' AND ({content_filter}) AND {date_expr} LIKE %s - ORDER BY {date_expr} DESC - LIMIT %s OFFSET %s""", - (f'{date_filter}%', limit, offset)) + cur = conn.execute(sql.format(extra=f"AND {d} LIKE ?"), + (f'{date_filter}%', limit, offset)) else: - cur.execute( - f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status, - v.summary_json, v.events_json, v.people_json, v.compute_provider, - v.processed_at, v.updated_at, - (SELECT COUNT(*) FROM sync_events se WHERE se.video_id = v.id) AS event_count - FROM sync_videos v - WHERE v.status='done' AND ({content_filter}) - ORDER BY {date_expr} DESC - LIMIT %s OFFSET %s""", - (limit, offset)) - return cur.fetchall() + cur = conn.execute(sql.format(extra=''), (limit, offset)) + return _rows(cur) finally: conn.close() -def get_sync_video(video_id: int) -> Optional[Dict]: +def get_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() + return _row(conn.execute("SELECT * FROM videos WHERE id=?", (video_id,))) finally: conn.close() -def delete_sync_video(video_id: int): - """删除本地镜像里的这条视频(events 先删,再删 videos)。 - - 增量同步(get_sync_delta)只做 upsert,感知不到 Oracle 那边的物理删除, - 所以删除操作不能走"转发 Oracle + trigger_now() 拉增量"这条老路,必须在 - Oracle 确认删除成功后由调用方显式清理本地镜像。 - """ +def get_events_for_video(video_id: int) -> List[Dict]: conn = get_conn() try: - cur = conn.cursor() - cur.execute("DELETE FROM sync_events WHERE video_id=%s", (video_id,)) - cur.execute("DELETE FROM sync_videos WHERE id=%s", (video_id,)) - conn.commit() - finally: - conn.close() - - -# ============================================================ -# 同步镜像: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, - person_appearances_json, is_attention_event, - updated_at, synced_at) - VALUES (%s,%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), - person_appearances_json=VALUES(person_appearances_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'), r.get('person_appearances_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( + return _rows(conn.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() + is_attention_event + FROM events WHERE video_id=? ORDER BY ts ASC""", (video_id,))) finally: conn.close() -def get_sync_people_clips(label: str, limit: int = 10) -> List[Dict]: - """某人物(label 或 canonical_name)出现过的运动片段列表。 - - 匹配:label 本身 + 其 canonical_name 下全部 label;命中 sync_events 的 - person_list_json(JSON_CONTAINS)→ 关联 video(优先运动片段 motion_)。 - 返回按 event_start_time 倒序,含 first_ts(该人物在片段内最早事件 ts, - 前端缩略图定位用)与 clip_events(片段内该人物事件数)。 - """ +def get_attention_events(limit: int = 200) -> List[Dict]: + """需关注事件(统计图表页用),按录制日期倒序。""" + conn = get_conn() + try: + return _rows(conn.execute( + """SELECT COALESCE(NULLIF(v.event_start_time,''), v.processed_at) AS ev_date, + e.person_list_json + FROM events e JOIN videos v ON e.video_id=v.id + WHERE e.is_attention_event = 1 + ORDER BY ev_date DESC LIMIT ?""", (limit,))) + finally: + conn.close() + + +def get_people_clips(label: str, limit: int = 10) -> List[Dict]: + """某人物出现过的运动片段列表。 + + 匹配 label 本身 + 同一 canonical_name 下的全部 label;返回含 first_ts + (该人物在片段内最早事件时间点,前端缩略图定位用)与 clip_events。 + """ + hit_e = _PERSON_HIT.format(col='e.person_list_json') + hit_e2 = _PERSON_HIT.format(col='e2.person_list_json') + hit_e3 = _PERSON_HIT.format(col='e3.person_list_json') conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) labels = {label} - cur.execute( - "SELECT label, canonical_name FROM sync_people WHERE label=%s OR canonical_name=%s", - (label, label)) - for r in cur.fetchall(): - cn = r.get('canonical_name') - if cn: - cur.execute("SELECT label FROM sync_people WHERE canonical_name=%s", (cn,)) - labels.update(x['label'] for x in cur.fetchall()) - clips = [] - seen = set() + for r in _rows(conn.execute( + "SELECT label, canonical_name FROM people WHERE label=? OR canonical_name=?", + (label, label))): + if r.get('canonical_name'): + labels.update(x['label'] for x in _rows(conn.execute( + "SELECT label FROM people WHERE canonical_name=?", (r['canonical_name'],)))) + clips, seen = [], set() for lb in sorted(labels): - cur.execute( - """SELECT DISTINCT v.id AS video_id, v.filename, v.event_start_time, - v.duration_sec, v.summary_json, v.camera_name, - (SELECT MIN(e2.ts) FROM sync_events e2 - WHERE e2.video_id=v.id AND e2.person_list_json IS NOT NULL - AND JSON_CONTAINS(e2.person_list_json, JSON_QUOTE(%s), '$')) AS first_ts, - (SELECT COUNT(*) FROM sync_events e3 - WHERE e3.video_id=v.id AND e3.person_list_json IS NOT NULL - AND JSON_CONTAINS(e3.person_list_json, JSON_QUOTE(%s), '$')) AS clip_events - FROM sync_events e - JOIN sync_videos v ON e.video_id=v.id - WHERE e.person_list_json IS NOT NULL - AND JSON_CONTAINS(e.person_list_json, JSON_QUOTE(%s), '$') - AND v.status='done'""", + cur = conn.execute( + f"""SELECT DISTINCT v.id AS video_id, v.filename, v.event_start_time, + v.duration_sec, v.summary_json, v.camera_name, + (SELECT MIN(e2.ts) FROM events e2 + WHERE e2.video_id=v.id AND {hit_e2}) AS first_ts, + (SELECT COUNT(*) FROM events e3 + WHERE e3.video_id=v.id AND {hit_e3}) AS clip_events + FROM events e JOIN videos v ON e.video_id=v.id + WHERE {hit_e} AND v.status='done'""", (lb, lb, lb)) - for row in cur.fetchall(): + for row in _rows(cur): if row['video_id'] not in seen: seen.add(row['video_id']) clips.append(row) @@ -304,397 +180,177 @@ def get_sync_people_clips(label: str, limit: int = 10) -> List[Dict]: conn.close() -def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]: +def query_events_for_person_date(person: str, date_str: str) -> List[Dict]: """问答上下文:某人在某天的事件。 - 说明: Oracle 事件 ts 为视频内相对时间点(如 00:01:23),不是绝对日期, - 因此按所属视频的录制日期(event_start_time,回退 processed_at)过滤, - 再按 person_list_json 命中人名。 - person 可为真名或抽象标签(Oracle 回灌上下文用真名,但历史标签也保留)。 + 事件 ts 是视频内的相对时间点(如 00:01:23),不是绝对日期,所以按所属视频的 + 录制日期过滤,再按 person_list_json 命中人名。person 可为真名或抽象标签。 """ + hit = _PERSON_HIT.format(col='e.person_list_json') 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 COALESCE(NULLIF(v.event_start_time,''), v.processed_at) LIKE %s - AND e.person_list_json IS NOT NULL - AND JSON_CONTAINS(e.person_list_json, JSON_QUOTE(%s), '$') - ORDER BY COALESCE(NULLIF(v.event_start_time,''), v.processed_at) ASC, e.ts ASC""", - (f'{date_str}%', person)) - return cur.fetchall() + return _rows(conn.execute( + f"""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 events e JOIN videos v ON e.video_id = v.id + WHERE COALESCE(NULLIF(v.event_start_time,''), v.processed_at) LIKE ? + AND {hit} + ORDER BY COALESCE(NULLIF(v.event_start_time,''), v.processed_at) ASC, e.ts ASC""", + (f'{date_str}%', person))) finally: conn.close() # ============================================================ -# 同步镜像:sync_people +# people # ============================================================ - -def upsert_sync_people(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 people 增量。 - - v2 修复 (2026-09-03):同 videos,以 label 为业务唯一键去重,命中就地 UPDATE - 保留 NAS 原 id(people.id 无外键引用,仅作展示排序),Oracle id 落 oracle_id 列 - 溯源;未命中插入(优先 Oracle id,主键冲突回退自增,避免 1062)。 - """ - if not rows: - return 0 - conn = get_conn() - n = 0 - try: - cur = conn.cursor() - for r in rows: - label = r.get('label') - oracle_id = r.get('id') - cur.execute("SELECT id FROM sync_people WHERE label=%s", (label,)) - row = cur.fetchone() - if row: - cur.execute( - """UPDATE sync_people SET - oracle_id=%s, canonical_name=%s, first_seen=%s, - appearances=%s, source=%s, features_json=%s, - display_uid=%s, updated_at=%s, synced_at=NOW() - WHERE id=%s""", - (oracle_id, r.get('canonical_name'), r.get('first_seen'), - r.get('appearances') or 0, r.get('source'), - r.get('features_json'), r.get('display_uid'), - r.get('updated_at'), row[0])) - else: - try: - cur.execute( - """INSERT INTO sync_people - (id, oracle_id, label, canonical_name, first_seen, - appearances, source, features_json, display_uid, - updated_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", - (oracle_id, oracle_id, label, r.get('canonical_name'), - r.get('first_seen'), r.get('appearances') or 0, - r.get('source'), r.get('features_json'), - r.get('display_uid'), r.get('updated_at'))) - except Exception: - cur.execute( - """INSERT INTO sync_people - (oracle_id, label, canonical_name, first_seen, - appearances, source, features_json, display_uid, - updated_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", - (oracle_id, label, r.get('canonical_name'), - r.get('first_seen'), r.get('appearances') or 0, - r.get('source'), r.get('features_json'), - r.get('display_uid'), r.get('updated_at'))) - n += 1 - conn.commit() - return n - finally: - conn.close() - - -def get_sync_people() -> List[Dict]: +def get_people() -> List[Dict]: conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) - cur.execute( + return _rows(conn.execute( "SELECT id, label, canonical_name, first_seen, appearances, source, " - "features_json, display_uid, updated_at " - "FROM sync_people ORDER BY id ASC") - return cur.fetchall() + "features_json, display_uid, updated_at FROM people ORDER BY id ASC")) finally: conn.close() -def upsert_sync_model_calls(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 model_calls 增量(幂等,重复覆盖)。""" - if not rows: - return 0 +def get_named_members() -> List[str]: + """已命名成员的真名列表(UI 下拉用)。""" conn = get_conn() - n = 0 try: - cur = conn.cursor() - for r in rows: - cur.execute( - """INSERT INTO sync_model_calls - (id, provider, model, video_id, filename, started_at, - duration_sec, success, error, created_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW()) - ON DUPLICATE KEY UPDATE - provider=VALUES(provider), - model=VALUES(model), - video_id=VALUES(video_id), - filename=VALUES(filename), - started_at=VALUES(started_at), - duration_sec=VALUES(duration_sec), - success=VALUES(success), - error=VALUES(error), - created_at=VALUES(created_at), - synced_at=NOW()""", - (r.get('id'), r.get('provider'), r.get('model'), - r.get('video_id'), r.get('filename'), r.get('started_at'), - r.get('duration_sec') or 0, 1 if r.get('success') else 0, - (r.get('error') or '')[:500], r.get('created_at'))) - n += 1 - conn.commit() - return n + return [r['canonical_name'] for r in _rows(conn.execute( + "SELECT DISTINCT canonical_name FROM people " + "WHERE canonical_name IS NOT NULL AND canonical_name != '' " + "ORDER BY canonical_name ASC"))] finally: conn.close() -def upsert_sync_identity_map(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 person_identity_map 增量(幂等,重复覆盖)。 - - v2 修复 (2026-08-29):不再以 Oracle 的 person_identity_map.id 作为 NAS 主键。 - 该 id 只是 Oracle 端代理主键,一旦 Oracle 库重建/恢复,id 会被复用(实测出现 - id=744 同时对应 NAS 上的 (1302,'人物E') 与 Oracle 现在的 (1073,'汤圆')), - 旧写法会先按 id 撞主键、再 UPDATE 成新组合,触发 uq_video_raw_uid 二次冲突 - 报 1062。现改为:NAS 本地自增 id 做主键,业务键 (video_id, raw_uid) 唯一, - Oracle 的 id 只落 oracle_id 列作溯源参考。 - """ - if not rows: - return 0 - conn = get_conn() - n = 0 - try: - cur = conn.cursor() - for r in rows: - cur.execute( - """INSERT INTO sync_identity_map - (oracle_id, video_id, raw_uid, canonical_name, source, updated_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s, NOW()) - ON DUPLICATE KEY UPDATE - oracle_id=VALUES(oracle_id), - canonical_name=VALUES(canonical_name), - source=VALUES(source), - updated_at=VALUES(updated_at), - synced_at=NOW()""", - (r.get('id'), r.get('video_id'), r.get('raw_uid'), - r.get('canonical_name'), r.get('source'), r.get('updated_at'))) - n += 1 - conn.commit() - return n - finally: - conn.close() +def get_known_members_context() -> str: + """人物清单文本,注入问答 Prompt,让模型用真名指代。""" + rows = get_people() + 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) -def get_sync_model_calls(limit: int = 200) -> List[Dict]: - """最近模型调用记录(前端统计展示)。""" +# ============================================================ +# model_calls +# ============================================================ +def get_model_calls(limit: int = 200) -> List[Dict]: conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) - cur.execute( + return _rows(conn.execute( "SELECT id, provider, model, video_id, filename, started_at, " "duration_sec, success, error, created_at " - "FROM sync_model_calls ORDER BY id DESC LIMIT %s", (limit,)) - return cur.fetchall() + "FROM model_calls ORDER BY id DESC LIMIT ?", (limit,))) finally: conn.close() -def get_sync_model_calls_stats() -> Dict: - """模型调用统计:按 provider+model 聚合成功/失败/平均耗时。""" +def get_model_calls_stats() -> List[Dict]: + """按 provider+model 聚合成功/失败/平均耗时。""" conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) - cur.execute( + return _rows(conn.execute( """SELECT provider, model, SUM(CASE WHEN success=1 THEN 1 ELSE 0 END) AS ok_cnt, SUM(CASE WHEN success=0 THEN 1 ELSE 0 END) AS fail_cnt, ROUND(AVG(duration_sec), 1) AS avg_duration, MAX(created_at) AS last_call - FROM sync_model_calls - GROUP BY provider, model ORDER BY provider, model""") - 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) + FROM model_calls GROUP BY provider, model ORDER BY provider, model""")) finally: conn.close() # ============================================================ -# 同步游标 +# 概览统计 # ============================================================ - -def get_sync_cursor() -> str: +def get_stats(date_str: str = None) -> Dict: + """视频数 / 事件数 / 关注事件数 / 出现人物数(可按日期过滤)。""" + d = _DATE_EXPR.format(a='sv') 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() - - -# ============================================================ -# 运动通知游标(MotionNotifier 增量去重用,复用 sync_cursor 表) -# ============================================================ - -def get_motion_cursor() -> str: - """返回上次推送到甲骨文的最大 SS 事件 id(字符串),无则返回空。""" - conn = get_conn() - try: - cur = conn.cursor() - cur.execute("SELECT `value` FROM sync_cursor WHERE `key`='motion_last_event_id'") - row = cur.fetchone() - return row[0] if row else '' - finally: - conn.close() - - -def set_motion_cursor(value: str): - conn = get_conn() - try: - cur = conn.cursor() - cur.execute( - """INSERT INTO sync_cursor (`key`, `value`) VALUES ('motion_last_event_id', %s) - ON DUPLICATE KEY UPDATE `value`=VALUES(`value`)""", - (value,)) - conn.commit() - finally: - conn.close() - - -# ============================================================ -# 统计 -# ============================================================ - -def get_attention_events(limit: int = 200) -> List[Dict]: - """需关注事件列表(供统计图表页展示日期 + 涉及人物),按录制日期倒序。""" - conn = get_conn() - try: - cur = conn.cursor(pymysql.cursors.DictCursor) - cur.execute( - """SELECT COALESCE(NULLIF(v.event_start_time,''), v.processed_at) AS ev_date, - e.person_list_json - FROM sync_events e JOIN sync_videos v ON e.video_id=v.id - WHERE e.is_attention_event = 1 - ORDER BY ev_date DESC LIMIT %s""", - (limit,)) - return cur.fetchall() - 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) - # 视频/事件/关注数(日期维度=录制时间 event_start_time,回退 processed_at) - _D = "COALESCE(NULLIF(sv.event_start_time,''), sv.processed_at)" if date_str: - cur.execute( - """SELECT - (SELECT COUNT(*) FROM sync_videos sv - WHERE sv.status='done' - AND (LEFT(sv.filename, 7) = 'motion_' - OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id=sv.id)) - AND {d} LIKE %s) AS videos, - (SELECT COUNT(*) FROM sync_events se - JOIN sync_videos sv ON se.video_id=sv.id - WHERE {d} 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 {d} LIKE %s) AS attention""".format(d=_D), - (f'{date_str}%', f'{date_str}%', f'{date_str}%')) + like = f'{date_str}%' + stat = _row(conn.execute( + f"""SELECT + (SELECT COUNT(*) FROM videos sv + WHERE sv.status='done' + AND (substr(sv.filename,1,7)='motion_' + OR EXISTS(SELECT 1 FROM events se WHERE se.video_id=sv.id)) + AND {d} LIKE ?) AS videos, + (SELECT COUNT(*) FROM events se JOIN videos sv ON se.video_id=sv.id + WHERE {d} LIKE ?) AS events, + (SELECT COALESCE(SUM(se.is_attention_event),0) FROM events se + JOIN videos sv ON se.video_id=sv.id + WHERE {d} LIKE ?) AS attention""", + (like, like, like))) or {} else: - cur.execute( + stat = _row(conn.execute( """SELECT - (SELECT COUNT(*) FROM sync_videos - WHERE status='done' - AND (LEFT(filename, 7) = 'motion_' - OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id=sync_videos.id))) 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 {} + (SELECT COUNT(*) FROM videos sv + WHERE sv.status='done' + AND (substr(sv.filename,1,7)='motion_' + OR EXISTS(SELECT 1 FROM events se WHERE se.video_id=sv.id))) AS videos, + (SELECT COUNT(*) FROM events) AS events, + (SELECT COALESCE(SUM(is_attention_event),0) FROM events) AS attention""")) 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: + # 出现人物数:逐 label/真名在事件里查有没有命中(沿用改写前的口径) + hits = 0 + for p in _rows(conn.execute("SELECT label, canonical_name FROM 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 + c = conn.execute("SELECT COUNT(*) AS c FROM events " + "WHERE person_list_json LIKE ?", (f'%{name}%',)).fetchone() + if c and c['c'] > 0: + hits += 1 + stat['people'] = hits return stat finally: conn.close() # ============================================================ -# chat_history(保留:问答历史) +# chat_history(本模块唯一写的表;2026-09-13 从 NAS MariaDB 迁入) # ============================================================ +def _ensure_chat_schema(conn): + global _chat_schema_ready + if _chat_schema_ready: + return + conn.executescript(""" + CREATE TABLE IF NOT EXISTS chat_history ( + chat_id INTEGER PRIMARY KEY AUTOINCREMENT, + user_question TEXT NOT NULL, + ai_answer TEXT NOT NULL, + context_summary TEXT, + queried_date TEXT, + queried_person TEXT, + created_at TEXT + ); + CREATE INDEX IF NOT EXISTS idx_chat_created ON chat_history(created_at); + """) + conn.commit() + _chat_schema_ready = True + def insert_chat_history(user_question: str, ai_answer: str, context_summary: str, queried_date: str, queried_person: str) -> int: - """插入对话记录""" conn = get_conn() try: - cur = conn.cursor() - cur.execute( + _ensure_chat_schema(conn) + cur = conn.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) - ) + (user_question, ai_answer, context_summary, queried_date, + queried_person, created_at) + VALUES (?, ?, ?, ?, ?, ?)""", + (user_question, ai_answer, context_summary, queried_date, + queried_person, _now())) conn.commit() return cur.lastrowid finally: @@ -702,24 +358,20 @@ def insert_chat_history(user_question: str, ai_answer: str, def get_chat_history(limit=20, offset=0, date_filter=None, person_filter=None) -> List[Dict]: - """获取对话历史(分页)""" conn = get_conn() try: - cur = conn.cursor(pymysql.cursors.DictCursor) - conditions = [] - params = [] + _ensure_chat_schema(conn) + conds, params = [], [] if date_filter: - conditions.append("queried_date = %s") + conds.append("queried_date = ?") params.append(date_filter) if person_filter: - conditions.append("queried_person = %s") + conds.append("queried_person = ?") params.append(person_filter) - where = f"WHERE {' AND '.join(conditions)}" if conditions else "" + where = f"WHERE {' AND '.join(conds)}" if conds else "" params.extend([limit, offset]) - cur.execute( - f"SELECT * FROM chat_history {where} ORDER BY created_at DESC LIMIT %s OFFSET %s", - params - ) - return cur.fetchall() + return _rows(conn.execute( + f"SELECT * FROM chat_history {where} ORDER BY created_at DESC LIMIT ? OFFSET ?", + params)) finally: conn.close() diff --git a/fam-core/src/fam_core/edge_client.py b/fam-core/src/fam_core/edge_client.py new file mode 100644 index 0000000..9b40cb7 --- /dev/null +++ b/fam-core/src/fam_core/edge_client.py @@ -0,0 +1,102 @@ +""" +Edge-Client - 访问同机 fam-edge 的薄客户端(2026-09-13 迁云后新增) + +迁云之前这些调用散在 `oracle_sync.py` 里(NAS 跨公网访问甲骨文),随镜像层一起 +删掉了。现在 fam-core 与 fam-edge 同机,全部走 127.0.0.1:不出公网、无 TLS、 +无隧道,token 仍然带着(fam-edge 侧的鉴权没变,且它同时对公网监听 :5000)。 + +只保留三类调用: + - push_name_correct / push_video_delete —— 写操作仍由 fam-edge 处理,因为它 + 除了改库还要做人物合并、删磁盘素材,不是单纯一条 SQL + - get_activity —— 服务状态页 + - fetch_image —— 事件帧图 / 人物头像(图源和裁剪都在 fam-edge) +""" +import requests + +from .config_loader import load_config +from .logger import setup_logger + +logger = setup_logger('fam-core.edge_client') + +_cfg = None + + +def _conf(): + global _cfg + if _cfg is None: + c = load_config().get('edge', {}) + _cfg = { + 'base_url': (c.get('base_url') or 'http://127.0.0.1:5000').rstrip('/'), + 'token': c.get('token') or '', + 'timeout': int(c.get('timeout', 60)), + } + return _cfg + + +def base_url(): + return _conf()['base_url'] + + +def token(): + return _conf()['token'] + + +def _post(rel: str, payload: dict): + """返回 (ok, err)。写操作统一用这个,失败时把原因带回给前端。""" + c = _conf() + body = dict(payload) + if c['token']: + body['token'] = c['token'] + try: + resp = requests.post(f"{c['base_url']}{rel}", json=body, timeout=(5, c['timeout'])) + except requests.RequestException as e: + logger.error(f"fam-edge {rel} 请求失败: {e}") + return False, str(e) + if resp.status_code != 200: + logger.warning(f"fam-edge {rel} 返回 {resp.status_code}: {resp.text[:160]}") + return False, f"HTTP {resp.status_code} {resp.text[:120]}" + return True, '' + + +def push_name_correct(label: str, canonical_name: str): + return _post('/api/oracle/people/correct', + {'label': label, 'canonical_name': canonical_name}) + + +def push_identity_correct(video_id: int, current_name: str, new_name: str): + return _post('/api/oracle/identity/correct', + {'video_id': video_id, 'current_name': current_name, 'new_name': new_name}) + + +def push_video_delete(video_id: int): + return _post('/api/oracle/video/delete', {'video_id': video_id}) + + +def get_activity(): + """服务状态页用:fam-edge 的队列 / 分割 / 模型活动快照。返回 (data, err)。""" + c = _conf() + try: + r = requests.get(f"{c['base_url']}/api/oracle/activity", + params={'token': c['token']}, timeout=(5, 15)) + except requests.RequestException as e: + return None, f"连接 fam-edge 失败: {e}" + if r.status_code != 200: + return None, f"fam-edge activity HTTP {r.status_code}" + return r.json(), None + + +def fetch_image(rel: str, params: dict): + """帧图 / 头像透传,失败返回 None(调用方给 404)。""" + c = _conf() + p = dict(params) + if c['token']: + p['token'] = c['token'] + try: + resp = requests.get(f"{c['base_url']}{rel}", params=p, timeout=(5, c['timeout'])) + except requests.RequestException as e: + logger.error(f"fam-edge {rel} 请求失败: {e}") + return None + if resp.status_code == 200 and resp.content: + return resp.content + logger.warning(f"fam-edge {rel} 返回 {resp.status_code}") + return None diff --git a/fam-core/src/fam_core/img_proxy.py b/fam-core/src/fam_core/img_proxy.py index 3edfd41..506d1a5 100644 --- a/fam-core/src/fam_core/img_proxy.py +++ b/fam-core/src/fam_core/img_proxy.py @@ -1,14 +1,12 @@ """ -Img-Proxy - NAS 端图片代理(图源在 Oracle,计算全部在 Oracle) +Img-Proxy - 图片代理(图源与裁剪都在 fam-edge) -浏览器不直连 Oracle(避免公网暴露 5000 端口与 token 外泄), -而是访问 NAS fam-core 的 /api/proxy/*,由 NAS 出网到 Oracle 拉取 jpeg 回传。 - -Oracle 侧已有磁盘缓存 / VLM 人物定位裁剪,NAS 端仅透传,不做图像计算。 +浏览器不直连 fam-edge(避免 token 外泄),而是访问 /api/proxy/*,由本服务转一手。 +迁云后这一跳是同机 127.0.0.1,纯透传,不做图像计算。 """ from flask import Blueprint, Response, request -from .oracle_sync import get_sync +from . import edge_client from .logger import setup_logger logger = setup_logger('fam-core.img_proxy') @@ -17,21 +15,7 @@ img_bp = Blueprint('img_proxy', __name__) def _fetch(rel: str, params: dict): - import requests - sync = get_sync() # 复用 oracle_sync 已解析好的 base_url/token/timeout,避免两处配置各读一份 - params = dict(params) - if sync.token: - params['token'] = sync.token - try: - resp = requests.get(f"{sync.base_url}{rel}", params=params, - timeout=(10, sync.timeout)) - if resp.status_code == 200 and resp.content: - return resp.content - logger.warning(f"Oracle {rel} 返回 {resp.status_code}: {resp.text[:120]}") - return None - except Exception as e: - logger.error(f"Oracle {rel} 请求失败: {e}") - return None + return edge_client.fetch_image(rel, params) @img_bp.route('/api/proxy/frame', methods=['GET']) diff --git a/fam-core/src/fam_core/member_manager/member_manager.py b/fam-core/src/fam_core/member_manager/member_manager.py index 58f8daf..808c8f4 100644 --- a/fam-core/src/fam_core/member_manager/member_manager.py +++ b/fam-core/src/fam_core/member_manager/member_manager.py @@ -14,7 +14,7 @@ from flask import Blueprint, request, jsonify from ..logger import setup_logger from .. import db_layer -from ..oracle_sync import get_sync +from .. import edge_client logger = setup_logger('fam-core.member_manager') @@ -24,7 +24,7 @@ member_bp = Blueprint('member_manager', __name__) @member_bp.route('/api/member/unnamed', methods=['GET']) def list_unnamed(): """列出未命名人物(canonical_name 为空)""" - members = db_layer.get_sync_people() + members = db_layer.get_people() result = [] for m in members: canonical = m.get('canonical_name') @@ -40,7 +40,7 @@ def list_unnamed(): @member_bp.route('/api/member/list', methods=['GET']) def list_members(): """列出所有人物(按 canonical_name 或 label 展示)""" - members = db_layer.get_sync_people() + members = db_layer.get_people() result = [] for m in members: canonical = m.get('canonical_name') @@ -58,7 +58,7 @@ def list_members(): @member_bp.route('/api/member/name', methods=['POST']) def name_member(): - """命名人物(回推 Oracle + 立即拉回本地镜像) + """命名人物(交给 fam-edge 落库 + 合并人物) 请求: {"label": "人物A", "canonical_name": "张三"} """ @@ -71,18 +71,13 @@ def name_member(): if not label or not canonical_name: return jsonify({"error": "缺少必填字段: label, canonical_name"}), 400 - logger.info(f"命名: {label} -> {canonical_name}(回推 Oracle)") - ok, err = get_sync().push_name_correct(label, canonical_name) + logger.info(f"命名: {label} -> {canonical_name}(回推 fam-edge)") + ok, err = edge_client.push_name_correct(label, canonical_name) if not ok: - return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502 + return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502 - # 立即拉回最新 people 镜像,前端无需等待下一个 30 分钟周期 - try: - get_sync().trigger_now() - except Exception as e: - logger.warning(f"命名后即时拉回失败(下一个周期会自动同步): {e}") - members = db_layer.get_sync_people() + members = db_layer.get_people() return jsonify({ "status": "ok", "label": label, @@ -114,24 +109,20 @@ def merge_member(): return jsonify({"error": "source 与 target 不能相同"}), 400 # 解析 target 的规范名 - members = {m['label']: m for m in db_layer.get_sync_people()} + members = {m['label']: m for m in db_layer.get_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) + logger.info(f"合并: {source} -> {canonical}(回推 fam-edge)") + ok, err = edge_client.push_name_correct(source, canonical) if not ok: - return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502 + return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502 - try: - get_sync().trigger_now() - except Exception as e: - logger.warning(f"合并后即时拉回失败(下一个周期会自动同步): {e}") - members = db_layer.get_sync_people() + members = db_layer.get_people() return jsonify({ "status": "ok", "source": source, @@ -163,15 +154,11 @@ def identity_correct(): if not video_id or not current_name or not new_name: return jsonify({"error": "缺少必填字段: video_id, current_name, new_name"}), 400 - logger.info(f"人物纠错: video_id={video_id} {current_name} -> {new_name}(回推 Oracle)") - ok, err = get_sync().push_identity_correct(video_id, current_name, new_name) + logger.info(f"人物纠错: video_id={video_id} {current_name} -> {new_name}(回推 fam-edge)") + ok, err = edge_client.push_identity_correct(video_id, current_name, new_name) if not ok: - return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502 + return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502 - try: - get_sync().trigger_now() - except Exception as e: - logger.warning(f"纠错后即时拉回失败(下一个周期会自动同步): {e}") return jsonify({ "status": "ok", "video_id": video_id, diff --git a/fam-core/src/fam_core/motion_bp.py b/fam-core/src/fam_core/motion_bp.py deleted file mode 100644 index 3108c93..0000000 --- a/fam-core/src/fam_core/motion_bp.py +++ /dev/null @@ -1,20 +0,0 @@ -""" -运动监测状态接口(fam-core) - -运动事件采集走 MotionNotifier 轮询 SS EventCenter.Event.List(真实 -event_id/start_time/duration),不再提供 /api/ss/webhook 接收端点 -(SS Webhook 行动规则已弃用,2026-08-25 移除)。 -""" -from flask import Blueprint, jsonify - -from .logger import setup_logger -from .motion_notifier.motion_notifier import get_motion_notifier - -logger = setup_logger('fam-core.motion_bp') - -motion_bp = Blueprint('motion_bp', __name__) - - -@motion_bp.route('/api/ss/status', methods=['GET']) -def ss_status(): - return jsonify(get_motion_notifier().status()), 200 diff --git a/fam-core/src/fam_core/motion_notifier/__init__.py b/fam-core/src/fam_core/motion_notifier/__init__.py deleted file mode 100644 index 5afb0ae..0000000 --- a/fam-core/src/fam_core/motion_notifier/__init__.py +++ /dev/null @@ -1,4 +0,0 @@ -"""MotionNotifier 包:NAS 端运动监测通知服务。""" -from .motion_notifier import MotionNotifier, get_motion_notifier - -__all__ = ["MotionNotifier", "get_motion_notifier"] diff --git a/fam-core/src/fam_core/oracle_sync.py b/fam-core/src/fam_core/oracle_sync.py deleted file mode 100644 index 619d26a..0000000 --- a/fam-core/src/fam_core/oracle_sync.py +++ /dev/null @@ -1,234 +0,0 @@ -""" -Oracle-Sync - NAS 端唯一后台线程 - -职责: - 1. 每 interval_sec(默认 1800s = 30 分钟)从甲骨文 FAM-Edge 拉取增量: - GET {base_url}/api/oracle/sync?since=&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.26.249: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 - # 拉取互斥锁:trigger_now(命名后即时拉回)与后台 _run 并发时只允许一个执行 - self._pull_lock = threading.Lock() - - # ------------------------------------------------------------------ - def _pull_once(self) -> bool: - """执行一次增量拉取(加锁防并发双拉)。返回是否成功。""" - with self._pull_lock: - return self._pull_once_locked() - - def _pull_once_locked(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 [] - model_calls = data.get('model_calls', []) or [] - identity_map = data.get('identity_map', []) 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) - n_calls = db_layer.upsert_sync_model_calls(model_calls) - n_identity = db_layer.upsert_sync_identity_map(identity_map) - - 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, n_calls, n_identity) - logger.info( - f"同步完成: videos+{n_videos} events+{n_events} people+{n_people} " - f"model_calls+{n_calls} identity_map+{n_identity} since={since!r} " - f"-> 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 push_identity_correct(self, video_id: int, current_name: str, new_name: str): - """回推事件时间轴/人物管理"这个人识别错了"纠错到 Oracle。 - - 跟 push_name_correct 的区别:这个按 (video_id, 当前展示名) 定位,只改这 - 一段视频里错认的那个人,不影响同名字符串在其他视频里的映射(人物 uid - 只在单次视频分析内稳定,跨视频复用同一字符串完全可能是不同真人)。 - - 返回 (success: bool, error: str) - """ - try: - video_id = int(video_id) - except (TypeError, ValueError): - return False, "video_id 必须是数字" - current_name = (current_name or '').strip() - new_name = (new_name or '').strip() - if not current_name or not new_name: - return False, "缺少 current_name / new_name" - try: - resp = requests.post( - f"{self.base_url}/api/oracle/identity/correct", - json={"video_id": video_id, "current_name": current_name, - "new_name": new_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 push_video_delete(self, video_id: int): - """回推事件时间轴"删除该视频"到 Oracle。返回 (success: bool, error: str)。""" - try: - video_id = int(video_id) - except (TypeError, ValueError): - return False, "video_id 必须是数字" - try: - resp = requests.post( - f"{self.base_url}/api/oracle/video/delete", - json={"video_id": video_id, "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, - } diff --git a/fam-core/src/fam_core/ui_api.py b/fam-core/src/fam_core/ui_api.py index 8e2de2f..08caaee 100644 --- a/fam-core/src/fam_core/ui_api.py +++ b/fam-core/src/fam_core/ui_api.py @@ -7,9 +7,8 @@ UI-API - Vue 前端数据接口(新架构 v3:Vue SPA 取代 Streamlit) `(删除视频会话),因为语义上属于 videos 这个资源,比塞进 member_manager.py 更清晰;写法沿用项目里"先回推 Oracle,成功后处理本地状态"的既有模式。 -/api/ui/service-status 需要代理 Oracle 的 /api/oracle/activity(浏览器不直连 -Oracle,避免 token 暴露),写法照抄 img_proxy.py 的模式:复用 oracle_sync.get_sync() -已解析好的 base_url/token,不再单独存一份配置。 +/api/ui/service-status 代理同机 fam-edge 的 /api/oracle/activity(浏览器不直连 +fam-edge,避免 token 暴露),走 edge_client 统一读 base_url/token。 """ import re from datetime import datetime @@ -19,7 +18,7 @@ from flask import Blueprint, request, jsonify from .logger import setup_logger from . import db_layer -from .oracle_sync import get_sync +from . import edge_client logger = setup_logger('fam-core.ui_api') @@ -50,7 +49,7 @@ def videos(): date_filter = request.args.get('date') or None page = max(0, request.args.get('page', 0, type=int)) page_size = 15 - rows = db_layer.get_sync_videos(limit=page_size, offset=page * page_size, + rows = db_layer.get_videos(limit=page_size, offset=page * page_size, date_filter=date_filter) return jsonify({"videos": _ser(rows), "page": page, "page_size": page_size}), 200 @@ -58,10 +57,10 @@ def videos(): @ui_bp.route('/api/ui/videos/', methods=['GET']) def video_detail(video_id): """单个视频会话详情 + 时间线事件列表(事件时间轴右侧)。""" - video = db_layer.get_sync_video(video_id) + video = db_layer.get_video(video_id) if not video: return jsonify({"error": "视频不存在"}), 404 - events = db_layer.get_sync_events_for_video(video_id) + events = db_layer.get_events_for_video(video_id) return jsonify({"video": _ser(video), "events": _ser(events)}), 200 @@ -69,16 +68,13 @@ def video_detail(video_id): def video_delete(video_id): """删除视频会话(事件时间轴"删除"入口)。 - 先回推 Oracle 物理删除(events + videos 行 + 磁盘文件),成功后再清理本地 - MariaDB 镜像——增量同步(get_sync_delta)只做 upsert 感知不到删除,不能像 - 命名纠错那样靠 trigger_now() 拉增量顺带清理,必须显式调 delete_sync_video。 - Oracle 回推失败时不清理本地镜像,避免"Oracle 还留着、NAS 却以为删了"的 - 数据不一致,让用户看错误提示后重试。 + 交给 fam-edge 处理:它删 events + videos 行,还要删磁盘上的运动片段文件。 + 迁云前这里还要再清一次 NAS 本地镜像(增量同步只 upsert,感知不到删除), + 现在没有镜像了,fam-edge 删完就是最终状态。 """ - ok, err = get_sync().push_video_delete(video_id) + ok, err = edge_client.push_video_delete(video_id) if not ok: return jsonify({"error": f"删除失败: {err}"}), 502 - db_layer.delete_sync_video(video_id) return jsonify({"status": "ok", "video_id": video_id}), 200 @@ -86,7 +82,7 @@ def video_delete(video_id): def stats(): """统计卡:视频/事件/人物/关注数(可选按日期过滤)。""" date_filter = request.args.get('date') or None - return jsonify(_ser(db_layer.get_sync_stats(date_filter))), 200 + return jsonify(_ser(db_layer.get_stats(date_filter))), 200 def _clean_person(s: str) -> str: @@ -100,7 +96,7 @@ def people(): 对应旧 Streamlit 版本 app.py 里的分组逻辑,原样搬到服务端。 """ - rows = db_layer.get_sync_people() + rows = db_layer.get_people() groups = {} for p in rows: key = p.get('canonical_name') or p['label'] @@ -128,7 +124,7 @@ def people(): def people_clips(): """某人物出现过的运动片段列表(人物卡「运动片段」区块)。 - 按 label/canonical_name 匹配 sync_events.person_list_json → 关联视频 + 按 label/canonical_name 匹配 events.person_list_json → 关联视频 (运动片段优先)。返回片段 video_id/filename/event_start_time/duration/ summary/camera_name + first_ts(缩略图定位)+ clip_events(片段内事件数)。 """ @@ -136,7 +132,7 @@ def people_clips(): if not label: return jsonify({"error": "缺少 label 参数"}), 400 limit = min(20, request.args.get('limit', 10, type=int) or 10) - clips = db_layer.get_sync_people_clips(label, limit) + clips = db_layer.get_people_clips(label, limit) return jsonify({"label": label, "clips": _ser(clips)}), 200 @@ -163,41 +159,25 @@ def attention_events(): @ui_bp.route('/api/ui/named-members', methods=['GET']) def named_members(): """已命名成员真名列表(AI 对话页快捷选择)。""" - return jsonify({"members": db_layer.get_sync_named_members()}), 200 + return jsonify({"members": db_layer.get_named_members()}), 200 @ui_bp.route('/api/ui/model-stats', methods=['GET']) def model_stats(): """云端模型调用统计:按模型聚合 + 最近调用明细。""" - agg = db_layer.get_sync_model_calls_stats() - calls = db_layer.get_sync_model_calls(limit=100) + agg = db_layer.get_model_calls_stats() + calls = db_layer.get_model_calls(limit=100) return jsonify({"aggregate": _ser(agg), "recent_calls": _ser(calls)}), 200 @ui_bp.route('/api/ui/service-status', methods=['GET']) def service_status(): - """服务状态页:NAS 同步状态 + Oracle 实时活动(代理,token 不下发浏览器)。""" - sync = get_sync() - nas_status = sync.status() + """服务状态页:fam-edge 的队列/分割/模型活动(代理,token 不下发浏览器)。 - oracle_data = None - oracle_error = None - if sync.base_url and sync.token: - import requests - try: - r = requests.get(f"{sync.base_url}/api/oracle/activity", - params={'token': sync.token}, timeout=15) - if r.status_code == 200: - oracle_data = r.json() - else: - oracle_error = f"Oracle activity HTTP {r.status_code}" - except Exception as e: - oracle_error = f"连接 Oracle 失败: {e}" - else: - oracle_error = "未配置 oracle_sync.base_url/token" - - return jsonify({ - "nas_sync": nas_status, - "oracle": oracle_data, - "oracle_error": oracle_error, - }), 200 + 迁云前这里还有一块 "NAS 同步状态"(镜像拉取的游标/周期/上次条数),随镜像层 + 一起删了——现在前端读的就是 fam-edge 写的那份库,没有"同步"这个中间状态。 + NAS 侧只剩运动事件推送,它的心跳在 fam-edge 的 service_activity 里, + 已经包含在 activity 快照中。 + """ + data, err = edge_client.get_activity() + return jsonify({"oracle": data, "oracle_error": err}), 200 diff --git a/fam-core/tests/test_db_layer.py b/fam-core/tests/test_db_layer.py new file mode 100644 index 0000000..b70ffd0 --- /dev/null +++ b/fam-core/tests/test_db_layer.py @@ -0,0 +1,155 @@ +"""db_layer 直读 SQLite 的单测(2026-09-13 迁云重写后新增)。 + +重点验证从 MySQL 翻译过来的几处方言:json_each 取代 JSON_CONTAINS、 +substr 取代 LEFT、以及历史脏数据(person_list_json 不是合法 JSON)不能 +把整条查询搞崩——这在 MySQL 下 JSON_CONTAINS 会直接报错,SQLite 下 +json_each 同样会,所以查询里加了 json_valid 前置判断。 +""" +import sqlite3 + +import pytest + +from fam_core import db_layer + +# fam-edge 建表语句的最小子集(列名与 oracle_db.py 保持一致) +_SCHEMA = """ +CREATE TABLE videos ( + id INTEGER PRIMARY KEY AUTOINCREMENT, filename TEXT UNIQUE, camera_name TEXT, + duration_sec REAL, event_start_time TEXT, status TEXT DEFAULT 'pending', + summary_json TEXT, events_json TEXT, people_json TEXT, compute_provider TEXT, + created_at TEXT, updated_at TEXT, processed_at TEXT); +CREATE TABLE events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, video_id INTEGER, ts TEXT, description TEXT, + person_list_json TEXT, person_appearances_json TEXT, is_attention_event INTEGER DEFAULT 0); +CREATE TABLE people ( + id INTEGER PRIMARY KEY AUTOINCREMENT, label TEXT UNIQUE, canonical_name TEXT, + first_seen TEXT, appearances INTEGER DEFAULT 0, source TEXT, features_json TEXT, + display_uid TEXT, updated_at TEXT); +CREATE TABLE model_calls ( + id INTEGER PRIMARY KEY AUTOINCREMENT, provider TEXT, model TEXT, video_id INTEGER, + filename TEXT, started_at TEXT, duration_sec REAL, success INTEGER, error TEXT, + created_at TEXT); +""" + + +@pytest.fixture +def db(tmp_path, monkeypatch): + path = tmp_path / "oracle.db" + conn = sqlite3.connect(path) + conn.executescript(_SCHEMA) + conn.executescript(""" + INSERT INTO videos (id, filename, camera_name, event_start_time, status, processed_at) + VALUES (1, 'motion_20260913_080000.mp4', '客厅', '2026-09-13 08:00:00', 'done', '2026-09-13 08:10:00'), + (2, 'whole_20260912_090000.mp4', '客厅', '2026-09-12 09:00:00', 'done', '2026-09-12 09:10:00'), + (3, 'whole_20260911_090000.mp4', '客厅', '2026-09-11 09:00:00', 'done', '2026-09-11 09:10:00'), + (4, 'motion_20260910_070000.mp4', '客厅', '2026-09-10 07:00:00', 'pending', NULL); + -- video 2 有事件(应展示),video 3 没有事件且不是 motion_(应被过滤掉) + INSERT INTO events (video_id, ts, description, person_list_json, is_attention_event) + VALUES (1, '00:00:05', '汤圆在客厅玩耍', '["汤圆"]', 0), + (1, '00:00:20', '有人靠近门口', '["人物B"]', 1), + (2, '00:01:00', '汤圆和奶奶', '["汤圆","奶奶"]', 0), + (2, '00:02:00', '脏数据事件', '不是合法JSON', 0); + INSERT INTO people (id, label, canonical_name, appearances, source) + VALUES (1, '人物A', '汤圆', 12, 'manual'), + (2, '人物B', NULL, 3, 'llm'), + (3, '汤圆', '汤圆', 5, 'manual'); + INSERT INTO model_calls (provider, model, success, duration_sec, created_at) + VALUES ('gemini', 'flash', 1, 10.0, '2026-09-13 08:10:00'), + ('gemini', 'flash', 0, 20.0, '2026-09-13 08:20:00'), + ('nvidia', 'vila', 1, 5.0, '2026-09-12 08:10:00'); + """) + conn.commit() + conn.close() + monkeypatch.setattr(db_layer, '_db_path', str(path)) + monkeypatch.setattr(db_layer, '_chat_schema_ready', False) + return path + + +# ---------------------------------------------------------------- videos +def test_get_videos_only_returns_content_sessions(db): + """只展示运动片段或含事件的会话:分割 0 段的整段素材不该淹没时间轴。""" + ids = [v['id'] for v in db_layer.get_videos()] + assert ids == [1, 2] # 3 无事件且非 motion_,4 未完成 + assert db_layer.get_videos()[0]['event_count'] == 2 + + +def test_get_videos_date_filter_and_paging(db): + assert [v['id'] for v in db_layer.get_videos(date_filter='2026-09-12')] == [2] + assert [v['id'] for v in db_layer.get_videos(limit=1, offset=1)] == [2] + + +def test_get_video_and_events(db): + assert db_layer.get_video(1)['filename'].startswith('motion_') + assert db_layer.get_video(999) is None + evs = db_layer.get_events_for_video(1) + assert [e['ts'] for e in evs] == ['00:00:05', '00:00:20'] + + +def test_get_attention_events(db): + rows = db_layer.get_attention_events() + assert len(rows) == 1 and rows[0]['ev_date'].startswith('2026-09-13') + + +# ---------------------------------------------------------------- 人物命中(JSON) +def test_people_clips_matches_label_and_its_aliases(db): + """'人物A' 的规范名是"汤圆",而事件里记的是"汤圆"——要能顺着别名找到。""" + clips = db_layer.get_people_clips('人物A') + assert {c['video_id'] for c in clips} == {1, 2} + assert clips[0]['video_id'] == 1 # 按录制时间倒序 + assert clips[0]['first_ts'] == '00:00:05' + assert clips[0]['clip_events'] == 1 + + +def test_person_hit_is_exact_not_substring(db): + """json_each 是精确匹配:查"人物"不该命中"人物B"。""" + assert db_layer.get_people_clips('人物') == [] + assert {c['video_id'] for c in db_layer.get_people_clips('人物B')} == {1} + + +def test_malformed_person_json_does_not_break_queries(db): + """video 2 里混了一条非 JSON 的脏数据,查询必须照常返回而不是整条报错。""" + rows = db_layer.query_events_for_person_date('汤圆', '2026-09-12') + assert len(rows) == 1 and rows[0]['ts'] == '00:01:00' + + +def test_query_events_for_person_date(db): + assert db_layer.query_events_for_person_date('汤圆', '2026-09-13')[0]['description'] == '汤圆在客厅玩耍' + assert db_layer.query_events_for_person_date('汤圆', '2026-09-01') == [] + + +# ---------------------------------------------------------------- people / model_calls +def test_people_and_named_members(db): + assert len(db_layer.get_people()) == 3 + assert db_layer.get_named_members() == ['汤圆'] # distinct 且非空 + ctx = db_layer.get_known_members_context() + assert '- 汤圆(标识:人物A)' in ctx and '- 人物B' in ctx + + +def test_model_calls_stats(db): + stats = {s['model']: s for s in db_layer.get_model_calls_stats()} + assert stats['flash']['ok_cnt'] == 1 and stats['flash']['fail_cnt'] == 1 + assert stats['flash']['avg_duration'] == 15.0 + assert len(db_layer.get_model_calls(limit=2)) == 2 + + +# ---------------------------------------------------------------- 统计 +def test_stats_全量与按日(db): + total = db_layer.get_stats() + assert total['videos'] == 2 and total['events'] == 4 and total['attention'] == 1 + day = db_layer.get_stats('2026-09-13') + assert day['videos'] == 1 and day['events'] == 2 + # 人物数:汤圆(含 label 人物A/汤圆两行都归一到"汤圆")+ 人物B + assert total['people'] >= 2 + + +# ---------------------------------------------------------------- chat_history +def test_chat_history_roundtrip_creates_table_on_demand(db): + """chat_history 是本模块唯一写的表,建表是懒加载的(库文件属于 fam-edge)。""" + chat_id = db_layer.insert_chat_history('今天有人来吗', '有,08:00 有人靠近门口', + '上下文', '2026-09-13', '汤圆') + assert chat_id == 1 + rows = db_layer.get_chat_history(limit=10) + assert rows[0]['user_question'] == '今天有人来吗' + assert rows[0]['created_at'] + assert db_layer.get_chat_history(person_filter='不存在的人') == [] + assert len(db_layer.get_chat_history(date_filter='2026-09-13')) == 1 diff --git a/fam-edge/src/fam_edge/api_gateway/api_gateway.py b/fam-edge/src/fam_edge/api_gateway/api_gateway.py index cb07491..bf5a09e 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -270,9 +270,16 @@ def activity(): disk_info["free_gb"] = round(_shutil.disk_usage(watch_path).free / (1024 ** 3), 1) except Exception as e: logger.warning(f"磁盘空间查询异常: {e}") + # NAS 推送链路(2026-09-13 起 NAS 上只剩 fam-notifier 这一个服务, + # 服务状态页的「NAS 运动推送」卡片读这里;心跳由它定期空 POST 刷新) + motion = { + "heartbeat_age_sec": db.get_motion_heartbeat_age_sec(), + "recent": db.get_recent_motion_events(5), + } return jsonify({ "queue": q_status, "db": db.get_queue_status(), + "motion": motion, "segment": segment, "disk": disk_info, "rclone": _last_activity('rclone'), diff --git a/fam-notifier/config/config.yaml.example b/fam-notifier/config/config.yaml.example new file mode 100644 index 0000000..430f6a4 --- /dev/null +++ b/fam-notifier/config/config.yaml.example @@ -0,0 +1,27 @@ +# fam-notifier 配置(NAS 上唯一保留的服务) +# 复制为 config.yaml 后填实际值;${VAR} 会从环境变量解析(start_notifier.sh 会 source .env) + +motion_notifier: + enabled: true + poll_enabled: true # 轮询主路径(默认开启) + dsm_host: "127.0.0.1" # Surveillance Station 就在本机(NAS),走回环即可 + dsm_port: 5000 + dsm_account: "${DSM_ACCOUNT}" + dsm_password: "${DSM_PASSWORD}" + camera_ids: [2] # 轮询关注的摄像头(Generic_ONVIF-001) + + # 推送目标:甲骨文 fam-edge。这是本服务唯一的出口,单向。 + oracle_base_url: "http://129.146.26.249:5000" + oracle_token: "${ORACLE_SYNC_TOKEN}" + timeout_sec: 10 # 单次 SS 请求超时 + + # 心跳:跟轮询 SS 无关,只是定期空 POST 一下甲骨文的 /api/ss/motion,证明 + # NAS->甲骨文这条推送链路本身还活着(enabled=true 就跑,不受 poll_enabled 影响)。 + # 甲骨文侧 dsm_motion_prefilter.max_heartbeat_age_sec(默认 900s)据此判断"无运动" + # 结论是否可信——这个心跳间隔要明显小于那个阈值,否则会被误判成链路已死。 + heartbeat_interval_sec: 300 + + # 轮询参数 + poll_interval_sec: 60 # 轮询间隔 + poll_window_hours: 2 # 每轮回看窗口(小时),覆盖轮询间隔内的新事件 + batch_size: 100 # 单批推送上限 diff --git a/fam-notifier/requirements.txt b/fam-notifier/requirements.txt new file mode 100644 index 0000000..befad9a --- /dev/null +++ b/fam-notifier/requirements.txt @@ -0,0 +1,3 @@ +# NAS 上唯一要装的依赖:一个 HTTP 客户端 + YAML 配置解析 +requests>=2.31.0 +PyYAML>=6.0 diff --git a/fam-notifier/scripts/start_notifier.sh b/fam-notifier/scripts/start_notifier.sh new file mode 100755 index 0000000..27af021 --- /dev/null +++ b/fam-notifier/scripts/start_notifier.sh @@ -0,0 +1,36 @@ +#!/bin/bash +# fam-notifier 启动脚本(NAS 端) +# 用法: bash scripts/start_notifier.sh +# +# NAS 上没有 systemd,挂了不会自启——重启方式见 docs/DEPLOY.md §1.2。 +set -e + +APP_DIR="$(cd "$(dirname "$0")/.." && pwd)" +CONFIG_FILE="$APP_DIR/config/config.yaml" + +if [ ! -f "$CONFIG_FILE" ]; then + echo "错误: 配置文件不存在: $CONFIG_FILE" + echo "请复制 config/config.yaml.example 为 config.yaml 并填入实际值" + exit 1 +fi + +cd "$APP_DIR" + +# 加载 DSM_ACCOUNT / DSM_PASSWORD / ORACLE_SYNC_TOKEN。 +# set -a 包裹:.env 里是裸赋值,不 export 的话子进程读不到(2026-09-01 踩过)。 +ENV_FILE="$APP_DIR/../.env" +if [ -f "$ENV_FILE" ]; then + set -a + source "$ENV_FILE" + set +a +fi + +if [ -d "venv" ]; then + source venv/bin/activate +fi + +# 包在 src/ 下,不加这行 `python -m fam_notifier` 找不到模块 +export PYTHONPATH="$APP_DIR/src${PYTHONPATH:+:$PYTHONPATH}" + +echo "启动 fam-notifier(轮询 Surveillance Station -> 推送甲骨文)..." +exec python -m fam_notifier diff --git a/fam-notifier/src/fam_notifier/__init__.py b/fam-notifier/src/fam_notifier/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-notifier/src/fam_notifier/__main__.py b/fam-notifier/src/fam_notifier/__main__.py new file mode 100644 index 0000000..59722d0 --- /dev/null +++ b/fam-notifier/src/fam_notifier/__main__.py @@ -0,0 +1,34 @@ +""" +fam-notifier 入口:`python -m fam_notifier` + +没有 Flask,没有对外端口——这个进程只做一件事:轮询 Surveillance Station, +把运动事件单向推给甲骨文的 fam-edge。状态查看去网站的「服务状态」页, +心跳和事件都记在甲骨文那边(service_activity / ss_motion_events)。 +""" +import sys +import time + +from .logger import setup_logger +from .notifier import MotionNotifier + +logger = setup_logger('fam-notifier') + + +def main(): + notifier = MotionNotifier() + if not notifier.enabled: + logger.error("motion_notifier.enabled=false,没什么可做的,退出") + return 1 + notifier.start() + logger.info("fam-notifier 已启动(Ctrl+C 退出)") + try: + while notifier.is_alive(): + time.sleep(5) + except KeyboardInterrupt: + logger.info("收到中断,停止轮询") + notifier.stop() + return 0 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/fam-notifier/src/fam_notifier/config_loader.py b/fam-notifier/src/fam_notifier/config_loader.py new file mode 100644 index 0000000..a884eba --- /dev/null +++ b/fam-notifier/src/fam_notifier/config_loader.py @@ -0,0 +1,32 @@ +""" +配置加载器 - 从 config.yaml 读取配置 +""" +import os +import re +import yaml + + +def _resolve_env_vars(value): + """递归解析字符串中的 ${ENV_VAR} 引用""" + if isinstance(value, str): + def replace_env(match): + env_name = match.group(1) + return os.environ.get(env_name, match.group(0)) + return re.sub(r'\$\{(\w+)\}', replace_env, value) + elif isinstance(value, dict): + return {k: _resolve_env_vars(v) for k, v in value.items()} + elif isinstance(value, list): + return [_resolve_env_vars(item) for item in value] + return value + + +def load_config(config_path=None): + """加载 YAML 配置文件,自动解析 ${ENV_VAR} 引用""" + if config_path is None: + config_path = os.path.join( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), + 'config', 'config.yaml' + ) + with open(config_path, 'r', encoding='utf-8') as f: + raw = yaml.safe_load(f) + return _resolve_env_vars(raw) diff --git a/fam-notifier/src/fam_notifier/logger.py b/fam-notifier/src/fam_notifier/logger.py new file mode 100644 index 0000000..b619400 --- /dev/null +++ b/fam-notifier/src/fam_notifier/logger.py @@ -0,0 +1,34 @@ +""" +日志工具 - 统一格式(stdout + 文件双写) +""" +import logging +import os +import sys + +# 日志目录:fam-notifier/logs/(相对 src 的上一级),失败则退化为仅 stdout +_LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), 'logs') + + +def setup_logger(name='fam-notifier', level=logging.INFO): + """配置并返回 logger(stdout + 文件双写)""" + logger = logging.getLogger(name) + if logger.handlers: + return logger + logger.setLevel(level) + formatter = logging.Formatter( + '%(asctime)s [%(name)s] %(levelname)s %(message)s', + datefmt='%Y-%m-%d %H:%M:%S' + ) + handler = logging.StreamHandler(sys.stdout) + handler.setFormatter(formatter) + logger.addHandler(handler) + # 文件输出(daemon 模式下 stdout 不可见,文件是唯一可追溯日志) + try: + os.makedirs(_LOG_DIR, exist_ok=True) + file_handler = logging.FileHandler( + os.path.join(_LOG_DIR, 'fam-notifier.log'), encoding='utf-8') + file_handler.setFormatter(formatter) + logger.addHandler(file_handler) + except OSError: + pass # 目录不可写时退化为仅 stdout + return logger diff --git a/fam-core/src/fam_core/motion_notifier/motion_notifier.py b/fam-notifier/src/fam_notifier/notifier.py similarity index 88% rename from fam-core/src/fam_core/motion_notifier/motion_notifier.py rename to fam-notifier/src/fam_notifier/notifier.py index c62a989..1b0f45e 100644 --- a/fam-core/src/fam_core/motion_notifier/motion_notifier.py +++ b/fam-notifier/src/fam_notifier/notifier.py @@ -1,5 +1,11 @@ """ -MotionNotifier - NAS 端运动监测通知服务(轮询主路径,2026-08-22 定型) +MotionNotifier - NAS 端运动监测通知服务(轮询主路径,2026-08-22 定型; +2026-09-13 从 fam-core 拆出,成为 NAS 上唯一保留的服务) + +拆出来的原因:摄像头插在 NAS 上(Surveillance Station 就是 NAS 本身),这部分 +搬不走;而其余所有东西——数据库、查询接口、AI 问答、登录——都已经迁到甲骨文。 +拆完之后 NAS 侧没有 Flask、没有 MariaDB、没有对外端口,只有这一个进程单向往 +甲骨文推事件,挂了重启即可,不影响网站。 职责: 在 NAS 本机**轮询**群晖 Surveillance Station 的运动侦测事件 @@ -12,13 +18,15 @@ MotionNotifier - NAS 端运动监测通知服务(轮询主路径,2026-08-22 - 推送后由甲骨文本地落库 ss_motion_events,供 video_processor 本地预过滤使用。 可靠性设计: - - 游标(MariaDB sync_cursor.motion_last_event_id)持久化;重启优先续用 DB 游标, + - 游标持久化在本地 JSON 文件(拆出前存 NAS MariaDB 的 sync_cursor 表, + 现在 NAS 上已经没有数据库了);重启优先续用游标, 停机期间的事件由窗口回看补推;仅首次部署时才初始化为 SS 当前最大 id。 - 推送失败的批次不推进游标,下一轮窗口回看重试(不丢事件)。 - 心跳线程(enabled 即跑,与轮询无关):定期空 POST /api/ss/motion,证明 NAS->Oracle 推送链路存活,供甲骨文侧判断"无运动"结论是否可信。 """ +import json import os import re import time @@ -27,11 +35,33 @@ from datetime import datetime, timezone import requests -from ..logger import setup_logger -from ..config_loader import load_config -from .. import db_layer +from .logger import setup_logger +from .config_loader import load_config -logger = setup_logger('fam-core.motion_notifier') +logger = setup_logger('fam-notifier') + +# 游标文件:拆出前这是 MariaDB 里的一行,现在 NAS 上没有数据库了。 +# 丢了也不致命——窗口回看会把最近的事件补推一遍,甲骨文侧按 event_id 幂等落库。 +_CURSOR_PATH = os.environ.get('FAM_NOTIFIER_CURSOR') or os.path.join( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), + 'data', 'cursor.json') + + +def _read_cursor() -> str: + try: + with open(_CURSOR_PATH, encoding='utf-8') as f: + return str(json.load(f).get('motion_last_event_id') or '') + except (OSError, ValueError): + return '' + + +def _write_cursor(value: str): + try: + os.makedirs(os.path.dirname(_CURSOR_PATH), exist_ok=True) + with open(_CURSOR_PATH, 'w', encoding='utf-8') as f: + json.dump({'motion_last_event_id': str(value)}, f) + except OSError as e: + logger.warning(f"游标写入失败(下轮窗口回看会补推,不丢事件): {e}") _MOTION_NOTIFIER = None @@ -235,15 +265,15 @@ class MotionNotifier: def _init_cursor(self): """启动时初始化 last_event_id。 - 优先使用 DB 已保存的游标:服务停机/重启期间产生的新事件,会在重启后 + 优先使用文件里保存的游标:服务停机/重启期间产生的新事件,会在重启后 被下一轮窗口回看捞到并补推(不丢事件)。 - 仅当 DB 无有效游标(首次部署)时,取 SS 当前最大事件 id 作为起点, + 仅当没有有效游标(首次部署)时,取 SS 当前最大事件 id 作为起点, 避免把历史事件全部回灌一遍。 """ - saved = db_layer.get_motion_cursor() + saved = _read_cursor() if saved and int(saved) > 0: self._last_event_id = int(saved) - logger.info(f"运动通知游标初始化(DB): last_event_id={self._last_event_id}") + logger.info(f"运动通知游标初始化(本地文件): last_event_id={self._last_event_id}") return now = int(datetime.now(timezone.utc).timestamp()) events = self._fetch_events(now - 3600, now) # 最近 1h 用于定位最大 id @@ -295,7 +325,7 @@ class MotionNotifier: if int(e.get('duration') or 0) <= 0: self._zero_dur_ids.add(int(e.get('id'))) self._last_event_id = cursor - db_layer.set_motion_cursor(str(self._last_event_id)) + _write_cursor(str(self._last_event_id)) def _run(self): logger.info(f"MotionNotifier 启动,轮询间隔 {self.poll_interval_sec}s," diff --git a/fam-notifier/tests/conftest.py b/fam-notifier/tests/conftest.py new file mode 100644 index 0000000..ecd43e0 --- /dev/null +++ b/fam-notifier/tests/conftest.py @@ -0,0 +1,7 @@ +import os +import sys + +# 让测试能直接 `from fam_notifier.xxx import yyy`,无需先 pip install -e . +_SRC = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), 'src') +if _SRC not in sys.path: + sys.path.insert(0, _SRC) diff --git a/fam-core/tests/test_motion_notifier.py b/fam-notifier/tests/test_notifier.py similarity index 87% rename from fam-core/tests/test_motion_notifier.py rename to fam-notifier/tests/test_notifier.py index a0962a7..985a591 100644 --- a/fam-core/tests/test_motion_notifier.py +++ b/fam-notifier/tests/test_notifier.py @@ -1,6 +1,6 @@ import time -from fam_core.motion_notifier.motion_notifier import MotionNotifier +from fam_notifier.notifier import MotionNotifier def _cfg(**overrides): @@ -22,7 +22,7 @@ def _cfg(**overrides): def _make_notifier(monkeypatch, **cfg_overrides): monkeypatch.setattr( - "fam_core.motion_notifier.motion_notifier.load_config", + "fam_notifier.notifier.load_config", lambda: _cfg(**cfg_overrides)) return MotionNotifier() @@ -48,7 +48,7 @@ def test_send_heartbeat_success_updates_state(monkeypatch): def fake_post(url, json=None, timeout=None): calls.append((url, json)) return _FakeResp(200, {"status": "ok", "received": 0, "stored": 0}) - monkeypatch.setattr("fam_core.motion_notifier.motion_notifier.requests.post", fake_post) + monkeypatch.setattr("fam_notifier.notifier.requests.post", fake_post) ok = n.send_heartbeat() assert ok is True assert n._last_heartbeat_at is not None @@ -63,7 +63,7 @@ def test_send_heartbeat_http_error_records_failure(monkeypatch): def fake_post(url, json=None, timeout=None): return _FakeResp(500, {}, text="boom") - monkeypatch.setattr("fam_core.motion_notifier.motion_notifier.requests.post", fake_post) + monkeypatch.setattr("fam_notifier.notifier.requests.post", fake_post) ok = n.send_heartbeat() assert ok is False assert n._last_heartbeat_error == "HTTP 500" @@ -75,7 +75,7 @@ def test_send_heartbeat_network_exception_records_failure(monkeypatch): def fake_post(url, json=None, timeout=None): raise _requests.RequestException("connection refused") - monkeypatch.setattr("fam_core.motion_notifier.motion_notifier.requests.post", fake_post) + monkeypatch.setattr("fam_notifier.notifier.requests.post", fake_post) ok = n.send_heartbeat() assert ok is False assert "connection refused" in n._last_heartbeat_error @@ -87,7 +87,7 @@ def test_start_runs_heartbeat_thread_even_when_poll_disabled(monkeypatch): has_motion_in_range_local() 会一直 fail-open,省配额的效果就没了。""" n = _make_notifier(monkeypatch, poll_enabled=False, heartbeat_interval_sec=3600) monkeypatch.setattr( - "fam_core.motion_notifier.motion_notifier.requests.post", + "fam_notifier.notifier.requests.post", lambda url, json=None, timeout=None: _FakeResp(200, {"stored": 0})) try: n.start() diff --git a/fam-ui/src/App.vue b/fam-ui/src/App.vue index 1ff251c..24a337c 100644 --- a/fam-ui/src/App.vue +++ b/fam-ui/src/App.vue @@ -1,11 +1,10 @@