[codex-analytics] add protocol-native turn timestamps (#16638)

---
[//]: # (BEGIN SAPLING FOOTER)
Stack created with [Sapling](https://sapling-scm.com). Best reviewed
with [ReviewStack](https://reviewstack.dev/openai/codex/pull/16638).
* #16870
* #16706
* #16659
* #16641
* #16640
* __->__ #16638
This commit is contained in:
rhan-oai
2026-04-06 16:22:59 -07:00
committed by GitHub
parent e88c2cf4d7
commit 756c45ec61
58 changed files with 1134 additions and 36 deletions
+128 -14
View File
@@ -127,6 +127,8 @@ use codex_protocol::protocol::RealtimeEvent;
use codex_protocol::protocol::ReviewDecision;
use codex_protocol::protocol::ReviewOutputEvent;
use codex_protocol::protocol::TokenCountEvent;
use codex_protocol::protocol::TurnAbortedEvent;
use codex_protocol::protocol::TurnCompleteEvent;
use codex_protocol::protocol::TurnDiffEvent;
use codex_protocol::request_permissions::PermissionGrantScope as CorePermissionGrantScope;
use codex_protocol::request_permissions::RequestPermissionProfile as CoreRequestPermissionProfile;
@@ -190,6 +192,9 @@ pub(crate) async fn apply_bespoke_event_handling(
items: Vec::new(),
error: None,
status: TurnStatus::InProgress,
started_at: payload.started_at,
completed_at: None,
duration_ms: None,
})
};
let notification = TurnStartedNotification {
@@ -201,14 +206,21 @@ pub(crate) async fn apply_bespoke_event_handling(
.await;
}
}
EventMsg::TurnComplete(_ev) => {
EventMsg::TurnComplete(turn_complete_event) => {
// All per-thread requests are bound to a turn, so abort them.
outgoing.abort_pending_server_requests().await;
let turn_failed = thread_state.lock().await.turn_summary.last_error.is_some();
thread_watch_manager
.note_turn_completed(&conversation_id.to_string(), turn_failed)
.await;
handle_turn_complete(conversation_id, event_turn_id, &outgoing, &thread_state).await;
handle_turn_complete(
conversation_id,
event_turn_id,
turn_complete_event,
&outgoing,
&thread_state,
)
.await;
}
EventMsg::SkillsUpdateAvailable => {
if let ApiVersion::V2 = api_version {
@@ -1704,7 +1716,14 @@ pub(crate) async fn apply_bespoke_event_handling(
thread_watch_manager
.note_turn_interrupted(&conversation_id.to_string())
.await;
handle_turn_interrupted(conversation_id, event_turn_id, &outgoing, &thread_state).await;
handle_turn_interrupted(
conversation_id,
event_turn_id,
turn_aborted_event,
&outgoing,
&thread_state,
)
.await;
}
EventMsg::ThreadRolledBack(_rollback_event) => {
let pending = {
@@ -1866,11 +1885,18 @@ async fn handle_turn_plan_update(
}
}
struct TurnCompletionMetadata {
status: TurnStatus,
error: Option<TurnError>,
started_at: Option<i64>,
completed_at: Option<i64>,
duration_ms: Option<i64>,
}
async fn emit_turn_completed_with_status(
conversation_id: ThreadId,
event_turn_id: String,
status: TurnStatus,
error: Option<TurnError>,
turn_completion_metadata: TurnCompletionMetadata,
outgoing: &ThreadScopedOutgoingMessageSender,
) {
let notification = TurnCompletedNotification {
@@ -1878,8 +1904,11 @@ async fn emit_turn_completed_with_status(
turn: Turn {
id: event_turn_id,
items: vec![],
error,
status,
error: turn_completion_metadata.error,
status: turn_completion_metadata.status,
started_at: turn_completion_metadata.started_at,
completed_at: turn_completion_metadata.completed_at,
duration_ms: turn_completion_metadata.duration_ms,
},
};
outgoing
@@ -2073,6 +2102,7 @@ async fn find_and_remove_turn_summary(
async fn handle_turn_complete(
conversation_id: ThreadId,
event_turn_id: String,
turn_complete_event: TurnCompleteEvent,
outgoing: &ThreadScopedOutgoingMessageSender,
thread_state: &Arc<Mutex<ThreadState>>,
) {
@@ -2083,22 +2113,40 @@ async fn handle_turn_complete(
None => (TurnStatus::Completed, None),
};
emit_turn_completed_with_status(conversation_id, event_turn_id, status, error, outgoing).await;
emit_turn_completed_with_status(
conversation_id,
event_turn_id,
TurnCompletionMetadata {
status,
error,
started_at: turn_summary.started_at,
completed_at: turn_complete_event.completed_at,
duration_ms: turn_complete_event.duration_ms,
},
outgoing,
)
.await;
}
async fn handle_turn_interrupted(
conversation_id: ThreadId,
event_turn_id: String,
turn_aborted_event: TurnAbortedEvent,
outgoing: &ThreadScopedOutgoingMessageSender,
thread_state: &Arc<Mutex<ThreadState>>,
) {
find_and_remove_turn_summary(conversation_id, thread_state).await;
let turn_summary = find_and_remove_turn_summary(conversation_id, thread_state).await;
emit_turn_completed_with_status(
conversation_id,
event_turn_id,
TurnStatus::Interrupted,
/*error*/ None,
TurnCompletionMetadata {
status: TurnStatus::Interrupted,
error: None,
started_at: turn_summary.started_at,
completed_at: turn_aborted_event.completed_at,
duration_ms: turn_aborted_event.duration_ms,
},
outgoing,
)
.await;
@@ -2871,6 +2919,9 @@ mod tests {
Arc::new(Mutex::new(ThreadState::default()))
}
const TEST_TURN_COMPLETED_AT: i64 = 1_716_000_456;
const TEST_TURN_DURATION_MS: i64 = 1_234;
async fn recv_broadcast_message(
rx: &mut mpsc::Receiver<OutgoingEnvelope>,
) -> Result<OutgoingMessage> {
@@ -2884,6 +2935,24 @@ mod tests {
}
}
fn turn_complete_event(turn_id: &str) -> TurnCompleteEvent {
TurnCompleteEvent {
turn_id: turn_id.to_string(),
last_agent_message: None,
completed_at: Some(TEST_TURN_COMPLETED_AT),
duration_ms: Some(TEST_TURN_DURATION_MS),
}
}
fn turn_aborted_event(turn_id: &str) -> TurnAbortedEvent {
TurnAbortedEvent {
turn_id: Some(turn_id.to_string()),
reason: codex_protocol::protocol::TurnAbortReason::Interrupted,
completed_at: Some(TEST_TURN_COMPLETED_AT),
duration_ms: Some(TEST_TURN_DURATION_MS),
}
}
fn command_execution_completion_item(command: &str) -> CommandExecutionCompletionItem {
CommandExecutionCompletionItem {
command: command.to_string(),
@@ -3648,10 +3717,25 @@ mod tests {
ThreadId::new(),
);
let thread_state = new_thread_state();
{
let mut state = thread_state.lock().await;
state.track_current_turn_event(&EventMsg::TurnStarted(
codex_protocol::protocol::TurnStartedEvent {
turn_id: event_turn_id.clone(),
started_at: Some(42),
model_context_window: None,
collaboration_mode_kind: Default::default(),
},
));
state.track_current_turn_event(&EventMsg::TurnComplete(turn_complete_event(
&event_turn_id,
)));
}
handle_turn_complete(
conversation_id,
event_turn_id.clone(),
turn_complete_event(&event_turn_id),
&outgoing,
&thread_state,
)
@@ -3663,6 +3747,9 @@ mod tests {
assert_eq!(n.turn.id, event_turn_id);
assert_eq!(n.turn.status, TurnStatus::Completed);
assert_eq!(n.turn.error, None);
assert_eq!(n.turn.started_at, Some(42));
assert_eq!(n.turn.completed_at, Some(TEST_TURN_COMPLETED_AT));
assert_eq!(n.turn.duration_ms, Some(TEST_TURN_DURATION_MS));
}
other => bail!("unexpected message: {other:?}"),
}
@@ -3696,6 +3783,7 @@ mod tests {
handle_turn_interrupted(
conversation_id,
event_turn_id.clone(),
turn_aborted_event(&event_turn_id),
&outgoing,
&thread_state,
)
@@ -3707,6 +3795,8 @@ mod tests {
assert_eq!(n.turn.id, event_turn_id);
assert_eq!(n.turn.status, TurnStatus::Interrupted);
assert_eq!(n.turn.error, None);
assert_eq!(n.turn.completed_at, Some(TEST_TURN_COMPLETED_AT));
assert_eq!(n.turn.duration_ms, Some(TEST_TURN_DURATION_MS));
}
other => bail!("unexpected message: {other:?}"),
}
@@ -3740,6 +3830,7 @@ mod tests {
handle_turn_complete(
conversation_id,
event_turn_id.clone(),
turn_complete_event(&event_turn_id),
&outgoing,
&thread_state,
)
@@ -3758,6 +3849,8 @@ mod tests {
additional_details: None,
})
);
assert_eq!(n.turn.completed_at, Some(TEST_TURN_COMPLETED_AT));
assert_eq!(n.turn.duration_ms, Some(TEST_TURN_DURATION_MS));
}
other => bail!("unexpected message: {other:?}"),
}
@@ -4000,7 +4093,14 @@ mod tests {
&thread_state,
)
.await;
handle_turn_complete(conversation_a, a_turn1.clone(), &outgoing, &thread_state).await;
handle_turn_complete(
conversation_a,
a_turn1.clone(),
turn_complete_event(&a_turn1),
&outgoing,
&thread_state,
)
.await;
// Turn 1 on conversation B
let b_turn1 = "b_turn1".to_string();
@@ -4014,11 +4114,25 @@ mod tests {
&thread_state,
)
.await;
handle_turn_complete(conversation_b, b_turn1.clone(), &outgoing, &thread_state).await;
handle_turn_complete(
conversation_b,
b_turn1.clone(),
turn_complete_event(&b_turn1),
&outgoing,
&thread_state,
)
.await;
// Turn 2 on conversation A
let a_turn2 = "a_turn2".to_string();
handle_turn_complete(conversation_a, a_turn2.clone(), &outgoing, &thread_state).await;
handle_turn_complete(
conversation_a,
a_turn2.clone(),
turn_complete_event(&a_turn2),
&outgoing,
&thread_state,
)
.await;
// Verify: A turn 1
let msg = recv_broadcast_message(&mut rx).await?;
@@ -6581,6 +6581,9 @@ impl CodexMessageProcessor {
items: vec![],
error: None,
status: TurnStatus::InProgress,
started_at: None,
completed_at: None,
duration_ms: None,
};
let response = TurnStartResponse { turn };
@@ -6917,6 +6920,9 @@ impl CodexMessageProcessor {
items,
error: None,
status: TurnStatus::InProgress,
started_at: None,
completed_at: None,
duration_ms: None,
}
}
@@ -9566,6 +9572,7 @@ mod tests {
state.track_current_turn_event(&EventMsg::TurnStarted(
codex_protocol::protocol::TurnStartedEvent {
turn_id: "turn-1".to_string(),
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
},
+3
View File
@@ -826,6 +826,9 @@ mod tests {
items: Vec::new(),
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: Some(0),
duration_ms: None,
},
})
));
+7 -1
View File
@@ -45,6 +45,7 @@ pub(crate) enum ThreadListenerCommand {
/// Per-conversation accumulation of the latest states e.g. error message while a turn runs.
#[derive(Default, Clone)]
pub(crate) struct TurnSummary {
pub(crate) started_at: Option<i64>,
pub(crate) file_change_started: HashSet<String>,
pub(crate) command_execution_started: HashSet<String>,
pub(crate) last_error: Option<TurnError>,
@@ -110,8 +111,13 @@ impl ThreadState {
}
pub(crate) fn track_current_turn_event(&mut self, event: &EventMsg) {
if let EventMsg::TurnStarted(payload) = event {
self.turn_summary.started_at = payload.started_at;
}
self.current_turn_history.handle_event(event);
if !self.current_turn_history.has_active_turn() {
if matches!(event, EventMsg::TurnAborted(_) | EventMsg::TurnComplete(_))
&& !self.current_turn_history.has_active_turn()
{
self.current_turn_history.reset();
}
}