Cog 容器运行时(Container Runtime)深度解析:Rust 父进程与 Python Worker 的双进程预测架构
【免费下载链接】cogContainers for machine learning项目地址: https://gitcode.com/GitHub_Trending/co/cog
Cog 将机器学习模型打包为生产级 OCI 镜像,而镜像启动后真正承担预测服务的就是容器运行时(Container Runtime)。本指南以 architecture/04-container-runtime.md 为核心,深入讲解 Cog 容器内部的双进程架构——Rust 编写的 HTTP 服务端(Axum)与 Python 编写的 Worker 子进程如何分工协作,覆盖进程角色、IPC 协议、健康状态机、Predictor 生命周期、预测请求全流程以及关键设计决策。读完本文,你将掌握 Cog 容器从启动到完成一次预测的完整链路,并能据此排查并发、超时、健康检查与日志问题。
概述:为什么容器内部需要两个进程
当 Cog 容器运行时,它执行的是双进程架构:一个 Rust 父进程(HTTP 服务器 + 编排器)和一个 Python Worker 子进程(预测执行)。这一设计把用户模型代码与 HTTP 服务器隔离,换来稳定性、资源管理和干净的关闭处理。
该运行时用 Rust 实现,HTTP 层使用 Axum,Python 集成使用 PyO3,最终以 Python wheel(coglet)形式分发。从架构全景看,容器运行时是 Model Source(用户写的 Runner 类)、Schema(OpenAPI 接口描述)与 Prediction API(HTTP 契约)三者交汇的落点——用户代码经 Schema 校验后,在这里被真正加载、执行并产出结果。
在crates/coglet/src/lib.rs中,coglet核心模块清晰地把这些职责拆分为service(预测生命周期管理)、orchestrator(Worker 子进程生命周期)、bridge(IPC 协议)、permit(并发控制)、transport(HTTP 传输)、prediction(预测状态机)与webhook(回调投递)等独立模块。
高层架构
整个运行时的数据流可以概括为:HTTP 请求进入 Axum 路由层 →PredictionService协调状态与并发许可 → 通过 Unix Socket 与管道驱动 Python Worker 执行预测:
源码佐证:PredictionService内部使用DashMap<String, PredictionEntry>作为活动预测的唯一事实来源,用RwLock<Health>维护健康状态,并通过watch::channel协调双向关闭(见 crates/coglet/src/service.rs)。
所有权模型:PredictionService 是唯一状态所有者
PredictionService是全部预测状态的唯一所有者,一切预测相关操作都流经它。一次预测的生命周期涉及三个关键对象:
- PredictionEntry(存放在并发
DashMap中)——预测状态的事实来源。持有Prediction状态机(经Arc共享)、一个取消令牌(CancellationToken)以及原始输入。在源码中,PredictionEntry还额外记录了cancel_on_stream_drop标志,用于流式订阅断开时是否自动取消(见 crates/coglet/src/service.rs)。 - PredictionSlot—— RAII 容器,把一次预测与一个并发许可(permit)配对。当 slot 被 drop 时,许可自动归还给
PermitPool。permit模块用 Rust typestate 模式在编译期强制状态迁移合法性:PermitInUse → PermitIdle(归还池)、PermitInUse → PermitPoisoned(丢弃),而PermitPoisoned → PermitIdle不存在任何方法,因此被毒化的槽位不可能被误归还(见 crates/coglet/src/permit/mod.rs)。 - PredictionHandle—— 返回给 HTTP 路由处理器。对同步请求,调用
sync_guard()会创建一个守卫(SyncPredictionGuard),当客户端连接断开时自动取消该预测。其Drop实现调用service.cancel(id),同时触发 Rust 侧的CancellationToken和编排器对 Worker 子进程的取消(见 crates/coglet/src/service.rs)。
Prediction结构体本身就是一个状态机:它的变更方法(set_processing、set_succeeded、append_log等)会以副作用的方式触发 webhook。这让 webhook 投递与状态迁移紧密耦合,而不是散落在各个调用点。PredictionStatus枚举定义了Starting / Processing / Succeeded / Failed / Canceled五种状态,其中后三种为终态(is_terminal())(见 crates/coglet/src/prediction.rs)。
进程角色
tini(PID 1)
- 是什么:极简 init 系统(约 30KB 二进制)。
- 为什么需要:正确地向子进程转发信号、回收僵尸进程。
- 入口:
ENTRYPOINT ["/sbin/tini", "--"]。
这一入口由 Dockerfile 生成器在构建期写入镜像:在 pkg/dockerfile/standard_generator.go 中,构建逻辑会下载指定版本的 tini 到/sbin/tini并设置ENTRYPOINT,确保容器 PID 1 是 tini 而非用户代码进程。
父进程(Rust HTTP 服务器)
- 入口:
CMD ["python", "-m", "cog.server.http"]—— 这个薄薄的 Python 启动器调用coglet.server.serve()。 - 职责:
- 端口 5000 上的 HTTP API(Axum)
- 请求校验
- 输入文件下载(从 URL 拉取)
- Webhook 投递(带重试与 trace 上下文传播)
- 输出文件上传
- 健康状态管理
- Worker 子进程生命周期管理
关于启动器:python/cog/server/http.py是一个 argparse 入口脚本,从PORT环境变量读取端口(默认 5000),从构建期写入的COG_PREDICT_TYPE_STUB(predict 模式)或COG_TRAIN_TYPE_STUB(train 模式)环境变量解析 predictor 引用,随后调用coglet.server.serve(predictor_ref, host, port, ...)(见 python/cog/server/http.py)。
Worker 子进程(Python)
- 启动方式:
python -c "import coglet; coglet.server._run_worker()" - 职责:
- 加载用户的 predictor 模块
- 启动时执行一次
setup() - 执行选定的
run()方法(老模型则执行 legacy 的predict()方法) - 通过基于 ContextVar 的日志路由捕获 stdout/stderr
- 通过 slot socket 向父进程回传事件
在 PyO3 绑定侧,crates/coglet-python/src/lib.rs暴露了serve()与_run_worker()两个模块入口;predictor.rs负责包装 Python predictor 类并检测同步/异步模式,worker_bridge.rs为 Python 实现PredictHandlertrait,log_writer.rs实现了基于 ContextVar 的、按 slot 路由的 stdout/stderr 日志写入。
为什么是双进程?
- 隔离:用户代码崩溃不会拖垮 HTTP 服务器
- 内存:模型加载拥有全新的地址空间
- CUDA:Worker 中可进行干净的 GPU 上下文初始化
- 稳定性:即使 Worker 崩溃,服务器仍继续运行(健康端点依然响应)
- 可观测:父进程独立跟踪 Worker 健康
Predictor 生命周期
Predictor 是一个单例:每个 Worker 进程只创建一个实例,且贯穿整个进程生命周期。
你可以依赖的保证:
setup()恰好执行一次,在任何预测被接受之前。用它来加载权重、初始化 GPU 上下文、预热缓存。如果它抛出异常,Worker 退出且健康状态变为SETUP_FAILED——没有重试。源码中SetupResult结构体记录了started_at、completed_at、status(starting/succeeded/failed)以及捕获的logs,便于在健康检查响应中暴露 setup 过程(见 crates/coglet/src/health.rs)。self状态跨所有run()调用持久存在。在setup()中把模型存到self.model,然后在每次run()中使用,这正是预期模式。没有 teardown 钩子。不存在
teardown()、cleanup()或__del__契约。容器关闭时进程直接退出。如果确实需要清理(比如冲刷日志缓冲区),请使用atexit。run()默认串行执行。COG_MAX_CONCURRENCY=1(默认值)时,run()绝不会被并发调用——每次调用都在下一次开始前完成。COG_MAX_CONCURRENCY > 1时,并发的run()调用共享self。异步 runner 在共享的 asyncio 事件循环上运行多个协程——并非真正的并行,而是在await点交错执行。如果模型把可变状态存在self上、且该状态可能在await边界被访问,需要格外小心。如果模型不适合并发调用,请把并发度保持在 1。Worker 崩溃是终结性的。如果 Worker 进程崩溃(段错误、OOM kill),运行时将使所有进行中的预测失败并停止接受新预测。HTTP 服务器保持存活(健康端点仍响应),但容器必须由外部重启——不存在自动的 Worker 重生。从源码看,Worker 在不可恢复错误时会发送
Fatal { reason }消息,父进程随即毒化所有 slot 并使所有进行中的预测失败(见 crates/coglet/src/bridge/protocol.rs)。
Worker 子进程协议
Rust 服务器与 Python Worker 之间通过两个通道通信。所有消息都是 JSON,一行一条。
控制通道(stdin/stdout)
面向 Worker 整体的生命周期消息。
父进程 → Worker:
| 消息 | 用途 |
|---|---|
Init { predictor_ref, num_slots, is_async, ... } | 引导 Worker——加载 predictor、创建 slots |
Cancel { slot } | 取消某个 slot 上正在运行的预测 |
Healthcheck { id } | 请求执行用户自定义健康检查 |
Shutdown | 优雅关闭 |
在源码中,ControlRequest::Init还携带transport_info(slot socket 传输信息)与is_train标志,且必须是 spawn 之后的第一条消息(见 crates/coglet/src/bridge/protocol.rs)。
Worker → 父进程:
| 消息 | 用途 |
|---|---|
Ready { slots, schema } | Worker 初始化完成,返回 slot ID 列表与 OpenAPI schema |
Log { source, data } | setup 阶段日志行(stdout 或 stderr) |
WorkerLog { target, level, message } | Worker 运行时自身(非用户代码)的结构化日志 |
Idle { slot } | slot 完成一次预测,重新可用 |
Cancelled { slot } | slot 上的预测被取消 |
Failed { slot, error } | slot 上的预测失败 |
Fatal { reason } | 不可恢复错误——Worker 正在关闭 |
DroppedLogs { count, interval_millis } | 因背压丢弃的日志消息数量 |
HealthcheckResult { id, status, error } | 用户自定义健康检查的结果 |
ShuttingDown | Worker 正在关闭 |
值得注意的源码细节:日志行在 Worker 侧会被截断保护——超过 4 MiB 的日志行会被裁剪并追加[**** LOG LINE TRUNCATED AT 4 MiB ****]提示(见 crates/coglet/src/bridge/protocol.rs)。
Slot 通道(每个 slot 一个 Unix socket)
承载单次预测的数据。每个 slot 使用独立的 socket,避免并发预测之间的队头阻塞(head-of-line blocking)。
父进程 → Worker:
| 消息 | 用途 |
|---|---|
Predict { id, input, input_file, output_dir } | 运行一次预测。input是内联 JSON;大载荷(>6MiB)时为null,input_file指向磁盘上的溢出文件 |
从源码看,SlotRequest::Predict还携带context字段(请求体中的dict[str, str]),预测器可通过current_scope().context访问(见 crates/coglet/src/bridge/protocol.rs)。
Worker → 父进程:
| 消息 | 用途 |
|---|---|
Log { source, data } | 来自run()的日志行 |
Output { output } | 产出的输出值(用于生成器/流式输出) |
FileOutput { filename, kind, mime_type } | run()产生的文件——按路径引用,由父进程上传 |
Metric { name, value, mode } | 自定义指标(mode:replace、increment或append) |
Done { id, output, predict_time, is_stream } | 预测成功完成 |
Failed { id, error } | 预测失败 |
Cancelled { id } | 预测被取消 |
关于FileOutputKind:它区分普通文件输出(FileType)与超出内联大小阈值的超大输出(Oversized)。关于流式输出,Worker 侧的OutputChunk消息带有序号index,父进程据此组装完整的输出序列(见 crates/coglet/src/bridge/protocol.rs)。
健康状态机
内部健康状态(Health枚举)与 HTTP 响应(HealthResponse)存在区分:Health枚举包括Unknown / Starting / Ready / Busy / SetupFailed / Defunct(见 crates/coglet/src/health.rs),而HealthResponse额外增加了一个瞬态状态UNHEALTHY——当用户自定义健康检查失败时返回,但不改变内部健康状态(见 crates/coglet/src/health.rs)。两者都按SCREAMING_SNAKE_CASE序列化(如READY、SETUP_FAILED),而SetupStatus则按小写序列化(starting/succeeded/failed),单元测试对此做了快照验证(见 crates/coglet/src/health.rs)。
另一个状态判断细节:HealthSnapshot提供is_ready()(state == Ready)与is_busy()(Ready但available_slots == 0)两个辅助方法,供传输层决定返回 200 还是 409/503(见 crates/coglet/src/service.rs)。
预测流程
同步请求(POST /predictions)
关键行为:SyncPredictionGuard在整个请求期间被持有。如果客户端连接断开,守卫被 drop,预测被自动取消。
异步请求(推荐:respond-async)
关键行为:不持有任何守卫。即使客户端断开,预测也会继续执行。
连接断开(同步模式)
一次预测的生命周期
跟随一次预测从 HTTP 请求到响应的完整路径:
- 请求到达Axum HTTP 层(
POST /predictions)。 - 输入被校验:在 Rust 边缘按 OpenAPI schema 校验——类型检查、必填字段、约束全部在 Python 看到任何东西之前完成。
InputValidator(见 crates/coglet/src/input_validation.rs)负责这一职责,相关 txtar 集成测试(如 integration-tests/tests/input_validation_before_start.txtar)验证了"校验先于启动"的行为。 - 获取 slot 许可:从
PermitPool获取。如果所有 slot 都忙,请求立即得到409 Conflict——没有排队。CreatePredictionError::AtCapacity对应这一失败路径(见 crates/coglet/src/service.rs)。 - 输入发送给 Worker:通过 slot 的 Unix socket,以
SlotRequest::Predict消息发送(内联 JSON,若超过 6 MiB 则溢出到临时文件)。 - URL 输入被下载:Worker 把任何
cog.PathURL 字段拉取到本地临时文件(使用线程池并行下载)。Predictor 收到的是本地文件路径,永远不会收到 URL。 - 调用
run(**kwargs):在单例 predictor 实例上执行。输入以原生 Python 类型到达——字符串、整数、pathlib.Path对象——而不是请求对象或原始 JSON。 - 输出经 slot socket 流回。对生成器,每次
yield立即发送一条Output消息——真正的流式,而非缓冲。对单一返回值,发送一条Output或FileOutput消息。 - 文件输出由父进程上传。
cog.Path返回值被上传到配置的存储(或内联响应时 base64 编码)。这对 Predictor 完全透明。upload_file实现位于 crates/coglet/src/orchestrator.rs,它 PUT 到签名端点并从Location头提取最终 URL。 - 组装响应。
Prediction状态机迁移到succeeded,slot 许可被释放,响应返回给客户端(异步请求则经 webhook 投递)。
出错时:如果run()抛出异常,Worker 发送Failed消息。预测被标记为failed,slot 回到 idle,runner 实例存活——它会正常处理下一个请求。只有进程级崩溃(段错误、OOM kill)才会销毁实例;之后会发生什么见上文 Predictor 生命周期。
调用路径
coglet 在运行 Cog 容器时是如何被调用的:
启动链路在源码中的落点:HTTP 传输层位于 crates/coglet/src/transport/(http/mod.rs、http/routes.rs、http/server.rs),编排器在 crates/coglet/src/orchestrator.rs(其流程注释清晰地描述了五步:spawn Worker → 发送 Init 等待 Ready → 用 slot socket 填充 PermitPool → 事件循环把响应路由给预测 → Worker 崩溃时使所有预测失败并关闭)。PyO3 入口在 crates/coglet-python/src/lib.rs。
关键设计决策
为什么用 Rust?
- 性能:请求处理上 Axum 比 Python HTTP 框架更快
- 稳定性:用户代码失败时服务器不会崩溃
- 资源管理:更好的背压与并发控制
- 内存安全:HTTP 层没有 Python GIL 竞争
为什么用 PyO3?
- ABI3 wheel:单个 wheel 兼容 Python 3.10–3.13
- 原生性能:直接 C API 调用,无序列化开销
- Predictor 代码不变:用户无需改动任何东西
- 即插即用:相同的 HTTP API、相同的行为
为什么用子进程(而非进程内)?
- 隔离:Python 崩溃/段错误不会杀死服务器
- CUDA 上下文:每个 Worker 有干净的 GPU 初始化
- 内存:模型加载拥有全新的地址空间
为什么用 slot(而非异步任务)?
- 可预测:并发预测数量固定
- 公平:许可机制防止饥饿
- 可观测:易于监控 slot 使用情况
- 简单:Worker 子进程内没有异步复杂度
从源码看,slot 机制还有一层防护:被毒化的 slot(例如其 mutex 被破坏、无法保证隔离)会通过SlotOutcome::Poisoned标记为Failed而非Idle,从类型上杜绝"中毒 slot 被误认为空闲"的错误(见 crates/coglet/src/bridge/protocol.rs)。
输入溢出(Input Spilling)
当一次预测的输入超过 6 MiB 时,内联通过 IPC socket 发送就太大了。此时父进程把它写入临时文件,并在input_file中发送文件路径(input置为 null)。Worker 读取文件、反序列化、然后删除文件,继续正常处理。这对 Predictor 代码完全透明。
源码细节:MAX_INLINE_IPC_SIZE = 6 MiB的取值依据是LengthDelimitedCodec默认帧上限为 8 MiB,6 MiB 为帧开销和其他消息字段留出了 2 MiB 安全余量(见 crates/coglet/src/bridge/protocol.rs)。SlotRequest::rehydrate_input负责恢复输入:读文件、立即删除溢出文件(即使在 JSON 损坏时也先删除,避免残留),再反序列化(见 crates/coglet/src/bridge/protocol.rs)。集成测试 integration-tests/tests/coglet_large_input.txtar 与 coglet_large_output.txtar 覆盖了大输入/大输出的端到端行为。
文件输出
当run()产生文件输出(cog.Path)时,Worker 发送带文件名和 MIME 类型的FileOutput消息。父进程负责上传文件(或对内联响应做 base64 编码)。Predict请求中的output_dir字段告诉 Worker 把输出文件写到哪个目录。FileOutputKind区分普通文件输出(FileType)与超出内联大小限制的超大输出(Oversized)。
自定义指标(Custom Metrics)
模型可以在 predict 方法中通过self.record_metric(name, value, mode)记录自定义指标。它们作为 slot 通道上的Metric消息发送。mode控制指标的聚合方式:
replace—— 覆盖任何已有值increment—— 加到当前值上(数值型)append—— 追加到列表
指标出现在预测响应的metrics对象中,与内建的predict_time并列。源码中MetricMode枚举定义了这三种模式(见 crates/coglet/src/bridge/protocol.rs),且支持点分路径键(如"timing.preprocess"),服务器会把它们解析为嵌套对象;crates/coglet/src/bridge/snapshots/下的slot_metric_replace/increment/append快照文件锁定了各模式的序列化格式。
用户自定义健康检查
模型可以实现一个自定义健康检查,与内建的健康状态机并行运行。父进程在控制通道发送Healthcheck { id };Worker 运行用户的健康检查并以HealthcheckResult { id, status, error }响应。
如果健康检查失败,HTTP/health-check端点返回UNHEALTHY——但这是瞬态的,不改变内部Health状态。模型保持READY并继续接受预测。相关集成测试包括 integration-tests/tests/healthcheck.txtar、healthcheck_unhealthy.txtar、healthcheck_async.txtar 与 healthcheck_timeout.txtar,覆盖了健康、不健康、异步与超时等场景。
环境变量
| 变量 | 默认值 | 用途 |
|---|---|---|
PORT | 5000 | HTTP 服务器端口 |
COG_LOG_LEVEL | INFO | 日志详细程度(若设置了RUST_LOG则被忽略) |
COG_MAX_CONCURRENCY | 1 | 并发预测 slot 数量 |
COG_SETUP_TIMEOUT | 无 | setup 超时秒数(0 被忽略) |
COG_THROTTLE_RESPONSE_INTERVAL | 0.5s | Webhook 响应节流间隔 |
LOG_FORMAT | json | 设为console获得人类可读的日志输出 |
日志相关的实现细节:crates/coglet-python/src/lib.rs中init_tracing优先读取RUST_LOG,否则按COG_LOG_LEVEL构造 EnvFilter(debug/warn/warning/error,其余默认info),并依据LOG_FORMAT选择 JSON 或 console 输出格式(见 crates/coglet-python/src/lib.rs)。webhook 节流间隔在 crates/coglet/src/webhook.rs 读取。
代码导航:Where to Look
coglet 核心(crates/coglet/src/):
- service.rs ——
PredictionService,中央协调器。从这里开始读。 - orchestrator.rs —— Worker 子进程的启动与生命周期
- bridge/ —— IPC 协议定义(
protocol.rs)与 Unix socket 传输 - permit/ —— 基于 slot 的并发控制(
PermitPool、PredictionSlot) - transport/http/ —— Axum HTTP 服务器与路由处理器
- prediction.rs —— 预测状态机,状态迁移时触发 webhook
coglet-python(crates/coglet-python/src/):
- lib.rs —— PyO3 模块入口:
serve()与_run_worker() predictor.rs—— 包装 Python predictor 类,处理同步/异步检测worker_bridge.rs—— 为 Python 实现PredictHandlertraitlog_writer.rs—— 基于 ContextVar 的、按预测 slot 路由的 stdout/stderr 写入
Python 启动器:python/cog/server/http.py —— 调用coglet.server.serve()的薄入口。
验证与测试:协议序列化快照位于 crates/coglet/src/bridge/snapshots/;端到端行为由 integration-tests/tests/ 下的 txtar 测试覆盖,其中与本文最相关的是healthcheck*.txtar、cancel_*_prediction.txtar、coglet_large_*、sse_*、setup_timeout_serial.txtar与sequential_state_leak.txtar(验证并发下self状态隔离)。
【免费下载链接】cogContainers for machine learning项目地址: https://gitcode.com/GitHub_Trending/co/cog
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考