Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions codex-rs/app-server-protocol/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ pub use protocol::common::*;
pub use protocol::event_mapping::*;
pub use protocol::item_builders::*;
pub use protocol::thread_history::*;
pub use protocol::thread_history_projection::*;
pub use protocol::v1::ApplyPatchApprovalParams;
pub use protocol::v1::ApplyPatchApprovalResponse;
pub use protocol::v1::ClientInfo;
Expand Down
1 change: 1 addition & 0 deletions codex-rs/app-server-protocol/src/protocol/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,5 +7,6 @@ pub mod item_builders;
mod mappers;
mod serde_helpers;
pub mod thread_history;
pub mod thread_history_projection;
pub mod v1;
pub mod v2;
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
//! Stateless projection from canonical paginated rollout records to thread-history changes.
//!
//! This module is only for the new paginated rollout format that persists canonical
//! `ItemCompleted(TurnItem)` records, not legacy event-only rollouts.

use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::RolloutItem;
use codex_protocol::protocol::RolloutLine;

use crate::protocol::thread_history::ThreadHistoryChangeSet;
use crate::protocol::thread_history::ThreadHistoryItemChange;
use crate::protocol::thread_history::ThreadHistoryTurnChange;
use crate::protocol::v2::ThreadItem;
use crate::protocol::v2::TurnError;
use crate::protocol::v2::TurnStatus;

/// Project one durable rollout line without reconstructing earlier history.
///
/// Callers that replay a JSONL suffix should invoke it once per line, in ordinal order, so storage
/// can preserve the first and latest timestamps for repeated item snapshots independently.
pub fn project_rollout_line(line: &RolloutLine) -> ThreadHistoryChangeSet {
match &line.item {
RolloutItem::EventMsg(EventMsg::TurnStarted(event)) => ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: event.turn_id.clone(),
status: TurnStatus::InProgress,
error: None,
started_at: event.started_at,
completed_at: None,
duration_ms: None,
}],
..Default::default()
},
RolloutItem::EventMsg(EventMsg::TurnComplete(event)) => ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: event.turn_id.clone(),
status: if event.error.is_some() {
TurnStatus::Failed
} else {
TurnStatus::Completed
},
error: event.error.as_ref().map(|error| TurnError {
message: error.message.clone(),
codex_error_info: error.codex_error_info.clone().map(Into::into),
additional_details: None,
}),
started_at: event.started_at,
completed_at: event.completed_at,
duration_ms: event.duration_ms,
}],
..Default::default()
},
RolloutItem::EventMsg(EventMsg::TurnAborted(event)) => {
let Some(turn_id) = event.turn_id.as_ref() else {
return ThreadHistoryChangeSet::default();
};
ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: turn_id.clone(),
status: TurnStatus::Interrupted,
error: None,
started_at: event.started_at,
completed_at: event.completed_at,
duration_ms: event.duration_ms,
}],
..Default::default()
}
}
RolloutItem::EventMsg(EventMsg::ItemCompleted(event)) => ThreadHistoryChangeSet {
changed_items: vec![ThreadHistoryItemChange {
turn_id: event.turn_id.clone(),
item: ThreadItem::from(event.item.clone()),
}],
..Default::default()
},
RolloutItem::SessionMeta(_)
| RolloutItem::ResponseItem(_)
| RolloutItem::InterAgentCommunication(_)
| RolloutItem::InterAgentCommunicationMetadata { .. }
| RolloutItem::Compacted(_)
| RolloutItem::TurnContext(_)
| RolloutItem::WorldState(_)
| RolloutItem::EventMsg(_) => ThreadHistoryChangeSet::default(),
}
}

#[cfg(test)]
#[path = "thread_history_projection_tests.rs"]
mod tests;
Original file line number Diff line number Diff line change
@@ -0,0 +1,211 @@
use codex_protocol::ThreadId;
use codex_protocol::items::AgentMessageContent;
use codex_protocol::items::AgentMessageItem;
use codex_protocol::items::TurnItem;
use codex_protocol::items::UserMessageItem;
use codex_protocol::protocol::CompactedItem;
use codex_protocol::protocol::ErrorEvent;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::ItemCompletedEvent;
use codex_protocol::protocol::RolloutItem;
use codex_protocol::protocol::RolloutLine;
use codex_protocol::protocol::TurnAbortReason;
use codex_protocol::protocol::TurnAbortedEvent;
use codex_protocol::protocol::TurnCompleteEvent;
use codex_protocol::protocol::TurnStartedEvent;
use codex_protocol::user_input::UserInput;
use pretty_assertions::assert_eq;

use super::*;
use crate::protocol::v2::ThreadItem;
use crate::protocol::v2::TurnError;

#[test]
fn projects_turn_lifecycle_without_prior_builder_state() {
let started = project(RolloutItem::EventMsg(EventMsg::TurnStarted(
TurnStartedEvent {
turn_id: "turn-1".to_string(),
trace_id: None,
started_at: Some(10),
model_context_window: None,
collaboration_mode_kind: Default::default(),
},
)));
let completed = project(RolloutItem::EventMsg(EventMsg::TurnComplete(
TurnCompleteEvent {
turn_id: "turn-1".to_string(),
last_agent_message: None,
error: None,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
time_to_first_token_ms: None,
},
)));

assert_eq!(started.changed_turns.len(), 1);
assert_eq!(started.changed_turns[0].turn_id, "turn-1");
assert_eq!(started.changed_turns[0].status, TurnStatus::InProgress);
assert_eq!(started.changed_turns[0].started_at, Some(10));
assert_eq!(
completed,
ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-1".to_string(),
status: TurnStatus::Completed,
error: None,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
}],
..Default::default()
}
);
}

#[test]
fn projects_failed_turn_completion_as_snapshot() {
let error = ErrorEvent {
message: "request failed".to_string(),
codex_error_info: None,
};

let changes = project(RolloutItem::EventMsg(EventMsg::TurnComplete(
TurnCompleteEvent {
turn_id: "turn-1".to_string(),
last_agent_message: None,
error: Some(error),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
time_to_first_token_ms: None,
},
)));

assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-1".to_string(),
status: TurnStatus::Failed,
error: Some(TurnError {
message: "request failed".to_string(),
codex_error_info: None,
additional_details: None,
}),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
}],
..Default::default()
}
);
}

#[test]
fn projects_completed_canonical_turn_items() {
let thread_id = ThreadId::default();
let user_item = TurnItem::UserMessage(UserMessageItem {
id: "user-1".to_string(),
client_id: None,
content: vec![UserInput::Text {
text: "hello".to_string(),
text_elements: Vec::new(),
}],
});
let agent_item = TurnItem::AgentMessage(AgentMessageItem {
id: "agent-1".to_string(),
content: vec![AgentMessageContent::Text {
text: "done".to_string(),
}],
phase: None,
memory_citation: None,
});

let user_changes = project(item_completed(thread_id, "turn-1", user_item.clone()));
let agent_changes = project(item_completed(thread_id, "turn-1", agent_item.clone()));

assert_eq!(
user_changes.changed_items,
vec![ThreadHistoryItemChange {
turn_id: "turn-1".to_string(),
item: ThreadItem::from(user_item),
}]
);
assert_eq!(
agent_changes.changed_items,
vec![ThreadHistoryItemChange {
turn_id: "turn-1".to_string(),
item: ThreadItem::from(agent_item),
}]
);
}

#[test]
fn ignores_legacy_abort_without_turn_id_and_context_only_records() {
let aborted = project(RolloutItem::EventMsg(EventMsg::TurnAborted(
TurnAbortedEvent {
turn_id: None,
reason: TurnAbortReason::Interrupted,
started_at: None,
completed_at: None,
duration_ms: None,
},
)));
let compacted = project(RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: None,
window_number: None,
first_window_id: None,
previous_window_id: None,
window_id: None,
}));

assert!(aborted.is_empty());
assert!(compacted.is_empty());
}

#[test]
fn projects_identified_turn_aborts() {
let changes = project(RolloutItem::EventMsg(EventMsg::TurnAborted(
TurnAbortedEvent {
turn_id: Some("turn-1".to_string()),
reason: TurnAbortReason::Interrupted,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
},
)));

assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-1".to_string(),
status: TurnStatus::Interrupted,
error: None,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
}],
..Default::default()
}
);
}

fn project(item: RolloutItem) -> ThreadHistoryChangeSet {
project_rollout_line(&RolloutLine {
timestamp: "2026-07-09T00:00:00.000Z".to_string(),
ordinal: Some(7),
item,
})
}

fn item_completed(thread_id: ThreadId, turn_id: &str, item: TurnItem) -> RolloutItem {
RolloutItem::EventMsg(EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: turn_id.to_string(),
item,
completed_at_ms: 123,
}))
}
Original file line number Diff line number Diff line change
Expand Up @@ -1015,6 +1015,7 @@ mod thread_processor_behavior_tests {

let line = RolloutLine {
timestamp: timestamp.clone(),
ordinal: None,
item: RolloutItem::SessionMeta(SessionMetaLine {
meta: session_meta.clone(),
git: None,
Expand Down Expand Up @@ -1082,6 +1083,7 @@ mod thread_processor_behavior_tests {

let line = RolloutLine {
timestamp,
ordinal: None,
item: RolloutItem::SessionMeta(SessionMetaLine {
meta: session_meta,
git: None,
Expand Down Expand Up @@ -1124,6 +1126,7 @@ mod thread_processor_behavior_tests {

let line = RolloutLine {
timestamp,
ordinal: None,
item: RolloutItem::SessionMeta(SessionMetaLine {
meta: session_meta,
git: None,
Expand Down
1 change: 1 addition & 0 deletions codex-rs/cli/src/doctor/thread_inventory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -819,6 +819,7 @@ mod tests {
let parsed_thread_id = ThreadId::from_string(thread_id).expect("thread id");
let rollout_line = RolloutLine {
timestamp: timestamp.to_string(),
ordinal: None,
item: RolloutItem::SessionMeta(codex_protocol::protocol::SessionMetaLine {
meta: codex_protocol::protocol::SessionMeta {
session_id: parsed_thread_id.into(),
Expand Down
Loading
Loading