Skip to content

Active goals can auto-continue concurrently across app-server processes #32793

Description

@sirouk

Codex version

Reproduced with codex-cli 0.144.3. The same process-local guards are present on main at 8b2c84d.

What happened?

Two app-server processes sharing one CODEX_HOME can independently resume the same thread while its persisted goal is active. Each process sees its own in-memory thread as idle and starts an automatic goal-continuation turn. The turns can overlap and execute tools concurrently in the same workspace.

I observed an existing CLI goal turn remain in flight while two later remote thread/resume requests each launched another real goal-continuation turn.

Reproduction outline

  1. Start app-server A with the goals feature enabled.
  2. Materialize a persistent thread, set an active goal, and keep its automatic continuation in flight.
  3. Start app-server B with the same CODEX_HOME.
  4. Send thread/resume for the same thread to app-server B.
  5. Observe another automatic continuation and a second turn/started / Responses request while A is still running.

Using a single managed daemon and attaching the TUI with codex --remote unix:// resume avoids the race, but independent app-server processes remain possible (for example a foreground Remote Control host or another Codex surface).

Expected behavior

A second process may resume/read the thread and receive its goal snapshot, but it must not start automatic goal work while another live process owns that thread's goal continuation. If the owner exits or crashes, a later process should be able to continue it.

Source trace

  • A cold thread/resume sends the response and then calls emit_resume_goal_snapshot_and_continue:
    let response_history = thread_history.clone();
    match self
    .thread_manager
    .resume_thread_with_history(
    config,
    thread_history,
    self.auth_manager.clone(),
    self.request_trace_context(&request_id).await,
    supports_openai_form_elicitation,
    )
    .await
    {
    Ok(NewThread {
    thread_id,
    thread: codex_thread,
    session_configured,
    ..
    }) => {
    if let Err(err) = Self::set_app_server_client_info(
    codex_thread.as_ref(),
    app_server_client_name,
    app_server_client_version,
    )
    .await
    {
    self.outgoing.send_error(request_id, err).await;
    return Ok(());
    }
    let instruction_sources = codex_thread.legacy_instruction_sources().await;
    let SessionConfiguredEvent { rollout_path, .. } = session_configured;
    let Some(rollout_path) = rollout_path else {
    let error =
    internal_error(format!("rollout path missing for thread {thread_id}"));
    self.outgoing.send_error(request_id, error).await;
    return Ok(());
    };
    // Auto-attach a thread listener when resuming a thread.
    log_listener_attach_result(
    self.ensure_conversation_listener(
    thread_id,
    request_id.connection_id,
    /*raw_events_enabled*/ false,
    )
    .await,
    thread_id,
    request_id.connection_id,
    "thread",
    );
    let mut thread = match self
    .load_thread_from_resume_source_or_send_internal(
    thread_id,
    codex_thread.as_ref(),
    &response_history,
    rollout_path.as_path(),
    resume_source_thread,
    include_turns,
    )
    .await
    {
    Ok(thread) => thread,
    Err(message) => {
    self.outgoing
    .send_error(request_id, internal_error(message))
    .await;
    return Ok(());
    }
    };
    thread.thread_source = codex_thread
    .config_snapshot()
    .await
    .thread_source
    .map(Into::into);
    self.thread_watch_manager
    .upsert_thread(thread.clone())
    .await;
    let thread_status = self
    .thread_watch_manager
    .loaded_status_for_thread(&thread.id)
    .await;
    set_thread_status_and_interrupt_stale_turns(
    &mut thread,
    thread_status,
    /*has_live_in_progress_turn*/ false,
    );
    let config_snapshot = codex_thread.config_snapshot().await;
    let sandbox = thread_response_sandbox_policy(
    &config_snapshot.permission_profile,
    config_snapshot.cwd().as_path(),
    );
    let active_permission_profile = thread_response_active_permission_profile(
    config_snapshot.active_permission_profile,
    );
    let token_usage_thread = include_turns.then(|| thread.clone());
    let mut initial_turns_page = if let Some(params) = initial_turns_page.as_ref() {
    match build_thread_resume_initial_turns_page(
    response_history.get_rollout_items(),
    thread.status.clone(),
    /*has_live_running_thread*/ false,
    /*active_turn*/ None,
    params,
    ) {
    Ok(page) => Some(page),
    Err(error) => {
    self.outgoing.send_error(request_id, error).await;
    return Ok(());
    }
    }
    } else {
    None
    };
    if redact_resume_payloads {
    redact_thread_resume_payloads(&mut thread.turns);
    if let Some(initial_turns_page) = initial_turns_page.as_mut() {
    redact_thread_resume_payloads(&mut initial_turns_page.data);
    }
    }
    let thread_originator = config_snapshot.originator.clone();
    let response = ThreadResumeResponse {
    thread,
    model: session_configured.model,
    model_provider: session_configured.model_provider_id,
    service_tier: session_configured.service_tier,
    cwd: session_configured.cwd,
    runtime_workspace_roots: config_snapshot.workspace_roots,
    instruction_sources,
    approval_policy: session_configured.approval_policy.into(),
    approvals_reviewer: session_configured.approvals_reviewer.into(),
    sandbox,
    active_permission_profile,
    reasoning_effort: session_configured.reasoning_effort,
    multi_agent_mode: MultiAgentMode::ExplicitRequestOnly,
    initial_turns_page,
    };
    let connection_id = request_id.connection_id;
    self.outgoing
    .send_response_with_thread_originator(request_id, response, thread_originator)
    .await;
    // `excludeTurns` is explicitly the cheap resume path, so avoid
    // rebuilding history only to attribute a replayed usage update.
    if let Some(token_usage_thread) = token_usage_thread {
    let token_usage_turn_id = latest_token_usage_turn_id_from_rollout_items(
    response_history.get_rollout_items(),
    token_usage_thread.turns.as_slice(),
    );
    // The client needs restored usage before it starts another turn.
    // Sending after the response preserves JSON-RPC request ordering while
    // still filling the status line before the next turn lifecycle begins.
    send_thread_token_usage_update_to_connection(
    &self.outgoing,
    connection_id,
    thread_id,
    &token_usage_thread,
    codex_thread.as_ref(),
    token_usage_turn_id,
    )
    .await;
    }
    self.thread_goal_processor
    .emit_resume_goal_snapshot_and_continue(thread_id, codex_thread.as_ref())
    .await;
  • That path emits the thread-idle lifecycle whenever goals are enabled:
    pub(crate) async fn emit_resume_goal_snapshot_and_continue(
    &self,
    thread_id: ThreadId,
    thread: &CodexThread,
    ) {
    if !self.config.features.enabled(Feature::Goals) {
    return;
    }
    self.emit_thread_goal_snapshot(thread_id).await;
    // App-server owns resume response and snapshot ordering, so wait until
    // those are sent before letting extensions react to the idle thread.
    thread.emit_thread_idle_lifecycle_if_idle().await;
  • GoalRuntimeHandle::continue_if_idle uses a per-runtime semaphore and calls the local thread's try_start_turn_if_idle:
    pub async fn restore_after_resume(&self) -> Result<(), String> {
    if !self.is_enabled() {
    return Ok(());
    }
    let goal = self
    .inner
    .state_dbs
    .thread_goals()
    .get_thread_goal(self.thread_id())
    .await
    .map_err(|err| err.to_string())?;
    match goal {
    Some(goal) if goal.status == codex_state::ThreadGoalStatus::Active => {
    self.inner
    .accounting_state
    .mark_idle_goal_active(goal.goal_id);
    self.inner.metrics.record_resumed();
    }
    Some(_) | None => self.inner.accounting_state.clear_active_goal(),
    }
    Ok(())
    }
    pub(crate) async fn continue_if_idle(&self) -> Result<(), String> {
    if !self.tools_visible() {
    self.inner.accounting_state.clear_active_goal();
    return Ok(());
    }
    // Hold this through the read/start window so external set/clear cannot
    // change the goal after we read it but before the continuation launches.
    let _goal_state_permit = self.goal_state_permit().await?;
    let Some(thread_manager) = self.inner.thread_manager.upgrade() else {
    tracing::debug!("skipping goal continuation because thread manager is unavailable");
    return Ok(());
    };
    let Ok(thread) = thread_manager.get_thread(self.inner.thread_id).await else {
    tracing::debug!("skipping goal continuation because live thread is unavailable");
    return Ok(());
    };
    let Some(goal) = self
    .inner
    .state_dbs
    .thread_goals()
    .get_thread_goal(self.thread_id())
    .await
    .map_err(|err| err.to_string())?
    else {
    self.inner.accounting_state.clear_active_goal();
    return Ok(());
    };
    if goal.status != codex_state::ThreadGoalStatus::Active {
    self.inner.accounting_state.clear_active_goal();
    return Ok(());
    }
    let item = continuation_steering_item(&protocol_goal_from_state(goal));
    if let Err(err) = thread.try_start_turn_if_idle(vec![item]).await {
    let reason = err.reason();
    tracing::debug!(
    ?reason,
    "skipping goal continuation because automatic idle work was rejected"
    );
    }
    let current_turn_is_goal_active = self
    .inner
    .accounting_state
    .current_turn_id()
    .is_some_and(|turn_id| {
    self.inner
    .accounting_state
    .turn_is_current_active_goal(turn_id.as_str())
    });
    if !current_turn_is_goal_active {
    self.inner.accounting_state.clear_active_goal();
    }
    Ok(())
    }
  • The thread map, active-turn check, mailbox, and goal semaphore are process-local. The shared thread_goals SQLite row has no owner or lease, so both processes can pass the idle gate.

A goal-specific cross-process advisory lock or ownership lease, released automatically on process death and held for the continuation lifetime, seems like the narrowest fix. This is related to multi-client attachment (#32445, #14722), but it is independently an execution-safety issue.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    app-serverIssues involving app server protocol or interfacesbugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions