Skip to content

refact(key-change): 运行键收敛为 int task_id + per-run 生命周期硬化 - #67

Merged
Jiadezhende merged 34 commits into
devfrom
refact/key-change
Jul 6, 2026
Merged

Jiadezhende merged 34 commits into
devfrom
refact/key-change

Conversation

@Jiadezhende

Copy link
Copy Markdown
Owner

背景

dev 上运行标识长期存在双身份:对外 source_ip 字符串既是客户端来源、又被当作内部路由键;CQ(ClientQueues)可跨 run 复用、无写入门控,重启/晚到结果易产生跨 run 串扰。本 PR 以"换键"为主线,把内部路由统一到 int task_id,并把 run 生命周期收敛为一次构造、状态机门控、单点编排,顺带完成 stream / inference / persistence 三侧的职责归位与死代码清理。

规模:34 commits / 89 files / +3998 −2767。

关键重构

1. 运行键 source_ip → int task_id(T5)

  • 注册表、per-task 锁、decoder、TemporalActor 全部改为 int task_id 键,删除 run_key str 副本,消除双身份。
  • source_ip 降级为纯诊断字段 + wire 边界适配(find_by_source_ip shim 保留,待前端换键后删)。
  • 业务端点 /ai/video、/api/terminate 双模兼容 task_id/client_id,前端可平滑迁移。

2. CQ per-run 不可变 + 状态机写门(T1/T2)

  • 一个 ClientQueues == 一个 run == 一个 (task_id, step_id),构造后身份不可变、不跨 run 复用。
  • 引入 ACTIVE → DRAINING → CLOSED 状态机,门控所有写入(帧/检测/告警);晚到结果携旧 CQ 句柄撞 DRAINING/CLOSED 即被拒,从机制上杜绝跨 run 串扰。
  • 身份改用 CQ 对象引用判定,拆机加对象身份 fence,防止 superseded run 被二次拆除。

3. 写回句柄化,去反查(T4)

  • L1 检测结果携 CQ 句柄写回(res.cq),取代 has_client(client_id) 反向查找;FeatureStore 增归属校验(owner fence),拒绝旧 CQ 写入,重启时 open_fresh 与 HLS step 目录清理对称。

4. 编排与职责归位

  • 新增 RunController 作为单一编排点,拆机顺序固定:to_draining → stop_stream → finalize_actor → persist_alarms/flush → close_feature_store → registry.remove。
  • persist_alarms 过闸编排从 persistence 移回 inference/temporal/alarm_sink(落库≠过闸编排),persistence 收敛为无状态落库。
  • HLS 分段落盘 PUSH → PULL,CQ 退回纯缓冲容器;persistence 生命周期上移。

5. stream / 清理

  • FFmpegDecoder 自持读循环合并双平台、RTSP-only 精简、子进程回收收敛;restart 同步停旧再起新,消除 ca_ready 双写窗口。
  • 删除 app/services/ai.py 等死代码,清理 CQ / persistence 无生产调用方的接口,frame_drop_total 接全真丢帧点。

影响面 & 兼容

  • 对外 wire 不变:枚举值、DB 列名、JSON 字段保持;前端可继续用 source_ip/client_id(双模)。
  • 新增测试:teardown 身份 fence、写回句柄 fence、聚合 mock 单一真源等。
  • 待办(本 PR 不含):前端 wire 换键(T6)完成后删除 find_by_source_ip shim;offline/ 消费端(Segmenter/Runner)仍为提案 stub。

验证

  • pytest tests/ 通过(含新增 fence/debounce 用例)。
  • 端到端(真实 RTSP)未在本 PR 跑,建议合并前按 DEPLOYMENT.md 走一次 start → 推理 → HLS → terminate 全链路。

🤖 Generated with Claude Code

yinweichen and others added 30 commits July 1, 2026 23:42
- 内核换 copy-on-write 不可变快照 + 单写锁:读全程无锁、枚举零拷贝
  (MappingProxyType),写极少走单锁复制换引用、不阻塞读
- 删 per-client 锁机器与 _task_to_client/_client_to_task 双向索引及 bind_task;
  get_client_by_task_id 改扫描当前快照
- 读写接口分离:读 get / has_client / snapshot / get_client_by_task_id,
  写 get_or_create / remove / remove_if(对象身份 fence,预留给拆除编排)
- 迁移 ~20 调用点(get_client→get/get_or_create,get_all_clients→snapshot,
  client_manager.remove_client→remove) + 同步测试 mock
- 保留 client_id 键(RunRegistry 改名/换 task_id 键留后续);全量 pytest 229 绿

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
T1 — CQ per-run 不可变:
- ClientQueues 一 run 一实例,身份构造注入不可变 (task_id/step_id/source_ip/stage);
  删 set_task/set_stage/clear_task_caches,getter 免锁直读,clear() 只释放 payload
- set_task 建**新** CQ 换槽 (ClientManager.set) + open_fresh 截断存储分区(重启 supersede)
- 配置单一出口 ClientConfig.cq_kwargs()(修 dead-kwargs);stream service 只取不建
- 旧 run settlement 落到捕获的 old_cq,排序不变式消解(actor 固化身份并入 T1)

T2 — CQ 状态机 + 写门 + close()(纯 queues.py):
- RunState ACTIVE/DRAINING/CLOSED + _state_lock(不与 7 payload 锁互嵌)
- 写时门控:帧/推理/检测非 ACTIVE 拒;rendered/temporal 放行清空;
  alarm 非对称(仅 CLOSED 拒,DRAINING 放行 settlement)
- close() 释放 payload 留身份小壳;clear() 接为 close() 别名

T3(编排,已落地部分):
- RunController 对称 start_run/stop_run 单一出口 + client_manager.lock_for
  per-client RLock 共用(api 经 asyncio.to_thread、HM 直调);api 委托、HM 孤儿停流持锁
- 收尾项(未做):stop_run 内 DRAINING 置位 + HM 路径 remove_if fence

tests: 迁移 stage_routing/alarm_increment 至不可变构造;新增 test_cq_immutable_run
+ test_cq_state_machine;pytest 242 passed

docs: 更新 20260628 运行身份/CQ 生命周期契约 + 落地现状

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
拆除单一出口补齐 _teardown_run 规格的两处缺口:

1. DRAINING 置位:stop_run 在停 decoder/落盘前先 cq.to_draining(),
   让 T2 的 DRAINING 门在生产链路真正触发——拆除期停生产者写
   (decoder 抽帧/结果写回/tick),放行 settlement 告警 + HLS flush;
   CLOSED 由 remove(cleanup=True)→close() 兜底。

2. HM 路径对象身份 fence:ReconnectState 加 cq 字段,_enter_reconnect_mode
   进入重连时捕获 cq_A;_handle_reconnecting_client 每 tick 比对当前槽位
   cq,异即弃本次重连;拆除链透传 expected=cq_A 至 stop_run,step 0 核对
   registry[cid] is expected 不符则整段放弃(skipped),step 3 用 remove_if。
   task-timeout/orphan 路径亦传 expected=cq;用户 start/terminate 不传
   (同锁内决策+执行,无 ABA),防 HM 过期决策误删健康新 run。

验证:pytest tests/ 247 passed (242 基线 + 3 teardown fence + 2 reconnect
fence);import app.main 无循环导入。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
写回句柄化(T4):DetectionTask/FrameInference 加 cq 句柄字段,dispatcher pop
时捕获、pool 透传;_write_back_results 改写 res.cq,删 has_client/get(client_id)
两处反查与 write-once 的 client_queues_map;顶层 is_active() 门 + stale_run 计数。
消除 dispatch→infer→write-back 期间换槽导致的跨 run 串台。

FeatureStore 归属校验:feature 腿是 T2 状态门唯一盖不住的(外部落盘 + 分区键
(task,step) 跨 run 共享,is_active() 到落盘之间状态漂移 + open_fresh 可插入截断)。
在 _JsonlBuffer 内以 owner=cq 对象引用校验分区归属,owner set/check 与所有文件
写/unlink 全在 store _lock 下串行(原锁外 _write 收进锁内)→ 无 TOCTOU;顺带消除
既有 flush-vs-enqueue 并发写交错。

FactLedger 解耦:FactLedger 是离线异步写的 store,不应由在线 manager 的 per-run
生命周期同步调度。从 InferenceManager 摘除其 instantiation + open_fresh/close/flush;
类/契约休眠预留,待离线 runner 自行 new+驱动。FeatureStore(在线写)生命周期不变。

pytest tests/ 255 passed(251+4),新增 test_writeback_handle_fence.py、
test_feature_store_owner_fence.py。import app.main 无循环导入。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
InferenceManager 越界处收敛(四角色表 Component owner 归位):

- 拆除口收敛为对称 per-run 接口 start_workflow/stop_workflow(原 set_task/
  remove_client):一把起/停本 run 全部 inference 自有组件(actor + feature 分区),
  stop_workflow 返回 settlement、不持久化;RunController 不感知 actor/feature 内部分解。
- persist_alarms + flush_residual_segments 迁入 PersistenceManager(删 alarm_sink.py),
  告警落库/HLS 切段归 persistence owner;persistence 包零 inference import。
- 告警别名前烧:ClientTemporalActor 构造期解析一次 _stage_alias,产出即烧进 alarm.stage,
  下游 persistence 直接读、不反向 import inference.naming。
- RunController.stop_run 按序调各 owner:stop_workflow → persist_alarms →
  清前端槽 + flush_residual_segments。
- 删 _persist_settlement_alarms/_flush_all_remaining_segments/_client_lifecycle_lock/
  persistence_manager 引用(互斥已由 lock_for 承接)。

persistence 生命周期上移:

- 新增 persistence.lifespan(),main.py 嵌套 health→persistence→ai(起序)/逆序停;
  persistence 先于 inference 起、后于 inference 停,承接停机结算告警+残段 flush 后再抽干队列。
- InferenceManager.start/stop 不再驱动 persistence 起停;停机残余结算走惰性 sink 调用。

仍按 client_id 键;CQ 构造上移/换 task_id 键留 T5。pytest tests/ 260 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
上一提交收敛 __init__ 时误删 self._actors={} 赋值(与 _client_lifecycle_lock
同段),导致 stop_workflow 运行时 AttributeError: 'InferenceManager' object has
no attribute '_actors'。现有测试全 mock/__new__ 绕过真构造故未暴露。

补真构造回归用例守卫 __init__ 必设属性 + stop_workflow 空跑。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- 删 CleaningTask VO:task_id/current_step/status/source_ip/stage 作不可变
  primitive 直构入 CQ;新增 run_key property(str(task_id))作路由键
- CQ 构造职责上移 RunController.start_run(编排者建 CQ→start_workflow(cq));
  InferenceManager 不再自建 CQ,_resolve_stage→公有 resolve_stage
- 边界垫片:/terminate、/ai/video wire 不变(?client_id=<source_ip>),
  经新增 ClientManager.find_by_source_ip(扫描匹配首个)解析回当前 run
- decoder 键/actor 线程名/lock_for 键/日志统一走 run_key;同 task_id 才抢占
  重启,不同 task_id 天然并发独立槽位
- 删 get_or_create/get_task/step_id_of 等复用机器;新增 rekey 垫片守卫测试

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
承接 b040187:把注册表/锁/decoder/actor/feature 的键类型从 str(task_id)
统一改为 int task_id 直取,删掉 ClientQueues.run_key property(消 str/int 双身份)。

- ClientManager: _runs/_task_locks/get/set/remove/remove_if/lock_for 全改 int 键;
  get_client_by_task_id 由扫描简化为 O(1) 直取(键即 task_id)
- RunController.start_run/stop_run: lock_for(task_id)、stop_run(task_id),不再 str(task_id)
- StreamService/FFmpegDecoder、HealthMonitor(_reconnecting_clients/_last_activity/
  ReconnectState)、InferenceManager._actors、DetectionTask/FrameInference.client_id、
  VisualizationWorker._last_rendered_ts 键类型同步改 int
- fix(admin): /clients/{client_id}/alarms 路径参数 str→int,否则 int 键注册表
  恒 miss 返 client_not_found(rekey 回归修复)
- 测试同步:身份构造 primitives、int 键断言

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
rekey 后 /overview 的 client_id = int task_id,但 /ai/video 仍按 source_ip
路由(生产前端契约:?client_id=<摄像头 source_ip>,见 docs/RTSP_FLOW.md)。
admin 实时监控此前把 task_id 塞进 ?client_id= → find_by_source_ip 恒 miss、
无画面。修复:

- _client_info 增补 source_ip 字段(overview/clients 载荷)
- admin 前端 startLive 改用所选 run 的 source_ip 连 WS;缺 source_ip 时提示不连
- 补 admin 页 ElMessage 解构(此前未引入)

lab 页与 lab.py 纯走 task_id/step_id + 文件系统/DB,不碰运行时注册表,无需改动。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
承接 int task_id rekey:把服务层内部的路由键形参/变量/字段名从 client_id
统一改名 task_id,消除"名叫 client_id、值却是 task_id"的语义漂移;面向人的
诊断字段(异常/告警 sink/结果 dict 的 client_id)语义锚定 source_ip,从捕获的
cq.source_ip 取值,不再拿路由键冒充。

- StreamService: start/stop/restart_stream、decoder 字典键、FFmpegDecoder
  形参/属性 client_id→task_id;get_all_client_ids→get_all_task_ids;异常
  client_id 改取 cq.source_ip
- GlobalHealthMonitor: reconnect/orphan/timeout 全链路 client_id→task_id;
  ReconnectState.client_id→task_id
- Dispatcher/Pool/Worker/Actor/models: DetectionTask/FrameInference.client_id
  →task_id;pool 错误上下文、actor persist_alarms、run_control 诊断字段改取
  source_ip
- ClientManager.remove/clear_all 返回 dict 键 client_id→task_id
- 测试与文档示例同步改名
- 新增 docs/update/20260704_RUNKEY_TASKID_LANDING.md:记录换键实际落地
  (source_ip→int task_id 单键)与 T6 wire 迁移暂缓的决策

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
异常身份升级(client_id → task_id+step_id+source_ip):
- AppError 及 7 子类删 client_id 字段, 改带 task_id+step_id(排障主键)+source_ip(辅助);
  __str__ 输出 [task=..][step=..][source_ip=..]; raise 现场从 cq 派生三元身份
  (decoder 加 _err_identity() 复用); 删无调用方的 get_client_id_from_exception。
- main.py 11 handler: 日志 extra 带三元坐标; 响应体 client_id wire 键保留(值=exc.source_ip),
  前端错误响应契约零改。

CQ 身份读取一致化:
- 删冗余 getter get_task_id/get_step_id/get_stage(T1 后身份即不可变公有属性, 直读即可,
  且旧 getter 只覆盖 3/6 字段、直连/getter 混用); 全仓 cq.get_X()→cq.X,
  内联 dispatcher._get_client_stage。

死代码清理:
- 删无调用方的 InferenceManager.get_result(连带 Frame import);
- 删 ClientManager.get_client_by_task_id(换键后与 get 等价, task.py 改直接 get(task_id))。

文档: RUNTIME_IDENTITY/ROUTING_BOUNDARY 标 T5 ✅ 已落地、T6 暂缓待前端;
RUNKEY 落地文档补异常身份升级 + 小清理记录。

pytest tests/ 264 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…step_id

- 删 status(恒 "running",冗余于 RunState) 与 current_step(仅比对/日志/展示用)
- step_id 收敛为唯一步骤身份:构造直收 int(下游 DBAlarm.BigInteger/落盘目录/
  FeatureStore 分区全链路 int),str→int 转换上移 RunController 边界一行 int()
- 删 _resolve_step_id 静态方法与 try/except 优雅 None 仪式(生产 current_step 恒数字)
- 幂等重启比对/重启日志改读 step_id;admin API+面板去 task_status、展示 step_id
- task_started_at 保留并补注释(health_monitor task_max_duration 看门狗消费者)
- _state_lock "0号锁" 黑话注释重写为"独立于 7 把 payload 锁、从不互嵌"

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- 删「总纲+引用」结构:RUNTIME_IDENTITY_BINDING(契约 hub,内容已再落户任务票)、
  MESSAGE_CONTEXT_PROPAGATION(从未实现的 superseded 方案)两篇整体删除
- 三篇任务级文档改为自包含:CLIENT_QUEUES(T1-T3)/CLIENT_ROUTING(T4-T6)/RUNKEY_LANDING
  删「契约前置/序列前置/相关」互引头,契约/概念内联进各文;仅保留一句「承接前序」依赖点名
  + 指向代码/对外 wire 文档(RTSP_FLOW)的真源链接
- _template.md 加「别搞总纲+引用」规则:每篇能独立读完,依赖用一句话点明不跳转取关键内容
- 顺手清 OFFLINE_PIPELINE 相关头里指向已删/换键工单的失效「后续」链接

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
测试基础设施重构(换键/生命周期重构后收口),行为等价、生产逻辑除 bug 修复外零改动。

- 新增 tests/factories.py(领域对象/CQ/推理消息的纯 builder 单一真源,可被
  integration_tests 复用)+ tests/conftest.py(factory-as-fixture + tmp_storage 共享 setup)
- 迁移 12 个测试文件的本地 _cq/_result/_make_record/_out/_frame 等 helper 到共享 factory,
  只换构造不改断言;静态核对 ClientQueues( 除 factories.py 外归零
- 接入 pytest-cov(.coveragerc + requirements),基线 54.7%;补桶 2 轻缺口后 57.3%:
  utils/context 27→96%、utils/decorators 13→76%、routers/admin 12→61%(+28 用例)
- I/O 边界(decoder/CUDA/WS/worker 循环)显式定为集成-only,不硬写 mock 单测
- 顺带修真 bug:_parse_metrics_json 用 families.get("*_total") 查 family,但 prometheus
  对 Counter 剥 _total 后缀,4/5 指标恒 miss;修 admin 4 处 lookup 键 + metrics.py 定义名
  去 _total 对齐约定(sample 名/对外 /metrics/PromQL 全不变,库自动补 _total)

全量 pytest 292 passed(迁移前 264 + 桶2 28)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
换键第一步:两个仍以 ?client_id=<source_ip> 调用的业务端点改为同时接受
task_id(新,首选)与 client_id(旧,=source_ip),task_id 优先直查运行键,
缺 task_id 时回落 find_by_source_ip 垫片。前端全切 task_id 后再下线 client_id。

- /api/terminate: 双可选参数;两参皆缺 → ValidationError;no-op 响应回带实际标识
- /ai/video: per-loop resolve() 闭包(get 锁定 run / find_by_source_ip 跟随);
  非法或两参皆缺 → 1008;日志统一 label(消 task_id 路径下 client_id=None)

admin/lab 前端本步不动(legacy 分支兜底),留后续直接切换。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
ClientQueues.append_ca_* 原本攒够 ca_segment_len 就直接调
persistence_manager.persist_hls_segment(懒 import 硬编码),让数据容器反向
依赖 persistence 服务。分段判定+落盘触发本应归 persistence。

- CQ:append_ca_raw/processed 收窄为纯缓冲(append+丢帧计数+刷 latest_raw),
  删段满触发/persistence 懒 import/step_id 日志;删 pop_n_ca_*,换带锁的
  take_raw_segment/take_processed_segment 供 sweeper 拉取。client 包不再 import persistence。
- persistence:新增 HLSSegmentSweeper(daemon,套 StorageCleanupWorker 模板),
  注入 client_manager.snapshot + persist_hls_segment,每 1s 扫活跃 CQ 拉走攒满整段。
  manager __init__ 接线,start/stop 生命周期(先于 hls_pool 停)。
- 配置:HLSConfig.sweep_interval_seconds=1.0 + yaml。

依赖方向翻正为 persistence→client 单向。运行期只落整段,残段仍由 teardown 的
flush_residual_segments 收尾,段大小/时机与旧 PUSH 等价。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Part 1 — frame_drop_total 变回“丢帧按成因”的单一 Prometheus 真源:
此前该指标仅 stale_run 一类会递增,/admin by_reason 误显“未丢帧”,而真正
load-bearing 的丢帧只落在 per-实例 int + 日志。现在在真丢帧点就地打 reason:
  - decoder: ingress_backpressure(入口背压)/ decode_error(解析异常,与背压拆开)
  - dispatcher: infer_backlog(推理阶段队列淘汰,计数放锁外)
  - queues: raw_backpressure / hls_backpressure(已随并发提交 3c84fa8 落地)
均为纯增量,不动现有 int(仍驱动 [INFER_PRESSURE]/[BACKPRESSURE] delta 日志)。

Part 2 — 退役 FrameDrop:已无生产者(唯一 raise 在 docstring 示例),被
return-False+counter 模型取代;Part 1 后 executor 的 FrameDrop→frame_drop_total
喂入彻底冗余。删 FrameDrop 类 + executor Action.DROP 路径与计数 + utils 引用;
删 test_exception_handling.py 中 6 个仅测 FrameDrop/DROP 的用例(19→13)。

顺带:清 service._inference_loop 里到不了的 FrameDrop/ModelInferenceError 死分支
(其 e.client_id 属性亦不存在),补 write-back 按确定 task_id 记 per-run 失败日志
(success=False 降级帧此前与“真没检到”不可分)。

验证:grep FrameDrop 归零;pytest 290 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
CQ 接口面收窄,消除重复方法与与不可变身份重复的入参:

- 删 to_status_dict(与 get_queue_depths 字节相同),调用点 inference/manager
  status() 改用 get_queue_depths。
- 删 get_latest_result(get_latest_rendered 的薄别名),routers/ai WS 推流改直调
  get_latest_rendered。
- append_alarm_record_with_gate 去掉冗余 task_id 入参,gate_key 改用 self.task_id
  (不可变身份免锁直读);persistence.persist_alarms 调用同步。
- get_task_alarm_message 去掉冗余 task_id 入参,payload 用 self.task_id;
  routers/task 调用同步。所有调用方传的 task_id 本就 == CQ 自身 task_id
  (注册表以 task_id 为键取 CQ),删参语义不变并消除传错键的口子。

测试同步去掉相应实参。pytest 290 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- decoder: 新增 _reader_loop(Windows/POSIX 统一阻塞读),删 _windows_reader_loop/
  on_stdout_ready/os.set_blocking;本地引用 stdout 规避 self.proc TOCTOU
- decoder: start() 秒退分支 + stop() 无条件 wait() 回收僵尸;stop() 锁外 join reader
- service: 删 selector 全套(sel/_selector_loop/run_once/register/unregister/cleanup),
  __init__ 不再起线程;restart_stream 改同步停旧 decoder,消除 ca_ready SPSC 双生产者窗口
- RTSP-only: 内联 _RTSP_INPUT_OPTS,删 protocol/protocol_opts/_build_protocol_opts/RTMP 分支
- fps: 后端内部停传(StartRequest.fps 保留于前端 wire 契约);ReconnectState 删 fps/protocol
- ffmpeg: _build_cmd 直读 settings.ffmpeg_path,去 import 期 FFMPEG_BIN 快照
- tests: 更新 test_reconnect_on_initial_failure 至新签名;全量 290 passed

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- 删模块级 try/except ImportError(会把 client 模块任意真实 ImportError 静默吞成
  None → 背压/队列悄然失效);改 _get_client_manager() 调用期导入并缓存:运行期各
  模块已初始化无顺序问题,真实 ImportError 直接冒出。app.services.client 不 import
  stream,无实际环,故惰性即可、无需 try/except
- 删 has_stream():全仓无调用方。孤儿检测已改用 get_all_task_ids()(看注册、不看
  is_alive),因重连设计需保留「死掉但仍注册」的 decoder;用 has_stream 的 is_alive
  反而会误清待重连 decoder,故其被有意取代
- 全量 290 passed

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…t 死接口

HLSPersistenceStrategy._dir_locks 随 (task_id, step_id) 单调增长且从不回收,
长跑内存慢泄漏。新增 release_dir_locks(task_id) 按 task 前缀批量剔除,由
RunController.stop_run 清 registry 之后调用——此时 CQ 已出 registry、sweeper
扫不到,不会再有该 task 的新段入队;在途残段若再取锁经 _get_dir_lock 按需
重建同一把,串行正确性不受影响。

顺带删掉 manager.flush_remaining → hls_pool.flush_client 这条全无调用方的
no-op 死链,原位改成 release_task_locks。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- service.py: 把 _start_stream_impl 内的 from app.settings import settings 上提到
  模块顶部,与 decoder.py 一致(settings 是模块级单例,顶层 import 无环、更清晰)
- decoder.stop(): 除 reader 外一并 join stderr 线程(daemon,pipe close 后自退),
  消除频繁 restart 时旧 stderr 线程悬空的轻微 churn,与 reader 回收对称
- 全量 290 passed

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
惰性 _get_client_manager() accessor 过度防御——已核实 client/manager.py 顶层仅
import .config/.queues、均不碰 stream/inference,client_manager 单例无环。改为顶层
from app.services.client.manager import client_manager,三处消费点直接引单例。
真实 ImportError 仍 boot 时 fail-fast,不吞。全量 290 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
ClientQueues.get_signals_10s 懒 import inference.naming.get_task_metric_map,把
「流名→metric 映射 + signals_10s 组织」这类展示逻辑塞进队列容器、反向读 inference
运行时状态——client 包唯一残留的 inference 运行时依赖。

- CQ:删 get_frontend_message(死代码)、get_signals_10s、get_task_alarm_message;
  新增纯访问器 get_slide_window_summary(按流名聚合,无 metric 映射)与
  get_alarm_snapshot(单 _alarm_lock 原子返回 (增量,max_seq),保住 max_seq>=max(seq) 不变式)。
  client 包不再有运行时 inference import。
- router task.py:新增 _build_signals_10s(流名→metric.value + 全量空模板)与
  _build_task_alarm_message(原子取告警+滑窗汇总+序列化);get_task_metric_map import
  收敛于此。_empty_alarm_payload 改用 _build_signals_10s({}),空模板/实时路径合一消除重复。
  端点 /task/message/{task_id} 改调装配函数。
- manager.py set() 注释修正:decoder/actor 不由 CQ 持有、由 RunController.stop_run
  显式拆除,不随换槽 GC(原注释误导)。

测试:test_alarm_increment 改测 get_slide_window_summary(流名键);新增
test_task_message_assembler(装配层 metric 映射/空模板/原子入口);test_task_message_api
改走 get_alarm_snapshot mock 路径;404 规避测试行为等价不变。pytest 293 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- actor/store/operator 文档 set_task()/remove_client() → start/stop_workflow()(改名后遗留)
- resolve_stage(current_step) 参数正名 step_id(与 docstring/恒等路由语义一致,调用方全 positional)
- get_batch_for_stage(max_size: int=None) → Optional[int],去掉 `# type: ignore`
- manager.py logger= 从 import 夹层移到 import 块之后

纯清理零行为改动;stage routing + api concurrency 测试全绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
顶层及各子包 __init__ 原平铺 re-export 28+ 符号(含 dispatcher/pool/service/
actor/visualization worker 等内部管件),但无任何调用方经包门面导入——路由/
run_control/manager/测试一律走深路径(.instance/.models/.manager/.feature.store
/.detection.detector/...),StageFactory 亦经全限定 class_path 用 importlib 实例化。

平铺 re-export 反而违背 instance.py 明载的"避免 import app.services.inference.*
触发 eager 构造":取轻量 .models.FrameInference 也会拉起 YOLO/cv2/workflows 重导入链。

改为纯包标记(仅留分层 docstring + 深路径导入约定),去掉全部 re-export 与 __all__。
零调用点改动;pkg / app.main 导入干净。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
- 删 get_ca_ready/raw/processed_capacity:三队列共用同一 maxlen,拆三个 getter 伪装
  「各队列容量可不同」;改存公有属性 self.ca_maxlen 直读(同 ca_segment_len 风格)
- 删 get_ca_raw_length / get_slide_window_latest:全仓零调用方
- 删 get_alarm_increment / get_alarm_max_seq:已被 get_alarm_snapshot 原子超集覆盖,
  仅测试残留;测试迁到 snapshot 单入口
- 生产调用点(decoder/dispatcher)改直读 .ca_maxlen

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
actor/manager 对 persistence 是单向 sink 依赖(persistence 不反 import inference,无环),
惰性 import 只是躲单例 import 期副作用;但 persistence_manager 本就是 package 导入期构造的
模块级单例、app 启动即建,run_control 亦早已顶层 import 之——故提顶层零风险、依赖关系更诚实。

- actor.py: _persist_alarms 内 lazy import → 顶层
- manager.py: stop() 内 lazy import → 顶层,同步去掉「惰性 import」陈述注释

import 无环验证通过;stage routing + api concurrency 全绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
yinweichen and others added 4 commits July 6, 2026 01:10
… 清死代码

PersistenceMetrics 全套(enqueued/completed/errors/queue_full)是只写不读的
影子计数:get_metrics() 全仓零调用方,未接 /health、未接 Prometheus(真指标走
app/utils/metrics.py 的 Counter 体系)。整套删除——models 的 dataclass、manager
的 self.metrics/get_metrics/所有 += 1、__init__ 导出。要观测再走 Prometheus,
不留平行 dict。

顺带:
- hls_strategy: raw/processed 两处 VideoWriter 改 try/finally,异常路径也释放
  原生编码器句柄(原先 write 抛异常时 release() 被跳过)
- manager: 删 vestigial _stop_event(创建即 set、无人读)+ 随之无用的 import threading
- config: 删重复定义的 raw_fps/processed_fps property(前一份引用不存在的
  self.hls.raw_fps,被 settings 单一真源版覆盖,是死代码)

全量 293 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
PersistenceManager.persist_alarms 早先把两件不同域的事捏在一起:client 侧过闸
(cq.append_alarm_record_with_gate,写 CQ 自有去重闸门+前端轮询的告警环形缓冲)
与 persistence 落库。结果 persistence 反向依赖 cq,违反"client 是零跨服务依赖
的中台 leaf、persistence 只做无状态落库"。

新增 inference/temporal/alarm_sink.py 作为编排入口(告警产出域):逐条 set mode
→ 过闸 → 存活者 persist_alarm(dict)。client_id/task_id/step_id 全由 cq 派生。
三个调用点(actor 实时 / inference.manager 停机结算 / run_control 拆除结算)改调
alarm_sink;删除 PersistenceManager.persist_alarms。

接缝归位:闸门归 client、落库归 persistence、编排归 alarm_sink,persistence 不再
import/依赖 cq(仅 flush_residual_segments 的 HLS 切段仍用 cq,属另一议题)。

测试:告警守卫移到 test_alarm_sink.py;test_persistence_sink.py 保留 flush 测试 +
test_persistence_manager_source_has_no_inference_import(搬迁后仍成立,守卫解耦)。
全量 293 passed。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…esh 对称

同 (task_id, step_id) 重启一次 run 时,HLS 段带唯一时间戳文件名不覆盖旧段,只会往同一 playlist 持续累计——旧行为仅靠 7 天 TTL 兜底回收,重启不清。

补一个与拆除侧 flush_residual_segments(cq) 对称的 per-run 起始钩子 persistence.start_run(cq):内部 rmtree 该 step 目录(HLS 磁盘无状态、无 owner-fence,逐段惰性重建)。RunController.start_run 在建新 CQ 后、start_workflow 前调用,与 inference.start_workflow(含 FeatureStore.open_fresh)两个 service 钩子并排——start/stop 两侧同构,run 层只编排、各 service 自封装自己的 storage supersede。全程持 lock_for(task_id),清目录时无活跃 worker。

转发链 manager.start_run -> hls_pool.purge_step_dir -> strategy.purge_step_dir 与既有 release_task_locks 同构。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

该 PR 以“换键”为主线,将后端运行时路由键统一收敛为 int task_id(替代历史上同时承担外部语义与内部键的 source_ip/client_id),并通过 CQ per-run 不可变 + ACTIVE/DRAINING/CLOSED 状态机写门 + 对象身份 fence + 单点 RunController 编排,系统性降低重启/晚到结果导致的跨 run 串扰风险;同时将告警持久化编排从 persistence 侧归位到 inference 域、HLS 分段落盘从 PUSH 改为 PULL、并补齐一组针对上述不变式的回归测试与测试构造基建。

Changes:

  • 运行键全面切换到 task_id:int,并为仍以 source_ip 访问的边界(/ai/video、/terminate 等)保留 shim 兼容。
  • 引入 RunController 作为 run 生命周期唯一编排点,强化 teardown 顺序与身份 fence;写回路径改为携带 res.cq 句柄投递,避免反查导致的串扰。
  • 持久化侧拆分职责:告警过闸编排回到 alarm_sink;HLS 分段落盘引入 HLSSegmentSweeper(PULL);新增/迁移大量测试并引入 tests/factories.py 单一真源。

Reviewed changes

Copilot reviewed 88 out of 89 changed files in this pull request and generated no comments.

Show a summary per file
File Description
tests/test_writeback_handle_fence.py 覆盖“写回只写 res.cq 句柄 + stale run 被门控并计数”的不变式
tests/test_temporal_debounce.py 将检测输出构造收敛到共享 factories,减少散点构造
tests/test_teardown_identity_fence.py 覆盖 stop_run DRAINING 先置位 + identity fence 防误拆新 run
tests/test_task_message_assembler.py 覆盖 router 装配层(signals_10s 模板/映射 + 原子快照入口)
tests/test_task_message_api.py 端点侧改用 client_manager.get(task_id) + 原子入口,更新断言
tests/test_rekey_source_ip_shim.py 覆盖“task_id 为键、source_ip 为被动字段 + find_by_source_ip shim”
tests/test_reconnect_on_initial_failure.py 健康监控重连逻辑对齐 task_id 键,并新增重连对象身份 fence 用例
tests/test_pipeline_drop_counters.py 测试构造收敛到 factories;ClientManager API 改为 snapshot;stage 属性直读
tests/test_persistence_sink.py 覆盖 flush_residual_segments 切段语义 + persistence 无 inference 反向依赖
tests/test_operator_framework.py CQ 构造/检测输出构造收敛;actor 构造键改为 task_id,并验证 stage 别名前烧
tests/test_offline_reservation.py FrameInference 构造收敛;离线直连 FeatureStore 场景明确 cq=None
tests/test_inference_stage_routing.py 将 stage 解析上移为 public resolve_stage,并补 init/stop_workflow 烟测
tests/test_hls_segment_sweeper.py 覆盖 HLSSegmentSweeper(PULL 分段落盘)核心语义与边界条件
tests/test_hls_eff_fps.py frames 构造收敛到 factories(去 numpy 手搓)
tests/test_feature_store_owner_fence.py 覆盖 FeatureStore owner fence:supersede 后拒旧 owner 迟到写/close
tests/test_exception_handling.py 异常处理测试对齐“RETRY/FATAL”与新异常身份字段,移除 FrameDrop 路径
tests/test_decorators.py 覆盖 decorators 的提取/清洗/透明性与 timing slow-warn 行为
tests/test_cq_state_machine.py 覆盖 CQ 状态机、写门、close 释放 payload、迟到写拒绝
tests/test_cq_immutable_run.py 覆盖“一 CQ 一 run 身份不可变 + 换槽 + open_fresh 截断”
tests/test_context.py 覆盖线程本地 context set/get/clear + ClientContext 嵌套与异常还原
tests/test_boundary_layers.py 边界层异常构造对齐 source_ip/task_id 字段变化
tests/test_alarm_sink.py 覆盖 alarm_sink 读 baked stage + gate reject 跳过落库
tests/test_alarm_increment.py 告警 gate/seq 测试对齐新 CQ 构造与 get_alarm_snapshot 原子入口
tests/test_admin_serialization.py 覆盖 admin 序列化/分位数/metrics JSON,并验证 drop counter 分支
tests/factories.py 新增:测试数据构造单一真源(domain/CQ/FrameInference/Alarm 等)
tests/conftest.py 新增:factory-as-fixture + tmp_storage 共享 setup
requirements.txt 增加 pytest-cov(覆盖率量化 opt-in)
requirements-cpu.txt 同上(CPU 依赖集)
docs/update/20260705_TEST_AGGREGATION_MOCK_UNIFY.md 记录测试构造收敛与覆盖率量化的设计与基线
docs/update/20260705_STREAM_READER_UNIFY_CLEANUP.md 记录 stream reader/decoder 清理与资源回收收敛说明
docs/update/20260704_RUNKEY_TASKID_LANDING.md 记录运行键切换到 task_id + source_ip shim 的落地说明
docs/update/20260701_SERVICE_INSTANTIATION_DESIGN.md 服务实例化/惰性边界的设计草案与清单
docs/update/20260628_CLIENT_QUEUES_LIFECYCLE_TASK.md CQ 生命周期工单文档更新(T1–T3)
docs/update/20260626_THREAD_INSTANCE_LIFECYCLE_AUDIT.md 线程生命周期审计文档小幅修订
docs/update/20260620_LAYERED_INFER_DATAFLOW.md 在线/离线分离与 FactLedger 生命周期归属的说明修订
docs/update/_template.md 变更记录模板补充约束(每篇自洽、避免“总纲引用”)
docs/CLEANUP_SIMPLIFICATION.md 文档中 API 名称对齐:get_all_task_ids
config/persistence_config.yaml 新增 HLS sweep_interval_seconds 配置项(PULL 模型)
app/utils/metrics.py Counter family 名对齐 Prometheus 规则(去 _total 后缀),并更新 drop reason 语义
app/utils/executor.py 执行器动作收敛为 RETRY/FATAL,移除 FrameDrop 路径与相关计数逻辑
app/utils/exceptions.py AppError 身份字段升级为 task_id/step_id/source_ip,删除 FrameDrop 与旧提取工具函数
app/utils/init.py 对齐异常导出与文档(移除 FrameDrop/get_client_id_from_exception)
app/static/admin/index.html admin 实时画面按 source_ip 连接 /ai/video;字段 current_step→step_id;补 ElMessage
app/services/run_control.py 新增:RunController 单点编排 start_run/stop_run(状态机门控 + identity fence + 顺序固定)
app/services/persistence/workers/segment_sweeper.py 新增:HLSSegmentSweeper 周期拉取 CQ 满段入队落盘(PULL)
app/services/persistence/workers/hls_worker.py 新增转发接口:release_dir_locks / purge_step_dir
app/services/persistence/strategies/hls_strategy.py 增加 dir lock 回收与 step 目录 purge;并确保 VideoWriter 异常路径 release
app/services/persistence/models.py 删除 PersistenceMetrics(持久化侧收敛为无状态落库)
app/services/persistence/manager.py 引入 segment_sweeper;新增 flush_residual_segments/start_run/release_task_locks;移除 metrics/stop_event
app/services/persistence/config.py 增加 sweep_interval_seconds 扁平访问器;fps 来源改由 settings 单一真源
app/services/persistence/init.py 新增 persistence.lifespan,生命周期不再由 inference 代管
app/services/inference/workflows/init.py 取消 re-export,避免 import 包触发 eager 拉起
app/services/inference/visualization/worker.py 可视化去重键改为 task_id;遍历 client_manager.snapshot;日志语义对齐 run
app/services/inference/visualization/init.py 取消 re-export,按需深路径导入
app/services/inference/temporal/operator.py 文档对齐:operator 实例创建点由 set_task→start_workflow
app/services/inference/temporal/alarm_sink.py 告警落库编排归位 inference 域:过闸 + 落库串联,直接读 baked stage
app/services/inference/temporal/actor.py actor 键改为 task_id;构造期解析 stage alias 并“前烧”到 alarm.stage;tick 线程超时跳过结算
app/services/inference/temporal/init.py 取消 re-export,包标记化并更新说明
app/services/inference/models.py DetectionTask/FrameInference 改为 task_id 键 + cq 句柄透传(写回句柄化)
app/services/inference/instance.py 新增 leaf 单例 inference_manager,避免 import inference.* 即 eager 构造
app/services/inference/feature/store.py 增加 FeatureStore owner fence(open_fresh/append/close),防 supersede 后迟到写串台
app/services/inference/feature/init.py 取消 re-export,并明确 FactLedger 在线侧休眠预留
app/services/inference/detection/service.py 写回改为 res.cq 句柄投递 + stale_run 计数;移除 client_id 反查与旧异常分支
app/services/inference/detection/pool.py 结果对象改为 task_id + cq 透传;错误上下文补 task/step/source_ip
app/services/inference/detection/dispatcher.py 遍历 snapshot;请求携带 cq 句柄;stage backlog 淘汰计数并上报 frame_drop_total(infer_backlog)
app/services/inference/detection/detector.py ModelInferenceError 身份字段对齐 task_id/step_id/source_ip
app/services/inference/detection/init.py 取消 re-export,按需深路径导入
app/services/inference/init.py 顶层取消 re-export,避免轻量导入触发重链路 eager import
app/services/health_monitor/types.py ReconnectState 关键字段对齐 task_id,并加入捕获 cq 的 identity fence 基准
app/services/client/config.py 新增 cq_kwargs 作为创建 CQ 的唯一配置出口,避免 dead-kwargs
app/services/ai.py 删除旧 ai facade(死代码清理)
app/services/init.py 去除对 services.ai 的导入
app/routers/task.py 将 signals_10s 映射与告警消息装配下沉到 router 层(CQ 只提供纯数据访问器)
app/routers/health.py health lifespan 改用 inference.instance.inference_manager(移除 ai.manager)
app/routers/ai.py WS /ai/video 支持 task_id 优先 + client_id(source_ip) 兼容;推理服务生命周期改为 inference_manager
app/routers/admin.py admin 载荷对齐 task_id 键并补 source_ip;修复 metrics family 取值(Counter 去 _total)
app/main.py lifespan 嵌套增加 persistence.lifespan;错误响应继续保留 wire 字段 client_id=source_ip
app/domain/task.py 删除 CleaningTask VO(运行态身份 primitives 直挂 CQ)
app/domain/alarm.py 文档对齐:stage 由 actor 前烧,mode/seq/timestamp 在落日志/持久化补全
.coveragerc 新增覆盖率配置(opt-in)

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

@Jiadezhende
Jiadezhende merged commit 0373afd into dev Jul 6, 2026
1 check passed
Jiadezhende pushed a commit that referenced this pull request Jul 6, 2026
以 dev 现码为准,把 refact/key-change(PR #67)一轮重构的结论融合进 docs/kb/,
21 文件更新 + 新增 SERVICE_RUN_CONTROL.md,统一更新时间 2026-07-06:

- Client/编排:int task_id 注册表(COW)、per-run 不可变 CQ + ACTIVE/DRAINING/CLOSED
  状态机 + 对象身份 fence;新增 RunController 起停编排(SERVICE_RUN_CONTROL.md)
- Inference:四子包分层(detection/feature/temporal/visualization)、Detector/Operator
  框架(analyze+judge 合并,subscribes 必填)、L1 单入口句柄写回、FeatureStore owner
  fence、online/offline 分离(offline 消费端标待实现)
- Persistence/Stream:HLS PULL、告警过闸移回 inference/temporal/alarm_sink、存储根
  settings.storage_base_dir 单一真源、keypoints 死写删;decoder 自持读循环、RTSP-only
- 数据模型:app/domain 分层(Alarm 吸收 AlarmRecord,metric 显式)、app/models.py
- API 双模键(task_id/client_id)、服务实例化 eager/lazy 分层、HealthMonitor 委托 RunController

未落地项按 KB 规范标注:前端 wire 换键(T6)shim 现状、offline 消费端待实现、
服务实例化草案待核验。刻意未动:SERVICE_LAB / BUSINESS_TRACEBACK_AND_LAB /
SERVICE_GATEWAY_MEDIAMTX / TESTING_MAP(与换键无耦合,另轮)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
@Jiadezhende
Jiadezhende deleted the refact/key-change branch September 3, 2026 16:53
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants