64 KiB
cmvr-edge-ai 架构与开发指南
本文描述当前代码已经实现的架构、配置语义和扩展边界。文中的“当前”指运行时 v1;规划能力会明确标注,避免把配置字段误认为已经具备的调度能力。
1. 目标与边界
cmvr-edge-ai 负责在单台机器人边缘计算设备上,把三类模块可靠地连接起来:
- 上游数据:cmvr-es 的相机、未来的麦克风音频流,以及平台下发的数据;
- 中间 AI:视频解码、检测、VAD、ASR、LLM、TTS、决策和策略;
- 下游输出: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 字段、非
normalpriority 和非空 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 内部契约层
协议边界进入运行时后统一变为:
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) |
所有组件共享以下生命周期:
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 时:
- 设置全局 shutdown event;
- 取消 Source task,并先调用 Source 的
stop(); - 让 Operator/Sink 在
shutdown_timeout_s内排空队列; - 超时后取消剩余 task、丢弃未处理消息;
- 按 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_push 与
passive_invoke 两组,并把“代码已注册”和“部署已就绪”作为不同状态。一个同时支持
两种模式的 capability 可以出现在两组中。
公开路由边界是 namespaced category,例如 detect.ppe、detect.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_workers 和 shutdown_timeout_s。边缘端应保持较小的 thread_workers;第三方推理库本身还可能创建线程,应同时限制 OpenMP/MKL/模型运行时线程数,避免过度订阅 CPU。
3.3 endpoints.<id>
| 字段 | 默认值 | 说明 |
|---|---|---|
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 支持:
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.<id>
| 字段 | 默认值 | 说明 |
|---|---|---|
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.<id>
| 字段 | 默认值 | 说明 |
|---|---|---|
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
- 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。最小
结构如下:
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 KiB~1 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、::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 从进程环境读取,不写入仓库):
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.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、白名单、参数限值、序号与限频,并建立批准边界 |
安全门参数:
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/<model-name>/vN/ 组织,每个版本目录同时
保存权重和独立 model card;Mobile Phone 只有 v1 PT 保留历史扁平布局:
- Construction PPE YOLOv8 v1;
- PPE YOLOv8n 6 Classes v1;
- People Talking YOLOv8x v1;
- YOLOv8n Mobile Phone v1。
- Construction PPE YOLOv8 ONNX v2;
- PPE YOLOv8n 6 Classes ONNX v2;
- People Talking YOLOv8x ONNX v2;
- YOLOv8n Mobile Phone ONNX v2。
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 参数:
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/RGB8packed buffer,并把RGB8转为 backend 使用的 BGR 顺序;它不运行 tracker,返回的track_id为None。
示例 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 模型都追赶相机帧率。
六类模型只包含正向装备类,没有 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 中重复:
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 被序列化为扁平对象:
{
"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 节点的危险参数;该节点仍必须直接连接前述安全门。此开关只应用于架空轮、受控台架等隔离环境:
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 时,请求外层包含:
{
"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:
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 不能声明 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:
[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
这里有四层进程内防线:
- 编译器要求
actuator的每个直接前驱都带safety_gate标签;安全门与 Sink 之间不能插入转换/透传节点,也不能增加旁路输入; - 安全门只接收
RobotCommand/v1,检查通过后才输出ApprovedRobotCommand/v1; - AGV Sink 的端口 schema 和运行时类型检查都只接受
ApprovedRobotCommand; - 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. 运维流程
推荐发布/启动顺序:
- 安装固定版本的核心包和可信插件包;
- 如果使用 cmvr-es gRPC 连接器,安装匹配 cmvr-es proto 版本的
cmvr-apiwheel,或运行 stub 生成脚本; - 在部署 YAML 中配置 endpoint、证书引用、模型路径和推理设备;凭据应通过部署平台的 secret 机制注入;
- 执行
cmvr-edge-ai plugins和cmvr-edge-ai models,保存实际插件、模型及标签清单; - 执行
cmvr-edge-ai validate -c ...; - 确认 cmvr-es 设备启用、设备 ID 正确、平台接口可达;
- 先运行回放/smoke,再连接真实传感器;
- 控制链先使用空载/限速/人工急停条件验证;
- 使用 SIGTERM 停止并给
shutdown_timeout_s留出排空时间。
被动调用服务还需要安装 server extra,在同一环境运行 Detect Client 时安装 http
extra;serve 会同时启动配置中全部启用的主动/被动 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、非
normalpriority 和显式保留 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 推荐演进顺序
- 真实流与模型验收:用录制数据和目标边缘设备验证 cmvr-es 编码连续性、PyAV 长时间恢复、PPE 精度/FPS/内存/显存和端到端告警;
- ONNX 现场验收与继续轻量化:用有授权的非方形相机图完成 PT(
rect=False)/ONNX parity,测量目标机 P50/P95/RSS;再决定 People-Talking 换小模型或现场校准 INT8; - 可靠告警投递:在现有 event ID 和有限重试之上增加有界持久 outbox、确认、恢复发送和容量/保留策略;
- 同人违规语义:增加 tracker 与 Worker/PPE 空间关联,验证 ID switch 后再启用
scope: track; - 补齐可观测性:导出 health、队列水位/丢弃、重连/重试、节点延迟和模型资源指标;
- 控制安全闭环:先在 cmvr-es 实现速度 lease/server-side deadman,再补齐设备状态输入、机器人型号限值、优先级与审计,然后才启用真实 AGV/机械臂动作;
- 接入音频双向流:proto 落地后实现麦克风 Source/扬声器 Sink,严格处理 chunk 顺序、背压和 discontinuity;
- 按测量结果优化并发:先使用现有显式 thread offload,再按测量结果增加常驻 model/process worker;确认复制瓶颈后才加入共享内存;
- 扩展协议:用相同内部契约实现 QUIC/UDP connector,不修改算法插件。
每一步都应先通过 validate、smoke、录制数据 replay 和资源峰值检查,再接入真实设备。