数据截至 (上游 commit 352f1bd7c1a0)
03 · Rollout Controller:把排队任务变成真实执行
本章讲什么: 第 01 章的账本里有 QUEUING 的任务、第 02 章的代理在等 agent 来调——中间缺一个「让 agent 真的跑起来」的角色。Controller 就是它,而且用了 K8s 世界经典的 reconciler(调谐器)模式:反复对比「期望状态」(账本)和「实际状态」(进程/Job),把差异抹平。
1. 通用骨架:两个实现,一个模式
入口 agl-controller(agentlightning/controller/__main__.py:46-52,Hydra 配置 agentlightning/config/controller.yaml)按 runner_type 选实现:
| 模式 | 执行单元 | 实现 | 适合 |
|---|---|---|---|
local | 短命 Python 子进程 | agentlightning/controller/local_reconciler.py | 单机开发调试 |
k8s | Kubernetes Job | agentlightning/controller/k8s_reconciler.py | 集群规模化、依赖隔离 |
两个实现共享同一个循环形状:
# 示意,两个 reconciler 的共同骨架
while 没收到停止信号:
await reconcile_once() # 拉取 queuing/running 的 rollout,对齐现实
await 等待 poll_interval 或停止信号
(依据:local_reconciler.py:116-126 的 _reconcile_loop 与 k8s_reconciler.py:149-161 的 _periodic_reconcile_loop 结构一致。)
核心纪律:Controller 自己不存任何状态(local 模式只有一个进程表),状态真相永远在 Gateway——每次决策都是「查账 → 对齐」。这换来很强的容错:Controller 崩了重启,下一轮 reconcile 照常收敛。
2. local 模式:子进程流水线
2.1 领任务与对齐(_reconcile_once)
每轮拉一批 state_in=queuing,running 的 rollout(local_reconciler.py:128-137),然后逐个对齐(:140-156 ):
| 账本状态 | 现实(进程表) | 动作 |
|---|---|---|
| QUEUING | 无进程 | 池子没满就 spawn(§2.2),满了先放着 |
| RUNNING | 无进程 | 孤儿回收:PATCH 成 FAILED(“local subprocess is not running”) |
| QUEUING | 进程活着 | PATCH 成 RUNNING(补记启动事实) |
| 任意 | 进程已退出 | 按退出码定成败:0→SUCCEEDED,非 0→FAILED(_finish_proc,:168-179) |
2.2 spawn:环境变量的「脐带」
_spawn_for(local_reconciler.py:197-245)是执行侧的对接核心。它构造的环境变量就是 agent 与框架的全部接口:
AGL_OPENAI_BASE_URL:{网关}/proxy/rollout/{id}/attempt/0/mode/{train|val}/openai/v1——第 02 章的归账 URL。AGL_EVENT_URL:{网关}/api/rollouts/{id}/attempt/0/events——agent 结束时报reward事件的地址。AGL_KEY:网关鉴权 key。
再加用户自定义的 env_map:把 rollout input 里的字段按路径(如 input.question)抠出来塞进环境变量(_build_env_from_map + _resolve_input_path,local_reconciler.py:58-86)。Calc-X 例子就靠它把题目和答案分别注入 QUESTION/RESULT(examples/calc_x/train_calc_agent.py:90-93)。
进程本身用 python -c "…_run_local_reconciler_worker(sys.argv[1])" 启动(local_reconciler.py:218-232),worker 函数按 模块:类 或 模块.类 路径 import 你的 Agent 类并调 run()(_run_local_reconciler_worker,:30-45)。两个工程细节:
start_new_session=True:子进程自成进程组——后面超时才能整组击杀。- spawn 成功后立即 PATCH 成 RUNNING 并记
last_attempt_id(:243)。
2.3 超时与关停
每轮 reconcile 末尾检查:进程存活时间超过 config.timeout_seconds(默认 3600,agentlightning/schemas.py:121)就 os.killpg(SIGKILL) 整组击杀、标 FAILED “timed out”(local_reconciler.py:158-166;击杀实现 _kill_process_group,:181-195,SIGKILL 后还等 5 秒确认退出)。收到 SIGTERM/SIGINT(controller/__main__.py:40-42 注册信号)则先做最后一轮 reconcile,再把所有活进程击杀并标 FAILED(_shutdown,:247-256)——关停也走账本,不会留下永远 RUNNING 的僵尸记录。
3. k8s 模式:Job 模板 + 双循环
3.1 Job 是怎么造出来的:build_job_spec
k8s 模式下每个 rollout 的执行配置里带一份完整的 Jinja2 Job 模板(RolloutK8sConfig.job_template,agentlightning/schemas.py:112-115)。build_job_spec(agentlightning/controller/k8s_reconciler.py:42-103)渲染并强制补齐安全字段:
| 步骤 | 内容 |
|---|---|
| 渲染 | job_name 和 input 作为模板变量;自定义 filter yaml_escape 用 json.dumps 安全转义任意值(:48-53) |
| 校验 | 必须渲染出恰好一个 kind: Job 的 YAML 文档,否则报错(:54-60) |
| 覆盖元数据 | 名字固定为 agl-rollout-{rollout_id}、打上 app.kubernetes.io/managed-by=agentlightning 等 label(:62-68)——这是后面 list/watch 的选择器 |
| 强制策略 | backoffLimit: 0、restartPolicy: Never(Job 失败就是失败,不许 K8s 自动重试,RL 数据要的是如实记录)、activeDeadlineSeconds 取超时(:70-76) |
| 注入脐带 | 遍历所有 container,写入与 local 模式同名同值的 AGL_OPENAI_BASE_URL/AGL_EVENT_URL/AGL_KEY(:87-102) |
真实模板长什么样(依据:examples/calc_x/job-template.yaml):
# 已精简,保留关键行
env:
- name: QUESTION
value: {{ input.question | yaml_escape }} # rollout 的 input 直接进模板
- name: RESULT
value: {{ input.result | yaml_escape }}
注意两套注入并存:任务数据走模板变量(用户写的),框架脐带走强制 env 覆盖(框架保证的)。
3.2 双循环:轮询 + watch
K8sReconciler.run(k8s_reconciler.py:128-141)同时跑两个任务:
_periodic_reconcile_loop(:149-161):每poll_interval(默认 5 秒,agentlightning/config/controller.yaml)拉一遍 QUEUING/RUNNING 的 rollout 和带managed-by=agentlightninglabel 的 Job,对齐(_reconcile_once,:163-242)——Job 没了而 rollout 还 RUNNING 就标 FAILED “Job disappeared”(孤儿回收,:185-186);Job 的Complete/Failedcondition 映射成 SUCCEEDED/FAILED。_watch_jobs_loop(:283-300):用 kr8s 的 watch API 订阅 Job 事件,Job 一完成立即 PATCH rollout 终态(_handle_job_event,:302-337)——把「发现完成」的延迟从轮询间隔压到事件级。watch 断了就 5 秒后重建(:298-300)。
轮询兜底 + watch 提速,两层互为保险。
3.3 创建限速
K8s API server 对创建速率敏感。_create_job(:244-279)用滑动窗口限速:过去 60 秒(JOB_CREATION_WINDOW_SECONDS,:34)内创建数达到 max_jobs_per_minute(默认 100)就推迟这个 rollout(不失败、下轮再试,:251-258)。创建失败也分情况:422/invalid(模板渲染出的 Job 不合法)直接把 rollout 标 FAILED;其他错误(网络等)只警告、下轮重试(:268-279)。
小结: Controller 是「状态哑、动作准」的调谐器——所有判断都基于 Gateway 账本,自己只维护进程表/Job 映射这类「现实快照」。local 模式用环境变量脐带对接任意 agent 类;k8s 模式用用户模板 + 框架强制字段保证隔离与如实记录。至此执行侧闭环:任务排队 → 进程/Job 启动 → agent 过代理调模型 → 报奖励 → 终态回账本。下一章看学习侧怎么消费这一切。