← paper & engineering

Mid-Turn Steering

Mid-Turn Steering

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

Codex的steer机制

Codex Steer:同一个 Turn 内追加输入

首先我不懂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())
    }