# 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. 分层架构 ```mermaid flowchart TB subgraph External["外部系统"] CE["cmvr-es
gRPC / future QUIC"] PLAT["业务平台
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:` 标签的连接器只能引用相同 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_ms` 和 `put_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 内部契约层 协议边界进入运行时后统一变为: ```text 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` 是冻结 dataclass,`attributes` 被复制为只读 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/v1`、`InferenceResponse/v1` 这样的内部 schema 名; JSON 文档中的 `schema_version` 使用上表 `cmvr.*` 名称。两者属于同一 DTO 的内部图契约 和外部 wire 版本,不应互相替换。 ### 2.3 组件层 运行时只有三类组件: | 组件 | 输入/输出 | 核心方法 | |---|---|---| | `Source` | 无输入,一个或多个输出端口 | `messages()` 返回 async iterator | | `Operator` | 一个或多个输入/输出端口 | `process(envelope, input_port)` | | `Sink` | 一个或多个输入端口,无输出 | `consume(envelope, input_port)` | 所有组件共享以下生命周期: ```text setup(context) -> start() -> 执行 -> health() -> stop() ``` 建议在 `__init__` 只解析轻量参数,在 `setup` 建立网络连接或加载模型,在 `stop` 释放资源并保证幂等。`ComponentContext` 提供 `pipeline_id`、`node_id`、全局 shutdown event,以及 endpoint、共享连接池和节点配置等 metadata。 Operator 可以返回: - `None`:过滤该消息; - `Envelope`:从默认 `output` 端口发送; - `Emission(port, envelope)`:从指定端口发送; - 上述对象的同步或异步 iterable:一进多出。 Source 也可以直接 yield `Envelope` 或 `Emission`。返回原始 bytes、模型对象或其他未包装值会在产生该值的节点附近报错。 ### 2.4 DAG 运行时 当前每个 Pipeline 在一个 asyncio event loop 中运行,每个节点一个 task。`execution.mode: async` 和 `inline` 在 v1 中都由该 task await 组件方法,Operator 和 Sink 每次只处理一个消息。v1 不会复制节点 task,也不提供配置驱动的并行调度。 每条有向边拥有独立队列。一个输出端口 fan-out 到多条边时,同一个只读 Envelope 会送入每个队列;队列之间的容量和丢弃计数独立。注意:路由会等待所有边的 `put()`,因此任意 fan-out 分支使用 `overflow: block` 且消费变慢时,仍会对共同生产者产生背压。希望“平台慢但控制链不停”时,平台分支必须选择合适的丢弃策略,或在后续加入独立 spool。 多个入边到同一节点时,运行时处理最先就绪的边,不保证不同边之间的全局顺序。单条边内部保持 FIFO(`drop_*` 导致的丢弃除外)。 节点自然结束后,其全部出边会关闭;消费者会读完已接受的数据后结束。收到 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,因此两种服务模式可以在同一 进程并存: ```mermaid 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."] --> Broker ``` `cmvr_edge_ai.capabilities` 统一登记 capability 的稳定 category、模型实现和支持的服务 模式;`GET /v1/models` 返回的 `ModelCatalog` 分成 `active_push` 与 `passive_invoke` 两组,并把“代码已注册”和“部署已就绪”作为不同状态。一个同时支持 两种模式的 capability 可以出现在两组中。 公开路由边界是 namespaced category,例如 `detect.ppe`、`detect.mobile_phone`,而不是 模型或 Pipeline 标识。客户端请求不能指定 `model_id`、Pipeline ID、权重路径或推理 设备;这些都由 `server.routes.` 和服务端 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_workers` 和 `shutdown_timeout_s`。边缘端应保持较小的 `thread_workers`;第三方推理库本身还可能创建线程,应同时限制 OpenMP/MKL/模型运行时线程数,避免过度订阅 CPU。 ### 3.3 `endpoints.` | 字段 | 默认值 | 说明 | |---|---:|---| | `transport` | 无 | 必填,转为小写;当前连接器支持 `grpc`、`http` | | `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 是否使用 TLS;HTTP endpoint 禁止设置该字段 | | `timeout_s` | `5.0` | 连接器请求超时;流式相机订阅不使用该 unary timeout | | `metadata` | `{}` | gRPC metadata 预留;当前 HTTP 连接器把它作为默认 header | | `options` | `{}` | 传输私有参数 | 当前 gRPC options 支持: ```yaml 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_url` 的 `https://` scheme 决定,证书验证默认开启;只有隔离测试环境才可显式设置 `options.verify_tls: false`。当前 gRPC `metadata` 尚未自动附加到 RPC,若需要鉴权应在连接器中显式实现。 ### 3.4 `pipelines.` | 字段 | 默认值 | 说明 | |---|---:|---| | `enabled` | `true` | 未指定 `--pipeline` 时是否启动 | | `priority` | `normal` | v1 必须为 `normal`;`low/high` 或数值优先级会因无调度器而被拒绝 | | `nodes` | 无 | 至少一个节点 | | `edges` | `[]` | 有向连接列表 | 仓库的主动推送配置为 `configs/active_detection.yaml`,当前只定义 `detection`。 CLI 显式传 `--pipeline detection` 时只编译并运行该链路;省略参数时会启动配置中所有 `enabled: true` 的 Pipeline。不能提供真实能力的 Talk 占位图已移除。 不要在活跃 Pipeline 中保留 `enabled: false` 节点;当前编译器会直接拒绝。要暂时关闭逻辑,请禁用整个 Pipeline 或从图和配置中移除该节点。 ### 3.5 `nodes.` | 字段 | 默认值 | 说明 | |---|---:|---| | `uses` | 无 | 版本化插件 ID,例如 `cmvr.grpc.camera_rgb_stream@1` | | `with` | `{}` | 原样传给插件工厂的参数;Python 模型内字段名为 `params` | | `execution` | 见下表 | 调度意图声明 | | `resources` | `{}` | 预留资源需求字段;v1 必须保持为空/默认值 | | `enabled` | `true` | 当前只允许活跃 Pipeline 中为 `true` | `execution` 字段: | 字段 | 默认值 | 允许值/限制 | |---|---:|---| | `mode` | `async` | v1 支持 `async`、`inline`;`thread/process/model_worker` 为保留值并被编译器拒绝 | | `concurrency` | `1` | v1 必须为 `1`,其他值被编译器拒绝 | | `max_in_flight` | `null` | 预留;非空值被编译器拒绝 | | `timeout_s` | `null` | 预留;非空值被编译器拒绝 | | `ordered` | `true` | v1 必须为 `true`,`false` 被编译器拒绝 | `async` 与 `inline` 当前没有调度差异,都是 event loop 中的单 task、单 in-flight 调用。配置模型保留其他枚举值是为了后续版本演进,不代表 v1 能执行它们;`validate` 会 fail closed,而不是静默退化到 async。 `resources` 模型预留了 `cpu_cores`、`memory_mb`、`gpu_device`、`gpu_memory_mb`、`exclusive_gpu`,但 v1 无法强制执行资源隔离,因此任何非空/非默认声明都会被编译器拒绝。不要依赖“只写说明但不生效”的资源配置;当前应在部署层和插件实现中限制 CPU/GPU。 ### 3.6 `edges[]` 与 QoS ```yaml - 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_latest`、`video_contiguous`、`audio_contiguous`、`request`、`telemetry` | | `capacity` | `1` | 队列最大消息数 | | `overflow` | `block` | `drop_oldest`、`drop_newest`、`block`、`reject`;`error` 是 `reject` 的兼容别名 | | `max_age_ms` | `null` | 已预留;设置后编译失败 | | `put_timeout_ms` | `null` | 已预留;设置后编译失败 | 溢出行为: | 策略 | 队列满时行为 | 适合场景 | |---|---|---| | `drop_oldest` | 丢掉最旧消息,再接受新消息 | 实时视频,优先处理最新画面 | | `drop_newest` | 保留已有消息,拒收这次新消息但不抛异常 | 需要保留已排队批次的低优先级数据 | | `block` | 生产者等待空位 | 必须连续的视频编码包、音频或请求;会传播背压 | | `reject` / `error` | 抛出队列满异常,Pipeline 失败 | 不能静默丢失且希望监督器介入的控制链 | v1 强制执行的 profile/overflow 矩阵: | profile | 允许的 overflow | 典型用途与注意事项 | |---|---|---| | `realtime_latest` | `drop_oldest`、`drop_newest` | 已解码的完整视频帧;常用容量 `1..2`,优先限制推理延迟 | | `video_contiguous` | `block`、`reject`、`error` | H264/H265 等 inter-frame 编码流在解码前不能静默丢包;此保证不覆盖上游在发送前已经丢失的数据 | | `audio_contiguous` | `block`、`reject`、`error` | 连续音频不能静默丢弃;阻塞预算必须明确,拒绝时应重建连续性 | | `request` | `block`、`reject`、`error` | 请求/控制;使用有限容量和消息 TTL,禁止执行积压旧动作 | | `telemetry` | 所有五种策略 | 平台上报;当前无磁盘 spool,按数据重要性选择阻塞、丢弃或失败 | 队列记录入队、出队、水位、丢弃、拒绝和强制关闭丢弃计数,当前可通过 Python `PipelineRuntime.edge_stats()` 获取,尚未暴露为服务指标。 ### 3.7 `server` `server` 是可选的进程级被动调用配置;只有 `cmvr-edge-ai serve` 会监听 HTTP。最小 结构如下: ```yaml 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 `;非 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 KiB~1 GiB | | `max_image_bytes` | `10485760` | Base64 解码后的单张输入图像上限,且不能大于请求体上限 | | `access_log` | `false` | 是否启用 Uvicorn access log | `server.routes.` 字段: | 字段 | 默认值 | 说明 | |---|---:|---| | `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`、`::1` 和 `localhost` 可在本机调试时不配置认证。任何其他 bind(包括 `0.0.0.0`、`::` 和普通 hostname)都必须设置 `bearer_token`,并默认要求同时提供 TLS 证书和私钥。只有部署在受信、隔离网络且明确接受 token 明文传输风险时,才能设置 `allow_insecure_remote: true` 代替 TLS。认证失败返回 `401 cmvr.inference-error/v1` 和 `WWW-Authenticate: Bearer`,健康探针保持无需认证。 安全的局域网监听示例(token 从进程环境读取,不写入仓库): ```yaml 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](../configs/server_detect.yaml)。它为 `detect.ppe` 与 `detect.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、白名单、参数限值、序号与限频,并建立批准边界 | 安全门参数: ```yaml 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` 可以列出还必须配置参数范围的 action;`set_velocity` 无条件要求 `argument_limits.set_velocity` 同时定义 `vx/vy/wz`。范围边界必须是有限数且 `min <= max`。过期、过快、重放、乱序、未授权、缺少参数或越界命令都会 fail-closed 丢弃并写审计 warning,不会因一条业务拒绝停止整条控制 Pipeline。通过检查后,安全门把原始 `RobotCommand` 包装为带 `approved_by/approved_at_ns` 的 `ApprovedRobotCommand`。 ### 4.2 检测模型、解码与触发规则 检测流水线内置三个模型无关插件: | 插件 ID | 输入 -> 输出 | 语义 | |---|---|---| | `media.video_decoder.pyav@1` | `ImageFrame/v1 -> ImageFrame/v1` | 用可选 PyAV 保持 H264/H265 decoder context,输出 packed `BGR8`;stream/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.sh` 按 `uv.lock` 安装 `grpc + http + video + image + server + onnx-cpu`;`--profile dev` 改用锁定的 `onnx-export-cpu`,额外提供 Torch/Ultralytics/ONNX 构建工具。旧 `@1` 回滚路径仍可 显式选择 `yolo` 或 `yolo-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//vN/` 组织,每个版本目录同时 保存权重和独立 model card;Mobile Phone 只有 v1 PT 保留历史扁平布局: - [Construction PPE YOLOv8 v1](../models/detection/construction-ppe-yolov8/v1/README.md); - [PPE YOLOv8n 6 Classes v1](../models/detection/ppe-6classes-yolov8n/v1/README.md); - [People Talking YOLOv8x v1](../models/detection/people-talking-yolov8x/v1/README.md); - [YOLOv8n Mobile Phone v1](../models/detection/yolov8n-mobile-phone/README.md)。 - [Construction PPE YOLOv8 ONNX v2](../models/detection/construction-ppe-yolov8/v2/README.md); - [PPE YOLOv8n 6 Classes ONNX v2](../models/detection/ppe-6classes-yolov8n/v2/README.md); - [People Talking YOLOv8x ONNX v2](../models/detection/people-talking-yolov8x/v2/README.md); - [YOLOv8n Mobile Phone ONNX v2](../models/detection/yolov8n-mobile-phone/v2/README.md)。 `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 的 `task`、`names`、`imgsz` 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` 参数: ```yaml 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` 属于 backend;ONNX 支持 `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_id` 为 `None`。 示例 detector 到 repeat gate 使用小容量 `drop_oldest` 队列以控制延迟和原始帧内存; 持续过载时丢弃的推理结果不会参与 `min_hits`。需要每个已完成推理都参与计数时可改 为 `block`,代价是背压和延迟向上游传播。 如果恢复六类模型,可把同一个 decoder 输出 fan-out 到两个 detector,而不是创建两个 独立 Pipeline: ```text 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 模型都追赶相机帧率。 六类模型只包含正向装备类,没有 `Person` 或 `No-*` 标签;它的 `DetectionResult/v1` 表示当前推理帧中“检测到了哪些装备”,不能单独推断某个人 缺少装备。该分支保持 `attach_frame: false`,绕过 repeat gate,把每个完成推理的 结果直接发送到 `/v1/ppe-detections`,所以不会生成 `DetectionAlert` 的 `rule_id`、 `hit_count`、event ID 或告警图字段;`cooldown_ms` 只是 repeat gate 的内部规则配置, 不属于告警 payload。输出边使用容量 16 的 `telemetry + drop_oldest`,模拟 平台持续变慢时会优先保留较新的结果;要求逐条可靠送达时应改用可接受背压的策略或 增加持久 outbox。 `detection.repeat_gate@1` 的 `rules` 每项包含 `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/auto`;`captured` 缺失会失败,`auto` 优先 captured 后回退 received。 告警图片是 repeat gate 的全局配置,不需要在每条 rule 中重复: ```yaml 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: true`。`jpeg_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 bytes;HTTP wire 中 `payload.image` 被序列化为扁平对象: ```json { "media_type": "image/jpeg", "width": 1280, "height": 720, "encoding": "base64", "data": "/9j/4AAQSk..." } ``` 关闭图片,或因第三方结果未附带帧、坏帧等原因渲染失败时,告警仍会发送且 `payload.image` 为 `null`;证据图失败不能阻断结构化告警。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 stop;WARNING/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_stop`、`clear_fault`、`pause_navigation`、`resume_navigation`、`cancel_navigation`、`stop_velocity_control`、`stop_mapping` 和 `set_velocity`; - `set_velocity` 默认禁用;只有显式设置 `unsafe_allow_unleased_velocity: true`,并在 Sink 再配置完整 `velocity_limits.vx/vy/wz` 后才会发出; - `set_velocity` 的 `vx/vy/wz` 必须是有限数,并同时通过安全门和 Sink 两层范围校验; - 同时检查 gRPC 调用异常与响应 `header.success`,并把 RPC timeout 限制在命令剩余 TTL 内; - 不自动重试运动命令;正常关闭时若曾成功发送速度,会尽力调用 `stopVelocityControl`。 下面只展示 actuator 节点的危险参数;该节点仍必须直接连接前述安全门。此开关只应用于架空轮、受控台架等隔离环境: ```yaml 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-safety:edge-ai 被 `SIGKILL`、崩溃、断电或与 cmvr-es 网络分区时,Python `stop()` 没有机会可靠执行。真实机器人必须由 cmvr-es 服务端持有速度 lease,并在 lease/heartbeat 超时后独立执行清零和停车。 ### 4.4 平台 HTTP 连接器 `platform.http_json_sink@1` 接收任意 schema: - 参数:`endpoint`、`path`(默认 `/`)、可选 `grpc_endpoint`、`failure_mode`、`max_attempts`(默认 3)、`retry_initial_s`、`retry_max_s` 和 `retry_statuses`; - POST Envelope 元数据、输入端口、attributes 和 payload; - 配置 `grpc_endpoint` 时,从被引用 gRPC endpoint 的 `target` 提取字面 IP,并在 Envelope JSON 顶层增加 `grpc_ip`;`source_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_es` 且 `endpoints.cmvr_es.target` 为 `192.168.0.119:50052` 时,请求外层包含: ```json { "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 只失败对应 caller,Pipeline 继续处理后续请求。客户端超时或取消后,已经被 source 取出的任务仍占用 该 route 的 broker capacity,直到模型返回 late result/failure 或服务关闭;这样不会把 仍在执行的旧推理伪装成空闲资源并继续接收新任务。尚未进入 DAG 的排队请求则可立即移除。 ## 5. 插件开发约定 ### 5.1 创建组件 插件工厂固定接收 `node_id: str` 与 `params: Mapping[str, Any]`,应在构造阶段检查无需 I/O 的必填参数。下面展示一个最小 Operator: ```python 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 ```python 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 不能声明 input,Sink 不能声明 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`: ```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: ```python 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 只暴露内部契约: ```text 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. 并发和资源模型 当前真实执行关系如下: ```mermaid 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. 控制安全模型 机器人控制链必须保持以下形态: ```text 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 plugins` 和 `cmvr-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` extra;`serve` 会同时启动配置中全部启用的主动/被动 Pipeline 和 HTTP API。 示例命令: ```bash 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_ms`、`put_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 和资源峰值检查,再接入真实设备。