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"),
        }
    }
}