diff --git a/crates/freshell-server/src/session_directory.rs b/crates/freshell-server/src/session_directory.rs
index 18490a793..95490a8ff 100644
--- a/crates/freshell-server/src/session_directory.rs
+++ b/crates/freshell-server/src/session_directory.rs
@@ -48,8 +48,10 @@ use axum::{
Json, Router,
};
use base64::Engine as _;
+use freshell_sessions::codex_history::CodexHistorySegment;
use freshell_sessions::codex_segments::CodexUnresolvedIdentity;
use freshell_sessions::directory_index::{IndexedSession, SessionIndex};
+use freshell_sessions::search::search_codex_history_segment;
// SESSION-07: the `userMessages`/`fullText` tier file-content search
// (`apply_file_search`, below) -- ports `server/session-directory/file-search.ts`.
use freshell_sessions::{search_session_file, FileSearchTier};
@@ -171,6 +173,9 @@ struct DirItem {
/// `fullText`. A composed Codex row carries every chronological segment;
/// other file-backed rows carry their single transcript. Never serialized.
source_files: Vec,
+ /// Selected Codex history bounds; absent for ordinary whole-file rows.
+ /// Captured with the row from one index generation and never serialized.
+ codex_history_segments: Option>,
/// STATUS-STRIP: live token usage (`SessionDirectoryItem.tokenUsage`,
/// `shared/read-models.ts`; Node's `CodingCliSession.tokenUsage`,
/// `coding-cli/types.ts:190`). Powers the fresh-agent strip's context
@@ -593,7 +598,20 @@ async fn session_directory(
} else {
indexed.source_file.clone().into_iter().collect()
};
- dir_item_from_indexed_with_source_files(indexed, source_files)
+ let mut item =
+ dir_item_from_indexed_with_source_files(indexed, source_files);
+ if indexed.provider == "codex" {
+ item.codex_history_segments = snapshot
+ .codex_history_segments
+ .get(&indexed.session_id)
+ .filter(|segments| {
+ !segments.is_empty()
+ && segments.last().map(|segment| &segment.path)
+ == indexed.source_file.as_ref()
+ })
+ .cloned();
+ }
+ item
})
.collect();
(items, snapshot.unresolved_codex_identities)
@@ -887,7 +905,7 @@ fn merge_unresolved_codex_identity_collisions(
*count = (*count).max(collision.duplicate_item_count);
}
for group in unresolved {
- if group.paths.len() < 2
+ if group.paths.is_empty()
|| codex_identity_is_canonically_soft_deleted(&group.session_id, overrides)
{
continue;
@@ -1117,6 +1135,7 @@ fn dir_item_from_indexed_with_source_files(
session_type: None,
title_source: idx.title_source.clone(),
source_files,
+ codex_history_segments: None,
token_usage: idx.token_usage.clone(),
// Provenance is overlay-derived (`apply_session_overrides`), never
// parsed from the transcript.
@@ -1283,6 +1302,7 @@ fn item_from_meta(
session_type: None,
title_source: meta.title_source.clone(),
source_files: source_file.into_iter().collect(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -1659,6 +1679,7 @@ fn build_live_terminal_session_item(
// `titleSource` either, `service.ts:110-130`).
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
// PARITY NOTE: Rust's `TerminalIdentity` carries no token usage, so a
// live-terminal-only row reports none here — unlike Node, whose
// `TerminalMeta` carries `tokenUsage`. Fresh-agent pane sessions are
@@ -1772,6 +1793,7 @@ mod join_tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -2235,11 +2257,11 @@ fn apply_file_search(
if !matches!(item.provider.as_str(), "claude" | "codex") {
continue;
}
- if item.source_files.is_empty() {
- continue;
- }
-
- for source_file in &item.source_files {
+ let source_count = item
+ .codex_history_segments
+ .as_ref()
+ .map_or(item.source_files.len(), Vec::len);
+ for source_index in 0..source_count {
if results.len() > limit {
break 'items;
}
@@ -2250,7 +2272,18 @@ fn apply_file_search(
}
scanned += 1;
- match search_session_file(source_file, &item.provider, query_text, tier) {
+ let searched = match &item.codex_history_segments {
+ Some(segments) => {
+ search_codex_history_segment(&segments[source_index], query_text, tier)
+ }
+ None => search_session_file(
+ &item.source_files[source_index],
+ &item.provider,
+ query_text,
+ tier,
+ ),
+ };
+ match searched {
Ok(Some(m)) => {
let mut matched = item.clone();
matched.matched_in = Some(m.matched_in.to_string());
@@ -2736,6 +2769,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: Some(freshell_sessions::meta::TokenSummary {
input_tokens: 10,
output_tokens: 5,
@@ -3010,6 +3044,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3229,6 +3264,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3284,6 +3320,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3351,6 +3388,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3427,6 +3465,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3498,6 +3537,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3550,6 +3590,7 @@ mod tests {
session_type: None,
title_source: title_source.map(str::to_string),
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -3694,6 +3735,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -6927,6 +6969,7 @@ mod tests {
session_type: None,
title_source: None,
source_files: Vec::new(),
+ codex_history_segments: None,
token_usage: None,
title_overridden: false,
provider_title: None,
@@ -6954,4 +6997,276 @@ mod tests {
assert_eq!(arr[0]["title"], json!("My Renamed Special Project"));
assert_eq!(arr[0]["matchedIn"], json!("title"));
}
+
+ fn write_codex_referenced_route_history(home: &Path) -> (PathBuf, PathBuf) {
+ let session_id = "b7936c10-4935-441c-837c-c1f33cafec2d";
+ let (older, newer) = codex_fixtures();
+ let mut older_records = older
+ .lines()
+ .map(|line| serde_json::from_str::(line).unwrap())
+ .collect::>();
+ older_records[0]["ordinal"] = json!(0);
+ for record in &mut older_records[1..] {
+ record["ordinal"] = json!(record["ordinal"].as_u64().unwrap() + 1);
+ }
+ let encode = |records: &[Value]| {
+ records
+ .iter()
+ .map(|record| format!("{record}\n"))
+ .collect::()
+ };
+ let end_byte_offset = encode(&older_records[..3]).len() as u64;
+ older_records[3]["type"] = json!("response_item");
+ older_records[3]["payload"] = json!({"type":"message", "role":"user", "content":[{"type":"input_text", "text":"Superseded user request"}]});
+ older_records[4]["type"] = json!("response_item");
+ older_records[4]["payload"] = json!({"type":"message", "role":"assistant", "content":[{"type":"output_text", "text":"Superseded assistant answer"}]});
+ let mut newer_records = newer
+ .lines()
+ .map(|line| serde_json::from_str::(line).unwrap())
+ .collect::>();
+ newer_records[0]["payload"]["history_base"] = json!({
+ "thread_id":session_id, "end_byte_offset":end_byte_offset, "end_ordinal_exclusive":3,
+ });
+ newer_records[0]["ordinal"] = json!(3);
+ newer_records[0]["payload"]["cli_version"] = json!("0.160.0");
+ for record in &mut newer_records[1..] {
+ record["ordinal"] = json!(record["ordinal"].as_u64().unwrap() + 4);
+ }
+ let codex_home = home.join(".codex");
+ let archive = codex_home.join("archived_sessions");
+ let sessions = codex_home.join("sessions");
+ std::fs::create_dir_all(&archive).unwrap();
+ std::fs::create_dir_all(&sessions).unwrap();
+ let root = archive.join(format!("rollout-2026-10-03T00-00-00-{session_id}.jsonl"));
+ let head = sessions.join(format!(
+ "rollout-2026-10-03T00-00-10-{session_id}_00000000-0000-4000-8000-000000000001.jsonl"
+ ));
+ std::fs::write(&root, encode(&older_records)).unwrap();
+ std::fs::write(&head, encode(&newer_records)).unwrap();
+ let db = rusqlite::Connection::open(codex_home.join("state_5.sqlite")).unwrap();
+ db.execute_batch("CREATE TABLE threads (id TEXT PRIMARY KEY, rollout_path TEXT, history_mode TEXT, archived INTEGER DEFAULT 0);").unwrap();
+ db.execute(
+ "INSERT INTO threads VALUES (?1, ?2, 'paginated', 0)",
+ rusqlite::params![session_id, head.to_str().unwrap()],
+ )
+ .unwrap();
+ (root, head)
+ }
+
+ #[tokio::test]
+ async fn codex_referenced_route_searches_archived_prefix_and_selected_head_with_boundaries() {
+ let home = unique_temp_dir();
+ let (root, head) = write_codex_referenced_route_history(&home);
+ let (app, index) = codex_session_directory_app(
+ &home,
+ None,
+ freshell_ws::identity::TerminalIdentityRegistry::new(),
+ );
+ let base = "/api/session-directory?priority=visible&includeNonInteractive=1";
+ let page = get_directory_page(&app, base).await;
+ assert_eq!(
+ page["items"].as_array().unwrap().len(),
+ 1,
+ "selected referenced history is one saved session: {page}"
+ );
+ assert!(
+ page.get("integrityError").is_none(),
+ "referenced history must not become a collision warning: {page}"
+ );
+ let snapshot = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ let session_id = "b7936c10-4935-441c-837c-c1f33cafec2d";
+ assert_eq!(snapshot.sessions[0].source_file.as_ref(), Some(&head));
+ assert_eq!(snapshot.codex_history_segments[session_id][0].path, root);
+ assert!(snapshot.codex_history_segments[session_id][0].end.is_some());
+ for (needle, tier, matched_in) in [
+ ("Older first request", "userMessages", "userMessage"),
+ ("Older assistant summary", "fullText", "assistantMessage"),
+ ("Continuation title", "userMessages", "userMessage"),
+ ("Continuation summary", "fullText", "assistantMessage"),
+ ] {
+ let result = get_directory_page(
+ &app,
+ &format!("{base}&query={}&tier={tier}", needle.replace(' ', "%20")),
+ )
+ .await;
+ assert_eq!(
+ result["items"].as_array().unwrap().len(),
+ 1,
+ "effective-history needle {needle:?}: {result}"
+ );
+ assert_eq!(result["items"][0]["matchedIn"], json!(matched_in));
+ assert!(result["items"][0]["snippet"]
+ .as_str()
+ .unwrap()
+ .contains(needle));
+ }
+ for tier in ["userMessages", "fullText"] {
+ let result =
+ get_directory_page(&app, &format!("{base}&query=Superseded&tier={tier}")).await;
+ assert!(
+ result["items"].as_array().unwrap().is_empty(),
+ "superseded tail must not leak through {tier}: {result}"
+ );
+ }
+ std::fs::remove_dir_all(home).unwrap();
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn codex_referenced_route_keeps_head_match_when_archived_prefix_disappears() {
+ let home = unique_temp_dir();
+ let (root, _) = write_codex_referenced_route_history(&home);
+ let (app, _) = codex_session_directory_app(
+ &home,
+ None,
+ freshell_ws::identity::TerminalIdentityRegistry::new(),
+ );
+ let base = "/api/session-directory?priority=visible&includeNonInteractive=1";
+ let warm = get_directory_page(&app, base).await;
+ assert_eq!(warm["items"].as_array().unwrap().len(), 1);
+ std::fs::remove_file(root).unwrap();
+ let result = get_directory_page(
+ &app,
+ &format!("{base}&query=Continuation%20title&tier=userMessages"),
+ )
+ .await;
+ assert_eq!(
+ result["items"].as_array().unwrap().len(),
+ 1,
+ "a later selected segment still matches: {result}"
+ );
+ assert_eq!(result["items"][0]["matchedIn"], json!("userMessage"));
+ assert_eq!(result["partial"], json!(true));
+ assert_eq!(result["partialReason"], json!("io_error"));
+ std::fs::remove_dir_all(home).unwrap();
+ }
+
+ #[tokio::test]
+ async fn codex_referenced_route_warns_for_a_single_unresolved_head_and_preserves_its_running_terminal(
+ ) {
+ let home = unique_temp_dir();
+ let (root, head) = write_codex_referenced_route_history(&home);
+ std::fs::remove_file(root).unwrap();
+ let session_id = "b7936c10-4935-441c-837c-c1f33cafec2d";
+ let identity = freshell_ws::identity::TerminalIdentityRegistry::new();
+ let (app, index) = codex_session_directory_app(&home, None, identity.clone());
+ let snapshot = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ assert_eq!(snapshot.unresolved_codex_identities.len(), 1);
+ assert_eq!(snapshot.unresolved_codex_identities[0].paths, vec![head]);
+ let base = "/api/session-directory?priority=visible&includeNonInteractive=1";
+ let hidden = get_directory_page(&app, base).await;
+ assert_eq!(
+ hidden["integrityError"]["collisionCount"],
+ json!(1),
+ "an unresolved referenced head must warn even with one physical file: {hidden}"
+ );
+ assert_eq!(hidden["partial"], json!(true));
+ assert!(
+ hidden["items"].as_array().unwrap().is_empty(),
+ "incomplete saved history stays quarantined: {hidden}"
+ );
+
+ identity.upsert(
+ "term-unresolved",
+ Some("codex"),
+ Some(session_id),
+ Some("/live-terminal"),
+ 400,
+ );
+ let running = get_directory_page(&app, base).await;
+ assert_eq!(running["integrityError"], hidden["integrityError"]);
+ let rows = running["items"].as_array().unwrap();
+ assert_eq!(rows.len(), 1);
+ assert_eq!(rows[0]["provider"], json!("codex"));
+ assert_eq!(rows[0]["sessionId"], json!(session_id));
+ assert_eq!(rows[0]["title"], json!("Codex CLI"));
+ assert_eq!(rows[0]["isRunning"], json!(true));
+ assert_eq!(rows[0]["runningTerminalId"], json!("term-unresolved"));
+ assert_eq!(rows[0]["projectPath"], json!("/live-terminal"));
+ assert_ne!(rows[0]["liveTerminalOnly"], json!(true));
+ std::fs::remove_dir_all(home).unwrap();
+ }
+
+ #[tokio::test]
+ async fn codex_referenced_route_quarantines_a_noncanonical_copy_and_recovers_after_removal() {
+ let home = unique_temp_dir();
+ let (_, head) = write_codex_referenced_route_history(&home);
+ let session_id = "b7936c10-4935-441c-837c-c1f33cafec2d";
+ let identity = freshell_ws::identity::TerminalIdentityRegistry::new();
+ let (app, index) = codex_session_directory_app(&home, None, identity.clone());
+ let base = "/api/session-directory?priority=visible&includeNonInteractive=1";
+ let initial = get_directory_page(&app, base).await;
+ let initial_rows = initial["items"].as_array().unwrap();
+ assert_eq!(initial_rows.len(), 1);
+ assert_eq!(
+ initial_rows[0]["firstUserMessage"],
+ json!("Older first request")
+ );
+ assert!(initial.get("integrityError").is_none());
+
+ // Discovery sees the same embedded identity even though this copied
+ // filename cannot identify a physical referenced-history segment.
+ let copy = head.parent().unwrap().join("copy.jsonl");
+ std::fs::copy(&head, ©).unwrap();
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ let conflicted = get_directory_page(&app, base).await;
+ assert_eq!(
+ conflicted["integrityError"]["kind"],
+ json!("identity_collision"),
+ "a noncanonical copy must not disappear from accepted history: {conflicted}"
+ );
+ assert_eq!(conflicted["integrityError"]["collisionCount"], json!(1));
+ assert_eq!(conflicted["partial"], json!(true));
+ assert!(
+ conflicted["items"].as_array().unwrap().is_empty(),
+ "saved history must remain quarantined while the copy exists: {conflicted}"
+ );
+
+ identity.upsert(
+ "term-copy-conflict",
+ Some("codex"),
+ Some(session_id),
+ Some("/live-terminal"),
+ 400,
+ );
+ let live = get_directory_page(&app, base).await;
+ assert_eq!(live["integrityError"], conflicted["integrityError"]);
+ let live_rows = live["items"].as_array().unwrap();
+ assert_eq!(live_rows.len(), 1);
+ assert_eq!(live_rows[0]["sessionId"], json!(session_id));
+ assert_eq!(live_rows[0]["title"], json!("Codex CLI"));
+ assert_eq!(live_rows[0]["isRunning"], json!(true));
+ assert_eq!(
+ live_rows[0]["runningTerminalId"],
+ json!("term-copy-conflict")
+ );
+ assert_eq!(live_rows[0]["projectPath"], json!("/live-terminal"));
+
+ std::fs::remove_file(copy).unwrap();
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ let recovered = get_directory_page(&app, base).await;
+ assert!(
+ recovered.get("integrityError").is_none(),
+ "copy removal must clear the warning: {recovered}"
+ );
+ let recovered_rows = recovered["items"].as_array().unwrap();
+ assert_eq!(recovered_rows.len(), 1);
+ assert_eq!(recovered_rows[0]["sessionId"], json!(session_id));
+ assert_eq!(
+ recovered_rows[0]["firstUserMessage"],
+ initial_rows[0]["firstUserMessage"]
+ );
+ assert_eq!(recovered_rows[0]["summary"], initial_rows[0]["summary"]);
+ assert_eq!(recovered_rows[0]["isRunning"], json!(true));
+ assert_eq!(
+ recovered_rows[0]["runningTerminalId"],
+ json!("term-copy-conflict")
+ );
+ std::fs::remove_dir_all(home).unwrap();
+ }
}
diff --git a/crates/freshell-sessions/Cargo.toml b/crates/freshell-sessions/Cargo.toml
index 9850104c6..2faee6c66 100644
--- a/crates/freshell-sessions/Cargo.toml
+++ b/crates/freshell-sessions/Cargo.toml
@@ -28,6 +28,8 @@ notify = "6"
# compiles the SQLite amalgamation so the parser is not coupled to the host's
# libsqlite3 version (deterministic across dev/CI/Windows).
rusqlite = { version = "0.31", features = ["bundled"] }
+# Canonical logical thread and physical rollout identifiers in Codex history.
+uuid = "1"
# SYNC-06 resume-resolve parity: the `shared/resume-input-parser.ts` port's
# candidate-extraction regexes + hint tables (`resume_input.rs`). `(?-u:\b)`
# keeps JS's ASCII \b word-boundary semantics. Same major already a direct
diff --git a/crates/freshell-sessions/src/codex_history.rs b/crates/freshell-sessions/src/codex_history.rs
new file mode 100644
index 000000000..6f7dd44fb
--- /dev/null
+++ b/crates/freshell-sessions/src/codex_history.rs
@@ -0,0 +1,824 @@
+//! Read-only reconstruction of Codex's selected, referenced rollout history.
+//!
+//! Logical thread identities stay stable after revert. A history_base points
+//! at a physical rollout UUID and inherits only its specified byte/ordinal
+//! prefix. Never concatenate superseded tails or guess a different DB head.
+
+use std::collections::{BTreeMap, HashMap, HashSet};
+use std::fs::{self, File};
+use std::hash::{Hash, Hasher};
+use std::io::{self, Read};
+use std::path::{Path, PathBuf};
+use std::time::UNIX_EPOCH;
+
+use chrono::NaiveDateTime;
+use rusqlite::{Connection, OpenFlags};
+use serde_json::Value;
+use uuid::Uuid;
+
+use crate::codex_segments::{
+ compose_codex_segments, scan_codex_file_bytes_evidence, supported_root_session_source,
+ CodexComposition, CodexFileEvidence, CodexSegmentEntry, CodexUnresolvedIdentity,
+};
+use crate::directory_index::{parse_codex_content, IndexedSession};
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+pub struct CodexHistoryBoundary {
+ pub end_byte_offset: u64,
+ pub end_ordinal_exclusive: u64,
+}
+
+/// A selected physical file, optionally restricted to an inherited prefix.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CodexHistorySegment {
+ pub path: PathBuf,
+ pub end: Option,
+}
+
+impl CodexHistorySegment {
+ /// Read exactly the retained range, without touching its superseded tail.
+ /// Validate it again because downstream search may run after discovery.
+ pub fn read(&self) -> io::Result> {
+ let bytes = self.read_bytes(None)?;
+ if self.end.is_some() {
+ validate_bytes(&bytes, self.end).map_err(|error| invalid_data(&error.detail))?;
+ }
+ Ok(bytes)
+ }
+
+ fn read_bytes(&self, snapshot_len: Option) -> io::Result> {
+ let mut bytes = Vec::new();
+ let file = File::open(&self.path)?;
+ let limit = match self.end {
+ Some(end) => end.end_byte_offset,
+ None => snapshot_len.unwrap_or(file.metadata()?.len()),
+ };
+ file.take(limit).read_to_end(&mut bytes)?;
+ if bytes.len() as u64 != limit {
+ return Err(invalid_data(
+ "source rollout is shorter than the selected byte range",
+ ));
+ }
+ Ok(bytes)
+ }
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+struct FileStamp {
+ path: PathBuf,
+ size: u64,
+ modified_nanos: u128,
+}
+
+impl FileStamp {
+ fn read(path: &Path) -> Result {
+ let metadata = fs::metadata(path)
+ .map_err(|error| HistoryError::at("source_unreadable", path, error))?;
+ let modified_nanos = metadata
+ .modified()
+ .and_then(|time| time.duration_since(UNIX_EPOCH).map_err(io::Error::other))
+ .map_err(|error| HistoryError::at("source_unreadable", path, error))?
+ .as_nanos();
+ Ok(Self {
+ path: path.to_owned(),
+ size: metadata.len(),
+ modified_nanos,
+ })
+ }
+}
+
+#[derive(Debug, Clone)]
+struct CachedHistory {
+ item: IndexedSession,
+ segments: Vec,
+ revision: u64,
+ dependencies: Vec,
+ // Record physical bindings too: a newly introduced duplicate must not
+ // evade validation merely because the selected bytes did not change.
+ bindings: Vec<(Uuid, PathBuf)>,
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+struct HistoryError {
+ reason: &'static str,
+ detail: String,
+}
+
+impl HistoryError {
+ fn new(reason: &'static str, detail: impl Into) -> Self {
+ Self {
+ reason,
+ detail: detail.into(),
+ }
+ }
+ fn at(reason: &'static str, path: &Path, error: impl std::fmt::Display) -> Self {
+ Self::new(reason, format!("{}: {error}", path.display()))
+ }
+}
+
+#[derive(Debug, Clone, Copy)]
+struct RolloutName {
+ logical: Uuid,
+ physical: Uuid,
+ timestamp: NaiveDateTime,
+}
+
+fn rollout_name(path: &Path) -> Option {
+ let name = path
+ .file_name()?
+ .to_str()?
+ .strip_prefix("rollout-")?
+ .strip_suffix(".jsonl")?;
+ let timestamp = NaiveDateTime::parse_from_str(name.get(..19)?, "%Y-%m-%dT%H-%M-%S").ok()?;
+ if name.get(19..20)? != "-" {
+ return None;
+ }
+ let ids = name.get(20..)?;
+ let (logical, physical) = ids.split_once('_').unwrap_or((ids, ids));
+ Some(RolloutName {
+ logical: Uuid::parse_str(logical).ok()?,
+ physical: Uuid::parse_str(physical).ok()?,
+ timestamp,
+ })
+}
+
+#[derive(Default)]
+struct Catalog {
+ paths: HashMap>,
+}
+
+impl Catalog {
+ fn add(&mut self, path: PathBuf) {
+ if let Some(name) = rollout_name(&path) {
+ let paths = self.paths.entry(name.physical).or_default();
+ if !paths.contains(&path) {
+ paths.push(path);
+ paths.sort();
+ }
+ }
+ }
+ fn unique(&self, physical: Uuid) -> Result {
+ match self.paths.get(&physical).map(Vec::as_slice) {
+ Some([path]) => Ok(path.clone()),
+ Some(_) => Err(HistoryError::new(
+ "ambiguous_rollout",
+ format!("physical rollout {physical} has multiple files"),
+ )),
+ None => Err(HistoryError::new(
+ "missing_reference",
+ format!("physical rollout {physical} was not found"),
+ )),
+ }
+ }
+}
+
+/// Stateful reconstruction cache owned by one Codex discovery source.
+#[derive(Debug)]
+pub struct CodexHistoryResolver {
+ codex_home: PathBuf,
+ cache: HashMap,
+ warnings: HashMap,
+ #[cfg(test)]
+ reads: usize,
+ #[cfg(test)]
+ after_read: Option,
+}
+
+impl CodexHistoryResolver {
+ pub fn new(codex_home: PathBuf) -> Self {
+ Self {
+ codex_home,
+ cache: HashMap::new(),
+ warnings: HashMap::new(),
+ #[cfg(test)]
+ reads: 0,
+ #[cfg(test)]
+ after_read: None,
+ }
+ }
+
+ pub fn compose(&mut self, entries: Vec) -> CodexComposition {
+ let db = self.database_selection();
+ // Fork/subagent ancestry is outside root continuation reconstruction.
+ // Preserve the old own-file projection for a single explicit child;
+ // the conservative composer still quarantines same-id child groups.
+ let child_ids: HashSet = entries
+ .iter()
+ .filter_map(|entry| {
+ let evidence = entry.evidence.as_ref()?;
+ explicit_child_lineage(evidence)
+ .then(|| evidence.session_id.clone())
+ .flatten()
+ })
+ .collect();
+ let mut referenced_ids: HashSet = entries
+ .iter()
+ .filter_map(|entry| {
+ let evidence = entry.evidence.as_ref()?;
+ (evidence
+ .lineage_markers
+ .iter()
+ .any(|marker| marker.name == "history_base")
+ || (evidence
+ .required_metadata
+ .get("history_mode")
+ .and_then(Value::as_str)
+ == Some("paginated")
+ && rollout_name(&entry.path)
+ .is_some_and(|name| name.logical != name.physical)))
+ .then(|| evidence.session_id.clone())
+ .flatten()
+ })
+ .collect();
+ if let Ok(rows) = &db {
+ for entry in &entries {
+ if let Some(id) = entry
+ .evidence
+ .as_ref()
+ .and_then(|evidence| evidence.session_id.as_ref())
+ {
+ if rows.get(id).is_some_and(|row| {
+ row.history_mode == "paginated"
+ && (row.archived
+ || rollout_name(&row.path)
+ .is_some_and(|name| name.logical != name.physical))
+ }) {
+ referenced_ids.insert(id.clone());
+ }
+ }
+ }
+ }
+ referenced_ids.retain(|id| !child_ids.contains(id));
+ let mut ordinary = Vec::new();
+ let mut groups = BTreeMap::>::new();
+ let mut catalog = Catalog::default();
+ for entry in entries {
+ catalog.add(entry.path.clone());
+ match entry
+ .evidence
+ .as_ref()
+ .and_then(|evidence| evidence.session_id.as_ref())
+ {
+ Some(id) if referenced_ids.contains(id) => {
+ groups.entry(id.clone()).or_default().push(entry)
+ }
+ _ => ordinary.push(entry),
+ }
+ }
+ self.cache.retain(|id, _| referenced_ids.contains(id));
+ self.warnings.retain(|id, _| referenced_ids.contains(id));
+ let mut result = compose_codex_segments(ordinary);
+ if groups.is_empty() {
+ return result;
+ }
+ let archive_result =
+ add_archive_paths(&self.codex_home.join("archived_sessions"), &mut catalog);
+ for (id, members) in groups {
+ let resolved =
+ self.selected_head(&id, &members, &db)
+ .and_then(|selected| match selected {
+ SelectedHead::Archived => Ok(None),
+ SelectedHead::Path(path) => validate_group_members(&id, &members)
+ .and_then(|_| archive_result.as_ref().map_err(Clone::clone))
+ .and_then(|_| self.resolve(&id, &path, &catalog).map(Some)),
+ });
+ match resolved {
+ Ok(Some(history)) => {
+ self.warnings.remove(&id);
+ result.segment_paths.insert(
+ id.clone(),
+ history
+ .segments
+ .iter()
+ .map(|segment| segment.path.clone())
+ .collect(),
+ );
+ result
+ .history_segments
+ .insert(id.clone(), history.segments.clone());
+ result
+ .history_revisions
+ .insert(id.clone(), history.revision);
+ result.items.push(history.item.clone());
+ self.cache.insert(id, history);
+ }
+ Ok(None) => {
+ self.warnings.remove(&id);
+ self.cache.remove(&id);
+ }
+ Err(error) => {
+ self.cache.remove(&id);
+ let signature = (error.reason.to_owned(), error.detail.clone());
+ if self.warnings.get(&id) != Some(&signature) {
+ tracing::warn!(event = "codex_history_unresolved", session_id = %id, reason = error.reason, detail = %error.detail,
+ "Codex selected history could not be reconstructed");
+ self.warnings.insert(id.clone(), signature);
+ }
+ let mut paths: Vec<_> =
+ members.iter().map(|member| member.path.clone()).collect();
+ paths.sort();
+ paths.dedup();
+ result.unresolved_identities.push(CodexUnresolvedIdentity {
+ session_id: id,
+ paths,
+ });
+ result
+ .items
+ .extend(members.into_iter().filter_map(|member| member.item));
+ }
+ }
+ }
+ result
+ .unresolved_identities
+ .sort_by(|left, right| left.session_id.cmp(&right.session_id));
+ result
+ }
+
+ fn selected_head(
+ &self,
+ id: &str,
+ members: &[CodexSegmentEntry],
+ db: &Result, HistoryError>,
+ ) -> Result {
+ if let Some(row) = db.as_ref().map_err(Clone::clone)?.get(id) {
+ if row.archived {
+ return Ok(SelectedHead::Archived);
+ }
+ if row.history_mode != "paginated" {
+ return Err(HistoryError::new(
+ "unsupported_history_mode",
+ "selected database row is not paginated",
+ ));
+ }
+ let path = if row.path.is_absolute() {
+ row.path.clone()
+ } else {
+ self.codex_home.join(&row.path)
+ };
+ if !path.is_file() {
+ return Err(HistoryError::at(
+ "selected_head_missing",
+ &path,
+ "database-selected rollout does not exist",
+ ));
+ }
+ return Ok(SelectedHead::Path(path));
+ }
+ let logical = Uuid::parse_str(id).map_err(|_| {
+ HistoryError::new("invalid_identity", "logical thread id is not a UUID")
+ })?;
+ members
+ .iter()
+ .filter_map(|member| {
+ let name = rollout_name(&member.path)?;
+ (name.logical == logical).then_some((
+ name.timestamp,
+ name.physical,
+ member.path.clone(),
+ ))
+ })
+ .max_by(|left, right| (left.0, left.1).cmp(&(right.0, right.1)))
+ .map(|(_, _, path)| SelectedHead::Path(path))
+ .ok_or_else(|| {
+ HistoryError::new(
+ "invalid_filename",
+ "selected history has no canonical rollout filename",
+ )
+ })
+ }
+
+ fn database_selection(&self) -> Result, HistoryError> {
+ let path = self.codex_home.join("state_5.sqlite");
+ if !path.exists() {
+ return Ok(HashMap::new());
+ }
+ let read_error = |error| HistoryError::at("state_db_unreadable", &path, error);
+ let connection = Connection::open_with_flags(
+ &path,
+ OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
+ )
+ .map_err(read_error)?;
+ let mut statement = connection
+ .prepare("SELECT id, rollout_path, history_mode, archived FROM threads")
+ .map_err(read_error)?;
+ let rows = statement
+ .query_map([], |row| {
+ Ok((
+ row.get::<_, String>(0)?,
+ DatabaseRow {
+ path: PathBuf::from(row.get::<_, String>(1)?),
+ history_mode: row.get(2)?,
+ archived: row.get::<_, i64>(3)? != 0,
+ },
+ ))
+ })
+ .map_err(read_error)?;
+ rows.map(|row| row.map_err(read_error)).collect()
+ }
+
+ fn resolve(
+ &mut self,
+ id: &str,
+ head: &Path,
+ catalog: &Catalog,
+ ) -> Result {
+ if let Some(cached) = self.cache.get(id) {
+ if cached.item.source_file.as_deref() == Some(head)
+ && cached
+ .bindings
+ .iter()
+ .all(|(physical, path)| catalog.unique(*physical).as_ref() == Ok(path))
+ && cached
+ .dependencies
+ .iter()
+ .all(|stamp| FileStamp::read(&stamp.path).as_ref() == Ok(stamp))
+ {
+ return Ok(cached.clone());
+ }
+ }
+ let logical = Uuid::parse_str(id).map_err(|_| {
+ HistoryError::new("invalid_identity", "logical thread id is not a UUID")
+ })?;
+ let mut path = head.to_owned();
+ let mut end = None;
+ let mut seen = HashSet::new();
+ let mut segments = Vec::new();
+ let mut bytes_by_segment = Vec::new();
+ let mut dependencies = Vec::new();
+ let mut bindings = Vec::new();
+ let root_creation_ms = loop {
+ let name = rollout_name(&path).ok_or_else(|| {
+ HistoryError::at(
+ "invalid_filename",
+ &path,
+ "not a canonical plain rollout filename",
+ )
+ })?;
+ if name.logical != logical {
+ return Err(HistoryError::at(
+ "identity_mismatch",
+ &path,
+ "reference belongs to another logical thread",
+ ));
+ }
+ if !seen.insert(name.physical) {
+ return Err(HistoryError::new(
+ "cycle",
+ "history references contain a cycle",
+ ));
+ }
+ let binding = catalog.unique(name.physical)?;
+ if binding != path {
+ return Err(HistoryError::at(
+ "identity_mismatch",
+ &path,
+ "selected path is not its physical rollout binding",
+ ));
+ }
+ let before = FileStamp::read(&path)?;
+ if end.is_some_and(|boundary: CodexHistoryBoundary| {
+ boundary.end_byte_offset > before.size
+ }) {
+ return Err(HistoryError::at(
+ "invalid_boundary",
+ &path,
+ "cutoff byte offset is past the source rollout",
+ ));
+ }
+ let segment = CodexHistorySegment {
+ path: path.clone(),
+ end,
+ };
+ #[cfg(test)]
+ {
+ self.reads += 1;
+ }
+ let bytes = segment
+ .read_bytes(Some(before.size))
+ .map_err(|error| HistoryError::at("source_unreadable", &path, error))?;
+ #[cfg(test)]
+ if let Some(after_read) = self.after_read {
+ after_read(&path);
+ }
+ let after = FileStamp::read(&path)?;
+ if before != after && after.size <= before.size {
+ return Err(HistoryError::at(
+ "source_changed",
+ &path,
+ "rollout changed during reconstruction",
+ ));
+ }
+ let header = validate_bytes(&bytes, end).map_err(|mut error| {
+ error.detail = format!("{}: {}", path.display(), error.detail);
+ error
+ })?;
+ validate_header_identity(&header.evidence, id, &path)?;
+ let base = header.base;
+ // Appends after the captured size do not invalidate a complete
+ // snapshot. Keep its earlier stamp so the next refresh rereads
+ // new records rather than caching old bytes as the newer file.
+ dependencies.push(before);
+ bindings.push((name.physical, path));
+ segments.push(segment);
+ bytes_by_segment.push(bytes);
+ let Some(base) = base else {
+ break header.creation_ms;
+ };
+ path = catalog.unique(base.physical)?;
+ end = Some(base.boundary);
+ };
+ segments.reverse();
+ bytes_by_segment.reverse();
+ let head_bytes = bytes_by_segment
+ .last()
+ .expect("a reconstructed history contains its head");
+ let head_header_end = head_bytes
+ .iter()
+ .position(|byte| *byte == b'\n')
+ .expect("validated header is newline terminated")
+ + 1;
+ let mut effective = head_bytes[..head_header_end].to_vec();
+ let mut hasher = std::collections::hash_map::DefaultHasher::new();
+ for bytes in &bytes_by_segment {
+ bytes.hash(&mut hasher);
+ let body_start = bytes
+ .iter()
+ .position(|byte| *byte == b'\n')
+ .expect("validated header is newline terminated")
+ + 1;
+ effective.extend_from_slice(&bytes[body_start..]);
+ }
+ let content =
+ std::str::from_utf8(&effective).expect("each selected segment was validated as UTF-8");
+ let mut item = parse_codex_content(content, head).ok_or_else(|| {
+ HistoryError::new(
+ "unrenderable_history",
+ "selected history has no displayable session metadata",
+ )
+ })?;
+ item.created_at = Some(root_creation_ms);
+ Ok(CachedHistory {
+ item,
+ segments,
+ revision: hasher.finish(),
+ dependencies,
+ bindings,
+ })
+ }
+}
+
+// Validate every candidate before consulting the selected-history cache: an
+// unclassified same-id copy must remain visible in quarantine diagnostics.
+// Canonical same-logical physical branches remain provider-selectable.
+fn validate_group_members(id: &str, members: &[CodexSegmentEntry]) -> Result<(), HistoryError> {
+ let logical = Uuid::parse_str(id)
+ .map_err(|_| HistoryError::new("invalid_identity", "logical thread id is not a UUID"))?;
+ let mut paths: Vec<_> = members.iter().map(|member| &member.path).collect();
+ paths.sort();
+ for path in paths {
+ let name = rollout_name(path).ok_or_else(|| {
+ HistoryError::at(
+ "invalid_filename",
+ path,
+ "same-id member does not have a canonical rollout filename",
+ )
+ })?;
+ if name.logical != logical {
+ return Err(HistoryError::at(
+ "identity_mismatch",
+ path,
+ "filename logical UUID does not match the embedded thread identity",
+ ));
+ }
+ }
+ Ok(())
+}
+
+fn explicit_child_lineage(evidence: &CodexFileEvidence) -> bool {
+ evidence
+ .lineage_markers
+ .iter()
+ .filter(|marker| marker.header_index == 0)
+ .any(|marker| match marker.name.as_str() {
+ "history_base" => false,
+ "is_subagent" => marker.value.as_bool() == Some(true),
+ _ => !marker.value.is_null(),
+ })
+}
+
+enum SelectedHead {
+ Path(PathBuf),
+ Archived,
+}
+struct DatabaseRow {
+ path: PathBuf,
+ history_mode: String,
+ archived: bool,
+}
+#[derive(Debug, Clone, Copy)]
+struct HistoryBase {
+ physical: Uuid,
+ boundary: CodexHistoryBoundary,
+}
+struct Header {
+ evidence: CodexFileEvidence,
+ base: Option,
+ creation_ms: i64,
+}
+
+fn validate_bytes(bytes: &[u8], end: Option) -> Result {
+ let content = std::str::from_utf8(bytes)
+ .map_err(|_| HistoryError::new("invalid_utf8", "selected prefix is not UTF-8"))?;
+ if !content.ends_with('\n') {
+ return Err(HistoryError::new(
+ "invalid_boundary",
+ "selected prefix does not end at a JSONL record boundary",
+ ));
+ }
+ let evidence = scan_codex_file_bytes_evidence(bytes);
+ if !evidence.scan_errors.is_empty() || !evidence.first_line_owned || evidence.header_count != 1
+ {
+ return Err(HistoryError::new(
+ "invalid_segment",
+ format!(
+ "selected prefix has invalid structural evidence: {:?}",
+ evidence.scan_errors
+ ),
+ ));
+ }
+ let first = content
+ .lines()
+ .next()
+ .ok_or_else(|| HistoryError::new("invalid_segment", "selected prefix is empty"))?;
+ let header: Value = serde_json::from_str(first)
+ .map_err(|_| HistoryError::new("invalid_segment", "invalid session header"))?;
+ let payload = &header["payload"];
+ let base = match payload.get("history_base") {
+ None | Some(Value::Null) => None,
+ Some(value) => {
+ let physical = value
+ .get("thread_id")
+ .and_then(Value::as_str)
+ .and_then(|id| Uuid::parse_str(id).ok())
+ .ok_or_else(|| {
+ HistoryError::new(
+ "invalid_reference",
+ "history_base has no physical rollout UUID",
+ )
+ })?;
+ let end_byte_offset = value
+ .get("end_byte_offset")
+ .and_then(Value::as_u64)
+ .filter(|offset| *offset > 0)
+ .ok_or_else(|| {
+ HistoryError::new(
+ "invalid_boundary",
+ "history_base byte offset is missing or zero",
+ )
+ })?;
+ let end_ordinal_exclusive = value
+ .get("end_ordinal_exclusive")
+ .and_then(Value::as_u64)
+ .filter(|ordinal| *ordinal > 0)
+ .ok_or_else(|| {
+ HistoryError::new(
+ "invalid_boundary",
+ "history_base ordinal is missing or zero",
+ )
+ })?;
+ Some(HistoryBase {
+ physical,
+ boundary: CodexHistoryBoundary {
+ end_byte_offset,
+ end_ordinal_exclusive,
+ },
+ })
+ }
+ };
+ let mut expected = base.map_or(0, |base| base.boundary.end_ordinal_exclusive);
+ for line in content.lines() {
+ let record: Value = serde_json::from_str(line)
+ .map_err(|_| HistoryError::new("invalid_segment", "invalid JSONL record"))?;
+ if record.get("ordinal").and_then(Value::as_u64) != Some(expected) {
+ return Err(HistoryError::new(
+ "invalid_ordinal",
+ format!("selected record must have ordinal {expected}"),
+ ));
+ }
+ if end.is_some_and(|end| expected >= end.end_ordinal_exclusive) {
+ return Err(HistoryError::new(
+ "invalid_boundary",
+ "byte prefix includes a record at or beyond its exclusive ordinal",
+ ));
+ }
+ expected = expected.checked_add(1).ok_or_else(|| {
+ HistoryError::new("invalid_ordinal", "selected record ordinal overflow")
+ })?;
+ }
+ if end.is_some_and(|end| expected != end.end_ordinal_exclusive) {
+ return Err(HistoryError::new(
+ "invalid_boundary",
+ "byte and ordinal boundaries select different prefixes",
+ ));
+ }
+ let creation_ms = header
+ .get("timestamp")
+ .and_then(Value::as_str)
+ .and_then(|timestamp| chrono::DateTime::parse_from_rfc3339(timestamp).ok())
+ .map(|timestamp| timestamp.timestamp_millis())
+ .ok_or_else(|| {
+ HistoryError::new(
+ "invalid_segment",
+ "session header has no valid outer timestamp",
+ )
+ })?;
+ Ok(Header {
+ evidence,
+ base,
+ creation_ms,
+ })
+}
+
+fn validate_header_identity(
+ evidence: &CodexFileEvidence,
+ id: &str,
+ path: &Path,
+) -> Result<(), HistoryError> {
+ if evidence.session_id.as_deref() != Some(id)
+ || evidence.first_header_id.as_ref().and_then(Value::as_str) != Some(id)
+ || evidence
+ .first_header_session_id
+ .as_ref()
+ .and_then(Value::as_str)
+ != Some(id)
+ {
+ return Err(HistoryError::at(
+ "identity_mismatch",
+ path,
+ "header id and session_id must match the logical thread",
+ ));
+ }
+ if evidence
+ .lineage_markers
+ .iter()
+ .any(|marker| marker.name != "history_base")
+ {
+ return Err(HistoryError::at(
+ "unsupported_lineage",
+ path,
+ "fork or subagent lineage cannot certify a root continuation",
+ ));
+ }
+ let metadata = &evidence.required_metadata;
+ if !["cwd", "cli_version", "originator"].iter().all(|key| {
+ metadata
+ .get(*key)
+ .and_then(Value::as_str)
+ .is_some_and(|value| !value.trim().is_empty())
+ }) || !metadata
+ .get("source")
+ .is_some_and(supported_root_session_source)
+ || metadata.get("thread_source").and_then(Value::as_str) != Some("user")
+ || metadata.get("history_mode").and_then(Value::as_str) != Some("paginated")
+ {
+ return Err(HistoryError::at(
+ "unsupported_metadata",
+ path,
+ "selected history is not a known paginated root user session",
+ ));
+ }
+ Ok(())
+}
+
+fn add_archive_paths(root: &Path, catalog: &mut Catalog) -> Result<(), HistoryError> {
+ let entries = match fs::read_dir(root) {
+ Ok(entries) => entries,
+ Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
+ Err(error) => return Err(HistoryError::at("archive_unreadable", root, error)),
+ };
+ for entry in entries {
+ let entry = entry.map_err(|error| HistoryError::at("archive_unreadable", root, error))?;
+ let path = entry.path();
+ let kind = entry
+ .file_type()
+ .map_err(|error| HistoryError::at("archive_unreadable", &path, error))?;
+ if kind.is_dir() {
+ add_archive_paths(&path, catalog)?;
+ } else if kind.is_file()
+ && path
+ .extension()
+ .is_some_and(|extension| extension == "jsonl")
+ {
+ catalog.add(path);
+ }
+ }
+ Ok(())
+}
+
+fn invalid_data(detail: &str) -> io::Error {
+ io::Error::new(io::ErrorKind::InvalidData, detail.to_owned())
+}
+
+#[cfg(test)]
+#[path = "codex_history_tests.rs"]
+mod tests;
diff --git a/crates/freshell-sessions/src/codex_history_tests.rs b/crates/freshell-sessions/src/codex_history_tests.rs
new file mode 100644
index 000000000..9b61ffcaf
--- /dev/null
+++ b/crates/freshell-sessions/src/codex_history_tests.rs
@@ -0,0 +1,1052 @@
+use super::*;
+use crate::directory_index::{CodexSource, SessionSource};
+use serde_json::json;
+
+const THREAD: &str = "11111111-1111-4111-8111-111111111111";
+const MIDDLE: &str = "22222222-2222-4222-8222-222222222222";
+const HEAD: &str = "33333333-3333-4333-8333-333333333333";
+
+struct Fixture {
+ home: PathBuf,
+}
+impl Fixture {
+ fn new() -> Self {
+ let home = std::env::temp_dir().join(format!("freshell-codex-history-{}", Uuid::new_v4()));
+ fs::create_dir_all(home.join("sessions/2026/10/01")).unwrap();
+ Self { home }
+ }
+ fn write(&self, timestamp: &str, physical: &str, content: &str) -> PathBuf {
+ let ids = if physical == THREAD {
+ THREAD.to_owned()
+ } else {
+ format!("{THREAD}_{physical}")
+ };
+ let path = self
+ .home
+ .join("sessions/2026/10/01")
+ .join(format!("rollout-{timestamp}-{ids}.jsonl"));
+ fs::write(&path, content).unwrap();
+ path
+ }
+ fn entries(&self, paths: &[PathBuf]) -> Vec {
+ let source = CodexSource::new(self.home.clone());
+ paths
+ .iter()
+ .map(|path| CodexSegmentEntry {
+ path: path.clone(),
+ item: source.parse(path),
+ evidence: Some(scan_codex_file_bytes_evidence(&fs::read(path).unwrap())),
+ })
+ .collect()
+ }
+ fn resolver(&self) -> CodexHistoryResolver {
+ CodexHistoryResolver::new(self.home.clone())
+ }
+ fn select(&self, path: &Path, archived: bool) {
+ let db = Connection::open(self.home.join("state_5.sqlite")).unwrap();
+ db.execute_batch("CREATE TABLE IF NOT EXISTS threads (id TEXT PRIMARY KEY, rollout_path TEXT NOT NULL, history_mode TEXT NOT NULL, archived INTEGER NOT NULL DEFAULT 0)").unwrap();
+ db.execute("INSERT OR REPLACE INTO threads(id,rollout_path,history_mode,archived) VALUES (?1,?2,'paginated',?3)", rusqlite::params![THREAD,path.to_str().unwrap(),archived as i64]).unwrap();
+ }
+ fn pair(&self) -> (PathBuf, PathBuf, String) {
+ // Retain only metadata; the old first user message is superseded.
+ let prefix = header(0, "2026-10-01T00:00:00Z", "/fixture/old", "0.159.2", None);
+ let root = self.write(
+ "2026-10-01T00-00-00",
+ THREAD,
+ &(prefix.clone() + &user(1, "2026-10-01T00:00:01Z", "Superseded question")),
+ );
+ let head = self.write(
+ "2026-10-01T00-01-00",
+ HEAD,
+ &(header(
+ 1,
+ "2026-10-01T00:01:00Z",
+ "/fixture/current",
+ "0.159.3",
+ Some(base(THREAD, &prefix, 1)),
+ ) + &user(2, "2026-10-01T00:01:01Z", "Current question")),
+ );
+ (root, head, prefix)
+ }
+}
+impl Drop for Fixture {
+ fn drop(&mut self) {
+ fs::remove_dir_all(&self.home).unwrap();
+ }
+}
+
+fn line(value: Value) -> String {
+ format!("{value}\n")
+}
+fn base(physical: &str, prefix: &str, end_ordinal: u64) -> Value {
+ json!({"thread_id":physical,"end_byte_offset":prefix.len(),"end_ordinal_exclusive":end_ordinal})
+}
+fn header(ordinal: u64, timestamp: &str, cwd: &str, version: &str, base: Option) -> String {
+ let mut payload = json!({"id":THREAD,"session_id":THREAD,"timestamp":"1999-01-01T00:00:00Z","cwd":cwd,"source":"vscode","thread_source":"user","cli_version":version,"originator":"codex-tui","history_mode":"paginated"});
+ if let Some(base) = base {
+ payload["history_base"] = base;
+ }
+ line(json!({"type":"session_meta","timestamp":timestamp,"ordinal":ordinal,"payload":payload}))
+}
+fn user(ordinal: u64, timestamp: &str, message: &str) -> String {
+ line(
+ json!({"type":"event_msg","timestamp":timestamp,"ordinal":ordinal,"payload":{"type":"user_message","message":message}}),
+ )
+}
+fn mutate_header(path: &Path, mutation: impl FnOnce(&mut Value)) {
+ let content = fs::read_to_string(path).unwrap();
+ let (header, body) = content.split_once('\n').unwrap();
+ let mut header: Value = serde_json::from_str(header).unwrap();
+ mutation(&mut header);
+ fs::write(path, line(header) + body).unwrap();
+}
+fn assert_unresolved(result: &CodexComposition) {
+ assert_eq!(result.unresolved_identities.len(), 1);
+ assert_eq!(result.unresolved_identities[0].session_id, THREAD);
+ assert!(!result.history_segments.contains_key(THREAD));
+ assert!(!result.history_revisions.contains_key(THREAD));
+}
+fn selected_text(result: &CodexComposition) -> String {
+ result.history_segments[THREAD]
+ .iter()
+ .map(|segment| String::from_utf8(segment.read().unwrap()).unwrap())
+ .collect()
+}
+
+#[test]
+fn reconstructs_header_only_prefix_and_uses_selected_metadata() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(&[root.clone(), head.clone()]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ let row = &result.items[0];
+ assert_eq!(row.first_user_message.as_deref(), Some("Current question"));
+ assert_eq!(row.cwd.as_deref(), Some("/fixture/current"));
+ assert_eq!(row.source_file.as_ref(), Some(&head));
+ assert_eq!(row.created_at, Some(1790812800000));
+ assert_eq!(
+ result.history_segments[THREAD],
+ vec![
+ CodexHistorySegment {
+ path: root,
+ end: Some(CodexHistoryBoundary {
+ end_byte_offset: prefix.len() as u64,
+ end_ordinal_exclusive: 1
+ })
+ },
+ CodexHistorySegment {
+ path: head,
+ end: None
+ },
+ ]
+ );
+ assert!(!selected_text(&result).contains("Superseded question"));
+}
+
+#[test]
+fn follows_physical_suffix_uuid_through_multiple_retained_prefixes() {
+ let fixture = Fixture::new();
+ let root_prefix = header(
+ 0,
+ "2026-10-01T00:00:00Z",
+ "/fixture/project",
+ "0.160.0",
+ None,
+ ) + &user(1, "2026-10-01T00:00:01Z", "Retained opening question");
+ let root = fixture.write(
+ "2026-10-01T00-00-00",
+ THREAD,
+ &(root_prefix.clone() + &user(2, "2026-10-01T00:00:02Z", "Superseded root tail")),
+ );
+ let middle_prefix = header(
+ 2,
+ "2026-10-01T00:01:00Z",
+ "/fixture/project",
+ "0.160.0",
+ Some(base(THREAD, &root_prefix, 2)),
+ ) + &user(3, "2026-10-01T00:01:01Z", "Retained middle question");
+ let middle = fixture.write(
+ "2026-10-01T00-01-00",
+ MIDDLE,
+ &(middle_prefix.clone() + &user(4, "2026-10-01T00:01:02Z", "Superseded middle tail")),
+ );
+ let head = fixture.write(
+ "2026-10-01T00-02-00",
+ HEAD,
+ &(header(
+ 4,
+ "2026-10-01T00:02:00Z",
+ "/fixture/project",
+ "0.160.0",
+ Some(base(MIDDLE, &middle_prefix, 4)),
+ ) + &user(5, "2026-10-01T00:02:01Z", "Current final question")),
+ );
+ let result =
+ fixture
+ .resolver()
+ .compose(fixture.entries(&[head.clone(), root.clone(), middle.clone()]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Retained opening question")
+ );
+ assert_eq!(result.segment_paths[THREAD], vec![root, middle, head]);
+ let text = selected_text(&result);
+ assert!(text.contains("Retained middle question"));
+ assert!(text.contains("Current final question"));
+ assert!(!text.contains("Superseded root tail"));
+ assert!(!text.contains("Superseded middle tail"));
+}
+
+#[test]
+fn ignores_malformed_superseded_tail() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ let mut invalid = prefix.into_bytes();
+ invalid.extend_from_slice(b"invalid JSON\n");
+ invalid.push(0xff);
+ invalid.extend_from_slice(b" missing final newline");
+ fs::write(&root, invalid).unwrap();
+ let result = fixture.resolver().compose(fixture.entries(&[root, head]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Current question")
+ );
+}
+
+#[test]
+fn resolves_archived_prefix_without_publishing_an_archived_row() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let archive = fixture.home.join("archived_sessions");
+ fs::create_dir_all(&archive).unwrap();
+ let archived = archive.join(root.file_name().unwrap());
+ fs::rename(root, &archived).unwrap();
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(std::slice::from_ref(&head)));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(result.segment_paths[THREAD], vec![archived, head]);
+}
+
+#[test]
+fn database_selection_overrides_newer_competing_branch() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ let branch = fixture.write(
+ "2026-10-01T00-02-00",
+ MIDDLE,
+ &(header(
+ 1,
+ "2026-10-01T00:02:00Z",
+ "/fixture/branch",
+ "0.160.0",
+ Some(base(THREAD, &prefix, 1)),
+ ) + &user(2, "2026-10-01T00:02:01Z", "Other branch question")),
+ );
+ fixture.select(&head, false);
+ let paths = [root, head.clone(), branch.clone()];
+ let mut resolver = fixture.resolver();
+ let first = resolver.compose(fixture.entries(&paths));
+ assert_eq!(first.items.len(), 1);
+ assert_eq!(
+ first.items[0].first_user_message.as_deref(),
+ Some("Current question")
+ );
+ fixture.select(&branch, false);
+ let second = resolver.compose(fixture.entries(&paths));
+ assert_eq!(second.items[0].source_file.as_ref(), Some(&branch));
+ assert_eq!(
+ second.items[0].first_user_message.as_deref(),
+ Some("Other branch question")
+ );
+ assert_ne!(
+ first.history_revisions[THREAD],
+ second.history_revisions[THREAD]
+ );
+}
+
+#[test]
+fn archived_database_selection_keeps_same_thread_files_hidden() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fixture.select(&head, true);
+ let result = fixture.resolver().compose(fixture.entries(&[root, head]));
+ assert!(result.items.is_empty());
+ assert!(result.unresolved_identities.is_empty());
+}
+
+#[test]
+fn archived_selection_keeps_noncanonical_copies_hidden_without_warnings() {
+ for warm_cache in [false, true] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let mut resolver = fixture.resolver();
+ if warm_cache {
+ let result = resolver.compose(fixture.entries(&[root.clone(), head.clone()]));
+ assert_eq!(result.items.len(), 1);
+ }
+ let copy = head.parent().unwrap().join("copy.jsonl");
+ fs::copy(&head, ©).unwrap();
+ fixture.select(&head, true);
+ let result = resolver.compose(fixture.entries(&[root, head, copy]));
+ assert!(result.items.is_empty());
+ assert!(result.unresolved_identities.is_empty());
+ assert!(result.history_segments.is_empty());
+ assert!(result.history_revisions.is_empty());
+ assert!(resolver.warnings.is_empty());
+ assert!(resolver.cache.is_empty());
+ }
+}
+
+#[test]
+fn stale_database_selection_never_falls_back_to_an_older_file() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fixture.select(&head, false);
+ fs::remove_file(&head).unwrap();
+ // Preserve the discovered reference-bearing member: selection remains
+ // authoritative even when the selected file disappears after discovery.
+ let head_entry = CodexSegmentEntry {
+ path: head.clone(),
+ item: None,
+ evidence: Some(scan_codex_file_bytes_evidence(
+ header(
+ 1,
+ "2026-10-01T00:01:00Z",
+ "/fixture/current",
+ "0.160.0",
+ Some(json!({"thread_id":THREAD,"end_byte_offset":1,"end_ordinal_exclusive":1})),
+ )
+ .as_bytes(),
+ )),
+ };
+ let mut entries = fixture.entries(&[root]);
+ entries.push(head_entry);
+ let mut resolver = fixture.resolver();
+ let result = resolver.compose(entries);
+ assert_unresolved(&result);
+ assert_eq!(resolver.warnings[THREAD].0, "selected_head_missing");
+}
+
+#[test]
+fn unreadable_database_does_not_choose_a_filesystem_branch() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fs::write(fixture.home.join("state_5.sqlite"), "not SQLite").unwrap();
+ let mut resolver = fixture.resolver();
+ let result = resolver.compose(fixture.entries(&[root, head]));
+ assert_unresolved(&result);
+ assert_eq!(resolver.warnings[THREAD].0, "state_db_unreadable");
+}
+
+#[test]
+fn missing_prefix_is_unresolved_even_for_one_renderable_head() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fs::remove_file(root).unwrap();
+ let mut resolver = fixture.resolver();
+ let result = resolver.compose(fixture.entries(&[head]));
+ assert_unresolved(&result);
+ assert_eq!(resolver.warnings[THREAD].0, "missing_reference");
+}
+
+#[test]
+fn invalid_byte_or_ordinal_cutoffs_remain_unresolved() {
+ for case in [
+ "past_eof",
+ "inside_record",
+ "ordinal_too_large",
+ "ordinal_too_small",
+ "zero_byte",
+ "zero_ordinal",
+ "missing_ordinal",
+ "invalid_physical_id",
+ ] {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ mutate_header(&head, |header| {
+ let base = &mut header["payload"]["history_base"];
+ match case {
+ "past_eof" => base["end_byte_offset"] = json!(999999),
+ "inside_record" => base["end_byte_offset"] = json!(prefix.len() - 1),
+ "ordinal_too_large" => base["end_ordinal_exclusive"] = json!(2),
+ "ordinal_too_small" => {
+ base["end_byte_offset"] = json!(fs::metadata(&root).unwrap().len())
+ }
+ "zero_byte" => base["end_byte_offset"] = json!(0),
+ "zero_ordinal" => base["end_ordinal_exclusive"] = json!(0),
+ "missing_ordinal" => {
+ base.as_object_mut()
+ .unwrap()
+ .remove("end_ordinal_exclusive");
+ }
+ "invalid_physical_id" => base["thread_id"] = json!("not-a-uuid"),
+ _ => unreachable!(),
+ }
+ });
+ let result = fixture.resolver().compose(fixture.entries(&[root, head]));
+ assert_unresolved(&result);
+ }
+}
+
+#[test]
+fn selected_records_require_contiguous_ordinals() {
+ for ordinal in [None, Some(0), Some(1), Some(3)] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let content = fs::read_to_string(&head).unwrap();
+ let (header, body) = content.split_once('\n').unwrap();
+ let mut record: Value = serde_json::from_str(body.trim()).unwrap();
+ match ordinal {
+ Some(ordinal) => record["ordinal"] = json!(ordinal),
+ None => {
+ record.as_object_mut().unwrap().remove("ordinal");
+ }
+ }
+ fs::write(&head, format!("{header}\n{}", line(record))).unwrap();
+ assert_unresolved(&fixture.resolver().compose(fixture.entries(&[root, head])));
+ }
+}
+
+#[test]
+fn selected_ancestor_requires_matching_logical_identity_and_root_metadata() {
+ for field in [
+ "id",
+ "session_id",
+ "history_mode",
+ "source",
+ "thread_source",
+ "parent_thread_id",
+ "cli_version",
+ ] {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ mutate_header(&root, |header| {
+ header["payload"][field] = match field {
+ "id" | "session_id" | "parent_thread_id" => json!(MIDDLE),
+ "cli_version" => json!(""),
+ _ => json!("unsupported"),
+ };
+ });
+ // Identity/metadata mutation changes prefix byte length; retain the
+ // actual entire header so this tests the header rather than a cut.
+ let updated = fs::read_to_string(&root)
+ .unwrap()
+ .lines()
+ .next()
+ .unwrap()
+ .to_owned()
+ + "\n";
+ assert_ne!(updated, prefix);
+ mutate_header(&head, |header| {
+ header["payload"]["history_base"]["end_byte_offset"] = json!(updated.len());
+ });
+ assert_unresolved(&fixture.resolver().compose(fixture.entries(&[root, head])));
+ }
+}
+
+#[test]
+fn prefix_search_reader_rejects_changed_invalid_boundaries() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(&[root.clone(), head]));
+ fs::write(root, "short\n").unwrap();
+ assert_eq!(
+ result.history_segments[THREAD][0]
+ .read()
+ .unwrap_err()
+ .kind(),
+ io::ErrorKind::InvalidData
+ );
+}
+
+#[test]
+fn memoizes_dependencies_and_preserves_revision_when_only_unused_tail_changes() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ let paths = [root.clone(), head];
+ let mut resolver = fixture.resolver();
+ let first = resolver.compose(fixture.entries(&paths));
+ assert_eq!(resolver.reads, 2);
+ let second = resolver.compose(fixture.entries(&paths));
+ assert_eq!(resolver.reads, 2);
+ assert_eq!(first.history_revisions, second.history_revisions);
+ fs::write(root, prefix + "a different malformed discarded tail\n").unwrap();
+ let third = resolver.compose(fixture.entries(&paths));
+ assert_eq!(resolver.reads, 4);
+ assert!(third.unresolved_identities.is_empty());
+ assert_eq!(first.history_revisions, third.history_revisions);
+ assert_eq!(first.items, third.items);
+}
+
+#[test]
+fn newly_ambiguous_physical_reference_invalidates_cached_history() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fixture.select(&head, false);
+ let mut resolver = fixture.resolver();
+ let first = resolver.compose(fixture.entries(&[root.clone(), head.clone()]));
+ assert!(first.unresolved_identities.is_empty());
+ let duplicate = fixture.write(
+ "2026-10-01T00-00-30",
+ THREAD,
+ &fs::read_to_string(&root).unwrap(),
+ );
+ let second = resolver.compose(fixture.entries(&[root, head, duplicate]));
+ assert_unresolved(&second);
+ assert_eq!(resolver.warnings[THREAD].0, "ambiguous_rollout");
+}
+
+#[test]
+fn warns_once_for_a_failure_and_rearms_after_recovery() {
+ // The tracing callsite interest cache is process-wide. Isolate capture
+ // from parallel tests installing/removing their own default dispatches.
+ const CHILD_ENV: &str = "FRESHELL_CODEX_HISTORY_LOG_CAPTURE_CHILD";
+ if std::env::var_os(CHILD_ENV).is_none() {
+ let output = std::process::Command::new(std::env::current_exe().unwrap())
+ .args([
+ "--exact",
+ "codex_history::tests::warns_once_for_a_failure_and_rearms_after_recovery",
+ "--nocapture",
+ ])
+ .env(CHILD_ENV, "1")
+ .output()
+ .unwrap();
+ assert!(
+ output.status.success(),
+ "isolated warning capture failed: {} {}",
+ String::from_utf8_lossy(&output.stdout),
+ String::from_utf8_lossy(&output.stderr)
+ );
+ return;
+ }
+ use tracing_subscriber::layer::SubscriberExt;
+ #[derive(Clone)]
+ struct Capture(std::sync::Arc>>>);
+ impl tracing_subscriber::layer::Layer for Capture {
+ fn on_event(
+ &self,
+ event: &tracing::Event<'_>,
+ _: tracing_subscriber::layer::Context<'_, S>,
+ ) {
+ struct Fields(BTreeMap);
+ impl tracing::field::Visit for Fields {
+ fn record_debug(
+ &mut self,
+ field: &tracing::field::Field,
+ value: &dyn std::fmt::Debug,
+ ) {
+ self.0.insert(field.name().into(), format!("{value:?}"));
+ }
+ fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
+ self.0.insert(field.name().into(), value.into());
+ }
+ }
+ let mut fields = Fields(BTreeMap::new());
+ event.record(&mut fields);
+ if event.metadata().level() == &tracing::Level::WARN
+ && fields.0.get("event").map(String::as_str) == Some("codex_history_unresolved")
+ {
+ self.0.lock().unwrap().push(fields.0);
+ }
+ }
+ }
+ let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
+ let subscriber = tracing_subscriber::registry().with(Capture(events.clone()));
+ let _guard = tracing::subscriber::set_default(subscriber);
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let root_bytes = fs::read(&root).unwrap();
+ fs::remove_file(&root).unwrap();
+ let mut resolver = fixture.resolver();
+ assert_unresolved(&resolver.compose(fixture.entries(std::slice::from_ref(&head))));
+ assert_unresolved(&resolver.compose(fixture.entries(std::slice::from_ref(&head))));
+ assert_eq!(events.lock().unwrap().len(), 1);
+ assert_eq!(events.lock().unwrap()[0]["reason"], "missing_reference");
+ assert_eq!(events.lock().unwrap()[0]["session_id"], THREAD);
+ fs::write(&root, root_bytes).unwrap();
+ let result = resolver.compose(fixture.entries(&[root.clone(), head.clone()]));
+ assert!(result.unresolved_identities.is_empty());
+ fs::remove_file(root).unwrap();
+ assert_unresolved(&resolver.compose(fixture.entries(&[head])));
+ assert_eq!(events.lock().unwrap().len(), 2);
+}
+
+#[test]
+fn valid_snapshot_survives_complete_or_partial_append_during_read() {
+ fn append_complete(path: &Path) {
+ if rollout_name(path).is_some_and(|name| name.physical.to_string() == HEAD) {
+ use std::io::Write;
+ let mut file = fs::OpenOptions::new().append(true).open(path).unwrap();
+ file.write_all(user(3, "2026-10-01T00:01:02Z", "A newly appended question").as_bytes())
+ .unwrap();
+ }
+ }
+ fn append_partial(path: &Path) {
+ if rollout_name(path).is_some_and(|name| name.physical.to_string() == HEAD) {
+ use std::io::Write;
+ fs::OpenOptions::new()
+ .append(true)
+ .open(path)
+ .unwrap()
+ .write_all(b"{\"type\":")
+ .unwrap();
+ }
+ }
+ for hook in [append_complete as fn(&Path), append_partial as fn(&Path)] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let paths = [root, head.clone()];
+ let mut resolver = fixture.resolver();
+ resolver.after_read = Some(hook);
+ let first = resolver.compose(fixture.entries(&paths));
+ assert!(first.unresolved_identities.is_empty());
+ assert_eq!(first.items.len(), 1);
+ assert_eq!(first.items[0].last_activity_at, 1790812861000);
+ // Finish any partial append before the next observation.
+ let old = fs::read_to_string(&head).unwrap();
+ if old.ends_with("{\"type\":") {
+ fs::write(
+ &head,
+ old.strip_suffix("{\"type\":").unwrap().to_owned()
+ + &user(3, "2026-10-01T00:01:02Z", "A newly appended question"),
+ )
+ .unwrap();
+ }
+ resolver.after_read = None;
+ let second = resolver.compose(fixture.entries(&paths));
+ assert!(second.unresolved_identities.is_empty());
+ assert_eq!(second.items[0].last_activity_at, 1790812862000);
+ assert_ne!(
+ first.history_revisions[THREAD],
+ second.history_revisions[THREAD]
+ );
+ assert_eq!(
+ resolver.reads, 4,
+ "new bytes must not share the old snapshot's cache stamp"
+ );
+ }
+}
+
+#[test]
+fn selected_token_usage_and_activity_exclude_superseded_tail() {
+ fn usage(ordinal: u64, timestamp: &str, total: u64) -> String {
+ line(
+ json!({"type":"event_msg","timestamp":timestamp,"ordinal":ordinal,"payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":total-1,"output_tokens":1,"total_tokens":total}}}}),
+ )
+ }
+ let fixture = Fixture::new();
+ let prefix = header(
+ 0,
+ "2026-10-01T00:00:00Z",
+ "/fixture/project",
+ "0.160.0",
+ None,
+ ) + &user(1, "2026-10-01T00:00:01Z", "Opening question")
+ + &usage(2, "2026-10-01T00:00:02Z", 100);
+ let root = fixture.write(
+ "2026-10-01T00-00-00",
+ THREAD,
+ &(prefix.clone()
+ + &user(3, "2026-10-01T23:59:58Z", "Superseded late question")
+ + &usage(4, "2026-10-01T23:59:59Z", 9999)),
+ );
+ let head = fixture.write(
+ "2026-10-01T00-01-00",
+ HEAD,
+ &(header(
+ 3,
+ "2026-10-01T00:01:00Z",
+ "/fixture/current",
+ "0.160.0",
+ Some(base(THREAD, &prefix, 3)),
+ ) + &user(4, "2026-10-01T00:01:01Z", "Current late question")),
+ );
+ let result = fixture.resolver().compose(fixture.entries(&[root, head]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(result.items[0].last_activity_at, 1790812861000);
+ assert_eq!(
+ result.items[0].token_usage.as_ref().unwrap().total_tokens,
+ 100
+ );
+ assert_eq!(
+ result.items[0].token_usage.as_ref().unwrap().input_tokens,
+ 99
+ );
+}
+
+#[test]
+fn filesystem_fallback_uses_filename_timestamp_then_physical_uuid() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ // MIDDLE has a lower UUID despite being added later to discovery.
+ let same_second = fixture.write(
+ "2026-10-01T00-01-00",
+ MIDDLE,
+ &(header(
+ 1,
+ "2026-10-01T00:01:00Z",
+ "/fixture/tie",
+ "0.160.0",
+ Some(base(THREAD, &prefix, 1)),
+ ) + &user(2, "2026-10-01T00:01:01Z", "Tie branch")),
+ );
+ let mut resolver = fixture.resolver();
+ let first =
+ resolver.compose(fixture.entries(&[head.clone(), root.clone(), same_second.clone()]));
+ assert_eq!(first.items.len(), 1);
+ assert_eq!(first.items[0].source_file.as_ref(), Some(&head));
+ let latest = fixture.write(
+ "2026-10-01T00-02-00",
+ MIDDLE,
+ &(header(
+ 1,
+ "2026-10-01T00:02:00Z",
+ "/fixture/latest",
+ "0.160.0",
+ Some(base(THREAD, &prefix, 1)),
+ ) + &user(2, "2026-10-01T00:02:01Z", "Newest branch")),
+ );
+ fs::remove_file(same_second).unwrap();
+ let second = resolver.compose(fixture.entries(&[latest.clone(), root, head]));
+ assert_eq!(second.items.len(), 1);
+ assert_eq!(second.items[0].source_file.as_ref(), Some(&latest));
+}
+
+#[test]
+fn cyclic_references_remain_unresolved() {
+ let fixture = Fixture::new();
+ let middle_prefix = header(
+ 5,
+ "2026-10-01T00:00:00Z",
+ "/fixture/project",
+ "0.160.0",
+ Some(json!({"thread_id":HEAD,"end_byte_offset":1,"end_ordinal_exclusive":5})),
+ ) + &user(6, "2026-10-01T00:00:01Z", "Retained question six")
+ + &user(7, "2026-10-01T00:00:02Z", "Retained question seven")
+ + &user(8, "2026-10-01T00:00:03Z", "Retained question eight")
+ + &user(9, "2026-10-01T00:00:04Z", "Retained question nine");
+ let head_header = header(
+ 10,
+ "2026-10-01T00:01:00Z",
+ "/fixture/project",
+ "0.160.0",
+ Some(base(MIDDLE, &middle_prefix, 10)),
+ );
+ let middle = fixture.write("2026-10-01T00-00-00", MIDDLE, &middle_prefix);
+ let head = fixture.write("2026-10-01T00-01-00", HEAD, &head_header);
+ let mut resolver = fixture.resolver();
+ assert_unresolved(&resolver.compose(fixture.entries(&[middle, head])));
+ assert_eq!(resolver.warnings[THREAD].0, "cycle");
+}
+
+#[test]
+fn standalone_replacement_without_history_base_supersedes_original_file() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fs::write(
+ &head,
+ header(
+ 0,
+ "2026-10-01T00:01:00Z",
+ "/fixture/replaced",
+ "0.160.0",
+ None,
+ ) + &user(1, "2026-10-01T00:01:01Z", "Entirely replaced conversation"),
+ )
+ .unwrap();
+ for use_database in [false, true] {
+ if use_database {
+ fixture.select(&head, false);
+ }
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(&[root.clone(), head.clone()]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Entirely replaced conversation")
+ );
+ assert_eq!(result.items[0].source_file.as_ref(), Some(&head));
+ // No prefix is inherited, so the selected replacement itself is the
+ // reconstruction root; superseded creation metadata is not inferred.
+ assert_eq!(result.items[0].created_at, Some(1790812860000));
+ assert_eq!(result.segment_paths[THREAD], vec![head.clone()]);
+ assert!(!selected_text(&result).contains("Superseded question"));
+ }
+}
+
+#[test]
+fn fresh_resolver_hides_original_when_database_selected_replacement_is_archived() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let archive = fixture.home.join("archived_sessions");
+ fs::create_dir_all(&archive).unwrap();
+ let archived = archive.join(head.file_name().unwrap());
+ fs::rename(head, &archived).unwrap();
+ fixture.select(&archived, true);
+ let result = fixture.resolver().compose(fixture.entries(&[root]));
+ assert!(result.items.is_empty());
+ assert!(result.unresolved_identities.is_empty());
+}
+
+#[test]
+fn cached_archived_dependencies_are_rechecked_without_active_file_changes() {
+ for operation in ["mutation", "removal", "move"] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let archive = fixture.home.join("archived_sessions");
+ fs::create_dir_all(&archive).unwrap();
+ let archived = archive.join(root.file_name().unwrap());
+ fs::rename(root, &archived).unwrap();
+ let active_entries = fixture.entries(std::slice::from_ref(&head));
+ let mut resolver = fixture.resolver();
+ let first = resolver.compose(active_entries.clone());
+ assert!(first.unresolved_identities.is_empty());
+ let moved = archive.join("moved").join(archived.file_name().unwrap());
+ match operation {
+ "mutation" => {
+ let content = fs::read_to_string(&archived)
+ .unwrap()
+ .replace("2026-10-01T00:00:00Z", "2026-10-01T00:00:05Z")
+ + "ignored superseded append\n";
+ fs::write(&archived, content).unwrap();
+ }
+ "removal" => fs::remove_file(&archived).unwrap(),
+ "move" => {
+ fs::create_dir_all(moved.parent().unwrap()).unwrap();
+ fs::rename(&archived, &moved).unwrap();
+ }
+ _ => unreachable!(),
+ }
+ let second = resolver.compose(active_entries);
+ match operation {
+ "mutation" => {
+ assert!(second.unresolved_identities.is_empty());
+ assert_eq!(second.items[0].created_at, Some(1790812805000));
+ assert_ne!(first.history_revisions, second.history_revisions);
+ }
+ "removal" => assert_unresolved(&second),
+ "move" => {
+ assert!(second.unresolved_identities.is_empty());
+ assert_eq!(second.segment_paths[THREAD], vec![moved, head]);
+ assert_eq!(first.history_revisions, second.history_revisions);
+ assert_ne!(first.history_segments, second.history_segments);
+ }
+ _ => unreachable!(),
+ }
+ }
+}
+
+#[test]
+fn truncation_during_read_is_diagnosed_instead_of_cached() {
+ fn truncate(path: &Path) {
+ if rollout_name(path).is_some_and(|name| name.physical.to_string() == HEAD) {
+ fs::write(path, b"truncated\n").unwrap();
+ }
+ }
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let mut resolver = fixture.resolver();
+ resolver.after_read = Some(truncate);
+ let result = resolver.compose(fixture.entries(&[root, head]));
+ assert_unresolved(&result);
+ assert_eq!(resolver.warnings[THREAD].0, "source_changed");
+}
+
+#[test]
+fn marker_free_non_paginated_single_files_preserve_existing_display() {
+ for mode in ["legacy", "full", "unsupported"] {
+ for use_database in [false, true] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fs::remove_file(root).unwrap();
+ fs::write(
+ &head,
+ header(
+ 0,
+ "2026-10-01T00:01:00Z",
+ "/fixture/project",
+ "0.160.0",
+ None,
+ ) + &user(
+ 1,
+ "2026-10-01T00:01:01Z",
+ "Existing single file conversation",
+ ),
+ )
+ .unwrap();
+ mutate_header(&head, |header| {
+ header["payload"]["history_mode"] = json!(mode);
+ });
+ if use_database {
+ fixture.select(&head, false);
+ let connection = Connection::open(fixture.home.join("state_5.sqlite")).unwrap();
+ connection
+ .execute(
+ "UPDATE threads SET history_mode=?1 WHERE id=?2",
+ [mode, THREAD],
+ )
+ .unwrap();
+ }
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(std::slice::from_ref(&head)));
+ assert!(
+ result.unresolved_identities.is_empty(),
+ "marker-free {mode} file must retain prior single-file behavior"
+ );
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Existing single file conversation")
+ );
+ assert!(!result.history_segments.contains_key(THREAD));
+ }
+ }
+}
+
+#[test]
+fn single_explicit_fork_or_subagent_reference_preserves_own_file_display() {
+ for marker in [
+ "forked_from_id",
+ "parent_thread_id",
+ "subagent_history_start_ordinal",
+ "is_subagent",
+ "source.subagent",
+ ] {
+ let fixture = Fixture::new();
+ let (_, head, _) = fixture.pair();
+ mutate_header(&head, |header| match marker {
+ "subagent_history_start_ordinal" => header["payload"][marker] = json!(1),
+ "is_subagent" => header["payload"][marker] = json!(true),
+ "source.subagent" => {
+ header["payload"]["source"] =
+ json!({"subagent":{"thread_spawn":{"parent_thread_id":MIDDLE,"depth":1}}})
+ }
+ _ => {
+ header["payload"][marker] = json!(MIDDLE);
+ header["payload"]["session_id"] = json!(MIDDLE);
+ }
+ });
+ let result = fixture
+ .resolver()
+ .compose(fixture.entries(std::slice::from_ref(&head)));
+ assert!(
+ result.unresolved_identities.is_empty(),
+ "single {marker} history retains previous own-file display"
+ );
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(result.items[0].source_file.as_ref(), Some(&head));
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Current question")
+ );
+ assert!(!result.history_segments.contains_key(THREAD));
+ assert!(!result.history_revisions.contains_key(THREAD));
+ }
+}
+
+#[test]
+fn multiple_same_id_fork_reference_files_stay_quarantined() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ mutate_header(&head, |header| {
+ header["payload"]["forked_from_id"] = json!(MIDDLE);
+ });
+ assert_unresolved(&fixture.resolver().compose(fixture.entries(&[root, head])));
+}
+
+#[test]
+fn false_subagent_flag_does_not_exempt_a_root_missing_reference() {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ fs::remove_file(root).unwrap();
+ mutate_header(&head, |header| {
+ header["payload"]["is_subagent"] = json!(false);
+ });
+ assert_unresolved(&fixture.resolver().compose(fixture.entries(&[head])));
+}
+
+#[test]
+fn superseded_child_header_does_not_poison_retained_root_prefix() {
+ let fixture = Fixture::new();
+ let (root, head, prefix) = fixture.pair();
+ let mut discarded: Value = serde_json::from_str(
+ header(1, "2026-10-01T00:00:01Z", "/fixture/child", "0.160.0", None).trim(),
+ )
+ .unwrap();
+ discarded["payload"]["parent_thread_id"] = json!(MIDDLE);
+ fs::write(&root, prefix + &line(discarded)).unwrap();
+ let result = fixture.resolver().compose(fixture.entries(&[root, head]));
+ assert!(result.unresolved_identities.is_empty());
+ assert_eq!(result.items.len(), 1);
+ assert_eq!(
+ result.items[0].first_user_message.as_deref(),
+ Some("Current question")
+ );
+ assert!(result.history_segments.contains_key(THREAD));
+}
+
+#[test]
+fn noncanonical_root_or_head_copy_quarantines_fresh_and_cached_reference_groups() {
+ assert_copy_quarantine_recovery("copy.jsonl", "invalid_filename");
+}
+
+#[test]
+fn foreign_logical_filename_copy_quarantines_fresh_and_cached_reference_groups() {
+ assert_copy_quarantine_recovery(
+ &format!("rollout-2026-10-01T00-03-00-{MIDDLE}.jsonl"),
+ "identity_mismatch",
+ );
+}
+
+fn assert_copy_quarantine_recovery(copy_name: &str, expected_reason: &str) {
+ for copy_head in [false, true] {
+ for warm_cache in [false, true] {
+ let fixture = Fixture::new();
+ let (root, head, _) = fixture.pair();
+ let originals = [root.clone(), head.clone()];
+ let mut resolver = fixture.resolver();
+ if warm_cache {
+ let accepted = resolver.compose(fixture.entries(&originals));
+ assert!(accepted.unresolved_identities.is_empty());
+ assert_eq!(accepted.items.len(), 1);
+ }
+ let copy = head.parent().unwrap().join(copy_name);
+ fs::copy(if copy_head { &head } else { &root }, ©).unwrap();
+ let candidates = [root.clone(), head.clone(), copy.clone()];
+ let rejected = resolver.compose(fixture.entries(&candidates));
+ assert_unresolved(&rejected);
+ let mut expected_paths = candidates.to_vec();
+ expected_paths.sort();
+ assert_eq!(
+ rejected.unresolved_identities[0].paths, expected_paths,
+ "do not lose copied paths from quarantine diagnostics"
+ );
+ assert_eq!(resolver.warnings[THREAD].0, expected_reason);
+ assert!(resolver.warnings[THREAD].1.contains(copy_name));
+ fs::remove_file(copy).unwrap();
+ let recovered = resolver.compose(fixture.entries(&originals));
+ assert!(recovered.unresolved_identities.is_empty());
+ assert_eq!(recovered.items.len(), 1);
+ assert_eq!(recovered.items[0].source_file.as_ref(), Some(&head));
+ assert_eq!(
+ recovered.items[0].first_user_message.as_deref(),
+ Some("Current question")
+ );
+ assert!(!resolver.warnings.contains_key(THREAD));
+ }
+ }
+}
diff --git a/crates/freshell-sessions/src/codex_segments.rs b/crates/freshell-sessions/src/codex_segments.rs
index 63aaf47c6..991fb9128 100644
--- a/crates/freshell-sessions/src/codex_segments.rs
+++ b/crates/freshell-sessions/src/codex_segments.rs
@@ -119,6 +119,11 @@ pub struct CodexComposition {
/// embedded id. This sidecar feeds bounded downstream search without
/// changing `source_file`'s existing meaning.
pub segment_paths: HashMap>,
+ /// Selected byte/ordinal ranges for referenced rollout histories. Ordinary
+ /// single-file and verified chronological histories retain whole-file paths.
+ pub history_segments: HashMap>,
+ /// Digest of only selected history bytes, excluding superseded tails.
+ pub history_revisions: HashMap,
}
/// One per-file input to [`compose_codex_segments`].
@@ -475,7 +480,13 @@ fn can_compose_group(session_id: &str, members: &[CodexSegmentEntry]) -> bool {
previous_end = Some(interval.end_nanos);
if let Some(previous) = matching_metadata {
- if previous != &evidence.required_metadata {
+ // A CLI update changes the recorder version, not the thread's
+ // identity. Still require a nonempty version in each header.
+ if REQUIRED_METADATA
+ .iter()
+ .filter(|key| **key != "cli_version")
+ .any(|key| previous.get(*key) != evidence.required_metadata.get(*key))
+ {
return false;
}
} else {
@@ -529,7 +540,7 @@ fn nonempty(value: Option<&String>) -> Option<&String> {
/// Codex sessions. Unknown and internal/subagent sources do not certify a
/// user-visible continuation; custom sources are supported only in their
/// serialized enum form and are compared exactly across members.
-fn supported_root_session_source(value: &Value) -> bool {
+pub(crate) fn supported_root_session_source(value: &Value) -> bool {
match value {
Value::String(source) => matches!(source.as_str(), "cli" | "vscode" | "exec" | "mcp"),
Value::Object(fields) => {
diff --git a/crates/freshell-sessions/src/directory_index.rs b/crates/freshell-sessions/src/directory_index.rs
index 32fe20b52..303829f79 100644
--- a/crates/freshell-sessions/src/directory_index.rs
+++ b/crates/freshell-sessions/src/directory_index.rs
@@ -41,6 +41,7 @@ use std::time::{Duration, Instant};
use tokio::sync::Mutex as AsyncMutex;
+use crate::codex_history::{CodexHistoryResolver, CodexHistorySegment};
use crate::codex_segments::{
compose_codex_segments, scan_codex_file_bytes_evidence, CodexComposition, CodexFileEvidence,
CodexSegmentEntry, CodexUnresolvedIdentity,
@@ -162,6 +163,10 @@ pub struct SessionIndexSnapshot {
/// Chronological source paths for accepted multi-file Codex rows, keyed
/// by their canonical embedded session id.
pub codex_segment_paths: Arc>>,
+ /// Effective bounded source segments, from this same published generation.
+ pub codex_history_segments: Arc>>,
+ /// Process-local content revisions for change notification.
+ pub codex_history_revisions: Arc>,
}
/// One discovered file: its absolute path plus the stat facts (`mtime`/`size`)
@@ -234,6 +239,12 @@ pub trait SessionSource: Send + Sync {
(self.parse(path), None)
}
+ /// Compose parsed Codex files using this provider's selected history.
+ /// Test sources and other providers retain the conservative file composer.
+ fn compose_codex_segments(&self, entries: Vec) -> CodexComposition {
+ compose_codex_segments(entries)
+ }
+
/// Batch C: direct-listed sources (opencode's single sqlite db, which
/// enumerates MANY sessions in ONE query rather than one file per
/// session) can't fit the per-file `discover`/`parse` cache — there's no
@@ -578,11 +589,15 @@ fn item_from_meta(
/// `ClaudeSource` joining `projects`.
pub struct CodexSource {
codex_home: PathBuf,
+ history: StdMutex,
}
impl CodexSource {
pub fn new(codex_home: PathBuf) -> Self {
- Self { codex_home }
+ Self {
+ history: StdMutex::new(CodexHistoryResolver::new(codex_home.clone())),
+ codex_home,
+ }
}
/// Convenience: discover + parse every currently-visible file in one
@@ -601,11 +616,15 @@ impl CodexSource {
}
})
.collect();
- compose_codex_segments(entries).items
+ self.compose_codex_segments(entries).items
}
}
impl SessionSource for CodexSource {
+ fn compose_codex_segments(&self, entries: Vec) -> CodexComposition {
+ self.history.lock().unwrap().compose(entries)
+ }
+
fn discover(&self) -> Vec {
discover_codex_sessions(&self.codex_home).unwrap_or_default()
}
@@ -713,7 +732,7 @@ fn parse_codex_file_with_evidence(
)
}
-fn parse_codex_content(content: &str, path: &Path) -> Option {
+pub(crate) fn parse_codex_content(content: &str, path: &Path) -> Option {
let meta = parse_codex_session_content(content);
meta.cwd.as_ref()?;
let fallback = extract_codex_session_id_from_filename(path);
@@ -1185,6 +1204,8 @@ struct CachedSnapshot {
/// Chronological source paths for accepted multi-file rows, keyed by
/// canonical embedded id. This remains internal to Rust consumers.
codex_segment_paths: Arc>>,
+ codex_history_segments: Arc>>,
+ codex_history_revisions: Arc>,
}
/// Fields copied together from one published generation. The identity list
@@ -1196,6 +1217,8 @@ struct SnapshotRead {
scan_failures: Vec,
unresolved_codex_identities: Arc>,
codex_segment_paths: Arc>>,
+ codex_history_segments: Arc>>,
+ codex_history_revisions: Arc>,
}
/// Bookkeeping for the persistent parse-cache's opportunistic-save gating
@@ -1519,6 +1542,8 @@ impl SessionIndex {
scan_failures: snapshot.scan_failures,
unresolved_codex_identities: snapshot.unresolved_codex_identities,
codex_segment_paths: snapshot.codex_segment_paths,
+ codex_history_segments: snapshot.codex_history_segments,
+ codex_history_revisions: snapshot.codex_history_revisions,
}
}
@@ -1658,6 +1683,8 @@ impl SessionIndex {
scan_failures: sorted_names(&c.scan_failures),
unresolved_codex_identities: Arc::clone(&c.unresolved_codex_identities),
codex_segment_paths: Arc::clone(&c.codex_segment_paths),
+ codex_history_segments: Arc::clone(&c.codex_history_segments),
+ codex_history_revisions: Arc::clone(&c.codex_history_revisions),
})
}
_ => None,
@@ -1966,6 +1993,8 @@ impl SessionIndex {
refreshed.amplifier_root_dirs,
refreshed.unresolved_codex_identities,
refreshed.codex_segment_paths,
+ refreshed.codex_history_segments,
+ refreshed.codex_history_revisions,
)
}
})
@@ -1977,6 +2006,8 @@ impl SessionIndex {
amplifier_root_dirs,
unresolved_codex_identities,
codex_segment_paths,
+ codex_history_segments,
+ codex_history_revisions,
) = match sweep_result {
Ok(result) => result,
Err(join_err) => {
@@ -2006,19 +2037,33 @@ impl SessionIndex {
let items = Arc::new(items);
let unresolved_codex_identities = Arc::new(unresolved_codex_identities);
let codex_segment_paths = Arc::new(codex_segment_paths);
+ let codex_history_segments = Arc::new(codex_history_segments);
+ let codex_history_revisions = Arc::new(codex_history_revisions);
let failure_names = sorted_names(&failures);
+ let codex_history_changed;
{
// ONE lock write publishes the snapshot AND its scan failures as
// a single generation — a reader (`cached_pair`) can never
// observe a failed-scan snapshot paired with a cleared failure
// set, nor a recovered snapshot paired with stale failures.
let mut guard = snapshot.lock().unwrap();
+ // SQLite selection or archived dependencies can change effective
+ // history without mutating a parsed-file cache entry. Keep their
+ // notification signal separate from persistence save accounting.
+ codex_history_changed = guard.as_ref().is_some_and(|previous| {
+ previous.codex_history_revisions != codex_history_revisions
+ || previous.codex_history_segments != codex_history_segments
+ || previous.codex_segment_paths != codex_segment_paths
+ || previous.unresolved_codex_identities != unresolved_codex_identities
+ });
*guard = Some(CachedSnapshot {
items: Arc::clone(&items),
fetched_at: Instant::now(),
scan_failures: failures,
unresolved_codex_identities,
codex_segment_paths,
+ codex_history_segments,
+ codex_history_revisions,
});
} // guard dropped here — never held across an .await.
// Self-correction report (amplifier watch-reduction design
@@ -2034,7 +2079,7 @@ impl SessionIndex {
if force_full {
*last_full_at.lock().unwrap() = Some(Instant::now());
}
- if changed > 0 {
+ if changed > 0 || codex_history_changed {
change_tx.send_modify(|gen| *gen += 1);
}
// Opportunistic persistence: gated (threshold/debounce) and, when
@@ -2250,6 +2295,8 @@ struct RefreshedSnapshot {
amplifier_root_dirs: Option>,
unresolved_codex_identities: Vec,
codex_segment_paths: HashMap>,
+ codex_history_segments: HashMap>,
+ codex_history_revisions: HashMap,
}
fn refresh_snapshot(
@@ -2570,7 +2617,13 @@ fn refresh_snapshot(
evidence: entry.codex_evidence.clone(),
})
.collect();
- let codex_composition: CodexComposition = compose_codex_segments(codex_entries);
+ let codex_composition: CodexComposition = match sources
+ .iter()
+ .find(|source| source.provider_name() == Some("codex"))
+ {
+ Some(source) => source.compose_codex_segments(codex_entries),
+ None => compose_codex_segments(codex_entries),
+ };
let mut items: Vec = cache
.values()
.filter(|entry| entry.source_name.as_deref() != Some("codex"))
@@ -2591,6 +2644,8 @@ fn refresh_snapshot(
amplifier_root_dirs,
unresolved_codex_identities: codex_composition.unresolved_identities,
codex_segment_paths: codex_composition.segment_paths,
+ codex_history_segments: codex_composition.history_segments,
+ codex_history_revisions: codex_composition.history_revisions,
}
}
@@ -3811,6 +3866,8 @@ pub(crate) mod tests {
scan_failures: HashSet::new(),
unresolved_codex_identities: Arc::new(Vec::new()),
codex_segment_paths: Arc::new(HashMap::new()),
+ codex_history_segments: Arc::new(HashMap::new()),
+ codex_history_revisions: Arc::new(HashMap::new()),
})));
let file_cache = Arc::new(StdMutex::new(HashMap::new()));
let direct_cache = Arc::new(StdMutex::new(HashMap::new()));
@@ -4538,9 +4595,9 @@ pub(crate) mod tests {
.replace("\"thread_source\":\"user\"", "\"thread_source\":\"subagent\""),
),
(
- "conflicting CLI version",
+ "empty CLI version",
older.clone(),
- valid_newer.replace("\"cli_version\":\"0.156.0\"", "\"cli_version\":\"0.155.0\""),
+ valid_newer.replace("\"cli_version\":\"0.156.0\"", "\"cli_version\":\"\""),
),
(
"conflicting originator",
@@ -8112,4 +8169,165 @@ pub(crate) mod tests {
);
std::fs::remove_dir_all(&home).ok();
}
+
+ struct MutableCodexCompositionSource {
+ composition: StdMutex,
+ }
+
+ impl SessionSource for MutableCodexCompositionSource {
+ fn discover(&self) -> Vec {
+ Vec::new()
+ }
+
+ fn parse(&self, _path: &Path) -> Option {
+ None
+ }
+
+ fn provider_name(&self) -> Option<&'static str> {
+ Some("codex")
+ }
+
+ fn compose_codex_segments(&self, _entries: Vec) -> CodexComposition {
+ let composition = self.composition.lock().unwrap();
+ CodexComposition {
+ items: composition.items.clone(),
+ unresolved_identities: composition.unresolved_identities.clone(),
+ segment_paths: composition.segment_paths.clone(),
+ history_segments: composition.history_segments.clone(),
+ history_revisions: composition.history_revisions.clone(),
+ }
+ }
+ }
+
+ #[tokio::test]
+ async fn codex_composition_sidecars_publish_and_notify_without_parse_cache_changes() {
+ use crate::codex_history::{CodexHistoryBoundary, CodexHistorySegment};
+
+ let home = unique_temp_dir("codex-sidecar-generation");
+ let (older, _) = codex_continuation_fixtures();
+ let path = home.join("rollout.jsonl");
+ std::fs::create_dir_all(&home).unwrap();
+ std::fs::write(&path, &older).unwrap();
+ let row = parse_codex_file(&path).unwrap();
+ let session_id = row.session_id.clone();
+ let segment = CodexHistorySegment {
+ path: path.clone(),
+ end: None,
+ };
+ let source = Arc::new(MutableCodexCompositionSource {
+ composition: StdMutex::new(CodexComposition {
+ items: vec![row.clone()],
+ segment_paths: HashMap::from([(session_id.clone(), vec![path.clone()])]),
+ history_segments: HashMap::from([(session_id.clone(), vec![segment])]),
+ history_revisions: HashMap::from([(session_id.clone(), 1)]),
+ ..Default::default()
+ }),
+ });
+ let index = SessionIndex::with_ttl_and_cache_path(
+ vec![source.clone()],
+ Duration::from_secs(3600),
+ None,
+ );
+ let initial = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ assert_eq!(initial.sessions.as_ref(), &vec![row]);
+ assert_eq!(initial.codex_history_revisions.get(&session_id), Some(&1));
+ let mut changes = index.subscribe_changes();
+ changes.borrow_and_update();
+
+ // An archive's selected contents changed, but its projected row and bounds did not.
+ source
+ .composition
+ .lock()
+ .unwrap()
+ .history_revisions
+ .insert(session_id.clone(), 2);
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ assert!(changes.has_changed().unwrap());
+ changes.borrow_and_update();
+ let changed = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ assert_eq!(changed.sessions, initial.sessions);
+ assert_eq!(changed.codex_history_revisions.get(&session_id), Some(&2));
+ assert_eq!(initial.codex_history_revisions.get(&session_id), Some(&1));
+
+ // A DB pointer selected a shorter history whose display metadata is unchanged.
+ let boundary = CodexHistoryBoundary {
+ end_byte_offset: 123,
+ end_ordinal_exclusive: 1,
+ };
+ source
+ .composition
+ .lock()
+ .unwrap()
+ .history_segments
+ .get_mut(&session_id)
+ .unwrap()[0]
+ .end = Some(boundary);
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ assert!(changes.has_changed().unwrap());
+ changes.borrow_and_update();
+ let bounded = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ assert_eq!(bounded.sessions, initial.sessions);
+ assert_eq!(
+ bounded.codex_history_segments[&session_id][0].end,
+ Some(boundary)
+ );
+ assert!(initial.codex_history_segments[&session_id][0].end.is_none());
+
+ source
+ .composition
+ .lock()
+ .unwrap()
+ .unresolved_identities
+ .push(CodexUnresolvedIdentity {
+ session_id: session_id.clone(),
+ paths: vec![path],
+ });
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ assert!(changes.has_changed().unwrap());
+ changes.borrow_and_update();
+
+ // A repeated DB notification with unchanged effective history must not wake clients.
+ index.mark_provider_dirty("codex");
+ index.wait_for_refresh_idle_for_test().await;
+ assert!(!changes.has_changed().unwrap());
+ assert!(index.file_cache.lock().unwrap().is_empty());
+ std::fs::remove_dir_all(&home).unwrap();
+ }
+
+ #[tokio::test]
+ async fn codex_source_composes_continuations_across_cli_version_updates() {
+ let home = unique_temp_dir("codex-cli-update-continuation");
+ let (older, newer) = codex_continuation_fixtures();
+ let newer = newer.replace("\"cli_version\":\"0.156.0\"", "\"cli_version\":\"0.159.3\"");
+ let (older_path, newer_path) = write_codex_pair(&home, &older, &newer);
+ let source = Arc::new(CodexSource::new(home.join(".codex")));
+ let direct = source.scan();
+ assert_eq!(direct.len(), 1);
+ assert_eq!(direct[0].source_file.as_ref(), Some(&newer_path));
+ assert_eq!(
+ direct[0].first_user_message.as_deref(),
+ Some("Older first request")
+ );
+ let index =
+ SessionIndex::with_ttl_and_cache_path(vec![source], Duration::from_secs(60), None);
+ let snapshot = index
+ .snapshot_with_failures_and_unresolved_codex_identities()
+ .await;
+ assert_eq!(snapshot.sessions.len(), 1);
+ assert!(snapshot.unresolved_codex_identities.is_empty());
+ assert_eq!(
+ snapshot.codex_segment_paths[&direct[0].session_id],
+ vec![older_path, newer_path]
+ );
+ std::fs::remove_dir_all(&home).unwrap();
+ }
}
diff --git a/crates/freshell-sessions/src/lib.rs b/crates/freshell-sessions/src/lib.rs
index 70832c032..c7177084a 100644
--- a/crates/freshell-sessions/src/lib.rs
+++ b/crates/freshell-sessions/src/lib.rs
@@ -20,6 +20,7 @@
pub mod amplifier;
pub mod amplifier_stub;
pub mod bundle_config;
+pub mod codex_history;
pub mod codex_locator;
pub mod codex_segments;
pub mod directory_index;
diff --git a/crates/freshell-sessions/src/provider_layout.rs b/crates/freshell-sessions/src/provider_layout.rs
index d9b181466..1a2c19e41 100644
--- a/crates/freshell-sessions/src/provider_layout.rs
+++ b/crates/freshell-sessions/src/provider_layout.rs
@@ -22,6 +22,14 @@ pub enum WatchMode {
NonRecursive,
}
+/// Whether a relevant change reparses one discovered file or recomposes
+/// the provider, including dependencies outside its discovery tree.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum WatchChangeScope {
+ File,
+ Provider,
+}
+
/// A provider's on-disk session layout. One trait, one source of truth.
pub trait ProviderLayout: Send + Sync {
/// Short provider name: "claude", "codex", "opencode", "amplifier".
@@ -40,6 +48,21 @@ pub trait ProviderLayout: Send + Sync {
/// Whether watcher targets should be watched recursively or not.
fn watch_mode(&self) -> WatchMode;
+ /// Override the mode for providers with both metadata and history roots.
+ fn watch_mode_for_base(&self, _home: &Path, _base: &Path) -> WatchMode {
+ self.watch_mode()
+ }
+
+ /// Ignore unrelated provider state, or select the refresh scope.
+ /// Existing providers retain their file-versus-directory behavior.
+ fn watch_change_scope(&self, _home: &Path, path: &Path) -> Option {
+ Some(if self.qualifies(path) {
+ WatchChangeScope::File
+ } else {
+ WatchChangeScope::Provider
+ })
+ }
+
/// Does `path` look like a session file for this provider? Used by the
/// session watcher to filter raw inotify events down to relevant changes.
/// Does NOT check whether the file exists or is readable — pure path shape.
@@ -120,13 +143,58 @@ impl ProviderLayout for CodexLayout {
}
fn watch_bases(&self, home: &Path) -> Vec {
- vec![self.session_root(home)]
+ // Keep the home watch permanently: it catches SQLite inode replacement
+ // and late history roots without recursive config/log churn.
+ vec![
+ home.to_path_buf(),
+ self.session_root(home),
+ home.join("archived_sessions"),
+ ]
}
fn watch_mode(&self) -> WatchMode {
WatchMode::Recursive
}
+ fn watch_mode_for_base(&self, home: &Path, base: &Path) -> WatchMode {
+ if base == home {
+ WatchMode::NonRecursive
+ } else {
+ WatchMode::Recursive
+ }
+ }
+
+ fn watch_change_scope(&self, home: &Path, path: &Path) -> Option {
+ let active = self.session_root(home);
+ let archived = home.join("archived_sessions");
+ if path == home || path == active || path == archived {
+ return Some(WatchChangeScope::Provider);
+ }
+ if path.starts_with(&active) {
+ return if self.qualifies(path) {
+ Some(WatchChangeScope::File)
+ } else if path.extension().is_none() {
+ Some(WatchChangeScope::Provider)
+ } else {
+ None
+ };
+ }
+ if path.starts_with(&archived) {
+ // These files are bounded ancestors, not independently listed rows.
+ return (self.qualifies(path) || path.extension().is_none())
+ .then_some(WatchChangeScope::Provider);
+ }
+ if path.parent() == Some(home)
+ && matches!(
+ path.file_name().and_then(|name| name.to_str()),
+ Some("state_5.sqlite" | "state_5.sqlite-wal")
+ )
+ {
+ return Some(WatchChangeScope::Provider);
+ }
+ None
+ }
+
fn qualifies(&self, path: &Path) -> bool {
path.extension().and_then(|s| s.to_str()) == Some("jsonl")
}
@@ -378,4 +446,90 @@ mod tests {
let hardcoded_root = home.join("projects");
assert_eq!(layout_root, hardcoded_root);
}
+
+ #[test]
+ fn codex_watch_modes_and_change_scopes_keep_active_rows_separate_from_references() {
+ let layout = CodexLayout;
+ let home = Path::new("/home/user/.codex");
+ let bases = layout.watch_bases(home);
+ assert_eq!(
+ bases,
+ vec![
+ home.to_path_buf(),
+ home.join("sessions"),
+ home.join("archived_sessions")
+ ]
+ );
+ assert_eq!(
+ layout.watch_mode_for_base(home, &bases[0]),
+ WatchMode::NonRecursive
+ );
+ assert_eq!(
+ layout.watch_mode_for_base(home, &bases[1]),
+ WatchMode::Recursive
+ );
+ assert_eq!(
+ layout.watch_mode_for_base(home, &bases[2]),
+ WatchMode::Recursive
+ );
+ for path in [
+ home.to_path_buf(),
+ home.join("sessions"),
+ home.join("sessions/2026/10"),
+ home.join("archived_sessions"),
+ home.join("archived_sessions/rollout.jsonl"),
+ home.join("state_5.sqlite"),
+ home.join("state_5.sqlite-wal"),
+ ] {
+ assert_eq!(
+ layout.watch_change_scope(home, &path),
+ Some(WatchChangeScope::Provider),
+ "provider dependency {path:?}"
+ );
+ }
+ assert_eq!(
+ layout.watch_change_scope(home, &home.join("sessions/2026/10/rollout.jsonl")),
+ Some(WatchChangeScope::File)
+ );
+ for path in [
+ home.join("config.toml"),
+ home.join("log/codex-tui.log"),
+ home.join("history.jsonl"),
+ home.join("auth.json"),
+ home.join("unrelated.sqlite"),
+ ] {
+ assert_eq!(
+ layout.watch_change_scope(home, &path),
+ None,
+ "unrelated Codex state {path:?}"
+ );
+ }
+ }
+
+ #[test]
+ fn provider_watch_scope_defaults_preserve_existing_providers() {
+ let home = Path::new("/home/user/.claude");
+ let claude = ClaudeLayout;
+ assert_eq!(
+ claude.watch_mode_for_base(home, &home.join("projects")),
+ claude.watch_mode()
+ );
+ assert_eq!(
+ claude.watch_change_scope(home, &home.join("projects/project/session.jsonl")),
+ Some(WatchChangeScope::File)
+ );
+ assert_eq!(
+ claude.watch_change_scope(home, &home.join("projects/project")),
+ Some(WatchChangeScope::Provider)
+ );
+ let opencode = OpencodeLayout;
+ assert_eq!(
+ opencode.watch_mode_for_base(home, home),
+ opencode.watch_mode()
+ );
+ assert_eq!(
+ opencode.watch_change_scope(home, &home.join("opencode.db-wal")),
+ Some(WatchChangeScope::File)
+ );
+ }
}
diff --git a/crates/freshell-sessions/src/search.rs b/crates/freshell-sessions/src/search.rs
index 8da06ea08..0e552128f 100644
--- a/crates/freshell-sessions/src/search.rs
+++ b/crates/freshell-sessions/src/search.rs
@@ -115,7 +115,27 @@ pub fn search_session_file(
tier: FileSearchTier,
) -> std::io::Result
}>
+
+ ,
+ )
+ const fallback = screen.getByText(text)
+ const observation = expectRenderedText(text).catch((error: unknown) => error)
+
+ await act(async () => renderer.resolve({ default: MarkdownRenderer }))
+
+ expect(fallback).not.toBeInTheDocument()
+ expect(screen.getByText(text)).toBeInTheDocument()
+ expect(await observation).toBeUndefined()
+ })
+})
+
function freshopencodeSnapshot(text: string, revision: number) {
return {
sessionType: 'freshopencode',
@@ -1041,7 +1083,7 @@ describe('new conversation close acceptance', () => {
const fixture = prepare(surface, scope)
let closing: ReturnType | undefined
try {
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
const reads = apiMock.getFreshAgentThreadSnapshot.mock.calls.length
const startNew = () => fireEvent.click(screen.getByRole('button', { name: 'Start new conversation' }))
if (timing === 'already pending') act(() => { closing = fixture.startClose() })
@@ -1076,7 +1118,7 @@ describe('new conversation close acceptance', () => {
const fixture = prepare(surfaces[2], 'tab')
let closing: ReturnType | undefined
try {
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
act(() => { closing = fixture.startClose() })
fireEvent.click(screen.getByRole('button', { name: 'Start new conversation' }))
await act(async () => fixture.stop.resolve(stopped))
@@ -1156,7 +1198,7 @@ describe('new conversation close acceptance', () => {
const acknowledgeKill = () => fixture.emit({ type: 'freshAgent.killed', sessionId: fixture.content.sessionRef!.sessionId,
sessionType: 'freshcodex', provider: 'codex', success: true })
try {
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
fireEvent.click(screen.getByRole('button', { name: 'Start new session', exact: true }))
expect(sentFreshAgentMessages('freshAgent.kill')).toHaveLength(1)
expect(fixture.getContent().createRequestId).toBe(fixture.content.createRequestId)
@@ -1185,12 +1227,12 @@ describe('new conversation close acceptance', () => {
it('refuses a late managed cleanup after the displayed conversation source changes', async () => {
const fixture = prepare(surfaces[2], 'tab')
try {
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
fireEvent.click(screen.getByRole('button', { name: 'Start new conversation' }))
const replacement = { ...fixture.content, createRequestId: 'replacement-close-race-create',
soulId: 'replacement-close-race-soul', soulIntentRevision: 18 }
await act(async () => fixture.store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: replacement })))
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
await act(async () => fixture.stop.resolve(stopped))
expect(fixture.getContent()).toMatchObject(replacement)
expect(screen.getByText(historyText)).toBeInTheDocument()
@@ -1205,7 +1247,7 @@ describe('new conversation close acceptance', () => {
const fixture = prepare(surfaces[2], 'tab', 'unmanaged')
const history = createDeferred()
try {
- expect(await screen.findByText(historyText)).toBeInTheDocument()
+ await expectRenderedText(historyText)
const reads = apiMock.getFreshAgentThreadSnapshot.mock.calls.length
apiMock.getFreshAgentThreadSnapshot.mockReturnValue(history.promise)
act(() => fixture.store.dispatch(requestPaneRefresh({ tabId: 'tab-1', paneId: 'pane-1' })))
@@ -1269,7 +1311,7 @@ describe('FreshAgentView', () => {
sessionType: 'freshcodex', provider: 'codex', sessionId: 'close-retained-thread',
}))
render()
- expect(await screen.findByText('Saved conversation remains here')).toBeInTheDocument()
+ await expectRenderedText('Saved conversation remains here')
expect(screen.queryByText(/Close failed/)).toBeNull()
await act(async () => { await store.dispatch(closeTab('tab-1')) })
const notice = await screen.findByText('Close failed: The pane could not be closed, so it was left open. Try again.')
@@ -1969,7 +2011,7 @@ describe('FreshAgentView', () => {
,
)
- expect(await screen.findByText('Visible transcript answer')).toBeInTheDocument()
+ await expectRenderedText('Visible transcript answer')
expect(screen.queryByText('Do not pin this session summary')).not.toBeInTheDocument()
})
@@ -2295,7 +2337,7 @@ describe('FreshAgentView', () => {
expect.objectContaining({ cwd: '/repo/from-ref' }),
)
})
- expect(await screen.findByText('Codex turn')).toBeInTheDocument()
+ await expectRenderedText('Codex turn')
})
it('restores a fresh-agent split pane remount without creating a replacement session', async () => {
@@ -7058,7 +7100,7 @@ describe('FreshAgentView', () => {
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
const answer = `Saved native ${provider === 'codex' ? 'Codex' : 'OpenCode'} answer`
- expect(await screen.findByText(answer)).toBeInTheDocument()
+ await expectRenderedText(answer)
apiMock.getFreshAgentThreadSnapshot.mockRejectedValue(new ApiError(404, message, { code: 'FRESH_AGENT_LOST_SESSION' }))
const recovering = { ...content, recoverySummary: { ...content.recoverySummary, recoveryState: 'recovering' as const } }
wsMock.send.mockClear()
@@ -7292,7 +7334,7 @@ describe('FreshAgentView', () => {
const text = native.turns.flatMap((turn) => turn.items).find((item) => item.kind === 'text') as {text: string}
expect(screen.queryByText(text.text)).not.toBeInTheDocument()
await act(async () => owned.resolve({ ...native, extensions: { [provider]: { nativeHistoryAvailable: true, ownerKind: 'vacant' } } }))
- expect(await screen.findByText(text.text)).toBeInTheDocument()
+ await expectRenderedText(text.text)
expect(composer).toBeDisabled()
expect(composer).toHaveValue('Owned outage draft')
expect(getFreshAgentPaneContent(store)).toEqual(before)
@@ -7444,7 +7486,7 @@ describe('FreshAgentView', () => {
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
apiMock.getFreshAgentThreadSnapshot.mockResolvedValue({ ...native, capabilities: { ...native.capabilities, send: true }, extensions: {} })
render()
- expect(await screen.findByText('Saved native Codex answer')).toBeInTheDocument()
+ await expectRenderedText('Saved native Codex answer')
const pending = createDeferred()
apiMock.getFreshAgentThreadSnapshot.mockReturnValue(pending.promise)
apiMock.getFreshAgentThreadSnapshot.mockClear()
@@ -7706,7 +7748,7 @@ describe('FreshAgentView', () => {
}
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Saved conversation before recovery')).toBeInTheDocument()
+ await expectRenderedText('Saved conversation before recovery')
expect(apiMock.getFreshAgentThreadSnapshot).toHaveBeenCalledWith(sessionType, provider, sessionId, expect.objectContaining({ soulId: 'saved-history-soul' }))
if (recoveryState === 'recovering') {
expect(screen.queryByTestId('managed-runtime-recovery-card')).not.toBeInTheDocument()
@@ -7736,7 +7778,7 @@ describe('FreshAgentView', () => {
durabilityState: 'resume_captured' as const, allocationState: 'verified_durable' as const } }
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText(`Saved native ${provider === 'claude' ? 'Claude' : provider === 'codex' ? 'Codex' : 'OpenCode'} answer`)).toBeInTheDocument()
+ await expectRenderedText(`Saved native ${provider === 'claude' ? 'Claude' : provider === 'codex' ? 'Codex' : 'OpenCode'} answer`)
expect(apiMock.getFreshAgentThreadSnapshot).toHaveBeenCalledWith(sessionType, provider, history.threadId, expect.objectContaining({ soulId: content.soulId }))
expect(getFreshAgentPaneContent(store)).toEqual(content)
expect(sentFreshAgentMessages('freshAgent.create')).toHaveLength(0)
@@ -7759,7 +7801,7 @@ describe('FreshAgentView', () => {
const tool = await screen.findByRole('button', { name: 'apply_patch tool call' })
expect(screen.getAllByRole('button', { name: 'apply_patch tool call' })).toHaveLength(1)
fireEvent.click(tool)
- expect(await screen.findByText(/Patch saved/)).toBeInTheDocument()
+ await expectRenderedText(/Patch saved/)
expect(screen.getByText(/\*\*\* Begin Patch/, { selector: 'pre' })).toBeInTheDocument()
expect(apiMock.getFreshAgentThreadSnapshot).toHaveBeenCalledWith('freshcodex', 'codex', history.threadId, expect.objectContaining({ soulId: content.soulId }))
expect(getFreshAgentPaneContent(store)).toEqual(content)
@@ -7782,13 +7824,13 @@ describe('FreshAgentView', () => {
createRequestId: 'revision-source-request', status: 'idle' as const, soulId: 'revision-source-soul', soulIntentRevision: 5 }
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Previously loaded live answer')).toBeInTheDocument()
+ await expectRenderedText('Previously loaded live answer')
const loaded = getFreshAgentPaneContent(store)
apiMock.getFreshAgentThreadSnapshot.mockResolvedValue(native)
wsMock.send.mockClear()
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: { ...loaded,
recoverySummary: { desiredState: 'running', recoveryState: 'blocked', durabilityState: 'resume_captured', allocationState: 'verified_durable' } } })))
- expect(await screen.findByText(`Saved native ${provider === 'codex' ? 'Codex' : 'OpenCode'} answer`)).toBeInTheDocument()
+ await expectRenderedText(`Saved native ${provider === 'codex' ? 'Codex' : 'OpenCode'} answer`)
expect(screen.queryByText('Previously loaded live answer')).not.toBeInTheDocument()
expect(getFreshAgentPaneContent(store).sessionId).toBe(native.threadId)
expect(wsMock.send).not.toHaveBeenCalledWith(expect.objectContaining({ type: expect.stringMatching(/^freshAgent\.|^pane\.reconcile/) }))
@@ -7797,7 +7839,7 @@ describe('FreshAgentView', () => {
turns: [{ ...live.turns[0], items: [{ id: 'resumed-text', kind: 'text', text: 'Resumed live answer' }] }] })
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: { ...loaded,
recoverySummary: { desiredState: 'running', recoveryState: 'healthy', durabilityState: 'resume_captured', allocationState: 'verified_durable' } } })))
- expect(await screen.findByText('Resumed live answer')).toBeInTheDocument()
+ await expectRenderedText('Resumed live answer')
})
it.each(([
@@ -7822,7 +7864,7 @@ describe('FreshAgentView', () => {
recoverySummary: { desiredState: 'running' as const, recoveryState,
durabilityState: 'resume_captured' as const, allocationState: 'verified_durable' as const } }
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: blocked })))
- expect(await screen.findByText('Saved native Codex answer')).toBeInTheDocument()
+ await expectRenderedText('Saved native Codex answer')
const identity = getFreshAgentPaneContent(store)
wsMock.send.mockClear()
await act(async () => {
@@ -7881,7 +7923,7 @@ describe('FreshAgentView', () => {
expect(wsMock.send).not.toHaveBeenCalledWith(expect.objectContaining({ type: expect.stringMatching(/^freshAgent\.|^pane\.reconcile/) }))
// The initial saved-history read is still useful and has no live actor authority.
await act(async () => resolveNative(native))
- expect(await screen.findByText('Saved native Codex answer')).toBeInTheDocument()
+ await expectRenderedText('Saved native Codex answer')
expect(getFreshAgentPaneContent(store)).toEqual(current)
})
@@ -7897,7 +7939,7 @@ describe('FreshAgentView', () => {
durabilityState: 'resume_captured' as const, allocationState: 'verified_durable' as const } }
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Saved native OpenCode answer')).toBeInTheDocument()
+ await expectRenderedText('Saved native OpenCode answer')
const empty = { ...native, revision: 0, latestTurnId: null, turns: [],
extensions: { opencode: { ownerKind: 'vacant', ownerEpoch: 1, ownerGeneration: 2 } } }
apiMock.getFreshAgentThreadSnapshot.mockResolvedValue(empty)
@@ -7932,7 +7974,7 @@ describe('FreshAgentView', () => {
turns: [{ id: 'new-live', turnId: 'new-live', role: 'assistant', summary: '', items: [{ id: 'new-live-text', kind: 'text', text: 'New resumed live answer' }] }] })
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: { ...content,
recoverySummary: { ...content.recoverySummary, recoveryState: 'healthy' } } })))
- expect(await screen.findByText('New resumed live answer')).toBeInTheDocument()
+ await expectRenderedText('New resumed live answer')
await act(async () => resolveNative(native))
expect(screen.getByText('New resumed live answer')).toBeInTheDocument()
expect(screen.queryByText('Saved native Codex answer')).not.toBeInTheDocument()
@@ -7959,7 +8001,7 @@ describe('FreshAgentView', () => {
fireEvent.click(screen.getByRole('button', { name: 'Retry recovery' }))
await waitFor(() => expect(getFreshAgentPaneContent(store).status).toBe('starting'))
await act(async () => resolveHistory(savedCodexNativeHistory))
- expect(await screen.findByText('Saved native Codex answer')).toBeInTheDocument()
+ await expectRenderedText('Saved native Codex answer')
expect(getFreshAgentPaneContent(store).status).toBe('starting')
expect(getFreshAgentPaneContent(store).resumeSessionId).toBeUndefined()
})
@@ -7980,7 +8022,7 @@ describe('FreshAgentView', () => {
next.turns[1].items[0] = { id: 'revision-answer', kind: 'text', text: 'Current revision history' }
apiMock.getFreshAgentThreadSnapshot.mockResolvedValue(next)
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: { ...content, soulIntentRevision: 2 } })))
- expect(await screen.findByText('Current revision history')).toBeInTheDocument()
+ await expectRenderedText('Current revision history')
await act(async () => resolveOld(savedCodexNativeHistory))
expect(screen.getByText('Current revision history')).toBeInTheDocument()
expect(screen.queryByText('Saved native Codex answer')).not.toBeInTheDocument()
@@ -8002,7 +8044,7 @@ describe('FreshAgentView', () => {
next.turns[1].items[0] = { id: 'next-answer', kind: 'text', text: 'Current soul answer' }
apiMock.getFreshAgentThreadSnapshot.mockResolvedValue(next)
act(() => store.dispatch(updatePaneContent({ tabId: 'tab-1', paneId: 'pane-1', content: { ...content, soulId: 'new-soul' } })))
- expect(await screen.findByText('Current soul answer')).toBeInTheDocument()
+ await expectRenderedText('Current soul answer')
await act(async () => resolveOld(savedCodexNativeHistory))
expect(screen.getByText('Current soul answer')).toBeInTheDocument()
expect(screen.queryByText('Saved native Codex answer')).not.toBeInTheDocument()
@@ -8065,7 +8107,7 @@ describe('FreshAgentView', () => {
}
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Already loaded durable conversation')).toBeInTheDocument()
+ await expectRenderedText('Already loaded durable conversation')
let resolveCold!: (snapshot: typeof cold) => void
apiMock.getFreshAgentThreadSnapshot.mockReturnValue(new Promise((resolve) => { resolveCold = resolve }))
wsMock.send.mockClear()
@@ -8387,7 +8429,7 @@ describe('FreshAgentView', () => {
durabilityState: 'resume_captured' as const, allocationState: 'verified_durable' as const } }
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Conversation retained before starting new')).toBeInTheDocument()
+ await expectRenderedText('Conversation retained before starting new')
expect(screen.getByRole('note', { name: 'Retained conversation history' })).toHaveTextContent(
'Showing retained conversation history. Older turns or large content were omitted from this view. The saved conversation has not been changed.')
expect(screen.getByRole('textbox', { name: 'Chat message input' })).toBeDisabled()
@@ -8410,7 +8452,7 @@ describe('FreshAgentView', () => {
}
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Conversation retained before starting new')).toBeInTheDocument()
+ await expectRenderedText('Conversation retained before starting new')
expect(apiMock.getFreshAgentThreadSnapshot).toHaveBeenCalledWith('freshcodex', 'codex', 'cold-lost-thread', expect.not.objectContaining({ soulId: expect.anything() }))
expect(getFreshAgentPaneContent(store)).toEqual(content)
expect(screen.getByRole('textbox', { name: 'Chat message input' })).toBeDisabled()
@@ -8441,7 +8483,7 @@ describe('FreshAgentView', () => {
}
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Conversation retained before starting new')).toBeInTheDocument()
+ await expectRenderedText('Conversation retained before starting new')
expect(screen.getByText('Close failed: Previous close was not confirmed')).toBeInTheDocument()
wsMock.send.mockClear()
@@ -8478,7 +8520,7 @@ describe('FreshAgentView', () => {
}
store.dispatch(initLayout({ tabId: 'tab-1', paneId: 'pane-1', content }))
render()
- expect(await screen.findByText('Conversation retained before starting new')).toBeInTheDocument()
+ await expectRenderedText('Conversation retained before starting new')
expect(screen.getByText('Close failed: Previous close was not confirmed')).toBeInTheDocument()
const retained = getFreshAgentPaneContent(store)
wsMock.send.mockClear()