cmvr_edge_ai/docs/architecture.md

64 KiB
Raw Permalink Blame History

cmvr-edge-ai 架构与开发指南

本文描述当前代码已经实现的架构、配置语义和扩展边界。文中的“当前”指运行时 v1规划能力会明确标注避免把配置字段误认为已经具备的调度能力。

1. 目标与边界

cmvr-edge-ai 负责在单台机器人边缘计算设备上,把三类模块可靠地连接起来:

  1. 上游数据cmvr-es 的相机、未来的麦克风音频流,以及平台下发的数据;
  2. 中间 AI视频解码、检测、VAD、ASR、LLM、TTS、决策和策略
  3. 下游输出cmvr-es 的 AGV/机械臂/扬声器控制,以及平台 HTTP/gRPC 接口。

框架刻意保持轻量:默认不依赖 Kafka、Redis、Celery 或 Kubernetes。一个进程内使用 asyncio 管理 I/O 和 DAG需要 CPU/GPU 隔离的插件自行接入有限线程池或常驻 Worker。这样可以在算力、内存和显存有限的边缘端按需增加复杂度。

核心设计约束:

  • 配置驱动YAML 只选择已注册插件、连接端口并声明 QoS/资源,不执行任意 Python 导入;
  • 协议隔离protobuf、HTTP JSON、未来 QUIC 包只在连接器中出现;
  • 内部契约稳定:算法通过 Envelope 和版本化 schema 交流;
  • 有界资源:所有边都有固定容量,满载行为必须显式配置;
  • 故障显式:节点异常默认终止 Pipeline只有连接器显式配置的有界重试和 log_and_drop 等降级策略可以改变该行为;
  • 控制失效安全:执行器只接收已批准命令,且其每个直接前驱都必须是安全门;持续运动还依赖 cmvr-es 服务端 lease/deadman。

2. 分层架构

flowchart TB
    subgraph External["外部系统"]
        CE["cmvr-es<br/>gRPC / future QUIC"]
        PLAT["业务平台<br/>HTTP / gRPC"]
    end

    subgraph Boundary["协议边界"]
        SRC["Source Connectors"]
        SNK["Sink Connectors"]
        POOL["共享 gRPC Channel / HTTP Client"]
    end

    subgraph Runtime["单进程 asyncio 运行时"]
        ENV["Envelope + 内部契约"]
        QUEUE["每条 Edge 的有界队列"]
        OPS["Operator DAG"]
        SAFE["Safety Gate"]
        LIFE["生命周期 / 健康状态 / 优雅退出"]
    end

    CFG["严格 YAML + 插件注册表 + DAG 编译器"] --> Runtime
    CE --> SRC
    PLAT --> SRC
    SRC --> ENV
    ENV --> QUEUE --> OPS
    OPS --> SAFE --> SNK
    OPS --> SNK
    SNK --> CE
    SNK --> PLAT
    POOL --- SRC
    POOL --- SNK
    LIFE --- SRC
    LIFE --- OPS
    LIFE --- SNK

2.1 配置层

config/loader.py 使用 yaml.safe_load 读取文件并递归展开环境变量,随后由 Pydantic 严格模型校验。所有模型均设置 extra="forbid",拼错字段会立即失败。

compiler.py 继续完成无法只靠字段类型判断的检查:

  • api_version 必须等于 cmvr.edge.ai/v1
  • 插件必须存在,工厂创建的对象必须与声明的 Source/Operator/Sink 种类一致;
  • Source 不能有入边Sink 不能有出边Operator/Sink 至少有一条入边;
  • 端口必须存在schema 必须相同或其中一端为通配符 *
  • 图必须无环,不能存在重复边;
  • transport:<name> 标签的连接器只能引用相同 transport 的 endpoint
  • 活跃 Pipeline 中不能包含 enabled: false 的节点;
  • 每个带 actuator 标签 Sink 的直接前驱都必须是带 safety_gate 标签的节点,安全门和执行器之间不能插入其他节点;
  • v1 只接受 execution.mode: async|inline 和默认单并发语义;保留的 mode/并发字段会在编译期被拒绝;
  • 显式保留 runtime 字段、非 normal priority 和非空 resources 会被拒绝,不会被静默忽略;
  • QoS profile 必须是 v1 支持的五种之一,且 overflow 必须符合该 profile 的语义;
  • 直接连接的 detection.model@1 -> detection.repeat_gate@1 会检查规则标签确实包含在 detector 本次选择的标签中;
  • 当前未实现的 max_age_msput_timeout_ms 不能设置。

validate 会实际调用插件工厂,因此也会检查构造函数所需参数、已注册模型 ID 和标签选择; 但不会执行组件的 setup()。ONNX 权重/manifest 是否存在且哈希匹配、网络连通性、 cmvr protobuf、PyAV/ONNX Runtime/Pillow 可选依赖和图 metadata要到 run 的 setup 或首个编码帧时检查PT 回滚 adapter 的 checkpoint 标签也在 setup 检查。启用 alert_image 后 Pillow 会在 repeat gate 的 setup 阶段提前检查。

2.2 内部契约层

协议边界进入运行时后统一变为:

Envelope[T]
├── payload: T
├── schema_name + schema_version
├── source_id + sequence
├── captured_at_ns + received_at_ns
├── deadline_ns
├── trace_id + session_id
└── attributes

Envelope 是冻结 dataclassattributes 被复制为只读 mapping。它只保证浅层不可变如果 payload 是 NumPy 数组等可变对象,发布后必须把它当作只读数据,才能安全地 fan-out 并为以后零拷贝/共享内存保留空间。

当前稳定内部 payload

Schema Python 类型 用途
ImageFrame/v1 ImageFrame 编码视频包或解码后的图像 buffer
AudioChunk/v1 AudioChunk 连续音频块,包含格式、采样率和 discontinuity 标记
DetectionResult/v1 DetectionResult 检测框、标签、置信度和推理时间;可临时附带对应解码帧供告警节点使用
DetectionAlert/v1 DetectionAlert 规则 ID、模型、scope、时间窗口、命中数、置信度、阈值帧检测框、可选 JPEG 告警图和唯一 event ID
InferenceRequest/v1 InferenceRequest 被动调用进入请求 Pipeline 后的统一输入;包含稳定 category、输入、业务参数和请求的 artifact roles
InferenceResponse/v1 InferenceResponse 被动调用 Pipeline 的统一相关响应包含状态、输出、artifact、实际模型引用、耗时、warning 或错误
TextEvent/v1 TextEvent ASR/LLM/TTS 链路中的文本事件
ChatTurn/v1 ChatTurn 带 session 和历史的对话输入
RobotCommand/v1 RobotCommand 带 TTL、序号和 command ID 的协议无关控制命令
ApprovedRobotCommand/v1 ApprovedRobotCommand 由可信安全门批准后交给执行器的控制命令

PluginSpec 中的 schema 是编译期契约字符串。当前运行时不会反射检查 payload 的 Python 类型;连接器和算法插件仍应使用 isinstance 或自己的严格模型在边界处失败。

HTTP 边界不复制一套临时字典协议,而是复用 cmvr_edge_ai.contracts 中严格、版本化的 wire DTO

schema_version Python 类型 HTTP 用途
cmvr.inference-request/v1 InferenceRequest POST /v1/inference 请求
cmvr.inference-response/v1 InferenceResponse 成功受理后的相关响应,包括业务失败状态
cmvr.inference-error/v1 InferenceErrorResponse 非 2xx 的稳定错误响应
cmvr.model-catalog/v1 ModelCatalog GET /v1/models 能力目录,分组列出主动推送与被动调用

Pipeline 端口使用 InferenceRequest/v1InferenceResponse/v1 这样的内部 schema 名; JSON 文档中的 schema_version 使用上表 cmvr.* 名称。两者属于同一 DTO 的内部图契约 和外部 wire 版本,不应互相替换。

2.3 组件层

运行时只有三类组件:

组件 输入/输出 核心方法
Source 无输入,一个或多个输出端口 messages() 返回 async iterator
Operator 一个或多个输入/输出端口 process(envelope, input_port)
Sink 一个或多个输入端口,无输出 consume(envelope, input_port)

所有组件共享以下生命周期:

setup(context) -> start() -> 执行 -> health() -> stop()

建议在 __init__ 只解析轻量参数,在 setup 建立网络连接或加载模型,在 stop 释放资源并保证幂等。ComponentContext 提供 pipeline_idnode_id、全局 shutdown event以及 endpoint、共享连接池和节点配置等 metadata。

Operator 可以返回:

  • None:过滤该消息;
  • Envelope:从默认 output 端口发送;
  • Emission(port, envelope):从指定端口发送;
  • 上述对象的同步或异步 iterable一进多出。

Source 也可以直接 yield EnvelopeEmission。返回原始 bytes、模型对象或其他未包装值会在产生该值的节点附近报错。

2.4 DAG 运行时

当前每个 Pipeline 在一个 asyncio event loop 中运行,每个节点一个 task。execution.mode: asyncinline 在 v1 中都由该 task await 组件方法Operator 和 Sink 每次只处理一个消息。v1 不会复制节点 task也不提供配置驱动的并行调度。

每条有向边拥有独立队列。一个输出端口 fan-out 到多条边时,同一个只读 Envelope 会送入每个队列;队列之间的容量和丢弃计数独立。注意:路由会等待所有边的 put(),因此任意 fan-out 分支使用 overflow: block 且消费变慢时,仍会对共同生产者产生背压。希望“平台慢但控制链不停”时,平台分支必须选择合适的丢弃策略,或在后续加入独立 spool。

多个入边到同一节点时,运行时处理最先就绪的边,不保证不同边之间的全局顺序。单条边内部保持 FIFOdrop_* 导致的丢弃除外)。

节点自然结束后,其全部出边会关闭;消费者会读完已接受的数据后结束。收到 SIGINT/SIGTERM 时:

  1. 设置全局 shutdown event
  2. 取消 Source task并先调用 Source 的 stop()
  3. 让 Operator/Sink 在 shutdown_timeout_s 内排空队列;
  4. 超时后取消剩余 task、丢弃未处理消息
  5. 按 setup 的反序释放组件和共享网络客户端。

任何未被组件显式处理的节点执行异常都会包装为带 node ID 的 NodeExecutionError,停止 Pipeline并由 CLI 以运行错误退出。相机 Source 和平台 HTTP Sink 内部实现了各自明确、 有界的重连/重试语义HTTP Sink 还可显式选择 log_and_drop,但框架仍不提供通用节点级 supervisor 或 retry policy。

2.5 主动推送与被动调用

运行时保留原有主动推送链路,同时增加共享 HTTP 入口的被动调用链路。serve 启动 配置中全部 enabled: true 的 Pipeline 和 HTTP API因此两种服务模式可以在同一 进程并存:

flowchart LR
    subgraph Active["主动推送 active_push"]
        Camera["相机 Source"] --> Decode["解码 / 模型 / 规则"] --> Platform["平台 Sink"]
    end
    subgraph Passive["被动调用 passive_invoke"]
        Client["HTTP Client"] --> API["POST /v1/inference"] --> Broker["category 请求队列"]
        Broker --> RequestSource["server.request_source@1"] --> Model["解码 / 模型"]
        Model --> Response["server.detection_response@1"] --> ResponseSink["server.response_sink@1"]
        ResponseSink --> Broker --> API --> Client
    end
    Routes["server.routes.<category>"] --> Broker

cmvr_edge_ai.capabilities 统一登记 capability 的稳定 category、模型实现和支持的服务 模式;GET /v1/models 返回的 ModelCatalog 分成 active_pushpassive_invoke 两组,并把“代码已注册”和“部署已就绪”作为不同状态。一个同时支持 两种模式的 capability 可以出现在两组中。

公开路由边界是 namespaced category例如 detect.ppedetect.mobile_phone,而不是 模型或 Pipeline 标识。客户端请求不能指定 model_id、Pipeline ID、权重路径或推理 设备;这些都由 server.routes.<category> 和服务端 Pipeline 拥有。响应和能力目录可以 报告实际模型元数据用于追踪,但该字段不是客户端选择器。

HTTP API 把每个 category 交给独立、有界的 invocation broker 队列,并等待相同请求的 相关响应。被动示例的图内边使用 request + block,不配置实时丢帧或 max_fps,保证 每个已接受请求最终得到响应或明确失败。未被 server.routes 引用的主动推送 Pipeline 不受被动调用图约束影响。

3. 配置参考

3.1 根字段

字段 必填 说明
api_version 当前只能是 cmvr.edge.ai/v1
runtime 进程级运行参数
endpoints 命名的外部服务连接信息
pipelines 至少一个命名 Pipeline
server 被动调用 HTTP 监听和稳定 category 到服务端 Pipeline/模型的路由

endpoint、pipeline 和 node 名称只能包含字母、数字、_-,且不能以数字开头。 端口引用固定使用 node.port。公开 category 使用单独的小写 namespaced 格式,详见 server.routes

3.2 runtime

字段 默认值 当前语义
thread_workers 4 创建共享 ThreadPoolExecutor,并通过 ComponentContext 提供给插件
max_processes 1 保留v1 配置中必须省略,显式设置(即使写默认值)也会被拒绝
process_start_method spawn 保留v1 配置中必须省略
shutdown_timeout_s 10.0 优雅排空的最长等待时间
health_bind null 保留v1 配置中必须省略,当前不监听健康检查端口
reserved_memory_mb 0 保留v1 配置中必须省略,当前不执行内存仲裁

v1 可配置的 runtime 字段只有 thread_workersshutdown_timeout_s。边缘端应保持较小的 thread_workers;第三方推理库本身还可能创建线程,应同时限制 OpenMP/MKL/模型运行时线程数,避免过度订阅 CPU。

3.3 endpoints.<id>

字段 默认值 说明
transport 必填,转为小写;当前连接器支持 grpchttp
target null gRPC 地址,例如 127.0.0.1:50052;若供 HTTP Sink 上报 grpc_ip,必须是单个字面 IP 与端口
base_url null HTTP 基础 URL
bind null 为未来 ingress/server 连接器预留
tls false 仅用于 gRPC channel 是否使用 TLSHTTP endpoint 禁止设置该字段
timeout_s 5.0 连接器请求超时;流式相机订阅不使用该 unary timeout
metadata {} gRPC metadata 预留;当前 HTTP 连接器把它作为默认 header
options {} 传输私有参数

当前 gRPC options 支持:

options:
  root_certificates: /path/to/ca.pem
  max_receive_mb: 32
  max_send_mb: 8
  channel_options:
    grpc.keepalive_time_ms: 20000

HTTP 是否使用 TLS 由 base_urlhttps:// scheme 决定,证书验证默认开启;只有隔离测试环境才可显式设置 options.verify_tls: false。当前 gRPC metadata 尚未自动附加到 RPC若需要鉴权应在连接器中显式实现。

3.4 pipelines.<id>

字段 默认值 说明
enabled true 未指定 --pipeline 时是否启动
priority normal v1 必须为 normallow/high 或数值优先级会因无调度器而被拒绝
nodes 至少一个节点
edges [] 有向连接列表

仓库的主动推送配置为 configs/active_detection.yaml,当前只定义 detection。 CLI 显式传 --pipeline detection 时只编译并运行该链路;省略参数时会启动配置中所有 enabled: true 的 Pipeline。不能提供真实能力的 Talk 占位图已移除。

不要在活跃 Pipeline 中保留 enabled: false 节点;当前编译器会直接拒绝。要暂时关闭逻辑,请禁用整个 Pipeline 或从图和配置中移除该节点。

3.5 nodes.<id>

字段 默认值 说明
uses 版本化插件 ID例如 cmvr.grpc.camera_rgb_stream@1
with {} 原样传给插件工厂的参数Python 模型内字段名为 params
execution 见下表 调度意图声明
resources {} 预留资源需求字段v1 必须保持为空/默认值
enabled true 当前只允许活跃 Pipeline 中为 true

execution 字段:

字段 默认值 允许值/限制
mode async v1 支持 asyncinlinethread/process/model_worker 为保留值并被编译器拒绝
concurrency 1 v1 必须为 1,其他值被编译器拒绝
max_in_flight null 预留;非空值被编译器拒绝
timeout_s null 预留;非空值被编译器拒绝
ordered true v1 必须为 truefalse 被编译器拒绝

asyncinline 当前没有调度差异,都是 event loop 中的单 task、单 in-flight 调用。配置模型保留其他枚举值是为了后续版本演进,不代表 v1 能执行它们;validate 会 fail closed而不是静默退化到 async。

resources 模型预留了 cpu_coresmemory_mbgpu_devicegpu_memory_mbexclusive_gpu,但 v1 无法强制执行资源隔离,因此任何非空/非默认声明都会被编译器拒绝。不要依赖“只写说明但不生效”的资源配置;当前应在部署层和插件实现中限制 CPU/GPU。

3.6 edges[] 与 QoS

- from: camera.frames
  to: decoder.frames
  qos:
    profile: video_contiguous
    capacity: 8
    overflow: block
- from: decoder.frames
  to: detector.frames
  qos:
    profile: realtime_latest
    capacity: 1
    overflow: drop_oldest
字段 默认值 当前语义
profile request v1 支持 realtime_latestvideo_contiguousaudio_contiguousrequesttelemetry
capacity 1 队列最大消息数
overflow block drop_oldestdrop_newestblockrejecterrorreject 的兼容别名
max_age_ms null 已预留;设置后编译失败
put_timeout_ms null 已预留;设置后编译失败

溢出行为:

策略 队列满时行为 适合场景
drop_oldest 丢掉最旧消息,再接受新消息 实时视频,优先处理最新画面
drop_newest 保留已有消息,拒收这次新消息但不抛异常 需要保留已排队批次的低优先级数据
block 生产者等待空位 必须连续的视频编码包、音频或请求;会传播背压
reject / error 抛出队列满异常Pipeline 失败 不能静默丢失且希望监督器介入的控制链

v1 强制执行的 profile/overflow 矩阵:

profile 允许的 overflow 典型用途与注意事项
realtime_latest drop_oldestdrop_newest 已解码的完整视频帧;常用容量 1..2,优先限制推理延迟
video_contiguous blockrejecterror H264/H265 等 inter-frame 编码流在解码前不能静默丢包;此保证不覆盖上游在发送前已经丢失的数据
audio_contiguous blockrejecterror 连续音频不能静默丢弃;阻塞预算必须明确,拒绝时应重建连续性
request blockrejecterror 请求/控制;使用有限容量和消息 TTL禁止执行积压旧动作
telemetry 所有五种策略 平台上报;当前无磁盘 spool按数据重要性选择阻塞、丢弃或失败

队列记录入队、出队、水位、丢弃、拒绝和强制关闭丢弃计数,当前可通过 Python PipelineRuntime.edge_stats() 获取,尚未暴露为服务指标。

3.7 server

server 是可选的进程级被动调用配置;只有 cmvr-edge-ai serve 会监听 HTTP。最小 结构如下:

server:
  enabled: true
  http:
    bind: 127.0.0.1
    port: 8081
    request_timeout_s: 30
    max_request_bytes: 16777216
    max_image_bytes: 10485760
    access_log: false
  routes:
    detect.ppe:
      pipeline: detect_ppe
      model_id: construction-ppe-yolov8@2
      queue_capacity: 4
      timeout_s: 30
      description: PPE detection

server 字段:

字段 默认值 说明
enabled false 是否启用被动 HTTP 服务;设为 true 时至少需要一条 route
http 见下表 监听、请求超时与请求体限制
routes {} 稳定 category 到 Pipeline 和模型的服务端映射

server.http 字段:

字段 默认值 说明
bind 127.0.0.1 监听 host/IP不接受 URL scheme 或 path
port 8080 HTTP 监听端口,范围 1..65535
bearer_token null Bearer 凭据;配置后除 /health/live/health/ready 外的入口都要求 Authorization: Bearer <token>;非 loopback 监听必须配置
tls_certfile null Uvicorn TLS 证书文件;必须与 tls_keyfile 同时配置
tls_keyfile null Uvicorn TLS 私钥文件;必须与 tls_certfile 同时配置
allow_insecure_remote false 是否明确允许非 loopback 地址通过明文 HTTP 传输 Bearer token仅限受信隔离网络
request_timeout_s 30.0 route 未单独设置超时时使用的请求总超时,范围 (0, 3600]
max_request_bytes 16777216 JSON 请求体上限,范围 1 KiB1 GiB
max_image_bytes 10485760 Base64 解码后的单张输入图像上限,且不能大于请求体上限
access_log false 是否启用 Uvicorn access log

server.routes.<category> 字段:

字段 默认值 说明
pipeline 必填;服务端拥有的、已启用的请求/响应 Pipeline ID
model_id 必填;服务端选择的版本化模型 ID
queue_capacity 4 该 category invocation broker 的有界请求容量
timeout_s null 可选 route 超时;省略时使用 server.http.request_timeout_s
description "" 能力目录中的部署说明

category 必须是含至少一个点的小写命名空间,例如 detect.ppe;每条 invocation Pipeline 只能被一个公开 category 引用。model_id 会出现在能力目录和响应中用于追踪, 但客户端请求中没有模型、Pipeline 或权重选择字段。

127.0.0.1::1localhost 可在本机调试时不配置认证。任何其他 bind包括 0.0.0.0:: 和普通 hostname都必须设置 bearer_token,并默认要求同时提供 TLS 证书和私钥。只有部署在受信、隔离网络且明确接受 token 明文传输风险时,才能设置 allow_insecure_remote: true 代替 TLS。认证失败返回 401 cmvr.inference-error/v1WWW-Authenticate: Bearer,健康探针保持无需认证。

安全的局域网监听示例token 从进程环境读取,不写入仓库):

server:
  enabled: true
  http:
    bind: 0.0.0.0
    bearer_token: env://CMVR_EDGE_AI_BEARER_TOKEN
    tls_certfile: /etc/cmvr-edge-ai/tls/server.crt
    tls_keyfile: /etc/cmvr-edge-ai/tls/server.key

配置模型与 DAG 编译器会 fail closed 地执行以下约束:

  • route 引用的 Pipeline 必须存在且启用,同一 Pipeline 不能映射到多个 category
  • 所有边都必须使用 request + block,避免已接受请求被 QoS 静默丢弃;
  • 每条被动 Pipeline 必须恰好包含一个 server.request_source@1 和一个 server.response_sink@1;请求 source 的 category 必须与 route 一致;
  • 必须恰好有一个节点输出 InferenceResponse/v1,其 category 也必须与 route 一致;
  • 必须恰好有一个模型所有者 Operator它通过 PluginSpec.route_model_param 声明模型参数, 编译器对所有模型类别校验该参数与 route 的 model_id 一致;
  • 每个 Operator 都必须声明 invocation_cardinality: exactly_one,运行时会先完整验证本次 调用确实只产生一个 emission再向下游路由返回 None、多条或惰性迭代异常都会立即 失败当前请求;
  • v1 请求图必须是 source 到 sink 的单线性路径;在引入显式 join/cardinality 语义前, 禁止 fan-out 后汇合导致同一个请求被重复处理;
  • detect.* Pipeline 必须恰好包含一个 detection.model@1,且不能设置会跳过请求的 max_fps
  • source 到 sink 必须存在有向路径,所有节点都必须位于某条有效请求-响应路径上,并且 不能存在绕过必需 response 节点或 detector 的旁路。

完整的双检测示例见 configs/server_detect.yaml。它为 detect.ppedetect.mobile_phone 分别创建独立 Pipeline所有图内边使用 request/block QoS。

4. 内置插件与连接器

运行 cmvr-edge-ai plugins 可查看当前实际注册结果。

4.1 基础插件

插件 ID 类型 端口 用途
core.sequence_source@1 Source output: * 有限模拟数据源,用于 smoke/replay
core.passthrough@1 Operator input: * -> output: * 占位和协议无关透传
core.log_sink@1 Sink input: * 把 Envelope 摘要和 payload 写日志
safety.robot_command_gate@1 Operator RobotCommand/v1 -> ApprovedRobotCommand/v1 TTL、白名单、参数限值、序号与限频并建立批准边界

安全门参数:

with:
  allowed_actions: [set_velocity, emergency_stop]
  allowed_devices: [src1100]
  min_interval_ms: 100
  max_ttl_ms: 5000
  enforce_monotonic_sequence: true
  argument_limits:
    set_velocity:
      vx: {min: -0.3, max: 0.3}
      vy: {min: -0.2, max: 0.2}
      wz: {min: -0.5, max: 0.5}

allowed_actions 必须非空;allowed_devices 为空表示不额外限定设备。max_ttl_ms 默认 5000拒绝有效期异常长的命令。require_limits_for 可以列出还必须配置参数范围的 actionset_velocity 无条件要求 argument_limits.set_velocity 同时定义 vx/vy/wz。范围边界必须是有限数且 min <= max。过期、过快、重放、乱序、未授权、缺少参数或越界命令都会 fail-closed 丢弃并写审计 warning不会因一条业务拒绝停止整条控制 Pipeline。通过检查后安全门把原始 RobotCommand 包装为带 approved_by/approved_at_nsApprovedRobotCommand

4.2 检测模型、解码与触发规则

检测流水线内置三个模型无关插件:

插件 ID 输入 -> 输出 语义
media.video_decoder.pyav@1 ImageFrame/v1 -> ImageFrame/v1 用可选 PyAV 保持 H264/H265 decoder context输出 packed BGR8stream/session/codec/尺寸变化、sequence gap 或解码错误后重置并等待关键帧
detection.model@1 ImageFrame/v1 -> DetectionResult/v1 加载一个已注册模型,在线程池中推理,按部署标签和阈值二次过滤;可附带对应解码帧
detection.repeat_gate@1 DetectionResult/v1 -> DetectionAlert/v1 按不同帧、时间窗口、scope 和 cooldown 把逐帧检测转换为平台告警,仅在规则触发时按需画框并编码 JPEG

核心包不会强制安装视觉运行时。解码器使用 video extra告警图使用独立的 image extra当前部署检测使用 onnx-cpu,只安装 NumPy、Pillow 和 ONNX Runtime不导入 Torch/Ultralytics。bash scripts/bootstrap.shuv.lock 安装 grpc + http + video + image + server + onnx-cpu--profile dev 改用锁定的 onnx-export-cpu,额外提供 Torch/Ultralytics/ONNX 构建工具。旧 @1 回滚路径仍可 显式选择 yoloyolo-cpu。CUDA、TensorRT 与 Jetson provider 必须按目标驱动 建立独立 source/lock 并重新验证,不能只改 YAML。

具体模型不再各自注册一套 DAG 插件。DetectionModelRegistry 保存 DetectionModelSpec

字段 说明
model_id 版本化稳定 ID例如 construction-ppe-yolov8@2
name 日志、CLI 和告警中的可读模型名称
supported_labels 有序、非空、无重复的标签;顺序必须与 backend 类别编号一致
backend backend 标识,例如 onnxruntime-yolov8
factory 接收 model_options 并返回 DetectionModel 的工厂

cmvr-edge-ai models 输出 model registry 的 ID、name、backend 和 labels。四个模型都保留 ultralytics-yolo @1 导出/回滚注册,并提供 onnxruntime-yolov8 @2;当前三个 部署 YAML 使用 @2。第三方包可通过 cmvr_edge_ai.detection_models entry point 暴露 DetectionModelSpec 或注册回调。

内置检测模型制品通常按 models/detection/<model-name>/vN/ 组织,每个版本目录同时 保存权重和独立 model cardMobile Phone 只有 v1 PT 保留历史扁平布局:

DetectionModelSpec.model_id 尾部的 @N 与制品目录的 vN 对应,例如 construction-ppe-yolov8@1 对应 construction-ppe-yolov8/v1/。这是注册表与制品 库的版本约定;当前 Mobile Phone 制品是导入布局的显式例外。运行时不会由模型 ID 自动推导权重路径,部署配置仍必须显式给出 model_options.weights。标签、训练来源、 评估、局限和许可信息由每个 model card 维护,架构文档只定义制品与运行时的边界。

四个 .pt 模型的 @1 版本继续保留;对应 @2 注册使用 onnxruntime-yolov8。该 backend 只接受静态 [1, 3, imgsz, imgsz] 输入和 [1, 4 + classes, anchors] 原始输出,在 NumPy 中执行 letterbox、类别筛选与 NMS 因此边缘运行时不导入 Torch/Ultralytics。加载阶段先校验 manifest、artifact SHA256、 内嵌模型身份,再校验 ONNX 的 tasknamesimgsz metadata、静态 shape、provider 和注册标签顺序,防止错误或被替换的制品静默运行。

scripts/export_detection_onnx.py 是构建期工具:固定 batch=1、dynamic=false、 nms=false、CPU FP32阻止 Ultralytics 自动安装未锁定依赖,清除训练机路径/时间戳, 并生成带来源/制品 SHA256 的 manifest。四份 v2 制品已经生成并完成真实 ORT smoke 现场发布仍必须补充有授权的非方形图片 parity 与目标硬件资源验收。People-Talking 的 YOLOv8x 计算量不会因为文件格式变化而消失,资源不足时应换小模型或使用现场校准 INT8。

detection.model@1 参数:

with:
  model: construction-ppe-yolov8@2
  detect_labels: [No-Helmet, No-Vest]
  confidence: 0.5
  label_confidence:
    No-Helmet: 0.6
  max_fps: 10
  inference_log_interval_s: 5
  attach_frame: true
  model_options:
    weights: models/detection/construction-ppe-yolov8/v2/model.onnx
    providers: [CPUExecutionProvider]
    intra_op_threads: 1
    inter_op_threads: 1
    imgsz: 640

weights 等相对文件路径按进程启动时的当前工作目录(cwd)解析,不是按 YAML 文件的所在目录解析。仓库内配置和文档命令以仓库根目录为 cwd 从其他目录启动时应使用绝对路径或在部署前将路径正规化。

  • detect_labels 省略时选择模型注册的全部标签;显式空列表、重复或未知标签会失败;
  • confidence 是全局阈值,label_confidence 可逐标签覆盖backend 接收所有选中标签中的最低阈值,通用 Operator 再逐框做严格后过滤;
  • max_fps 是推理启动频率上限,跳过的帧不会产生 DetectionResult
  • inference_log_interval_s 省略时不输出周期推理日志;配置后第一帧立即输出,之后按 周期聚合 window_frames/window_detections/hit_labels/avg_inference_ms。模型加载成功 始终输出一次 INFO。这样可以区分“权重已加载但没有输入帧”“持续推理但没有命中”与 “检测到目标”,同时避免逐帧日志拖慢边缘端;
  • attach_frame 默认 false;启用时 DetectionResult 临时引用对应的解码帧,供 repeat gate 在阈值帧上生成告警图片。由于它携带未压缩图像detector 到 repeat gate 的队列容量必须保持较小;
  • model_options 属于 backendONNX 支持 weights/providers/intra_op_threads/inter_op_threads/imgsz/iou/max_det/agnostic_nms PT 回滚 adapter 支持 weights/device/imgsz/iou/half/max_det/agnostic_nms
  • YOLO 只接受解码后的 BGR8/RGB8 packed buffer并把 RGB8 转为 backend 使用的 BGR 顺序;它不运行 tracker返回的 track_idNone

示例 detector 到 repeat gate 使用小容量 drop_oldest 队列以控制延迟和原始帧内存; 持续过载时丢弃的推理结果不会参与 min_hits。需要每个已完成推理都参与计数时可改 为 block,代价是背压和延迟向上游传播。

如果恢复六类模型,可把同一个 decoder 输出 fan-out 到两个 detector而不是创建两个 独立 Pipeline

camera -> decoder -> construction detector -> repeat gate -> alert API
                  `-> six-class detector -----------------> result API

该可选双分支方案的两条 decoder.frames 出边使用独立的 realtime_latest + drop_oldest 容量 1 队列, 路由的是同一个只读 ImageFrame 引用。这样只建立一个相机 gRPC stream 和一个 PyAV decoder如果把两个分支拆成两个 Pipeline即使 gRPC channel 可以复用,也会创建 两个订阅 RPC 和两个 decoder 实例。fan-out 只复用输入,两个 YOLO 模型仍分别加载并 执行推理,共享进程级有界线程池。示例把六类 YOLOv8n 限制为 5 FPS避免在尚未完成 目标硬件测量前让两个 CPU 模型都追赶相机帧率。

六类模型只包含正向装备类,没有 PersonNo-* 标签;它的 DetectionResult/v1 表示当前推理帧中“检测到了哪些装备”,不能单独推断某个人 缺少装备。该分支保持 attach_frame: false,绕过 repeat gate把每个完成推理的 结果直接发送到 /v1/ppe-detections,所以不会生成 DetectionAlertrule_idhit_count、event ID 或告警图字段;cooldown_ms 只是 repeat gate 的内部规则配置, 不属于告警 payload。输出边使用容量 16 的 telemetry + drop_oldest,模拟 平台持续变慢时会优先保留较新的结果;要求逐条可靠送达时应改用可接受背压的策略或 增加持久 outbox。

detection.repeat_gate@1rules 每项包含 id/labels/min_hits/window_ms/cooldown_ms/min_confidence/scope。一个规则/scope 在同一 (source_id, sequence) 帧最多增加一次 hit不按框数量累加window_ms 内达到 min_hits 才输出告警。触发后清空命中窗口,cooldown_ms 内的帧不累计,冷却结束后重新计数。time_source 支持 captured/received/autocaptured 缺失会失败,auto 优先 captured 后回退 received。

告警图片是 repeat gate 的全局配置,不需要在每条 rule 中重复:

with:
  time_source: received
  alert_image:
    enabled: true
    jpeg_quality: 85
  rules:
    - id: no-helmet
      labels: [No-Helmet]
      min_hits: 3
      window_ms: 2000

alert_image.enabled 默认 false;启用时上游 detector 必须设置 attach_frame: truejpeg_quality 默认 85只接受 1..95 的整数。Pillow 只在 规则达到阈值并准备发出告警时绘制 bounding boxes、标签和置信度并编码 JPEG 不会为每个 DetectionResult 产生图片。

scope: source 按相机和规则隔离状态;scope: track 还按 track_id 隔离,同一 track ID 在不同相机之间不会混合。没有 track ID 的检测会被 track 规则忽略。当前 Construction PPE 和 People-Talking 告警分支都使用 source;要判断同一个人,需要先增加 tracker 和人员/检测框关联节点。

每次触发创建 DetectionAlert.event_id。告警的 detections 和图片中的 bounding boxes 只来自达到 min_hits 的阈值帧,以控制内存和 HTTP payload窗口内更早帧只参与 labels/max_confidence/first_seen/last_seen/hit_count 的聚合。内部可选图片为 JPEG bytesHTTP wire 中 payload.image 被序列化为扁平对象:

{
  "media_type": "image/jpeg",
  "width": 1280,
  "height": 720,
  "encoding": "base64",
  "data": "/9j/4AAQSk..."
}

关闭图片,或因第三方结果未附带帧、坏帧等原因渲染失败时,告警仍会发送且 payload.imagenull证据图失败不能阻断结构化告警。Base64 比原始 JPEG 额外增加约三分之一体积,平台入口和反向代理必须设置相应的请求体上限。 同一阈值帧同时触发多条规则或多个 track 时,节点对所有相关告警框取并集并只编码 一次,再让这些告警共享同一个不可变 JPEG从而限制瞬时 CPU 与内存开销。 非有限坐标、反向/退化框、非法置信度或 track ID 的 detection 会在规则计数前被 忽略并增加 invalid_detections 健康计数,防止 NaN/Inf 破坏整个 HTTP JSON。

4.3 cmvr-es gRPC 连接器

cmvr.grpc.camera_rgb_stream@1

  • 参数:endpoint、必填 device_id、可选 pixel_format(默认 BGR8)、 stream_log_interval_s(默认 30 秒),以及 reconnect/reconnect_initial_s/reconnect_max_s/max_reconnect_attempts
  • 每个首次连接和重连 session 都先调用 unary CameraService.StartCamera,复用 endpoint 的 timeout_s,并检查业务反馈 header.success/error_message;只有成功后才调用 CameraService.GetRGBImageStream。这两步不能省略:前者打开物理相机,后者在服务端 启动编码流;
  • 发送一次订阅请求并保持 request side兼容当前 cmvr-es 行为和未来真正双向协议;
  • color_frame 转为 ImageFrame/v1,保留 codec、关键帧、内参、序号和时间戳
  • 默认在 Start RPC、断流、流 RPC 或业务失败后指数退避重连;每次尝试使用新的 session_id,收到一帧后重置连续失败计数;max_reconnect_attempts 省略表示持续尝试;
  • INFO 日志覆盖 source configured、StartCamera 请求/成功、stream opening、首帧、周期 progress、stream close 和 source stopWARNING/ERROR 覆盖业务反馈错误、断流、重连与 重试耗尽。progress 即使尚未收到首帧也会定时输出,并报告帧数、吞吐、关键帧、 codec、尺寸和最后一帧年龄但不会输出图像 bytes
  • stop 时主动 cancel stream并且等待重连期间也能被 shutdown event 立即打断。默认 不调用设备级 StopCamera,避免中断同一相机的其他客户端;
  • 不做 H264 解码,编码包必须交给有状态解码器插件。

重连不能修复 cmvr-es 生产端已经丢失的编码包。当前 cmvr-es 服务端使用 getLatestEncodedFrame 取得最新数据,负载或时序竞争可能在 edge-ai 看到数据之前跳过 H264/H265 参考包;本地 video_contiguous 只保护已经入队的数据,而且 camera Source 的本地 sequence 不能可靠表示这种上游跳包。decoder 遇到 FFmpeg 错误会重置并等待关键帧,但上线前仍应验证服务端提供连续 access unit或增加原始/JPEG/显式 discontinuity 的 AI 接口。

cmvr.grpc.agv_command_sink@1

  • 参数:endpoint 和固定、非空的 device_id;一个 Sink 只绑定一台设备;
  • 输入必须为 ApprovedRobotCommand/v1,直接传入原始 RobotCommand 会被拒绝;
  • 支持 emergency_stopclear_faultpause_navigationresume_navigationcancel_navigationstop_velocity_controlstop_mappingset_velocity
  • set_velocity 默认禁用;只有显式设置 unsafe_allow_unleased_velocity: true,并在 Sink 再配置完整 velocity_limits.vx/vy/wz 后才会发出;
  • set_velocityvx/vy/wz 必须是有限数,并同时通过安全门和 Sink 两层范围校验;
  • 同时检查 gRPC 调用异常与响应 header.success,并把 RPC timeout 限制在命令剩余 TTL 内;
  • 不自动重试运动命令;正常关闭时若曾成功发送速度,会尽力调用 stopVelocityControl

下面只展示 actuator 节点的危险参数;该节点仍必须直接连接前述安全门。此开关只应用于架空轮、受控台架等隔离环境:

actuator:
  uses: cmvr.grpc.agv_command_sink@1
  with:
    endpoint: cmvr_es
    device_id: src1100
    unsafe_allow_unleased_velocity: true
    velocity_limits:
      vx: {min: -0.3, max: 0.3}
      vy: {min: -0.2, max: 0.2}
      wz: {min: -0.5, max: 0.5}

unsafe_allow_unleased_velocity 必须是 YAML 布尔字面量,字符串 "true"/"false" 会被拒绝。该 override 不是 crash-safetyedge-ai 被 SIGKILL、崩溃、断电或与 cmvr-es 网络分区时Python stop() 没有机会可靠执行。真实机器人必须由 cmvr-es 服务端持有速度 lease并在 lease/heartbeat 超时后独立执行清零和停车。

4.4 平台 HTTP 连接器

platform.http_json_sink@1 接收任意 schema

  • 参数:endpointpath(默认 /)、可选 grpc_endpointfailure_modemax_attempts(默认 3retry_initial_sretry_max_sretry_statuses
  • POST Envelope 元数据、输入端口、attributes 和 payload
  • 配置 grpc_endpoint 时,从被引用 gRPC endpoint 的 target 提取字面 IP并在 Envelope JSON 顶层增加 grpc_ipsource_id 仍表示相机/数据源;
  • bytes 转换为 {encoding: base64, data: ...}
  • EncodedImage 特判为 media_type/width/height/encoding/data 同层的扁平对象,避免 payload.image.data 再嵌套一层;
  • payload 有非空 event_id 时以它作为 Idempotency-Key,否则使用 trace_id;平台仍必须真正实现按键去重;
  • 对连接/超时类错误以及默认 408/425/429/500/502/503/504 执行有限指数退避;不在 allowlist 的状态不会重试;最终失败由 failure_mode 决定:raise 抛异常, log_and_drop 输出 WARNING、丢弃当前报告并继续
  • 没有持久化 outbox/spool。有限重试只存在于当前进程内崩溃、断电或重启不会恢复尚未投递的告警。

例如 grpc_endpoint: cmvr_esendpoints.cmvr_es.target192.168.0.119:50052 时,请求外层包含:

{
  "schema": "DetectionAlert/v1",
  "source_id": "wrist_cam",
  "grpc_ip": "192.168.0.119",
  "input_port": "input",
  "payload": {"event_id": "..."},
  "attributes": {}
}

为避免 DNS 多地址与运行时漂移,grpc_endpoint 引用的 target 必须是单个字面 IPv4/IPv6 加端口hostname、Unix socket 和多地址 target 会在 Sink setup 时被拒绝。

同一个通用 Sink 可以按节点配置不同的 endpoint 和 path。当前检测配置把 DetectionAlert/v1 发到 ppe_alert_platform/v1/detection-alerts。若恢复六类模型, 可把每次推理的 DetectionResult/v1 发到独立 endpoint它与告警复用 Envelope JSON 外层,但 payload 契约不同,且没有 event ID幂等键会回退为 trace_id。Sink 不支持 通过 YAML 重命名或重排 payload 字段,如果目标平台要求自定义 wire contract应增加 平台专用转换节点或 Sink。

运行时对未处理异常采用 fail-fast。HTTP Sink 默认 failure_mode: raise,最终发送失败 会终止它所在的 Pipeline当前告警 Sink 显式使用 log_and_drop,所以平台离线只会 产生 WARNING 并丢弃对应告警,不会停止相机和检测。该模式不是可靠投递机制。

4.5 被动调用边界插件与 HTTP API

被动 Pipeline 使用四个内置边界/转换插件:

插件 ID 输入 -> 输出 用途
server.request_source@1 无 -> InferenceRequest/v1 从对应 category 的 broker 取出已接受请求
media.image_decoder.pillow@1 InferenceRequest/v1 -> ImageFrame/v1 解码内联 JPEG/PNG并执行媒体类型、尺寸和像素数检查
server.detection_response@1 DetectionResult/v1 -> InferenceResponse/v1 规范化检测框、实际模型信息和按需生成的 annotated/original JPEG artifact
server.response_sink@1 InferenceResponse/v1 -> 无 按内部 request ID + invocation token 完成 broker waiter把结果交还 HTTP 请求

cmvr-edge-ai serve 通过 cmvr_edge_ai.server 启动这些 Pipeline 和共享 HTTP API。 POST /v1/inference 只接受 application/json请求体、单张解码图像、category、输入类型、 artifact role 和业务参数都在进入 Pipeline 前校验。GET /v1/models 返回能力及动态部署 状态;GET /health/live 只表示进程可响应,GET /health/ready 还要求应用与所需 Pipeline 就绪。入站 HTTP 支持由 server extra 提供,调用端 SDK 位于 cmvr_edge_ai.client.detect,使用 http extra。

单条解码、模型或响应转换异常会通过 Envelope 中不可复用的 invocation token 只失败对应 callerPipeline 继续处理后续请求。客户端超时或取消后,已经被 source 取出的任务仍占用 该 route 的 broker capacity直到模型返回 late result/failure 或服务关闭;这样不会把 仍在执行的旧推理伪装成空闲资源并继续接收新任务。尚未进入 DAG 的排队请求则可立即移除。

5. 插件开发约定

5.1 创建组件

插件工厂固定接收 node_id: strparams: Mapping[str, Any],应在构造阶段检查无需 I/O 的必填参数。下面展示一个最小 Operator

from collections.abc import Mapping
from typing import Any

from cmvr_edge_ai.contracts import TextEvent
from cmvr_edge_ai.core import Emission, Envelope, Operator


class UppercaseOperator(Operator):
    def __init__(self, node_id: str, params: Mapping[str, Any]) -> None:
        self._node_id = node_id
        self._prefix = str(params.get("prefix", ""))

    async def process(
        self,
        envelope: Envelope[Any],
        input_port: str = "input",
    ) -> Emission:
        del input_port
        event = envelope.payload
        if not isinstance(event, TextEvent):
            raise TypeError(
                f"{self._node_id} expected TextEvent, got {type(event).__name__}"
            )
        result = TextEvent(
            text=self._prefix + event.text.upper(),
            role=event.role,
            final=event.final,
        )
        return Emission(
            "output",
            envelope.with_payload(
                result,
                schema_name="TextEvent",
                schema_version=1,
            ),
        )

使用 with_payload() 可以保留 trace、session、sequence、采集时间和 deadline只更新载荷与接收时间。除非开始了一次新的独立请求不要随意生成新的 trace_id

5.2 注册 PluginSpec

from cmvr_edge_ai.plugins import PluginKind, PluginRegistry, PluginSpec


def register_plugins(registry: PluginRegistry) -> None:
    registry.register(
        PluginSpec(
            plugin_id="example.text_upper@1",
            kind=PluginKind.OPERATOR,
            factory=UppercaseOperator,
            inputs={"input": "TextEvent/v1"},
            outputs={"output": "TextEvent/v1"},
            description="Uppercase a text event",
            tags=frozenset({"domain:dialogue"}),
        )
    )

规则:

  • ID 必须包含版本,例如 @1;不兼容的端口或参数变更发布新版本;
  • Source 不能声明 inputSink 不能声明 output
  • 端口名必须与组件实际发送/接收的端口一致;
  • 使用精确 schema只在真正协议无关的调试/路由节点使用 *
  • 执行器 Sink 加 actuator 标签,并只声明/接受 ApprovedRobotCommand/v1;安全节点加 safety_gate 标签并完成 RobotCommand -> ApprovedRobotCommand 转换;不要为了通过编译给普通变换节点冒充安全标签;连接器加 transport:grpc 等标签;
  • 同一个 ID 重复注册会失败。

若 Operator 要进入 server.routes 引用的被动 Pipeline还必须声明 invocation_cardinality=InvocationCardinality.EXACTLY_ONE。其中恰好一个模型所有者还要 设置 route_model_param(例如 "model_id"),让编译器验证服务 route 与实际 adapter 使用同一模型。当前被动 v1 不接受可能返回 0/N 条的插件,也不接受分支/汇合图;主动上报 Pipeline 不受这组约束。

5.3 通过 entry point 发布

第三方插件包的 pyproject.toml

[project]
name = "cmvr-edge-ai-example-plugin"
version = "0.1.0"
dependencies = ["cmvr-edge-ai>=0.1,<0.2"]

[project.entry-points."cmvr_edge_ai.plugins"]
example = "my_cmvr_plugin.plugins:register_plugins"

entry point 可以直接导出一个 PluginSpec,也可以像示例一样导出接收 registry 的回调。安装包后,默认 registry 会在 CLI 启动时发现它。只安装可信插件entry point 是 Python 可执行代码,虽然 YAML 本身不能任意导入模块,已安装插件仍拥有当前进程权限。

5.4 阻塞计算与模型

禁止直接在 process() 中执行长时间阻塞 I/O 或 Python/模型推理,否则整个 event loop 的 gRPC、HTTP、队列和其他 Pipeline 都会停顿。

框架提供 workers.run_blocking()。插件在 setup() 中取得应用共享的有界线程池,再显式 offload

from cmvr_edge_ai.workers import run_blocking

async def setup(self, context):
    self._thread_executor = context.metadata["thread_executor"]

async def process(self, envelope, input_port="input"):
    result = await run_blocking(
        self._blocking_infer,
        envelope.payload,
        executor=self._thread_executor,
    )
    return envelope.with_payload(
        result,
        schema_name="DetectionResult",
        schema_version=1,
    )

线程 offload 适合释放 GIL 的推理库或阻塞 SDK。不要省略 executor 后假定 runtime.thread_workers 仍然生效;未指定 executor 时会落到宿主 event loop 的默认 executor。

纯 Python CPU 密集任务、大模型或需要故障隔离的模型应由插件维护常驻进程/模型 Worker不要每帧创建进程或重复加载模型。框架提供 PersistentProcessWorker 作为小型基础设施:它使用一个常驻子进程、有界请求/结果队列、串行关联和可选超时,并把子进程异常还原为 RemoteWorkerError。请求一旦超时或在途取消Worker 会标记为 poisoned 并拒绝后续 submit避免迟到结果被误配插件必须 stop 后新建 Worker。插件仍负责在 setup/start 创建和启动它、在 stop 回收它,并确保 spawn 模式下 factory/payload 可序列化。

需要不同 Python/native 依赖的外部项目使用 FramedSubprocessWorker。它以 4 字节大端 JSON header 长度、JSON header、blob_lengths 和 raw blobs 传输数据,启动时要求 ready 握手,并串行关联请求。超时、取消、破损帧会终止当前进程组,下一次请求重新 启动,避免迟到响应串到新请求。gauge.analog_reader@1 使用该边界把 Python 3.8 的 ETHZ Analog Gauge Reader 与 Python 3.10 主服务隔离worker 的 stdout 专供协议, 第三方模型日志重定向到 stderr。部署与独立 uv.lock 位于 server/gauge/worker/

在 v1 配置中,节点仍必须使用 execution.mode: async(默认)或 inline,且保持默认单并发字段。thread|process|model_worker 以及非默认并发字段会在编译期被拒绝。上面的线程/进程 helper 是组件内部显式调用的实现细节,框架不会根据 YAML 自动 offload。插件需要自行限制 in-flight 数量并在 stop() 回收 Worker。若跨进程发送大图像优先传编码数据确认复制成为瓶颈后再实现固定大小共享内存池并只通过 IPC 传 slot/shape/dtype/时间戳。

5.5 开发新协议连接器

新增 gRPC、HTTP、UDP 或 QUIC 支持时,应把 wire format 的解析和错误映射放在 connectors//transports/,对 DAG 只暴露内部契约:

wire message -> 校验 -> 内部 payload -> Envelope -> AI Operators
AI result -> typed command/event -> wire message -> external service

不要让算法插件导入 *_pb2。连接器应:

  • setup 获取/建立连接,在 stop 取消 stream 并释放资源;
  • 为每条消息保留设备 ID、序号、采集时间、trace/session
  • 区分 transport failure 与业务响应失败;
  • 明确超时、重连、幂等和背压策略;
  • 对流式编码数据保持状态,不能把任意 H264 packet 当作独立图片;
  • 为音频丢包/重连产生 discontinuity,不能静默拼接不连续 PCM
  • 对控制接口使用显式 action 到 RPC 的白名单映射,禁止从配置反射任意方法名。
  • 新增控制 Sink 时只接受 ApprovedRobotCommand,并要求安全门成为图上的直接前驱;不要兼容接收原始 RobotCommand

gRPC transport 已提供有界的 BidiRequestStream 请求迭代器,可供未来真正的双向音频 RPC 使用。它限制为单消费者,正常 close 会排空已接受请求并原子唤醒被背压阻塞的 producer若 gRPC 放弃迭代器则丢弃无法再发送的请求。它只解决 request side 的背压与关闭,不包含任何音频 proto、重连、响应解析或扬声器策略这些仍属于具体连接器。

6. 并发和资源模型

当前真实执行关系如下:

flowchart LR
    MAIN["一个 Python 进程"] --> LOOP["一个 asyncio event loop"]
    LOOP --> P1["Pipeline A: 每节点一个 task"]
    LOOP --> P2["Pipeline B: 每节点一个 task"]
    LOOP --> TP["共享有界 ThreadPoolExecutor"]
    P1 -. "插件显式 workers.run_blocking" .-> TP
    P2 -. "插件显式 workers.run_blocking" .-> TP
    P1 -. "插件自行实现" .-> MW["可选常驻进程/Model Worker"]

多个 Pipeline 共用同一 event loop、gRPC ChannelPool、HttpClientPool 和应用线程池。同名 endpoint 会复用客户端;同一 endpoint ID 如果请求了冲突的 settings 会报错。线程池只通过 ComponentContext.metadata["thread_executor"] 提供;execution.mode: thread 在 v1 是非法保留值,不会自动调用线程池。

共享 gRPC channel 不等于共享流式 RPC 或 Source 实例。两个 Pipeline 各自配置同一 camera node 时会建立两个订阅并重复解码;同一物理输入需要供多个模型使用时,应优先 在一个 Pipeline 内从 decoder 的输出端口 fan-out。每个 fan-out 分支仍应使用独立的 小容量队列,避免慢模型把实时图像积压在内存中。

边缘端容量规划建议:

  • 优先减少输入帧率、分辨率和队列容量,不用无界缓存换吞吐;
  • 同一 GPU 上尽量复用一个常驻模型实例,避免每 Pipeline 重复加载;
  • 编码视频在 decoder 前使用 video_contiguous,不能用 latest-frame 丢包;只有 decoder 输出完整图像后,才能在推理前使用 realtime_latest 主动丢旧帧;当前运行时不会根据 deadline 自动丢弃;
  • 音频需要连续性,不能照搬视频的 drop_oldest
  • 平台上报与机器人控制使用不同边,并为非关键遥测选择非阻塞溢出策略;
  • 通过实际端到端延迟、queue high watermark、drop counter 和显存峰值决定是否引入进程/共享内存。

7. 控制安全模型

机器人控制链必须保持以下形态:

AI 结果 -> Policy -> RobotCommand/v1
                    -> Safety Gate
                    -> ApprovedRobotCommand/v1
                    -> 类型化执行器 Sink

这里有四层进程内防线:

  1. 编译器要求 actuator 的每个直接前驱都带 safety_gate 标签;安全门与 Sink 之间不能插入转换/透传节点,也不能增加旁路输入;
  2. 安全门只接收 RobotCommand/v1,检查通过后才输出 ApprovedRobotCommand/v1
  3. AGV Sink 的端口 schema 和运行时类型检查都只接受 ApprovedRobotCommand
  4. AGV Sink 只支持明确列出的 action不会根据配置字符串反射调用任意 gRPC 方法,也不会自动重试运动命令。

ApprovedRobotCommand 是进程内类型边界,不是加密签名。只有可信插件可以被安装并赋予 safety_gate 标签;恶意或被篡改的 Python 插件仍拥有当前进程权限,因此插件供应链、包版本锁定和文件完整性属于安全模型的一部分。

当前安全门已经检查:

  • payload 必须是 RobotCommand
  • Envelope deadline 和 RobotCommand.valid_until_ns
  • action 白名单;
  • 可选 device 白名单;
  • 为 action 配置的每个参数必须存在、是有限数并位于 [min, max]
  • set_velocity 必须为 vx/vy/wz 配置完整范围;require_limits_for 可把同一规则扩展到其他 action
  • 默认按 (device, action) 拒绝 sequence 重放和乱序;
  • 同一 (device, action) 的最小发送间隔。

投入真实机器人前仍必须由机器人型号相关 Policy/Safety 插件补齐:

  • 除已配置参数范围外的速度、位置、加速度和 jerk 联合约束;
  • 当前模式、故障、急停、定位质量、电量及障碍物状态;
  • command ID 去重和多来源优先级仲裁;
  • 急停独立高优先级路径和权限控制;
  • 审计日志、认证、TLS 和密钥管理。

最重要的边界位于进程外:持续速度命令必须由 cmvr-es 服务端发放短期 lease并由独立 deadman 在续租/heartbeat 超时后清零速度。当前 Sink 会在调用前检查 TTL并把客户端 RPC timeout 截断到剩余 TTL但当前 cmvr-es 请求没有承载 valid_until/command_id/sequence,服务端无法据此拒绝网络中迟到的命令。unsafe_allow_unleased_velocity、客户端 TTL、安全门以及 Sink stop() 都不能覆盖进程崩溃、SIGKILL、断电或网络分区。cmvr-es 未实现 server-side lease/deadman 前,set_velocity 只允许在隔离测试环境使用,禁止用于真实机器人持续运动。

运动命令默认不重试。只有外部 API 提供明确幂等契约,并完成 command ID 去重后,才能为特定 action 增加有界重试。

8. 运维流程

推荐发布/启动顺序:

  1. 安装固定版本的核心包和可信插件包;
  2. 如果使用 cmvr-es gRPC 连接器,安装匹配 cmvr-es proto 版本的 cmvr-api wheel或运行 stub 生成脚本;
  3. 在部署 YAML 中配置 endpoint、证书引用、模型路径和推理设备凭据应通过部署平台的 secret 机制注入;
  4. 执行 cmvr-edge-ai pluginscmvr-edge-ai models,保存实际插件、模型及标签清单;
  5. 执行 cmvr-edge-ai validate -c ...
  6. 确认 cmvr-es 设备启用、设备 ID 正确、平台接口可达;
  7. 先运行回放/smoke再连接真实传感器
  8. 控制链先使用空载/限速/人工急停条件验证;
  9. 使用 SIGTERM 停止并给 shutdown_timeout_s 留出排空时间。

被动调用服务还需要安装 server extra在同一环境运行 Detect Client 时安装 http extraserve 会同时启动配置中全部启用的主动/被动 Pipeline 和 HTTP API。

示例命令:

uv run --no-sync cmvr-edge-ai models
uv run --no-sync cmvr-edge-ai validate -c configs/active_detection.yaml --pipeline detection
uv run --no-sync cmvr-edge-ai run -c configs/active_detection.yaml \
  --pipeline detection \
  --log-level INFO \
  --log-format json
uv run --no-sync cmvr-edge-ai validate -c configs/server_detect.yaml
uv run --no-sync cmvr-edge-ai serve -c configs/server_detect.yaml \
  --log-level INFO \
  --log-format text
curl -sS http://127.0.0.1:8081/v1/models
CMVR_DETECT_BASE_URL=http://127.0.0.1:8081 \
  uv run --no-sync python client/detect/example.py \
  /path/to/image.jpg detect.ppe

当前日志可输出文本或单行 JSON。日志中不要写入音频原始数据、图像 base64、认证 metadata 或用户隐私内容;生产插件应只记录 trace ID、schema、耗时、尺寸、丢弃计数和经过脱敏的错误信息。

9. 当前限制与演进顺序

9.1 当前限制

  • cmvr-es 音频双向流 proto 尚未实现;框架只有 AudioChunk 契约,不提供 Talk 占位配置;
  • 检测链路已经提供 PyAV H264/H265 解码、双 PPE YOLO 注册/推理、单次解码 fan-out 和重复触发规则VAD、ASR、LLM、TTS 仍需插件提供;
  • UDP/QUIC transport/connector 尚未实现;
  • execution v1 只支持 async/inline 默认单并发;其他 mode 和非默认并发字段会在编译期拒绝,线程/进程/模型 offload 必须由插件显式实现;
  • 非空 resources、非 normal priority 和显式保留 runtime 字段都会被拒绝; runtime.health_bind 仍不启动独立服务,但 serve 已提供 /health/live/health/ready,详细队列指标尚未导出;
  • 无热更新、overlay、通用节点级 supervisor/restart、持久化 outbox/spool相机和 HTTP 的重连/重试是 connector 内部的局部策略;
  • 无自动 deadline 丢弃、max_age_msput_timeout_ms
  • 无共享内存池;跨进程 Worker 由插件负责;
  • Analog Gauge worker 当前固定为 Linux x86_64/Python 3.8 CPU 环境,首次加载约需 1.2 GiB 级别内存;它是可工作的兼容适配,不是轻量模型;
  • 配置只支持单个 YAML 加环境变量,不支持 include/merge
  • 当前 cmvr-es 相机成功响应通常未填写 header.timestamp,因此连接器的 captured_at_ns 可能为 null;精确采集时延需要后续在帧协议中增加设备采集时钟,不能用客户端接收时间冒充;
  • cmvr-es 当前通过 getLatestEncodedFrame 取得最新编码数据,可能在 edge-ai 收到之前跳过 inter-frame 参考包;本地连续队列和重连无法补回上游丢失内容,必须验证/改造服务端流语义;
  • HTTP Sink 对告警使用 event_id、对普通逐帧结果回退使用 trace_id 作为幂等键,并 提供有限退避和 raise/log_and_drop 终态策略,但没有持久投递;log_and_drop 的 WARNING 即表示该报告已丢失,进程崩溃或断电也可能丢失内存中的输出;
  • 内置 YOLO adapter 不产生 track_id,所以 scope: track 需要外部 tracker 及人员/PPE 关联节点;
  • PPE .pt 只能来自可信制品源商业部署还需核对权重说明、Ultralytics runtime/checkpoint 的 AGPL/商业许可,以及训练数据和权重分发权利;
  • 四份 FP32 ONNX 已生成且配置已切到 @2,但仓库没有有授权的现场正例图片;当前 smoke 只证明制品完整和可执行,不能替代非方形现场图 parity、精度与资源 gate
  • FP32 ONNX 文件约为 PT 的两倍People-Talking 单文件约 273 MB 且仍是 YOLOv8x ONNX 化减少运行依赖,不等于降低网络 FLOPs 或存储体积;
  • AGV set_velocity 默认禁用unsafe override 和正常 shutdown stop 都不具备崩溃安全性,不能替代 cmvr-es server-side lease/deadman、机器人本体限位、急停和功能安全系统。

9.2 推荐演进顺序

  1. 真实流与模型验收:用录制数据和目标边缘设备验证 cmvr-es 编码连续性、PyAV 长时间恢复、PPE 精度/FPS/内存/显存和端到端告警;
  2. ONNX 现场验收与继续轻量化:用有授权的非方形相机图完成 PT(rect=False)/ONNX parity测量目标机 P50/P95/RSS再决定 People-Talking 换小模型或现场校准 INT8
  3. 可靠告警投递:在现有 event ID 和有限重试之上增加有界持久 outbox、确认、恢复发送和容量/保留策略;
  4. 同人违规语义:增加 tracker 与 Worker/PPE 空间关联,验证 ID switch 后再启用 scope: track
  5. 补齐可观测性:导出 health、队列水位/丢弃、重连/重试、节点延迟和模型资源指标;
  6. 控制安全闭环:先在 cmvr-es 实现速度 lease/server-side deadman再补齐设备状态输入、机器人型号限值、优先级与审计然后才启用真实 AGV/机械臂动作;
  7. 接入音频双向流proto 落地后实现麦克风 Source/扬声器 Sink严格处理 chunk 顺序、背压和 discontinuity
  8. 按测量结果优化并发:先使用现有显式 thread offload再按测量结果增加常驻 model/process worker确认复制瓶颈后才加入共享内存
  9. 扩展协议:用相同内部契约实现 QUIC/UDP connector不修改算法插件。

每一步都应先通过 validate、smoke、录制数据 replay 和资源峰值检查,再接入真实设备。