Cog 容器运行时(Container Runtime)深度解析:Rust 父进程与 Python Worker 的双进程预测架构

Cog 容器运行时(Container Runtime)深度解析:Rust 父进程与 Python Worker 的双进程预测架构 Cog 容器运行时Container Runtime深度解析Rust 父进程与 Python Worker 的双进程预测架构【免费下载链接】cogContainers for machine learning项目地址: https://gitcode.com/GitHub_Trending/co/cogCog 将机器学习模型打包为生产级 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 层使用 AxumPython 集成使用 PyO3最终以 Python wheelcoglet形式分发。从架构全景看容器运行时是 Model Source用户写的 Runner 类、SchemaOpenAPI 接口描述与 Prediction APIHTTP 契约三者交汇的落点——用户代码经 Schema 校验后在这里被真正加载、执行并产出结果。在crates/coglet/src/lib.rs中coglet核心模块清晰地把这些职责拆分为service预测生命周期管理、orchestratorWorker 子进程生命周期、bridgeIPC 协议、permit并发控制、transportHTTP 传输、prediction预测状态机与webhook回调投递等独立模块。高层架构整个运行时的数据流可以概括为HTTP 请求进入 Axum 路由层 →PredictionService协调状态与并发许可 → 通过 Unix Socket 与管道驱动 Python Worker 执行预测源码佐证PredictionService内部使用DashMapString, PredictionEntry作为活动预测的唯一事实来源用RwLockHealth维护健康状态并通过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。进程角色tiniPID 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 APIAxum请求校验输入文件下载从 URL 拉取Webhook 投递带重试与 trace 上下文传播输出文件上传健康状态管理Worker 子进程生命周期管理关于启动器python/cog/server/http.py是一个 argparse 入口脚本从PORT环境变量读取端口默认 5000从构建期写入的COG_PREDICT_TYPE_STUBpredict 模式或COG_TRAIN_TYPE_STUBtrain 模式环境变量解析 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 实现PredictHandlertraitlog_writer.rs实现了基于 ContextVar 的、按 slot 路由的 stdout/stderr 日志写入。为什么是双进程隔离用户代码崩溃不会拖垮 HTTP 服务器内存模型加载拥有全新的地址空间CUDAWorker 中可进行干净的 GPU 上下文初始化稳定性即使 Worker 崩溃服务器仍继续运行健康端点依然响应可观测父进程独立跟踪 Worker 健康Predictor 生命周期Predictor 是一个单例每个 Worker 进程只创建一个实例且贯穿整个进程生命周期。你可以依赖的保证setup()恰好执行一次在任何预测被接受之前。用它来加载权重、初始化 GPU 上下文、预热缓存。如果它抛出异常Worker 退出且健康状态变为SETUP_FAILED——没有重试。源码中SetupResult结构体记录了started_at、completed_at、statusstarting/succeeded/failed以及捕获的logs便于在健康检查响应中暴露 setup 过程见 crates/coglet/src/health.rs。self状态跨所有run()调用持久存在。在setup()中把模型存到self.model然后在每次run()中使用这正是预期模式。没有 teardown 钩子。不存在teardown()、cleanup()或__del__契约。容器关闭时进程直接退出。如果确实需要清理比如冲刷日志缓冲区请使用atexit。run()默认串行执行。COG_MAX_CONCURRENCY1默认值时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、创建 slotsCancel { slot }取消某个 slot 上正在运行的预测Healthcheck { id }请求执行用户自定义健康检查Shutdown优雅关闭在源码中ControlRequest::Init还携带transport_infoslot socket 传输信息与is_train标志且必须是 spawn 之后的第一条消息见 crates/coglet/src/bridge/protocol.rs。Worker → 父进程消息用途Ready { slots, schema }Worker 初始化完成返回 slot ID 列表与 OpenAPI schemaLog { source, data }setup 阶段日志行stdout 或 stderrWorkerLog { 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 }用户自定义健康检查的结果ShuttingDownWorker 正在关闭值得注意的源码细节日志行在 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时为nullinput_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 }自定义指标modereplace、increment或appendDone { 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状态机迁移到succeededslot 许可被释放响应返回给客户端异步请求则经 webhook 投递。出错时如果run()抛出异常Worker 发送Failed消息。预测被标记为failedslot 回到 idlerunner 实例存活——它会正常处理下一个请求。只有进程级崩溃段错误、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 竞争为什么用 PyO3ABI3 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 MiB6 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覆盖了健康、不健康、异步与超时等场景。环境变量变量默认值用途PORT5000HTTP 服务器端口COG_LOG_LEVELINFO日志详细程度若设置了RUST_LOG则被忽略COG_MAX_CONCURRENCY1并发预测 slot 数量COG_SETUP_TIMEOUT无setup 超时秒数0 被忽略COG_THROTTLE_RESPONSE_INTERVAL0.5sWebhook 响应节流间隔LOG_FORMATjson设为console获得人类可读的日志输出日志相关的实现细节crates/coglet-python/src/lib.rs中init_tracing优先读取RUST_LOG否则按COG_LOG_LEVEL构造 EnvFilterdebug/warn/warning/error其余默认info并依据LOG_FORMAT选择 JSON 或 console 输出格式见 crates/coglet-python/src/lib.rs。webhook 节流间隔在 crates/coglet/src/webhook.rs 读取。代码导航Where to Lookcoglet 核心crates/coglet/src/service.rs ——PredictionService中央协调器。从这里开始读。orchestrator.rs —— Worker 子进程的启动与生命周期bridge/ —— IPC 协议定义protocol.rs与 Unix socket 传输permit/ —— 基于 slot 的并发控制PermitPool、PredictionSlottransport/http/ —— Axum HTTP 服务器与路由处理器prediction.rs —— 预测状态机状态迁移时触发 webhookcoglet-pythoncrates/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),仅供参考