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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 20 additions & 20 deletions codex-rs/rollout/src/session_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ use codex_protocol::protocol::SessionMetaLine;
use codex_protocol::protocol::SessionSource;
use serde::Deserialize;
use serde::Serialize;
use tokio::io::AsyncBufReadExt;

const SESSION_INDEX_FILE: &str = "session_index.jsonl";
static SESSION_INDEX_LOCK: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::new(()));
Expand Down Expand Up @@ -130,26 +129,27 @@ pub async fn find_thread_names_by_ids(
return Ok(HashMap::new());
}

let file = tokio::fs::File::open(&path).await?;
let reader = tokio::io::BufReader::new(file);
let mut lines = reader.lines();
let mut names = HashMap::with_capacity(thread_ids.len());

while let Some(line) = lines.next_line().await? {
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let Ok(entry) = serde_json::from_str::<SessionIndexEntry>(trimmed) else {
continue;
};
let name = entry.thread_name.trim();
if !name.is_empty() && thread_ids.contains(&entry.id) {
names.insert(entry.id, name.to_string());
let mut remaining_ids = thread_ids.clone();
tokio::task::spawn_blocking(move || {
let mut scanner = ReverseJsonlScanner::new(File::open(path)?)?;
let mut names = HashMap::with_capacity(remaining_ids.len());
while let Some(outcome) = scanner.scan_next::<SessionIndexEntry>()? {
let ScanOutcome::Parsed(entry) = outcome else {
continue;
};
let name = entry.thread_name.trim();
// The first nonempty name seen for an id is its latest usable name.
if !name.is_empty() && remaining_ids.remove(&entry.id) {
names.insert(entry.id, name.to_string());
if remaining_ids.is_empty() {
break;
}
}
}
}

Ok(names)
Ok(names)
})
.await
.map_err(std::io::Error::other)?
}

/// Locate the readable rollout with the newest modification time for a recorded thread name.
Expand Down
53 changes: 53 additions & 0 deletions codex-rs/rollout/src/session_index_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,59 @@ async fn find_thread_names_by_ids_prefers_latest_entry() -> std::io::Result<()>
Ok(())
}

#[tokio::test]
async fn find_thread_names_by_ids_skips_unusable_names_and_reads_across_chunks()
-> std::io::Result<()> {
let temp = TempDir::new()?;
let path = session_index_path(temp.path());
let renamed = ThreadId::new();
let older = ThreadId::new();
let missing = ThreadId::new();
let mut contents = String::new();
for (id, thread_name) in [
(renamed, "original".to_string()),
(older, " older name ".to_string()),
(ThreadId::new(), "unrelated".repeat(/*n*/ 10_000)),
(renamed, " 最新 café ".to_string()),
(renamed, " \t ".to_string()),
] {
let entry = SessionIndexEntry {
id,
thread_name,
updated_at: "2024-01-01T00:00:00Z".to_string(),
};
contents.push_str(&serde_json::to_string(&entry)?);
contents.push_str("\n\nnot json\n");
}
// A partial final write must not hide the latest complete name.
contents.push_str("{\"id\":");
std::fs::write(path, contents)?;

for (ids, expected) in [
(
HashSet::from([renamed]),
HashMap::from([(renamed, "最新 café".to_string())]),
),
(
HashSet::from([renamed, older]),
HashMap::from([
(renamed, "最新 café".to_string()),
(older, "older name".to_string()),
]),
),
(
HashSet::from([renamed, older, missing]),
HashMap::from([
(renamed, "最新 café".to_string()),
(older, "older name".to_string()),
]),
),
] {
assert_eq!(find_thread_names_by_ids(temp.path(), &ids).await?, expected);
}
Ok(())
}

#[test]
fn scan_index_finds_latest_match_among_mixed_entries() -> std::io::Result<()> {
let temp = TempDir::new()?;
Expand Down
Loading