Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
10b8501
refact: 清理service中的ai.py,规范化各个模块导出方式
Jul 1, 2026
829f0e5
refactor(client): ClientManager 重构为 COW 中台 + 读写接口分离
Jul 1, 2026
9727b4f
refactor(client): CQ per-run 不可变 + 状态机写门 + RunController 单一编排
Jul 1, 2026
7c684ac
refactor(run): T3 收尾 — stop_run 封闸(DRAINING)置位 + HM 路径对象身份 fence
Jul 1, 2026
505248c
refact(infer): T4 写回句柄化 + FeatureStore 归属校验 + FactLedger 解耦在线调度
Jul 2, 2026
b4ed8e3
refactor(infer): InferenceManager↔persistence 职责收敛 + persistence 生命周期上移
Jul 2, 2026
a73eaac
fix(infer): 恢复 InferenceManager.__init__ 漏设的 self._actors
Jul 2, 2026
b040187
refactor(client): T5 运行键改造 — 注册表键 source_ip→run_key(str(task_id))
Jul 3, 2026
4c41c7b
refactor(client): 运行键收敛为 int task_id — 删 run_key str 副本,消双身份
Jul 3, 2026
f5d67d6
fix(admin): 实时画面按 source_ip 连 /ai/video(注册表改 int task_id 后的路由修复)
Jul 4, 2026
671d688
refactor(client): 内部路由键命名对齐 client_id→task_id + 换键落地文档
Jul 4, 2026
fcb1f8a
refactor(client): 换键收尾 — 异常身份升级 + CQ getter 一致化 + 死代码清理
Jul 4, 2026
0a6217e
refactor(client): ClientQueues 身份字段收敛 — 删 current_step/status,单一 int …
Jul 4, 2026
4a5c25a
docs(update): 换键工单去总纲化 — 任务级自包含 + 删冗余 hub/superseded 两篇
Jul 5, 2026
c13b0b2
test: 聚合 mock 构造为单一真源 + 覆盖率量化 + 顺带修 admin metrics bug
Jul 5, 2026
f7f39c5
feat(api): 业务端点 task_id/client_id 双模兼容 — /ai/video + /api/terminate
Jul 5, 2026
3c84fa8
refactor(client/persistence): HLS 分段落盘 PUSH→PULL,CQ 退回纯缓冲容器
Jul 5, 2026
2fd854a
observ(metrics): frame_drop_total 接全真丢帧点 + 退役 vestigial FrameDrop 异常
Jul 5, 2026
bf2d8d3
refactor(client): 清理 ClientQueues 死接口与冗余 task_id 入参
Jul 5, 2026
44d65c2
refactor(stream): decoder 自持读循环合并双平台 + RTSP-only 清理 + 子进程回收收敛
Jul 5, 2026
88efbd7
refactor(stream): client_manager 惰性导入 + 删死代码 has_stream
Jul 5, 2026
a8c775e
fix(persistence): _dir_locks 任务拆除时回收 + 退役 flush_remaining/flush_clien…
Jul 5, 2026
ca61977
docs(update): stream 清理记录补收尾节(惰性导入 + 删 has_stream)
Jul 5, 2026
531e40c
refactor(stream): settings 顶层导入统一 + stop() 对称 join stderr 线程
Jul 5, 2026
236431b
refactor(stream): client_manager 收敛为直接顶层导入(与 run_control 一致)
Jul 5, 2026
87f208d
refactor(client): CQ 卸掉前端消息装配职责,展示逻辑下沉 router 装配层
Jul 5, 2026
831dbcd
chore(inference): 低危清理 — stale 注释/参数名/签名/logger 归位
Jul 5, 2026
8ce2a46
refactor(inference): __init__ 收敛为纯包标记,不再平铺 re-export 内部管件
Jul 5, 2026
ad8dcbe
refactor(client): CQ 死接口清理 — 删 7 个无生产调用方的访问器,容量统一走 ca_maxlen 直读
Jul 5, 2026
123c843
refactor(inference): persistence sink 惰性 import 提顶层(诚实化)
Jul 5, 2026
a55d4ad
fix: TemporalActor停机后检查,若线程未及时退出,则跳过结算告警
Jul 5, 2026
ac23774
refactor(persistence): 删只写不读的 PersistenceMetrics + VideoWriter 泄漏兜底 +…
Jul 5, 2026
60818ad
refactor(alarm): persist_alarms 过闸编排移回 alarm_sink,persistence 只留无状态落库
Jul 5, 2026
451d342
feat(persistence): 重启同 task supersede 清 HLS step 目录,与 feature open_fr…
Jul 6, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions .coveragerc
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
# 覆盖率配置(opt-in:仅在 `pytest --cov=app` 时生效,不改默认测试行为)。
# 用法:pytest tests/ --cov=app --cov-report=term-missing [--cov-report=html]
[run]
source = app
branch = True
omit =
app/main.py ; 进程启动壳(lifespan 编排),走集成/端到端而非单测
*/__pycache__/*

[report]
show_missing = True
skip_covered = False
precision = 1

[html]
directory = htmlcov
3 changes: 2 additions & 1 deletion app/domain/alarm.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,8 @@ class Alarm:
"""告警:由 Operator.judge()(实时上升沿)/ finalize()(结算)产出。

核心字段(alarm_type/level/message/metric/metadata)由产出方(Operator)填;
mode/stage/seq/timestamp 在落 alarm_log 时由 alarm_sink / 环形缓冲补全(故带默认值)。
stage 由 temporal actor 产出时烧入别名,mode/seq/timestamp 在落 alarm_log 时由
persistence sink / 环形缓冲补全(故均带默认值)。
metric 由产出方显式填,下游持久化直接读 alarm.metric,不靠文案反推。
(原 AlarmRecord 已并入此类——同一份告警从产出到落缓冲只有一种形态。)
"""
Expand Down
15 changes: 0 additions & 15 deletions app/domain/task.py

This file was deleted.

56 changes: 36 additions & 20 deletions app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
from fastapi.responses import JSONResponse, Response

from app.routers import admin, ai, api, health, lab, media, task, traceback as traceback_router
from app.services import persistence
from app.utils import (
AppError,
ConflictError,
Expand Down Expand Up @@ -57,17 +58,20 @@ async def lifespan(app: FastAPI):
shutdown_event = asyncio.Event()
app.state.shutdown_event = shutdown_event

# 按照服务模块启动生命周期管理
# 按照服务模块启动生命周期管理(起序 = 嵌套顺序,停序 = 逆序):
# 1. 健康监控服务(依赖 client_manager, stream_service, inference_manager)
# 2. AI 推理服务
# 2. 持久化服务(平级服务,须先于 inference 起、后于 inference 停,
# 以承接 inference.stop() 的结算告警 + HLS 残段 flush 后再抽干队列)
# 3. AI 推理服务
async with health.lifespan():
async with ai.lifespan():
try:
yield
finally:
# yield 返回时立即通知 WebSocket 退出,不等待后续清理
# 否则:WebSocket 等 shutdown_event → 清理等 WebSocket → 死锁
shutdown_event.set()
async with persistence.lifespan():
async with ai.lifespan():
try:
yield
finally:
# yield 返回时立即通知 WebSocket 退出,不等待后续清理
# 否则:WebSocket 等 shutdown_event → 清理等 WebSocket → 死锁
shutdown_event.set()



Expand Down Expand Up @@ -123,7 +127,7 @@ async def stream_error_handler(request: Request, exc: StreamConnectionError):
logger.warning(
"[BoundaryLayer3] Stream connection error: %s", exc,
extra={
"client_id": exc.client_id,
"client_id": exc.source_ip, # wire 键保留,值=source_ip
"url": str(request.url),
"method": request.method,
},
Expand All @@ -133,7 +137,7 @@ async def stream_error_handler(request: Request, exc: StreamConnectionError):
content={
"error": "Stream unavailable",
"detail": str(exc),
"client_id": exc.client_id,
"client_id": exc.source_ip, # wire 键保留,值=source_ip
},
)

Expand All @@ -151,7 +155,9 @@ async def ffmpeg_error_handler(request: Request, exc: FFmpegError):
logger.error(
"[BoundaryLayer3] FFmpeg error: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"url": str(request.url),
},
)
Expand All @@ -160,7 +166,7 @@ async def ffmpeg_error_handler(request: Request, exc: FFmpegError):
content={
"error": "FFmpeg error",
"detail": str(exc),
"client_id": exc.client_id,
"client_id": exc.source_ip, # wire 键保留,值=source_ip
},
)

Expand All @@ -178,7 +184,9 @@ async def database_error_handler(request: Request, exc: DatabaseError):
logger.error(
"[BoundaryLayer3] Database error: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"url": str(request.url),
},
)
Expand All @@ -205,7 +213,9 @@ async def inference_error_handler(request: Request, exc: ModelInferenceError):
logger.error(
"[BoundaryLayer3] Model inference error: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"url": str(request.url),
},
)
Expand All @@ -214,7 +224,7 @@ async def inference_error_handler(request: Request, exc: ModelInferenceError):
content={
"error": "Inference failed",
"detail": str(exc),
"client_id": exc.client_id,
"client_id": exc.source_ip, # wire 键保留,值=source_ip
},
)

Expand All @@ -232,7 +242,9 @@ async def persistence_error_handler(request: Request, exc: PersistenceError):
logger.error(
"[BoundaryLayer3] Persistence error: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"url": str(request.url),
},
)
Expand Down Expand Up @@ -316,7 +328,9 @@ async def conflict_error_handler(request: Request, exc: ConflictError):
logger.warning(
"[BoundaryLayer3] Resource conflict: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"resource_type": exc.resource_type,
"resource_id": exc.resource_id,
"url": str(request.url),
Expand All @@ -327,7 +341,7 @@ async def conflict_error_handler(request: Request, exc: ConflictError):
content={
"error": "Resource conflict",
"detail": str(exc),
"client_id": exc.client_id,
"client_id": exc.source_ip, # wire 键保留,值=source_ip
"resource_type": exc.resource_type,
"resource_id": exc.resource_id,
},
Expand All @@ -347,7 +361,9 @@ async def cleansight_exception_handler(request: Request, exc: AppError):
logger.error(
"[BoundaryLayer3] CleanSight exception: %s", exc,
extra={
"client_id": exc.client_id,
"task_id": exc.task_id,
"step_id": exc.step_id,
"source_ip": exc.source_ip,
"url": str(request.url),
},
)
Expand Down
51 changes: 27 additions & 24 deletions app/routers/admin.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,13 @@
# 内部工具函数
# ---------------------------------------------------------------------------

def _client_info(client_id: str, client_queues) -> dict:
task = client_queues.get_task()
def _client_info(client_id: int, client_queues) -> dict:
depths = client_queues.get_queue_depths()
task_id = task.task_id if task else None
task_status = task.status if task else None
current_step = task.current_step if task else None
return {
"client_id": client_id,
"task_id": task_id,
"task_status": task_status,
"current_step": current_step,
"client_id": client_id, # 注册表键 = task_id(int)
"task_id": client_queues.task_id,
"source_ip": client_queues.source_ip, # /ai/video 按 source_ip 路由,前端据此连 WS
"step_id": client_queues.step_id,
"queue_depths": depths,
}

Expand Down Expand Up @@ -71,8 +67,8 @@ def _parse_metrics_json() -> dict:
}
result["infer_latency_ms"] = latency_result

# 2. 推理失败 Counter
fail_fam = families.get("infer_failure_total")
# 2. 推理失败 Counter(family 名去 _total 后缀:prometheus 对 Counter 剥 _total)
fail_fam = families.get("infer_failure")
if fail_fam:
by_type: dict = {}
total_fail = 0
Expand All @@ -84,8 +80,8 @@ def _parse_metrics_json() -> dict:
total_fail += sample.value
result["infer_failure_total"] = {"total": int(total_fail), "by_type": by_type}

# 3. 帧丢弃 Counter
drop_fam = families.get("frame_drop_total")
# 3. 帧丢弃 Counter(family 名去 _total)
drop_fam = families.get("frame_drop")
if drop_fam:
by_reason: dict = {}
total_drop = 0
Expand All @@ -97,14 +93,14 @@ def _parse_metrics_json() -> dict:
total_drop += sample.value
result["frame_drop_total"] = {"total": int(total_drop), "by_reason": by_reason}

# 4. GPU OOM Counter
oom_fam = families.get("gpu_oom_total")
# 4. GPU OOM Counter(family 名去 _total)
oom_fam = families.get("gpu_oom")
if oom_fam:
total_oom = sum(s.value for s in oom_fam.samples if s.name.endswith("_total"))
result["gpu_oom_total"] = int(total_oom)

# 5. 重试 Counter
retry_fam = families.get("retry_total")
# 5. 重试 Counter(family 名去 _total)
retry_fam = families.get("retry")
if retry_fam:
by_op: dict = {}
total_retry = 0
Expand Down Expand Up @@ -148,8 +144,12 @@ def _quantile(sorted_buckets: list, q: float, total: float) -> float:

@router.get("/overview")
def get_overview():
"""聚合仪表盘:活跃客户端、队列深度。前端每 3s 轮询。"""
all_clients = client_manager.get_all_clients()
"""聚合仪表盘:活跃 run(任务)、队列深度。前端每 3s 轮询。

换键后注册表键即 task_id,一条目 = 一个活跃 run;响应键 `clients`/`client_id`
沿用旧名(admin 页 wire,值为 task_id),语义已是 run/任务。
"""
all_clients = client_manager.snapshot()
clients_info = [_client_info(cid, q) for cid, q in all_clients.items()]
total_queued = sum(
d["queue_depths"].get("ca_ready", 0)
Expand All @@ -167,17 +167,20 @@ def get_overview():

@router.get("/clients")
def get_clients():
"""活跃客户端列表(轻量),供前端下拉框使用。"""
all_clients = client_manager.get_all_clients()
"""活跃 run(任务)列表(轻量),供前端下拉框使用。

一 run 一 CQ = `registry[task_id]`;响应/路径的 `clients`·`client_id` 为 admin 页 wire 旧名(值=task_id)。
"""
all_clients = client_manager.snapshot()
return [_client_info(cid, q) for cid, q in all_clients.items()]


@router.get("/clients/{client_id}/alarms")
def get_client_alarms(client_id: str, n: int = Query(20, ge=1, le=100)):
"""从内存告警日志读取该客户端最近 n 条告警(不走 DB)。"""
def get_client_alarms(client_id: int, n: int = Query(20, ge=1, le=100)):
"""从内存告警日志读取该 run(task_id)最近 n 条告警(不走 DB)。"""
if not client_manager.has_client(client_id):
return {"client_id": client_id, "alarms": [], "error": "client_not_found"}
cq = client_manager.get_client(client_id)
cq = client_manager.get(client_id)
alarms = cq.get_recent_alarms(n=n)
return {
"client_id": client_id,
Expand Down
Loading