mirror of
https://github.com/openai/codex.git
synced 2026-09-28 08:43:01 +08:00
Route rollout reads through the canonical JSON decoder (#42378)
## Why Directly deserializing the flattened `RolloutLine` envelope can reject nested decimal values, preventing affected paginated sessions from resuming. ## What changed - Add canonical string, byte, and reverse-scanner helpers that decode rollout records through `serde_json::Value` before decoding the flattened item. - Route rollout readers across session discovery, history, migration, search, thread storage, and transcript previews through those helpers. - Remove `Deserialize` from `RolloutLine` so new readers cannot bypass the canonical persistence decoder. ## Testing Add coverage that resumes a paginated rollout after a token-count record with a decimal rate-limit value and verifies that ordinal sequencing continues. GitOrigin-RevId: 49abac1e0751c073daa5a93a840d8a483fd2d013
This commit is contained in:
BIN
Binary file not shown.
@@ -27,7 +27,6 @@ use codex_config::types::AuthCredentialsStoreMode;
|
||||
use codex_protocol::models::ContentItem;
|
||||
use codex_protocol::models::ResponseItem;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use core_test_support::responses;
|
||||
use core_test_support::skip_if_no_network;
|
||||
use pretty_assertions::assert_eq;
|
||||
@@ -334,7 +333,7 @@ fn replace_attribution_fragment_with_legacy(
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(|line| {
|
||||
let mut line = serde_json::from_str::<RolloutLine>(line)?;
|
||||
let mut line = codex_rollout::parse_rollout_line(line)?;
|
||||
if let RolloutItem::ResponseItem(response_item) = &mut line.item
|
||||
&& let ResponseItem::Message { role, content, .. } = &mut response_item.item
|
||||
&& role == "developer"
|
||||
|
||||
@@ -1778,7 +1778,7 @@ async fn assert_thread_fork_freezes_active_paginated_turn_as_interrupted(
|
||||
.expect("fork history base");
|
||||
let child_rollout = std::fs::read_to_string(forked_path.as_path())?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<Result<Vec<_>, _>>()?;
|
||||
assert!(matches!(
|
||||
child_rollout.as_slice(),
|
||||
|
||||
@@ -44,7 +44,6 @@ use codex_protocol::protocol::MultiAgentVersion;
|
||||
use codex_protocol::protocol::SessionSource as CoreSessionSource;
|
||||
use codex_protocol::protocol::SubAgentSource;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use codex_rollout::append_rollout_item_to_path;
|
||||
use codex_rollout::read_session_meta_line;
|
||||
use codex_state::DirectionalThreadSpawnEdgeStatus;
|
||||
@@ -218,7 +217,7 @@ fn set_rollout_cwd(path: &Path, cwd: &Path) -> Result<()> {
|
||||
let first_line = lines
|
||||
.first_mut()
|
||||
.ok_or_else(|| anyhow::anyhow!("rollout at {} is empty", path.display()))?;
|
||||
let mut rollout_line: RolloutLine = serde_json::from_str(first_line)?;
|
||||
let mut rollout_line = codex_rollout::parse_rollout_line(first_line)?;
|
||||
let RolloutItem::SessionMeta(mut session_meta_line) = rollout_line.item else {
|
||||
return Err(anyhow::anyhow!(
|
||||
"rollout at {} does not start with session metadata",
|
||||
|
||||
@@ -38,7 +38,6 @@ use codex_protocol::config_types::Settings;
|
||||
use codex_protocol::openai_models::ReasoningEffort;
|
||||
use codex_protocol::protocol::EventMsg;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use codex_rollout::read_session_meta_line;
|
||||
use codex_utils_absolute_path::AbsolutePathBuf;
|
||||
use pretty_assertions::assert_eq;
|
||||
@@ -108,7 +107,7 @@ async fn thread_revert_preserves_fork_cutoff_after_cold_resume() -> Result<()> {
|
||||
let inherited_revert_cutoff =
|
||||
std::fs::read_to_string(parent.path.as_ref().expect("parent rollout"))?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.find_map(|line| match line.item {
|
||||
|
||||
@@ -5,7 +5,6 @@ use super::Config;
|
||||
use super::DoctorCheck;
|
||||
use super::DoctorIssue;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::protocol::InternalSessionSource;
|
||||
use codex_protocol::protocol::SessionSource;
|
||||
use codex_protocol::protocol::SubAgentSource;
|
||||
@@ -543,7 +542,7 @@ async fn thread_id_from_rollout(path: &Path) -> RolloutThreadId {
|
||||
Err(_) => continue,
|
||||
};
|
||||
if item_type == "session_meta" {
|
||||
return match serde_json::from_str::<RolloutLine>(line.trim()) {
|
||||
return match codex_rollout::parse_rollout_line(line.trim()) {
|
||||
Ok(line) => match line.item {
|
||||
RolloutItem::SessionMeta(session_meta) => {
|
||||
RolloutThreadId::Id(session_meta.meta.id.to_string())
|
||||
@@ -560,7 +559,7 @@ async fn thread_id_from_rollout(path: &Path) -> RolloutThreadId {
|
||||
};
|
||||
}
|
||||
if !has_legacy_item {
|
||||
has_legacy_item = serde_json::from_str::<RolloutLine>(line.trim()).is_ok();
|
||||
has_legacy_item = codex_rollout::parse_rollout_line(line.trim()).is_ok();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -745,6 +744,7 @@ where
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::ThreadId;
|
||||
use codex_utils_absolute_path::test_support::PathExt;
|
||||
use pretty_assertions::assert_eq;
|
||||
@@ -862,7 +862,7 @@ mod tests {
|
||||
fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id);
|
||||
let contents = std::fs::read_to_string(&path).expect("rollout file");
|
||||
let mut rollout_line =
|
||||
serde_json::from_str::<RolloutLine>(contents.trim()).expect("rollout line");
|
||||
codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line");
|
||||
let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else {
|
||||
panic!("expected session metadata");
|
||||
};
|
||||
@@ -981,7 +981,7 @@ mod tests {
|
||||
fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id);
|
||||
let contents = std::fs::read_to_string(&metadata_path).expect("rollout file");
|
||||
let mut rollout_line =
|
||||
serde_json::from_str::<RolloutLine>(contents.trim()).expect("rollout line");
|
||||
codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line");
|
||||
let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else {
|
||||
panic!("expected session metadata");
|
||||
};
|
||||
@@ -1114,7 +1114,7 @@ mod tests {
|
||||
fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id);
|
||||
let contents = std::fs::read_to_string(&path).expect("rollout file");
|
||||
let mut rollout_line =
|
||||
serde_json::from_str::<RolloutLine>(contents.trim()).expect("rollout line");
|
||||
codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line");
|
||||
let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else {
|
||||
panic!("expected session metadata");
|
||||
};
|
||||
@@ -1148,7 +1148,7 @@ mod tests {
|
||||
fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id);
|
||||
let contents = std::fs::read_to_string(&path).expect("rollout file");
|
||||
let mut rollout_line =
|
||||
serde_json::from_str::<RolloutLine>(contents.trim()).expect("rollout line");
|
||||
codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line");
|
||||
let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else {
|
||||
panic!("expected session metadata");
|
||||
};
|
||||
|
||||
@@ -22,7 +22,6 @@ use codex_extension_api::empty_extension_registry;
|
||||
use codex_features::Feature;
|
||||
use codex_history::CompactedItem;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::AuthManager;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_protocol::AgentPath;
|
||||
@@ -1262,7 +1261,7 @@ async fn spawn_agent_fork_from_paginated_parent_uses_model_context_prefix() {
|
||||
let lines = std::fs::read_to_string(&rollout_path)
|
||||
.expect("read child rollout")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<RolloutLine>(line).expect("parse rollout line"))
|
||||
.map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line"))
|
||||
.collect::<Vec<_>>();
|
||||
let RolloutItem::SessionMeta(meta_line) = &lines[0].item else {
|
||||
panic!("child rollout should start with session metadata");
|
||||
@@ -1466,7 +1465,7 @@ async fn spawn_agent_fork_drops_inherited_token_usage_state() {
|
||||
let lines = std::fs::read_to_string(&rollout_path)
|
||||
.expect("read child rollout")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<RolloutLine>(line).expect("parse rollout line"))
|
||||
.map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line"))
|
||||
.collect::<Vec<_>>();
|
||||
assert!(
|
||||
!lines.iter().any(|line| {
|
||||
|
||||
@@ -3,7 +3,6 @@ use codex_core::StartThreadOptions;
|
||||
use codex_core::SuspendTurnOutcome;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
@@ -152,7 +151,7 @@ async fn root_turn_suspension_preserves_unfinished_turn_history() {
|
||||
let items = rollout
|
||||
.lines()
|
||||
.map(|line| {
|
||||
serde_json::from_str::<RolloutLine>(line)
|
||||
codex_rollout::parse_rollout_line(line)
|
||||
.expect("parse durable rollout")
|
||||
.item
|
||||
})
|
||||
|
||||
@@ -8,7 +8,6 @@ use codex_exec_server::LOCAL_ENVIRONMENT_ID;
|
||||
use codex_exec_server::REMOTE_ENVIRONMENT_ID;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_home::CodexHomeUserInstructionsProvider;
|
||||
use codex_protocol::config_types::TrustLevel;
|
||||
use codex_protocol::models::PermissionProfile;
|
||||
@@ -98,7 +97,7 @@ fn remove_agents_md_world_state_section(rollout_path: &Path) -> Result<()> {
|
||||
let mut removed_section = false;
|
||||
let retained = rollout
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.map(|mut line| {
|
||||
|
||||
@@ -2,7 +2,6 @@ use anyhow::Result;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_models_manager::model_info::model_info_from_slug;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
@@ -1037,7 +1036,7 @@ async fn cold_resume_refreshes_legacy_collaboration_snapshot_once(
|
||||
let legacy_rollout = std::fs::read_to_string(&rollout_path)?
|
||||
.lines()
|
||||
.map(|original_line| {
|
||||
let mut line = serde_json::from_str::<RolloutLine>(original_line)?;
|
||||
let mut line = codex_rollout::parse_rollout_line(original_line)?;
|
||||
if let RolloutItem::WorldState(world_state) = &mut line.item
|
||||
&& let Some(snapshot) = world_state.state.get_mut("collaboration_mode")
|
||||
{
|
||||
|
||||
@@ -6,7 +6,6 @@ use codex_core::compact::SUMMARY_PREFIX;
|
||||
use codex_core::config::Config;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_model_provider_info::ModelProviderInfo;
|
||||
use codex_model_provider_info::built_in_model_providers;
|
||||
@@ -367,7 +366,7 @@ fn replacement_history_from_rollout(path: &Path) -> Result<Vec<Value>> {
|
||||
.map(str::trim)
|
||||
.filter(|line| !line.is_empty())
|
||||
{
|
||||
let entry: RolloutLine = serde_json::from_str(line)?;
|
||||
let entry = codex_rollout::parse_rollout_line(line)?;
|
||||
if let RolloutItem::Compacted(compacted) = entry.item
|
||||
&& let Some(items) = compacted.replacement_history
|
||||
{
|
||||
@@ -723,7 +722,7 @@ async fn summarize_context_three_requests_and_instructions(
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let Ok(entry): Result<RolloutLine, _> = serde_json::from_str(trimmed) else {
|
||||
let Ok(entry) = codex_rollout::parse_rollout_line(trimmed) else {
|
||||
continue;
|
||||
};
|
||||
match entry.item {
|
||||
@@ -1049,7 +1048,7 @@ async fn manual_compact_records_durable_and_local_token_usage() {
|
||||
let rollout_items = fs::read_to_string(rollout_path)
|
||||
.expect("read rollout")
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.map(|line| line.item)
|
||||
.collect::<Vec<_>>();
|
||||
let records = rollout_items
|
||||
@@ -3430,7 +3429,7 @@ async fn pre_sampling_compact_recovers_comp_hash_after_resume() {
|
||||
let rollout = fs::read_to_string(&rollout_path).expect("read rollout");
|
||||
let persisted_comp_hash = rollout
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.find_map(|line| match line.item {
|
||||
RolloutItem::TurnContext(context) => context.comp_hash,
|
||||
_ => None,
|
||||
@@ -3713,7 +3712,7 @@ async fn auto_compact_persists_rollout_entries() {
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let Ok(entry): Result<RolloutLine, _> = serde_json::from_str(trimmed) else {
|
||||
let Ok(entry) = codex_rollout::parse_rollout_line(trimmed) else {
|
||||
continue;
|
||||
};
|
||||
match entry.item {
|
||||
|
||||
@@ -16,7 +16,6 @@ use codex_features::Feature;
|
||||
use codex_history::CodexHarnessMetadata;
|
||||
use codex_history::InitialHistory;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_login::auth::AgentIdentityAuth;
|
||||
use codex_login::auth::AgentIdentityAuthRecord;
|
||||
@@ -404,7 +403,7 @@ fn annotate_retained_user_in_rollout(path: &Path, retained_text: &str) -> Result
|
||||
let mut rollout = fs::read_to_string(path)?
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
rollout
|
||||
.iter_mut()
|
||||
@@ -432,7 +431,7 @@ fn assert_compacted_user_metadata(path: &Path, retained_text: &str) -> Result<()
|
||||
let replacement_history = fs::read_to_string(path)?
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.rev()
|
||||
@@ -669,7 +668,7 @@ async fn remote_compact_v2_retains_only_client_developer_messages_when_enabled(
|
||||
codex.shutdown_and_wait().await?;
|
||||
let replacement_history = fs::read_to_string(&rollout_path)?
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.filter_map(|line| match line.item {
|
||||
RolloutItem::Compacted(compacted) => compacted.replacement_history,
|
||||
_ => None,
|
||||
@@ -738,7 +737,7 @@ async fn remote_compact_v2_records_usage_before_output_validation() -> Result<()
|
||||
|
||||
let record = fs::read_to_string(&rollout_path)?
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.filter_map(|line| match line.item {
|
||||
RolloutItem::TokenUsageRecord(record) => Some(record),
|
||||
_ => None,
|
||||
@@ -3349,7 +3348,7 @@ async fn remote_compact_persists_replacement_history_in_rollout() -> Result<()>
|
||||
.map(str::trim)
|
||||
.filter(|l| !l.is_empty())
|
||||
{
|
||||
let Ok(entry) = serde_json::from_str::<RolloutLine>(line) else {
|
||||
let Ok(entry) = codex_rollout::parse_rollout_line(line) else {
|
||||
continue;
|
||||
};
|
||||
if let RolloutItem::Compacted(compacted) = entry.item
|
||||
|
||||
@@ -7,7 +7,6 @@ use std::path::PathBuf;
|
||||
use anyhow::Result;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_protocol::config_types::ServiceTier;
|
||||
use codex_protocol::models::PermissionProfile;
|
||||
@@ -791,7 +790,7 @@ fn replacement_history_from_rollout(path: &Path) -> Result<Value> {
|
||||
.map(str::trim)
|
||||
.filter(|line| !line.is_empty())
|
||||
{
|
||||
let Ok(entry) = serde_json::from_str::<RolloutLine>(line) else {
|
||||
let Ok(entry) = codex_rollout::parse_rollout_line(line) else {
|
||||
continue;
|
||||
};
|
||||
if let RolloutItem::Compacted(compacted) = entry.item
|
||||
|
||||
@@ -18,7 +18,6 @@ use codex_core::config::Config;
|
||||
use codex_core::spawn::CODEX_SANDBOX_NETWORK_DISABLED_ENV_VAR;
|
||||
use codex_history::CodexHarnessMetadata;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
use codex_protocol::config_types::ModeKind;
|
||||
use codex_protocol::config_types::Settings;
|
||||
@@ -94,7 +93,7 @@ fn seed_first_checkpoint_harness_metadata(path: &Path, retained_text: &str) -> R
|
||||
let mut lines = std::fs::read_to_string(path)?
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<Result<Vec<_>, _>>()?;
|
||||
let replacement_history = lines
|
||||
.iter_mut()
|
||||
@@ -127,7 +126,7 @@ fn assert_latest_checkpoint_retains_harness_metadata(
|
||||
let replacement_history = std::fs::read_to_string(path)?
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.rev()
|
||||
|
||||
@@ -7,7 +7,6 @@ use codex_core::parse_turn_item;
|
||||
use codex_history::InitialHistory;
|
||||
use codex_history::ResumedHistory;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::ThreadId;
|
||||
use codex_protocol::items::TurnItem;
|
||||
use codex_protocol::mcp::ClientMcpExtensions;
|
||||
@@ -311,7 +310,7 @@ fn read_rollout_items(path: &std::path::Path) -> Vec<RolloutItem> {
|
||||
let parse_json_message = format!("failed to parse rollout JSON line `{line}`");
|
||||
let v: serde_json::Value = serde_json::from_str(line).expect(&parse_json_message);
|
||||
let parse_line_message = format!("failed to parse rollout line `{line}`");
|
||||
let rl: RolloutLine = serde_json::from_value(v).expect(&parse_line_message);
|
||||
let rl = codex_rollout::decode_rollout_line(v).expect(&parse_line_message);
|
||||
match rl.item {
|
||||
RolloutItem::SessionMeta(_) => {}
|
||||
other => items.push(other),
|
||||
|
||||
@@ -22,7 +22,6 @@ use codex_extension_api::ToolStartInput;
|
||||
use codex_features::CurrentTimeSource;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_protocol::ThreadId;
|
||||
use codex_protocol::config_types::ApprovalsReviewer;
|
||||
@@ -618,7 +617,7 @@ async fn guardian_session_prewarms_and_is_reused_for_first_review(
|
||||
test.codex.shutdown_and_wait().await?;
|
||||
let guardian_rollout = fs::read_to_string(guardian_rollout_path)?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?;
|
||||
assert_eq!(
|
||||
guardian_rollout.iter().find_map(|line| match &line.item {
|
||||
|
||||
@@ -16,7 +16,6 @@ use codex_core::config::Constrained;
|
||||
use codex_core::config::ThreadStoreConfig;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_model_provider_info::ModelProviderInfo;
|
||||
use codex_model_provider_info::built_in_model_providers;
|
||||
use codex_plugin::PluginHookSource;
|
||||
@@ -1145,7 +1144,7 @@ fn rollout_hook_prompt_texts(text: &str) -> Result<Vec<String>> {
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let rollout: RolloutLine = serde_json::from_str(trimmed).context("parse rollout line")?;
|
||||
let rollout = codex_rollout::parse_rollout_line(trimmed).context("parse rollout line")?;
|
||||
if let RolloutItem::ResponseItem(envelope) = rollout.item
|
||||
&& let ResponseItem::Message { role, content, .. } = envelope.item
|
||||
&& role == "user"
|
||||
|
||||
@@ -4,7 +4,6 @@ use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
use codex_protocol::config_types::ModeKind;
|
||||
use codex_protocol::config_types::Settings;
|
||||
@@ -50,7 +49,7 @@ fn find_user_message_with_image(text: &str) -> Option<ResponseItem> {
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let rollout: RolloutLine = match serde_json::from_str(trimmed) {
|
||||
let rollout = match codex_rollout::parse_rollout_line(trimmed) {
|
||||
Ok(rollout) => rollout,
|
||||
Err(_) => continue,
|
||||
};
|
||||
@@ -321,7 +320,7 @@ async fn resumed_history_only_emits_resize_notices_for_new_images() -> anyhow::R
|
||||
|
||||
let mut rollout_lines = fs::read_to_string(&rollout_path)?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?;
|
||||
let historical_content = rollout_lines
|
||||
.iter_mut()
|
||||
@@ -501,7 +500,7 @@ async fn resumed_history_only_emits_resize_notices_for_new_images() -> anyhow::R
|
||||
.lines()
|
||||
.skip(existing_rollout_lines)
|
||||
.filter_map(|line| {
|
||||
let RolloutItem::ResponseItem(envelope) = serde_json::from_str::<RolloutLine>(line)
|
||||
let RolloutItem::ResponseItem(envelope) = codex_rollout::parse_rollout_line(line)
|
||||
.expect("new rollout line should deserialize")
|
||||
.item
|
||||
else {
|
||||
|
||||
@@ -4,7 +4,6 @@ use anyhow::Ok;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
use codex_protocol::config_types::ModeKind;
|
||||
use codex_protocol::config_types::Settings;
|
||||
@@ -361,7 +360,7 @@ async fn web_search_item_is_emitted() -> anyhow::Result<()> {
|
||||
let rollout = std::fs::read_to_string(rollout_path)?;
|
||||
let persisted_completion = rollout
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.find_map(|line| match line.item {
|
||||
|
||||
@@ -13,7 +13,6 @@ use codex_extension_api::ThreadLifecycleContributor;
|
||||
use codex_extension_api::ThreadStartInput;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_mcp::CODEX_APPS_MCP_SERVER_NAME;
|
||||
use codex_mcp::McpResourceClient;
|
||||
use codex_protocol::capabilities::CapabilityRootLocation;
|
||||
@@ -1003,7 +1002,7 @@ async fn initially_empty_deferred_tool_world_state_is_not_rendered_or_persisted(
|
||||
let world_states = tokio::fs::read_to_string(rollout_path)
|
||||
.await?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
@@ -1046,7 +1045,7 @@ async fn deferred_tool_world_state_survives_resume_without_duplicate_updates() -
|
||||
let persisted_tools = tokio::fs::read_to_string(&rollout_path)
|
||||
.await?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
|
||||
@@ -8,7 +8,6 @@ use codex_core::TurnInputRequest;
|
||||
use codex_core::config::Config;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::approvals::ElicitationRequest;
|
||||
use codex_protocol::config_types::ApprovalsReviewer;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
@@ -1229,7 +1228,7 @@ async fn apps_default_writes_prompts_for_writes_but_not_reads() -> Result<()> {
|
||||
let persisted_hints = tokio::fs::read_to_string(rollout_path)
|
||||
.await?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
|
||||
@@ -6,7 +6,6 @@ use codex_core::TurnInputRequest;
|
||||
use codex_core::config::Constrained;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_models_manager::bundled_models_response;
|
||||
use codex_models_manager::manager::RefreshStrategy;
|
||||
@@ -442,7 +441,7 @@ async fn model_change_appends_model_instructions_developer_message() -> Result<(
|
||||
let rollout_path = test.codex.rollout_path().expect("rollout path");
|
||||
let model_states = std::fs::read_to_string(rollout_path)?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
|
||||
@@ -12,7 +12,6 @@ use codex_extension_items::ExtensionItem;
|
||||
use codex_extension_items::sleep::SleepItem;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_protocol::AgentPath;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
@@ -721,7 +720,7 @@ async fn any_new_input_interrupts_sleep() {
|
||||
.expect("read rollout");
|
||||
let persisted_sleep_items = rollout
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.filter_map(|line| match line.item {
|
||||
RolloutItem::EventMsg(EventMsg::ItemCompleted(event)) => match event.item {
|
||||
TurnItem::Extension(ExtensionItem::Sleep(item)) => Some(item),
|
||||
|
||||
@@ -38,7 +38,6 @@ use codex_extension_api::WorldStateContributionInput;
|
||||
use codex_extension_api::WorldStateSectionContribution;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_http_client::HttpClientFactory;
|
||||
use codex_http_client::OutboundProxyPolicy;
|
||||
use codex_network_proxy::NetworkProxyConfig;
|
||||
@@ -1078,7 +1077,7 @@ async fn deferred_executor_promotes_primary_environment_when_startup_completes()
|
||||
let rollout = fs::read_to_string(test.codex.rollout_path().context("rollout path")?)?;
|
||||
let world_state_patch = rollout
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
@@ -2766,7 +2765,7 @@ async fn deferred_executor_compaction_preserves_then_updates_environment_once()
|
||||
let rollout = fs::read_to_string(rollout_path)?;
|
||||
let world_state_items = rollout
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<serde_json::Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.filter_map(|line| match line.item {
|
||||
|
||||
@@ -7,7 +7,6 @@ use codex_core::find_thread_path_by_id_str;
|
||||
use codex_exec_server::CreateDirectoryOptions;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_protocol::config_types::ApprovalsReviewer;
|
||||
use codex_protocol::config_types::CollaborationMode;
|
||||
@@ -209,7 +208,7 @@ async fn review_op_emits_lifecycle_and_review_output() {
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.find_map(|line| {
|
||||
let rollout_line: RolloutLine = serde_json::from_str(line).expect("rollout line");
|
||||
let rollout_line = codex_rollout::parse_rollout_line(line).expect("rollout line");
|
||||
match rollout_line.item {
|
||||
RolloutItem::SessionMeta(session_meta) => Some(session_meta.meta.id.to_string()),
|
||||
_ => None,
|
||||
@@ -251,7 +250,7 @@ async fn review_op_emits_lifecycle_and_review_output() {
|
||||
continue;
|
||||
}
|
||||
let v: serde_json::Value = serde_json::from_str(line).expect("jsonl line");
|
||||
let rl: RolloutLine = serde_json::from_value(v).expect("rollout line");
|
||||
let rl = codex_rollout::decode_rollout_line(v).expect("rollout line");
|
||||
if let RolloutItem::ResponseItem(envelope) = rl.item
|
||||
&& let ResponseItem::Message { role, content, .. } = envelope.item
|
||||
{
|
||||
@@ -763,8 +762,8 @@ async fn review_uses_updated_turn_permissions_and_approval_policy() {
|
||||
let review_session_cwd = review_rollout
|
||||
.lines()
|
||||
.find_map(|line| {
|
||||
let rollout_line: RolloutLine =
|
||||
serde_json::from_str(line).expect("review rollout line should be valid");
|
||||
let rollout_line = codex_rollout::parse_rollout_line(line)
|
||||
.expect("review rollout line should be valid");
|
||||
match rollout_line.item {
|
||||
RolloutItem::SessionMeta(session_meta) => Some(session_meta.meta.cwd),
|
||||
_ => None,
|
||||
@@ -775,8 +774,8 @@ async fn review_uses_updated_turn_permissions_and_approval_policy() {
|
||||
let review_context = review_rollout
|
||||
.lines()
|
||||
.filter_map(|line| {
|
||||
let rollout_line: RolloutLine =
|
||||
serde_json::from_str(line).expect("review rollout line should be valid");
|
||||
let rollout_line = codex_rollout::parse_rollout_line(line)
|
||||
.expect("review rollout line should be valid");
|
||||
match rollout_line.item {
|
||||
RolloutItem::TurnContext(turn_context) => Some(turn_context),
|
||||
_ => None,
|
||||
@@ -1229,7 +1228,7 @@ async fn review_input_isolated_from_parent_history() {
|
||||
continue;
|
||||
}
|
||||
let v: serde_json::Value = serde_json::from_str(line).expect("jsonl line");
|
||||
let rl: RolloutLine = serde_json::from_value(v).expect("rollout line");
|
||||
let rl = codex_rollout::decode_rollout_line(v).expect("rollout line");
|
||||
if let RolloutItem::ResponseItem(envelope) = rl.item
|
||||
&& let ResponseItem::Message { role, content, .. } = envelope.item
|
||||
&& role == "user"
|
||||
|
||||
@@ -241,7 +241,7 @@ async fn compaction_checkpoints_settings_changed_during_its_model_request() -> R
|
||||
let rollout_path = test.session_configured.rollout_path.expect("rollout path");
|
||||
let rollout: Vec<RolloutLine> = std::fs::read_to_string(rollout_path)?
|
||||
.lines()
|
||||
.map(serde_json::from_str)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<std::result::Result<_, _>>()?;
|
||||
let checkpoint = rollout
|
||||
.iter()
|
||||
|
||||
@@ -6,7 +6,6 @@ use codex_core::config::AgentRoleConfig;
|
||||
use codex_core::config::Config;
|
||||
use codex_features::Feature;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_login::CodexAuth;
|
||||
use codex_models_manager::manager::RefreshStrategy;
|
||||
use codex_models_manager::manager::SharedModelsManager;
|
||||
@@ -454,7 +453,7 @@ async fn multi_agent_v2_cold_resume_refreshes_legacy_usage_hints_once(
|
||||
let mut removed_usage_hint_presence = false;
|
||||
let legacy_rollout = std::fs::read_to_string(&rollout_path)?
|
||||
.lines()
|
||||
.map(serde_json::from_str::<RolloutLine>)
|
||||
.map(codex_rollout::parse_rollout_line)
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?
|
||||
.into_iter()
|
||||
.map(|mut line| {
|
||||
|
||||
@@ -2,7 +2,6 @@
|
||||
|
||||
use anyhow::Result;
|
||||
use codex_history::RolloutItem;
|
||||
use codex_history::RolloutLine;
|
||||
use codex_protocol::SessionId;
|
||||
use codex_protocol::protocol::TokenUsageRecord;
|
||||
use core_test_support::responses::ev_assistant_message;
|
||||
@@ -21,7 +20,7 @@ fn token_usage_records(path: &std::path::Path) -> Vec<TokenUsageRecord> {
|
||||
std::fs::read_to_string(path)
|
||||
.expect("read rollout")
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| codex_rollout::parse_rollout_line(line).ok())
|
||||
.filter_map(|line| match line.item {
|
||||
RolloutItem::TokenUsageRecord(record) => Some(record),
|
||||
_ => None,
|
||||
|
||||
@@ -1567,7 +1567,7 @@ async fn parse_latest_turn_context_cwd(path: &Path) -> Option<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let reader = codex_rollout::open_rollout_seekable_reader(&path).ok()?;
|
||||
let mut scanner = codex_rollout::ReverseJsonlScanner::new(reader).ok()?;
|
||||
while let Some(outcome) = scanner.scan_next::<RolloutLine>().ok()? {
|
||||
while let Some(outcome) = scanner.scan_next_rollout_line().ok()? {
|
||||
if let codex_rollout::ScanOutcome::Parsed(RolloutLine {
|
||||
item: RolloutItem::TurnContext(item),
|
||||
..
|
||||
|
||||
@@ -232,7 +232,11 @@ impl From<CompactedItem> for ResponseItem {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Serialize, Deserialize, Clone, JsonSchema)]
|
||||
/// One persisted rollout JSONL record.
|
||||
///
|
||||
/// This intentionally does not implement Deserialize: JSONL readers must use
|
||||
/// codex_rollout's canonical parser so nested decimal values survive the flattened envelope.
|
||||
#[derive(Serialize, Clone, JsonSchema)]
|
||||
pub struct RolloutLine {
|
||||
pub timestamp: String,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
|
||||
@@ -53,7 +53,11 @@ fn response_item_rollout_line_preserves_shape() -> Result<()> {
|
||||
},
|
||||
});
|
||||
|
||||
let line = serde_json::from_value::<RolloutLine>(legacy_line.clone())?;
|
||||
let line = RolloutLine {
|
||||
timestamp: "2025-01-03T12:00:00.000Z".to_string(),
|
||||
ordinal: Some(7),
|
||||
item: serde_json::from_value(legacy_line.clone())?,
|
||||
};
|
||||
let RolloutItem::ResponseItem(envelope) = &line.item else {
|
||||
panic!("expected response item");
|
||||
};
|
||||
@@ -93,8 +97,8 @@ fn response_item_envelope_stores_metadata_beside_rollout_payload() -> Result<()>
|
||||
);
|
||||
assert_eq!(serialized["payload"].get("metadata"), None);
|
||||
|
||||
let restored = serde_json::from_value::<RolloutLine>(serialized)?;
|
||||
let RolloutItem::ResponseItem(envelope) = restored.item else {
|
||||
let restored = serde_json::from_value(serialized)?;
|
||||
let RolloutItem::ResponseItem(envelope) = restored else {
|
||||
panic!("expected response item");
|
||||
};
|
||||
assert_eq!(
|
||||
@@ -157,7 +161,7 @@ fn response_item_envelope_preserves_harness_authored_configuration_provenance()
|
||||
#[test]
|
||||
/// Keeps future metadata fields from making older binaries reject persisted items.
|
||||
fn response_item_envelope_ignores_unknown_harness_metadata_fields() -> Result<()> {
|
||||
let line = serde_json::from_value::<RolloutLine>(json!({
|
||||
let line = serde_json::from_value(json!({
|
||||
"timestamp": "2025-01-03T12:00:00.000Z",
|
||||
"ordinal": 7,
|
||||
"type": "response_item",
|
||||
@@ -174,7 +178,7 @@ fn response_item_envelope_ignores_unknown_harness_metadata_fields() -> Result<()
|
||||
},
|
||||
}))?;
|
||||
|
||||
let RolloutItem::ResponseItem(envelope) = line.item else {
|
||||
let RolloutItem::ResponseItem(envelope) = line else {
|
||||
panic!("expected response item");
|
||||
};
|
||||
assert_eq!(envelope.metadata, Some(CodexHarnessMetadata::default()));
|
||||
|
||||
@@ -46,7 +46,9 @@ pub(crate) use codex_protocol::protocol;
|
||||
/// Remove it once Serde supports format-specific buffering.
|
||||
pub fn decode_rollout_line(value: Value) -> serde_json::Result<RolloutLine> {
|
||||
let Value::Object(mut fields) = value else {
|
||||
return serde_json::from_value(value);
|
||||
return Err(serde_json::Error::custom(
|
||||
"rollout line must be a JSON object",
|
||||
));
|
||||
};
|
||||
let timestamp = fields
|
||||
.remove("timestamp")
|
||||
@@ -66,6 +68,16 @@ pub fn decode_rollout_line(value: Value) -> serde_json::Result<RolloutLine> {
|
||||
})
|
||||
}
|
||||
|
||||
/// Parses a persisted JSONL rollout record through the canonical JSON decoder.
|
||||
pub fn parse_rollout_line(line: &str) -> serde_json::Result<RolloutLine> {
|
||||
serde_json::from_str::<Value>(line).and_then(decode_rollout_line)
|
||||
}
|
||||
|
||||
/// Parses persisted JSONL rollout record bytes through the canonical JSON decoder.
|
||||
pub fn parse_rollout_line_bytes(bytes: &[u8]) -> serde_json::Result<RolloutLine> {
|
||||
serde_json::from_slice::<Value>(bytes).and_then(decode_rollout_line)
|
||||
}
|
||||
|
||||
pub const SESSIONS_SUBDIR: &str = "sessions";
|
||||
pub const ARCHIVED_SESSIONS_SUBDIR: &str = "archived_sessions";
|
||||
pub static INTERACTIVE_SESSION_SOURCES: LazyLock<Vec<SessionSource>> = LazyLock::new(|| {
|
||||
|
||||
@@ -20,7 +20,6 @@ use super::SESSIONS_SUBDIR;
|
||||
use super::compression;
|
||||
use super::rollout_file_name::RolloutFileName;
|
||||
use crate::RolloutItem;
|
||||
use crate::RolloutLine;
|
||||
use crate::protocol::EventMsg;
|
||||
use crate::state_db;
|
||||
use codex_file_search as file_search;
|
||||
@@ -1128,7 +1127,7 @@ async fn read_head_summary(path: &Path, head_limit: usize) -> io::Result<HeadTai
|
||||
}
|
||||
lines_scanned += 1;
|
||||
|
||||
let parsed: Result<RolloutLine, _> = serde_json::from_str(trimmed);
|
||||
let parsed = crate::parse_rollout_line(trimmed);
|
||||
let rollout_line = match parsed {
|
||||
Ok(rollout_line) => rollout_line,
|
||||
Err(_) => {
|
||||
@@ -1241,7 +1240,7 @@ pub async fn read_head_for_summary(path: &Path) -> io::Result<Vec<serde_json::Va
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
if let Ok(rollout_line) = serde_json::from_str::<RolloutLine>(trimmed) {
|
||||
if let Ok(rollout_line) = crate::parse_rollout_line(trimmed) {
|
||||
match rollout_line.item {
|
||||
RolloutItem::SessionMeta(session_meta_line) => {
|
||||
if let Ok(value) = serde_json::to_value(session_meta_line) {
|
||||
@@ -1300,7 +1299,7 @@ pub async fn read_session_meta_line(path: &Path) -> io::Result<SessionMetaLine>
|
||||
if trimmed.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let Ok(rollout_line) = serde_json::from_str::<RolloutLine>(trimmed) else {
|
||||
let Ok(rollout_line) = crate::parse_rollout_line(trimmed) else {
|
||||
if let Ok(value) = serde_json::from_str::<serde_json::Value>(trimmed) {
|
||||
crate::recorder::reject_unknown_thread_history_mode(&value)?;
|
||||
}
|
||||
|
||||
@@ -10,7 +10,6 @@ use codex_protocol::protocol::HistoryPosition;
|
||||
use codex_protocol::protocol::ThreadHistoryMode;
|
||||
|
||||
use crate::RolloutItem;
|
||||
use crate::RolloutLine;
|
||||
use crate::reverse_jsonl_scanner::ReverseJsonlScanner;
|
||||
use crate::reverse_jsonl_scanner::ScanOutcome;
|
||||
|
||||
@@ -67,7 +66,7 @@ pub(crate) fn ordinal_state_for_rollout(
|
||||
|
||||
let mut scanner = ReverseJsonlScanner::new(file)?;
|
||||
let record = loop {
|
||||
match scanner.scan_next::<RolloutLine>()? {
|
||||
match scanner.scan_next_rollout_line()? {
|
||||
Some(ScanOutcome::Parsed(record)) => break record,
|
||||
Some(ScanOutcome::Rejected(_)) => continue,
|
||||
None => {
|
||||
@@ -111,7 +110,7 @@ fn read_history_metadata(
|
||||
if line.trim().is_empty() {
|
||||
continue;
|
||||
}
|
||||
let record: RolloutLine = serde_json::from_str(line.as_str()).map_err(|error| {
|
||||
let record = crate::parse_rollout_line(line.as_str()).map_err(|error| {
|
||||
io::Error::other(format!(
|
||||
"failed to parse first rollout record at {}: {error}",
|
||||
path.display()
|
||||
|
||||
@@ -14,11 +14,14 @@ use codex_protocol::protocol::AgentMessageEvent;
|
||||
use codex_protocol::protocol::AskForApproval;
|
||||
use codex_protocol::protocol::EventMsg;
|
||||
use codex_protocol::protocol::HistoryPosition;
|
||||
use codex_protocol::protocol::RateLimitSnapshot;
|
||||
use codex_protocol::protocol::RateLimitWindow;
|
||||
use codex_protocol::protocol::SandboxPolicy;
|
||||
use codex_protocol::protocol::SessionMeta;
|
||||
use codex_protocol::protocol::SessionMetaLine;
|
||||
use codex_protocol::protocol::SessionSource;
|
||||
use codex_protocol::protocol::ThreadHistoryMode;
|
||||
use codex_protocol::protocol::TokenCountEvent;
|
||||
use codex_protocol::protocol::TurnContextItem;
|
||||
use codex_protocol::protocol::UserMessageEvent;
|
||||
use codex_protocol::security_risk::SecurityRiskScore;
|
||||
@@ -103,7 +106,7 @@ fn read_rollout_lines(path: &Path) -> std::io::Result<Vec<RolloutLine>> {
|
||||
fs::read_to_string(path)?
|
||||
.lines()
|
||||
.filter(|line| !line.trim().is_empty())
|
||||
.map(|line| serde_json::from_str(line).map_err(std::io::Error::other))
|
||||
.map(|line| crate::parse_rollout_line(line).map_err(std::io::Error::other))
|
||||
.collect()
|
||||
}
|
||||
|
||||
@@ -687,7 +690,7 @@ async fn recorder_materializes_on_flush_with_pending_items() -> std::io::Result<
|
||||
vec![Some(0), Some(1), Some(2)]
|
||||
);
|
||||
let first_line = text.lines().next().expect("session metadata line");
|
||||
let session_meta: RolloutLine = serde_json::from_str(first_line)?;
|
||||
let session_meta = crate::parse_rollout_line(first_line)?;
|
||||
let RolloutItem::SessionMeta(session_meta) = session_meta.item else {
|
||||
panic!("expected session metadata in rollout");
|
||||
};
|
||||
@@ -995,6 +998,68 @@ async fn resumed_paginated_rollout_continues_after_ordinal_gap() -> std::io::Res
|
||||
recorder.shutdown().await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumed_paginated_rollout_continues_after_decimal_token_count() -> std::io::Result<()> {
|
||||
let home = TempDir::new().expect("temp dir");
|
||||
let config = test_config(home.path());
|
||||
let thread_id = ThreadId::new();
|
||||
let recorder = RolloutRecorder::new(
|
||||
&config,
|
||||
RolloutRecorderParams::new(
|
||||
thread_id,
|
||||
/*forked_from_id*/ None,
|
||||
/*parent_thread_id*/ None,
|
||||
SessionSource::Exec,
|
||||
/*thread_source*/ None,
|
||||
"test_originator".to_string(),
|
||||
BaseInstructions::default(),
|
||||
Vec::new(),
|
||||
)
|
||||
.with_history_mode(ThreadHistoryMode::Paginated),
|
||||
)
|
||||
.await?;
|
||||
let rollout_path = recorder.rollout_path().to_path_buf();
|
||||
recorder
|
||||
.record_canonical_items(&[RolloutItem::EventMsg(EventMsg::TokenCount(
|
||||
TokenCountEvent {
|
||||
info: None,
|
||||
rate_limits: Some(RateLimitSnapshot {
|
||||
limit_id: None,
|
||||
limit_name: None,
|
||||
primary: Some(RateLimitWindow {
|
||||
used_percent: 0.0,
|
||||
window_minutes: Some(60),
|
||||
resets_at: Some(1_800_000_000),
|
||||
}),
|
||||
secondary: None,
|
||||
credits: None,
|
||||
individual_limit: None,
|
||||
spend_control_reached: None,
|
||||
plan_type: None,
|
||||
rate_limit_reached_type: None,
|
||||
normal_model_slug: None,
|
||||
}),
|
||||
},
|
||||
))])
|
||||
.await?;
|
||||
recorder.persist().await?;
|
||||
recorder.shutdown().await?;
|
||||
|
||||
let resumed =
|
||||
RolloutRecorder::new(&config, RolloutRecorderParams::resume(rollout_path.clone())).await?;
|
||||
resumed
|
||||
.record_canonical_items(&[agent_message_item("after-resume")])
|
||||
.await?;
|
||||
resumed.shutdown().await?;
|
||||
|
||||
let ordinals = read_rollout_lines(&rollout_path)?
|
||||
.into_iter()
|
||||
.map(|line| line.ordinal)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(ordinals, vec![Some(0), Some(1), Some(2)]);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumed_paginated_rollout_repairs_unsafe_tail() -> std::io::Result<()> {
|
||||
let valid_unterminated = serde_json::to_string(&RolloutLine {
|
||||
@@ -1034,7 +1099,7 @@ async fn resumed_paginated_rollout_repairs_unsafe_tail() -> std::io::Result<()>
|
||||
assert!(contents.ends_with('\n'), "{name} tail should be terminated");
|
||||
let ordinals = contents
|
||||
.lines()
|
||||
.filter_map(|line| serde_json::from_str::<RolloutLine>(line).ok())
|
||||
.filter_map(|line| crate::parse_rollout_line(line).ok())
|
||||
.map(|line| line.ordinal)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
|
||||
@@ -4,6 +4,9 @@ use std::io::Seek;
|
||||
use std::io::SeekFrom;
|
||||
|
||||
use serde::de::DeserializeOwned;
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::RolloutLine;
|
||||
|
||||
const READ_CHUNK_SIZE: usize = 64 * 1024;
|
||||
|
||||
@@ -128,6 +131,17 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
/// Scans the next rollout record through the canonical persisted JSON decoder.
|
||||
pub fn scan_next_rollout_line(&mut self) -> io::Result<Option<ScanOutcome<RolloutLine>>> {
|
||||
Ok(self.scan_next::<Value>()?.map(|outcome| match outcome {
|
||||
ScanOutcome::Parsed(value) => match crate::decode_rollout_line(value) {
|
||||
Ok(line) => ScanOutcome::Parsed(line),
|
||||
Err(error) => ScanOutcome::Rejected(error),
|
||||
},
|
||||
ScanOutcome::Rejected(error) => ScanOutcome::Rejected(error),
|
||||
}))
|
||||
}
|
||||
|
||||
fn finish_record<T>(&mut self) -> Option<ScanOutcome<T>>
|
||||
where
|
||||
T: DeserializeOwned,
|
||||
|
||||
@@ -18,7 +18,6 @@ use super::SESSIONS_SUBDIR;
|
||||
use super::compression;
|
||||
use crate::ResponseItemEnvelope;
|
||||
use crate::RolloutItem;
|
||||
use crate::RolloutLine;
|
||||
|
||||
const MATCH_CONTEXT_BEFORE_CHARS: usize = 48;
|
||||
const MATCH_CONTEXT_AFTER_CHARS: usize = 96;
|
||||
@@ -248,7 +247,7 @@ fn case_insensitive_literal_regex(search_term: impl AsRef<str>) -> io::Result<Re
|
||||
}
|
||||
|
||||
fn content_match_snippet(jsonl_line: &str, search_term: &Regex) -> Option<String> {
|
||||
let rollout_line = serde_json::from_str::<RolloutLine>(jsonl_line.trim()).ok()?;
|
||||
let rollout_line = crate::parse_rollout_line(jsonl_line.trim()).ok()?;
|
||||
let text = conversation_text_from_item(&rollout_line.item)?;
|
||||
excerpt_around_match(text.as_str(), search_term)
|
||||
}
|
||||
|
||||
@@ -59,7 +59,7 @@ fn rollout_line_decoder_preserves_canonical_json_compatibility() -> Result<()> {
|
||||
|
||||
for encoded in cases {
|
||||
let value = serde_json::from_str::<serde_json::Value>(encoded)?;
|
||||
let decoded = crate::decode_rollout_line(value.clone())?;
|
||||
let decoded = crate::parse_rollout_line(encoded)?;
|
||||
let mut expected = value;
|
||||
if expected["type"] != "response_item" {
|
||||
expected
|
||||
|
||||
@@ -7,7 +7,6 @@ use codex_rollout::ModelContextScan;
|
||||
use codex_rollout::ModelContextScanProgress;
|
||||
use codex_rollout::ReverseJsonlScanner;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use codex_rollout::ScanOutcome;
|
||||
|
||||
use super::LocalThreadStore;
|
||||
@@ -137,7 +136,7 @@ fn scan_model_context_from_lineage_blocking(
|
||||
Some(end_byte_offset) => ReverseJsonlScanner::new_at(file, end_byte_offset)?,
|
||||
None => ReverseJsonlScanner::new(file)?,
|
||||
};
|
||||
while let Some(outcome) = scanner.scan_next::<RolloutLine>()? {
|
||||
while let Some(outcome) = scanner.scan_next_rollout_line()? {
|
||||
let ScanOutcome::Parsed(line) = outcome else {
|
||||
continue;
|
||||
};
|
||||
|
||||
@@ -599,8 +599,8 @@ fn rollout_end_byte_offset(path: &Path, end_ordinal_exclusive: u64) -> u64 {
|
||||
let contents = std::fs::read(path).expect("read rollout");
|
||||
let mut byte_offset = 0_u64;
|
||||
for line in contents.split_inclusive(|byte| *byte == b'\n') {
|
||||
let parsed: RolloutLine =
|
||||
serde_json::from_slice(line).expect("parse rollout line for byte offset");
|
||||
let parsed = codex_rollout::parse_rollout_line_bytes(line)
|
||||
.expect("parse rollout line for byte offset");
|
||||
if parsed.ordinal == Some(end_ordinal_exclusive) {
|
||||
return byte_offset;
|
||||
}
|
||||
|
||||
@@ -327,7 +327,7 @@ fn rollout_end_byte_offset(path: &Path, end_ordinal_exclusive: u64) -> u64 {
|
||||
let end_byte_offset = bytes
|
||||
.split_inclusive(|byte| *byte == b'\n')
|
||||
.take_while(|line| {
|
||||
serde_json::from_slice::<RolloutLine>(line)
|
||||
codex_rollout::parse_rollout_line_bytes(line)
|
||||
.expect("parse rollout fixture")
|
||||
.ordinal
|
||||
.expect("paginated rollout ordinal")
|
||||
|
||||
@@ -39,7 +39,7 @@ pub(super) fn parse_legacy_rollout_value(mut value: Value) -> Result<Option<Roll
|
||||
normalize_legacy_rate_limit_resets(&mut value);
|
||||
normalize_legacy_review_entry(&mut value);
|
||||
normalize_legacy_command_cwd(&mut value)?;
|
||||
serde_json::from_value(value)
|
||||
codex_rollout::decode_rollout_line(value)
|
||||
.map(Some)
|
||||
.map_err(|error| error.to_string())
|
||||
}
|
||||
|
||||
@@ -23,12 +23,12 @@ fn parses_legacy_numeric_event_payloads_through_value() {
|
||||
"info": null,
|
||||
"rate_limits": {
|
||||
"primary": {
|
||||
"used_percent": 2,
|
||||
"used_percent": 0.0,
|
||||
"window_minutes": 300,
|
||||
"resets_at": 1_770_414_841,
|
||||
},
|
||||
"secondary": {
|
||||
"used_percent": 12,
|
||||
"used_percent": 12.5,
|
||||
"window_minutes": 10_080,
|
||||
"resets_at": 1_770_698_702,
|
||||
},
|
||||
|
||||
@@ -23,7 +23,6 @@ use std::os::unix::fs::PermissionsExt;
|
||||
|
||||
use codex_protocol::ThreadId;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use tokio::fs::File;
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
@@ -170,7 +169,7 @@ pub(super) async fn rewrite_subagent_history_boundary(
|
||||
.read_until(b'\n', &mut head_bytes)
|
||||
.await
|
||||
.map_err(migration_error)?;
|
||||
let mut head = serde_json::from_slice::<RolloutLine>(&head_bytes).map_err(migration_error)?;
|
||||
let mut head = codex_rollout::parse_rollout_line_bytes(&head_bytes).map_err(migration_error)?;
|
||||
let RolloutItem::SessionMeta(session_meta) = &mut head.item else {
|
||||
return Err(migration_error(
|
||||
"staged rollout head is not session metadata",
|
||||
|
||||
@@ -259,7 +259,7 @@ fn read_rollout(path: &Path) -> Vec<RolloutLine> {
|
||||
fs::read_to_string(path)
|
||||
.expect("read migrated rollout")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str(line).expect("parse migrated rollout"))
|
||||
.map(|line| codex_rollout::parse_rollout_line(line).expect("parse migrated rollout"))
|
||||
.collect()
|
||||
}
|
||||
|
||||
|
||||
@@ -1392,7 +1392,7 @@ fn rollout_end_byte_offset(path: &std::path::Path, end_ordinal_exclusive: u64) -
|
||||
let end_byte_offset = bytes
|
||||
.split_inclusive(|byte| *byte == b'\n')
|
||||
.take_while(|line| {
|
||||
serde_json::from_slice::<RolloutLine>(line)
|
||||
codex_rollout::parse_rollout_line_bytes(line)
|
||||
.expect("parse rollout fixture")
|
||||
.ordinal
|
||||
.expect("paginated rollout ordinal")
|
||||
|
||||
@@ -504,7 +504,7 @@ async fn paginated_realtime_items_materialize_separately_in_rollout_order() {
|
||||
let expected_rows = fs::read_to_string(rollout_path.as_path())
|
||||
.expect("read canonical rollout")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<RolloutLine>(line).expect("parse rollout line"))
|
||||
.map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line"))
|
||||
.filter_map(|line| match line.item {
|
||||
RolloutItem::RealtimeItem(item) => Some((
|
||||
item.id,
|
||||
@@ -2675,7 +2675,7 @@ fn rollout_line_byte_offsets(path: &std::path::Path, ordinal: u64) -> (i64, i64)
|
||||
let mut start_byte_offset = 0;
|
||||
for line in bytes.split_inclusive(|byte| *byte == b'\n') {
|
||||
let end_byte_offset = start_byte_offset + line.len();
|
||||
if serde_json::from_slice::<RolloutLine>(line)
|
||||
if codex_rollout::parse_rollout_line_bytes(line)
|
||||
.ok()
|
||||
.and_then(|line| line.ordinal)
|
||||
== Some(ordinal)
|
||||
|
||||
@@ -22,7 +22,6 @@ use codex_protocol::ThreadId;
|
||||
use codex_protocol::protocol::EventMsg;
|
||||
use codex_rollout::ReverseJsonlScanner;
|
||||
use codex_rollout::RolloutItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use codex_rollout::ScanOutcome;
|
||||
|
||||
const MAX_TRANSCRIPT_PREVIEW_LINES: usize = 6;
|
||||
@@ -162,7 +161,7 @@ fn scan_legacy_transcript_preview(
|
||||
)?
|
||||
.with_max_record_bytes(MAX_LEGACY_TRANSCRIPT_PREVIEW_SCAN_BYTES);
|
||||
loop {
|
||||
let outcome = match scanner.scan_next::<RolloutLine>() {
|
||||
let outcome = match scanner.scan_next_rollout_line() {
|
||||
Ok(Some(outcome)) => outcome,
|
||||
Ok(None) => break,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::UnexpectedEof => return Ok(None),
|
||||
|
||||
@@ -18,6 +18,7 @@ use codex_protocol::protocol::AgentMessageEvent;
|
||||
use codex_protocol::protocol::ThreadRolledBackEvent;
|
||||
use codex_protocol::protocol::UserMessageEvent;
|
||||
use codex_rollout::CompactedItem;
|
||||
use codex_rollout::RolloutLine;
|
||||
use core_test_support::responses;
|
||||
use pretty_assertions::assert_eq;
|
||||
use tempfile::tempdir;
|
||||
|
||||
Reference in New Issue
Block a user