RegularTask 看似只是调用 run_turn,实际承担三道边界:尽早发 TurnStarted、消费一次性的启动预热、在完成竞态处再次检查 pending input。少任何一道,都会出现 UI 延迟、连接浪费或用户补充输入丢失。
入口先尝试 Steer,再决定创建任务
UserInput handler 先构造候选 TurnContext,再调用 steer_input。有活跃 RegularTask 时输入进入旧 Turn;只有明确返回 NoActiveTurn,才用候选 Context 创建新的 RegularTask。Review/Compact 活跃时则报不可 steer,而不是静默替换。
async def user_input_or_turn(submission):
candidate = await new_turn_with_sub_id(submission.id, submission.settings)
result = await steer_input(candidate.input, expected_id=submission.turn_id)
if result.kind == "steered":
return result.active_turn_id
if result.kind == "no_active_turn":
await start_task(candidate, candidate.input, RegularTask())
return candidate.sub_id
raise ActiveTurnNotSteerable(result.task_kind)
TurnStarted 为什么在等待预热之前发送
首次 Session 初始化可能后台预热 ModelClientSession。RegularTask 的 run 先内联发送 TurnStarted 并重置 server reasoning inclusion,然后才等待预热 resolution。这样客户端不会把连接准备时间误判为“提交没有生效”。
预热结果有三种:Ready 交给第一次 run_turn;Unavailable 使用新建 ModelClientSession;Cancelled 则仍运行 hooks 和记录输入,然后返回 None。预热句柄只消费一次,后续循环把 None 传入。
async def regular_run(ctx, input, cancel):
await send(TurnStarted(ctx.id))
reset_server_reasoning_flag()
warm = await consume_startup_prewarm(cancel)
if warm.is_cancelled:
await hooks_and_record(input)
return None
next_input = input
client = warm.session_if_ready
while True:
answer = await run_turn(ctx, next_input, take_once(client), cancel.child())
if not await input_queue.has_pending_input(active_turn):
return answer
next_input = []
为什么 RegularTask 外层还有一重循环
run_turn 自己已经在每次采样后检查 pending input,但存在一个极窄竞态:它完成最后检查准备返回后,用户输入才被 steer 进 TurnState。外层 RegularTask 在 run_turn 返回后再检查一次,就能用同一 TurnContext 再进一轮,而不是留下无人消费的输入。
第二轮传空 next_input,因为新输入已在 InputQueue;若再次把原始 input 传入,会重复记录第一条用户消息。ModelClientSession 是否复用也有细节:预热 session 只作为第一次 run_turn 的初始客户端;每次 run_turn 内部会在多次采样和 retry 间复用自己的 session。
错误和取消怎样回到统一生命周期
RegularTask 用 ? 把 run_turn 的 TurnAborted 传播给 task wrapper。wrapper 先尝试 flush,只有 token 未被取消时才调用 on_task_finished。显式 interrupt 通常已经由 abort 路径取得终态所有权,因此取消后的后台 task 不能再发一个 TurnComplete。
评论
登录后即可评论