mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
goals: keep pause transitions explicit (#23088)
## Problem This addresses several user-reported cases where active goals were paused even though the user had not explicitly asked for that transition: - the guardian approval-review circuit breaker interrupted a turn and implicitly paused the goal - a shutdown in one app-server instance could pause a goal while a second instance was still actively running the same thread - steering-style interrupts could also pause the goal even though they are meant to redirect work, not stop the goal lifecycle The common problem was that core treated `TurnAbortReason::Interrupted` as an implicit request to transition the persisted goal to `paused`. That made unrelated interrupt paths mutate goal state as a side effect, and in the multi-app-server case it allowed stale process teardown to pause a live goal owned by another running client. After this change, transitioning a goal to `paused` is always an explicit action performed by a client or another intentional goal-state mutation. It is never an implicit transition triggered by generic interrupt handling. Refs #22884. ## What changed - Remove the goal runtime path that paused active goals after interrupted task aborts. - Drop the now-unused abort reason from `GoalRuntimeEvent::TaskAborted`. - Update the focused regression coverage so an interrupted active goal still accounts usage but remains `active`.
This commit is contained in:
committed by
GitHub
Unverified
parent
ae03d073b3
commit
55f6bbc667
@@ -26,13 +26,11 @@ use codex_otel::GOAL_USAGE_LIMITED_METRIC;
|
||||
use codex_protocol::config_types::ModeKind;
|
||||
use codex_protocol::models::ContentItem;
|
||||
use codex_protocol::models::ResponseInputItem;
|
||||
use codex_protocol::protocol::Event;
|
||||
use codex_protocol::protocol::EventMsg;
|
||||
use codex_protocol::protocol::ThreadGoal;
|
||||
use codex_protocol::protocol::ThreadGoalStatus;
|
||||
use codex_protocol::protocol::ThreadGoalUpdatedEvent;
|
||||
use codex_protocol::protocol::TokenUsage;
|
||||
use codex_protocol::protocol::TurnAbortReason;
|
||||
use codex_protocol::protocol::validate_thread_goal_objective;
|
||||
use codex_rollout::state_db::reconcile_rollout;
|
||||
use codex_thread_store::LocalThreadStore;
|
||||
@@ -156,7 +154,6 @@ pub(crate) enum GoalRuntimeEvent<'a> {
|
||||
MaybeContinueIfIdle,
|
||||
TaskAborted {
|
||||
turn_context: Option<&'a TurnContext>,
|
||||
reason: TurnAbortReason,
|
||||
},
|
||||
UsageLimitReached {
|
||||
turn_context: &'a TurnContext,
|
||||
@@ -337,8 +334,8 @@ impl Session {
|
||||
/// starts capture the active goal and token baseline, tool completions
|
||||
/// account usage and may inject budget steering, completion accounting
|
||||
/// suppresses that steering, external mutations account best-effort before
|
||||
/// changing state, interrupts pause active goals, thread resumes restore
|
||||
/// runtime state for already-active goals, explicit maybe-continue events
|
||||
/// changing state, thread resumes restore runtime state for already-active
|
||||
/// goals, explicit maybe-continue events
|
||||
/// start idle goal continuation turns, and continuation turns with no counted
|
||||
/// autonomous activity suppress the next automatic continuation until
|
||||
/// user/tool/external activity resets it.
|
||||
@@ -390,12 +387,8 @@ impl Session {
|
||||
self.maybe_continue_goal_if_idle_runtime().await;
|
||||
Ok(())
|
||||
}),
|
||||
GoalRuntimeEvent::TaskAborted {
|
||||
turn_context,
|
||||
reason,
|
||||
} => Box::pin(async move {
|
||||
self.handle_thread_goal_task_abort(turn_context, reason)
|
||||
.await;
|
||||
GoalRuntimeEvent::TaskAborted { turn_context } => Box::pin(async move {
|
||||
self.handle_thread_goal_task_abort(turn_context).await;
|
||||
Ok(())
|
||||
}),
|
||||
GoalRuntimeEvent::UsageLimitReached { turn_context } => Box::pin(async move {
|
||||
@@ -947,11 +940,7 @@ impl Session {
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_thread_goal_task_abort(
|
||||
&self,
|
||||
turn_context: Option<&TurnContext>,
|
||||
reason: TurnAbortReason,
|
||||
) {
|
||||
async fn handle_thread_goal_task_abort(&self, turn_context: Option<&TurnContext>) {
|
||||
if let Some(turn_context) = turn_context {
|
||||
self.take_thread_goal_continuation_turn(&turn_context.sub_id)
|
||||
.await;
|
||||
@@ -974,12 +963,6 @@ impl Session {
|
||||
accounting.turn = None;
|
||||
}
|
||||
}
|
||||
|
||||
if reason == TurnAbortReason::Interrupted
|
||||
&& let Err(err) = self.pause_active_thread_goal_for_interrupt().await
|
||||
{
|
||||
tracing::warn!("failed to pause active thread goal after interrupt: {err}");
|
||||
}
|
||||
}
|
||||
|
||||
async fn account_thread_goal_progress(
|
||||
@@ -1187,57 +1170,6 @@ impl Session {
|
||||
}
|
||||
}
|
||||
|
||||
async fn pause_active_thread_goal_for_interrupt(&self) -> anyhow::Result<()> {
|
||||
if should_ignore_goal_for_mode(self.collaboration_mode().await.mode) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if !self.enabled(Feature::Goals) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let _continuation_guard = self
|
||||
.goal_runtime
|
||||
.continuation_lock
|
||||
.acquire()
|
||||
.await
|
||||
.context("goal continuation semaphore closed")?;
|
||||
let Some(state_db) = self.state_db_for_thread_goals().await? else {
|
||||
return Ok(());
|
||||
};
|
||||
self.account_thread_goal_wall_clock_usage(
|
||||
&state_db,
|
||||
codex_state::ThreadGoalAccountingMode::ActiveStatusOnly,
|
||||
TerminalMetricEmission::Emit,
|
||||
)
|
||||
.await?;
|
||||
let Some(goal) = state_db
|
||||
.thread_goals()
|
||||
.pause_active_thread_goal(self.conversation_id)
|
||||
.await?
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let goal = protocol_goal_from_state(goal);
|
||||
*self.goal_runtime.budget_limit_reported_goal_id.lock().await = None;
|
||||
self.goal_runtime
|
||||
.accounting
|
||||
.lock()
|
||||
.await
|
||||
.wall_clock
|
||||
.clear_active_goal();
|
||||
self.send_event_raw(Event {
|
||||
id: uuid::Uuid::new_v4().to_string(),
|
||||
msg: EventMsg::ThreadGoalUpdated(ThreadGoalUpdatedEvent {
|
||||
thread_id: self.conversation_id,
|
||||
turn_id: None,
|
||||
goal,
|
||||
}),
|
||||
})
|
||||
.await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn usage_limit_active_thread_goal_for_turn(
|
||||
&self,
|
||||
turn_context: &TurnContext,
|
||||
|
||||
@@ -3355,7 +3355,6 @@ impl Session {
|
||||
pub async fn interrupt_task(self: &Arc<Self>) {
|
||||
info!("interrupt received: abort current task, if any");
|
||||
let had_active_turn = self.active_turn.lock().await.is_some();
|
||||
// Even without an active task, interrupt handling pauses any active goal.
|
||||
self.abort_all_tasks(TurnAbortReason::Interrupted).await;
|
||||
if !had_active_turn {
|
||||
self.cancel_mcp_startup().await;
|
||||
|
||||
@@ -8243,7 +8243,7 @@ async fn abort_empty_active_turn_preserves_pending_input() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_accounts_active_goal_before_pausing() -> anyhow::Result<()> {
|
||||
async fn interrupt_accounts_active_goal_without_pausing() -> anyhow::Result<()> {
|
||||
let (sess, tc, _rx, _codex_home) = make_goal_session_and_context_with_rx().await;
|
||||
sess.set_thread_goal(
|
||||
tc.as_ref(),
|
||||
@@ -8273,7 +8273,7 @@ async fn interrupt_accounts_active_goal_before_pausing() -> anyhow::Result<()> {
|
||||
.await?
|
||||
.expect("goal should remain persisted after interrupt");
|
||||
assert_eq!(
|
||||
codex_protocol::protocol::ThreadGoalStatus::Paused,
|
||||
codex_protocol::protocol::ThreadGoalStatus::Active,
|
||||
goal.status
|
||||
);
|
||||
assert_eq!(70, goal.tokens_used);
|
||||
@@ -8283,6 +8283,34 @@ async fn interrupt_accounts_active_goal_before_pausing() -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn shutdown_without_active_turn_keeps_active_goal_active() -> anyhow::Result<()> {
|
||||
let (sess, tc, _rx, _codex_home) = make_goal_session_and_context_with_rx().await;
|
||||
sess.set_thread_goal(
|
||||
tc.as_ref(),
|
||||
SetGoalRequest {
|
||||
objective: Some("Keep improving the benchmark".to_string()),
|
||||
status: None,
|
||||
token_budget: None,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
|
||||
assert!(sess.active_turn.lock().await.is_none());
|
||||
assert!(handlers::shutdown(&sess, "shutdown".to_string()).await);
|
||||
|
||||
let goal = sess
|
||||
.get_thread_goal()
|
||||
.await?
|
||||
.expect("goal should remain persisted after shutdown");
|
||||
assert_eq!(
|
||||
codex_protocol::protocol::ThreadGoalStatus::Active,
|
||||
goal.status
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn active_goal_continuation_runs_again_after_no_tool_turn() -> anyhow::Result<()> {
|
||||
let server = start_mock_server().await;
|
||||
|
||||
@@ -513,7 +513,6 @@ impl Session {
|
||||
&& let Err(err) = self
|
||||
.goal_runtime_apply(GoalRuntimeEvent::TaskAborted {
|
||||
turn_context: turn_context.as_deref(),
|
||||
reason: reason.clone(),
|
||||
})
|
||||
.await
|
||||
{
|
||||
@@ -561,7 +560,6 @@ impl Session {
|
||||
if let Err(err) = self
|
||||
.goal_runtime_apply(GoalRuntimeEvent::TaskAborted {
|
||||
turn_context: turn_context.as_deref(),
|
||||
reason: reason.clone(),
|
||||
})
|
||||
.await
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user