Mid-Turn Steering
估计熟练使用codex以及claude的人都会在执行任务中出现这样一个问题,即当你任务执行一半了,但是发现不对,或者有少许偏移,但是已经执行了半天,此时shut是浪费token, 所以此时Agent需要输入指令以纠偏,steer功能便是以此为目的诞生的
Codex的steer机制

首先我不懂rust但是学过c++, 所以大部分还是能看懂,基本上codex实现是将当前的input直接插入到上下文当中,这就意味着,codex没有对当前的执行进行隐式打断,我觉得这一点非常重要,我做业务的时候选择了隐式的打断方式,即将steer指令放入pending的队列,然后打断主Agent的执行,copy上下文到当前新一轮的Agent Run中,然后实现的继续执行,当时设计的初衷是为了解耦,而两种截然不同的方式也必然带来了额外的逻辑,
我的steer机制
延续了之前的机制,即如下
no-steer:用户发起请求-> 选择orchestrator -> 创建Agent运行事件接口创建AgentRunEvent -> 把事件写入AgentEvent表,标记为start,输入的消息进行入msg库 -> 事件ID返回前端存储起来,-> push agent event id,进入异步任务队列 -> 队列收到agent启动的id后,触发agent的消息生成函数
steer: 用户发起请求-> 选择orchestrator -> 检查当前thread下运行的agent事件是否running? -> running -> 则先放入队列 -> 结束当前的main和subdag,即将任务优雅的导向end节点-> 同步消息和结束事件到sql -> 再次循环的检查是否runnig -> pending拿出事件,接受之前已经结束的checkpointer进行继续执行
pub async fn steer_input(
&self,
input: Vec<UserInput>,
additional_context: BTreeMap<String, AdditionalContextEntry>,
expected_turn_id: Option<&str>,
client_user_message_id: Option<String>,
responsesapi_client_metadata: Option<HashMap<String, String>>,
) -> Result<String, SteerInputError> {
let mut active = self.active_turn.lock().await;
let Some(active_turn) = active.as_mut() else {
return Err(SteerInputError::NoActiveTurn(input));
};
let Some(active_task) = active_turn.task.as_ref() else {
return Err(SteerInputError::NoActiveTurn(input));
};
let active_turn_id = &active_task.turn_context.sub_id;
if let Some(expected_turn_id) = expected_turn_id
&& expected_turn_id != active_turn_id
{
return Err(SteerInputError::ExpectedTurnMismatch {
expected: expected_turn_id.to_string(),
actual: active_turn_id.clone(),
});
}
match active_task.kind {
crate::state::TaskKind::Regular => {}
crate::state::TaskKind::Review => {
return Err(SteerInputError::ActiveTurnNotSteerable {
turn_kind: NonSteerableTurnKind::Review,
});
}
crate::state::TaskKind::Compact => {
return Err(SteerInputError::ActiveTurnNotSteerable {
turn_kind: NonSteerableTurnKind::Compact,
});
}
}
if input.is_empty() {
return Err(SteerInputError::EmptyInput);
}
let additional_context_input = {
let mut state = self.state.lock().await;
state.additional_context.merge(additional_context)
};
if let Some(responsesapi_client_metadata) = responsesapi_client_metadata {
active_task
.turn_context
.turn_metadata_state
.set_responsesapi_client_metadata(responsesapi_client_metadata);
}
let mut pending_input = additional_context_input
.into_iter()
.map(ResponseItem::from)
.map(TurnInput::ResponseItem)
.collect::<Vec<_>>();
pending_input.push(TurnInput::UserInput {
content: input,
client_id: client_user_message_id,
});
self.input_queue
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
pending_input,
)
.await;
Ok(active_turn_id.clone())
}