Engine 这一个接口。单 rank、单副本的引擎在你的进程里 inline 运行;换成本机多卡,step()、submit()、close() 一个都不变,变的只是 dispatcher 后面的 executor。
执行层级
ReplicaExecutor 永远代表一个完整的模型副本,不论这个副本是一个 rank 还是很多个。dispatcher 在副本之间选路,从不把副本里的 rank 当成独立的 worker。
inline 路径没有进程 supervisor,也没有 rendezvous 服务,单卡机器人控制器要的正是这份清净。
模型并行
ParallelConfig 描述一个副本怎么切。外层的 pipeline 和 cfg 对所有 domain 生效;dense、attention、MoE 各自描述在同一个 rank 池上的 TP/DP/CP/EP 切法:
dense_tp、attention_cp、moe_ep 和 world(P.all_reduce(x, group="dense_tp"))。在一个 pipeline/CFG 分区里,三个 domain 覆盖的 rank 数必须相同:
Entry.parallel_domains 里声明自己实现了哪些 domain,还可以在 Entry.validate_parallel() 里把规则收得更紧。默认只接受单 rank,所以插件跑不了的拓扑会在 worker 启动前就失败,而不是悄悄跑出几份互不相干的模型。以 Cosmos3 为例,它接受 dense 与 attention 相同的 TP 加 CFG,其余一概拒绝。拓扑只从这些对象读取;RANK、WORLD_SIZE、LOCAL_RANK 描述的是物理进程,不是模型布局。
服务副本
副本是整个模型的复制品,由 dispatcher 在它们之间路由。副本数放在DeploymentConfig 里,不放进 ParallelConfig,副本之间也不会建立 process group:
模式选择
DeploymentConfig() 默认 mode="auto":
单卡路径完全不需要部署配置:
devices 里的条目是启动进程可见顺序里的逻辑索引,可以和 shell 掩码叠加:在 CUDA_VISIBLE_DEVICES=4,5,6,7 下,同样的 devices=(0, 1, 2, 3) 落在物理 GPU 4 到 7 上。插件要在 parallel_domains 里声明 dense 或 attention,Cosmos3 两个都声明了。本机 worker 用 Python 的 spawn 启动,会重新 import 你的主模块,所以引擎要放在 if __name__ == "__main__": 里构造。
跨机器用外部 launcher(torchrun --nnodes ...)配 DeploymentConfig.external(),每个 rank 一个进程。这种模式下引擎不会广播请求:每个 rank 都要以同样的顺序、用同样的请求调用 step(),因为所有 rank 参与同一组集合通信,请求只到一个 rank 会让其他 rank 永远等下去。要么每个 rank 都读同一份输入,要么在一个 rank 上收到后先用 torch.distributed 广播。只有 output rank 返回结果,其余返回 None。跨机器的请求路由属于引擎之上的 gateway。
请求路由
每个请求整个交给一个健康的副本,选在途请求最少的那个,平手时轮转。该副本的每个 rank 都会收到请求,只有 output rank 作答。引擎从不在副本间切分、聚合或合并。如果一个 batch 里的样本彼此独立,你可以自己切:每个副本submit() 一份,把各个 future 的结果搬到同一个设备后按顺序拼起来。
两类失败靠异常类型区分。请求级失败,比如插件拒绝某个 payload,在 inline 和 external executor 上按插件自己的异常抛出,副本依然健康;在 managed worker 里同样的失败会带倒整个 group,因为共享集合通信的 rank 无法在请求中途恢复,所以提交前先校验请求。phyai.EngineUnavailableError 表示后端已经服务不了了:worker 退出、硬超时触发,或者引擎已关闭。重建引擎,或者让上层把它摘出轮询。
Engine.mode 告诉你最终选了哪个后端。Engine.submit() 返回标准的 concurrent.futures.Future;只有 inline 上还在排队的请求可以取消。
Worker 与 tensor 所有权
worker 原样继承父进程的设备可见性,各自用torch.cuda.set_device 绑定自己的索引,和 sglang、vllm 的放置方式一样,所以 cuda:1 在每个进程里都指同一块 GPU。所有 worker 报告就绪,启动才算完成。
结果原样返回。CUDA tensor 以 CUDA-IPC view 的形式跨进程:零拷贝,仍在 worker 的 GPU 上,只在 worker group 存活期间有效。要留住它们,就在关闭引擎前 .cpu() 或 .clone()。请求可以带任意本机设备上的 tensor,CPU tensor 走共享内存。
故障处理是 fail-fast。worker 退出、响应格式不对,或 WorkerSupervisorConfig.execution_timeout_s 到期,都会让整个 group 失败,之后的请求抛 EngineUnavailableError。PhyAI 不重试已经开始执行的请求,也不重启 worker,这归拉起引擎的那一层管,不论它是 systemd、Kubernetes 还是 RL 框架。超时从提交那一刻开始计时,排队的时间也算在内。
关闭
close() 幂等:取消排队中的 inline 工作,用 EngineUnavailableError 结束未完成的 managed future,停掉 worker,释放 process group。把 Engine 当 context manager 用,这些都会自动发生。
