File size: 6,015 Bytes
afa0cbf | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 | //! Attempts one new continuation turn immediately after daemon thread restoration.
//! Completed or superseded work is discarded; recovery never waits for a later idle event.
//! Permissions and the saved local executor selection must match the resumed runtime.
use super::ThreadRequestProcessor;
use codex_app_server_protocol::ThreadEnvironment;
use codex_app_server_transport::daemon_recovery::InterruptedTurn;
use codex_core::TurnInput;
use codex_core::TurnInputRequest;
use codex_core::TurnInputSubmission;
use codex_core::TurnStartOptions;
use codex_core::context::ContextualUserFragment;
use codex_core::context::InternalContextSource;
use codex_core::context::InternalModelContextFragment;
use codex_protocol::ThreadId;
use codex_protocol::protocol::EnvironmentConfigState;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::TurnAbortReason;
use codex_protocol::protocol::TurnAbortedEvent;
use codex_rollout::RolloutItem;
impl ThreadRequestProcessor {
pub(crate) async fn continue_daemon_turn(&self, thread_id: &str, saved: InterruptedTurn) {
let Ok(thread_id) = ThreadId::from_string(thread_id) else {
return;
};
let Ok(thread) = self.thread_manager.get_thread(thread_id).await else {
return;
};
let Some(path) = thread.rollout_path() else {
return;
};
let history = match codex_rollout::RolloutRecorder::load_rollout_items(&path).await {
Ok((history, _, _)) => history,
Err(err) => {
tracing::warn!(%thread_id, %err, "failed to read interrupted turn history");
return;
}
};
// The snapshot has complete recorded input, but the old process may have
// completed or replaced its turn during the shutdown grace period.
let Some(start) = history
.iter()
.rposition(|item| matches!(item, RolloutItem::EventMsg(EventMsg::TurnStarted(_))))
else {
return;
};
if !matches!(&history[start], RolloutItem::EventMsg(EventMsg::TurnStarted(event))
if event.turn_id == saved.turn_id)
|| history[start..].iter().any(|item| match item {
RolloutItem::EventMsg(EventMsg::TurnComplete(event)) => {
event.turn_id == saved.turn_id
}
RolloutItem::EventMsg(EventMsg::TurnAborted(event)) => event
.turn_id
.as_ref()
.is_none_or(|turn_id| turn_id == &saved.turn_id),
_ => false,
})
{
return;
}
let previous_context = history.iter().rev().find_map(|item| match item {
RolloutItem::TurnContext(context)
if context.turn_id.as_ref() == Some(&saved.turn_id) =>
{
Some(context)
}
_ => None,
});
let Some(previous_context) = previous_context else {
return;
};
let config = thread.config_snapshot().await;
let [environment] = config.environment_selections() else {
return;
};
if environment.environment_id != codex_exec_server::LOCAL_ENVIRONMENT_ID
|| environment.config != EnvironmentConfigState::FromThread
|| saved.local_environment.as_ref() != Some(&ThreadEnvironment::from(environment))
{
return;
}
// Recovery must not override a stricter saved or newly configured policy.
if previous_context.permission_profile() != config.permission_profile {
return;
}
let continuation = ContextualUserFragment::into(InternalModelContextFragment::new(
InternalContextSource::from_static("daemon_recovery"),
"The server restarted and interrupted the previous turn. Continue the unfinished work from the saved conversation. Check the current state before repeating actions that may already have completed.",
));
let request = TurnInputRequest::new(TurnInput::ResponseItem(continuation)).on_start(
TurnStartOptions {
turn_trigger: Some("daemon_recovery".to_string()),
final_output_json_schema: saved.output_schema,
service_tier: saved.service_tier,
cyber_access_program: saved.cyber_access_program,
root_turn_id: previous_context.root_turn_id.clone(),
..Default::default()
},
);
// Close the old turn in persisted history before exposing a new running turn.
if let Err(err) = thread
.append_rollout_items(&[RolloutItem::EventMsg(EventMsg::TurnAborted(
TurnAbortedEvent {
turn_id: Some(saved.turn_id.clone()),
reason: TurnAbortReason::Interrupted,
started_at: None,
completed_at: None,
duration_ms: None,
},
))])
.await
{
tracing::warn!(%thread_id, %err, "failed to close interrupted turn history");
return;
}
match thread.continue_turn_if_idle(request, saved.turn_id).await {
Ok(TurnInputSubmission::Started { .. }) => {
crate::extensions::send_thread_warning(
&self.outgoing,
&self.thread_state_manager,
thread_id,
"Resuming interrupted work".to_string(),
)
.await;
}
Ok(TurnInputSubmission::NotSubmitted { reason }) => {
tracing::debug!(%thread_id, ?reason, "recovery continuation was not started");
}
Ok(TurnInputSubmission::Steered { .. }) => unreachable!("continuation cannot steer"),
Err(err) => tracing::warn!(%thread_id, %err, "recovery continuation failed"),
}
}
}
|