perf(codex): cache parent rollout timelines across fork cutoffs (#5626)

* perf(codex): cache parent rollout timelines

* fix(codex): harden parent timeline cache

* fix(codex): tighten replay cache invalidation

---------

Co-authored-by: Ayanami <ay@nami.ltd>
Co-authored-by: SaladDay <92240037+SaladDay@users.noreply.github.com>
This commit is contained in:
ayanamislover
2026-07-29 16:49:46 +01:00
committed by GitHub
parent 30409878bd
commit 56fb46c093
2 changed files with 349 additions and 38 deletions
+5 -1
View File
@@ -91,7 +91,11 @@ webkit2gtk = { version = "2.0.1", features = ["v2_16"] }
[target.'cfg(target_os = "windows")'.dependencies]
winreg = "0.52"
windows-sys = { version = "0.61", features = ["Win32_Globalization", "Win32_UI_Shell"] }
windows-sys = { version = "0.61", features = [
"Win32_Globalization",
"Win32_Storage_FileSystem",
"Win32_UI_Shell",
] }
[target.'cfg(all(target_os = "windows", target_arch = "aarch64"))'.dependencies]
rquickjs = { version = "0.8", features = ["bindgen"] }
+344 -37
View File
@@ -29,9 +29,17 @@ use rust_decimal::Decimal;
use std::collections::HashMap;
use std::fs;
use std::io::{BufRead, BufReader};
#[cfg(unix)]
use std::os::unix::fs::MetadataExt;
#[cfg(windows)]
use std::os::windows::io::AsRawHandle;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::SystemTime;
#[cfg(windows)]
use windows_sys::Win32::Storage::FileSystem::{
FileIdInfo, GetFileInformationByHandleEx, FILE_ID_INFO,
};
const CODEX_THREAD_REQUEST_ID_PREFIX: &str = "codex_session:thread-v1";
@@ -72,6 +80,115 @@ struct TokenUsageSignature {
last: Option<TokenCountersSignature>,
}
#[derive(Debug)]
struct TimestampedTokenSignature {
timestamp: DateTime<Utc>,
signature: TokenUsageSignature,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ParentFileStamp {
modified_nanos: i64,
size: u64,
#[cfg(unix)]
device: u64,
#[cfg(unix)]
inode: u64,
#[cfg(windows)]
volume_serial: u64,
#[cfg(windows)]
file_id: [u8; 16],
}
impl ParentFileStamp {
fn from_file(file: &fs::File) -> Option<Self> {
let metadata = file.metadata().ok()?;
#[cfg(windows)]
let (volume_serial, file_id) = windows_file_identity(file)?;
Some(Self {
modified_nanos: metadata_modified_nanos(&metadata),
size: metadata.len(),
#[cfg(unix)]
device: metadata.dev(),
#[cfg(unix)]
inode: metadata.ino(),
#[cfg(windows)]
volume_serial,
#[cfg(windows)]
file_id,
})
}
}
#[cfg(windows)]
fn windows_file_identity(file: &fs::File) -> Option<(u64, [u8; 16])> {
let mut information = FILE_ID_INFO::default();
// SAFETY: `file` owns a live handle for this call, and `information` is a
// valid writable FILE_ID_INFO buffer of the size passed to Windows.
let succeeded = unsafe {
GetFileInformationByHandleEx(
file.as_raw_handle(),
FileIdInfo,
std::ptr::addr_of_mut!(information).cast(),
std::mem::size_of::<FILE_ID_INFO>() as u32,
)
} != 0;
succeeded.then_some((
information.VolumeSerialNumber,
information.FileId.Identifier,
))
}
#[derive(Debug)]
struct ParentTokenTimeline {
events: Vec<TimestampedTokenSignature>,
max_timestamp: Option<DateTime<Utc>>,
has_token_without_timestamp: bool,
}
impl ParentTokenTimeline {
fn signatures_before(
&self,
parent_path: &Path,
cutoff: DateTime<Utc>,
) -> Result<Vec<TokenUsageSignature>, String> {
if self.has_token_without_timestamp {
return Err(format!(
"父 rollout {} 的 token_count 缺少有效 timestamp",
parent_path.display()
));
}
if self
.max_timestamp
.is_none_or(|timestamp| timestamp < cutoff)
{
return Err(format!(
"父 rollout {} 尚未写到 child fork 时刻",
parent_path.display()
));
}
Ok(self
.events
.iter()
.filter(|event| event.timestamp <= cutoff)
.map(|event| event.signature.clone())
.collect())
}
}
#[derive(Debug)]
struct CachedParentTimeline {
stamp: ParentFileStamp,
timeline: Arc<ParentTokenTimeline>,
}
#[derive(Debug)]
struct CachedReplayPrefix {
modified: i64,
size: u64,
prefix: usize,
}
#[derive(Debug)]
struct ParsedTokenEvent {
line_offset: i64,
@@ -116,8 +233,8 @@ struct PendingEntry {
#[derive(Debug, Default)]
struct CodexReplayCaches {
parent_signatures: HashMap<(PathBuf, i64), Vec<TokenUsageSignature>>,
replay_prefixes: HashMap<(PathBuf, i64, u64), usize>,
parent_timelines: HashMap<PathBuf, CachedParentTimeline>,
replay_prefixes: HashMap<PathBuf, CachedReplayPrefix>,
pending: HashMap<PathBuf, PendingEntry>,
}
@@ -757,20 +874,28 @@ fn parent_signatures_before(
parent_path: &Path,
cutoff: DateTime<Utc>,
) -> Result<Vec<TokenUsageSignature>, String> {
let cache_key = (parent_path.to_path_buf(), cutoff.timestamp_micros());
if let Ok(caches) = replay_caches().lock() {
if let Some(signatures) = caches.parent_signatures.get(&cache_key) {
return Ok(signatures.clone());
}
}
let file = fs::File::open(parent_path)
.map_err(|error| format!("无法打开父 rollout {}: {error}", parent_path.display()))?;
let mut signatures = Vec::new();
let mut max_timestamp: Option<DateTime<Utc>> = None;
let stamp = ParentFileStamp::from_file(&file);
let cached_timeline = stamp.and_then(|stamp| {
replay_caches().lock().ok().and_then(|caches| {
caches
.parent_timelines
.get(parent_path)
.filter(|entry| entry.stamp == stamp)
.map(|entry| Arc::clone(&entry.timeline))
})
});
if let Some(timeline) = cached_timeline {
return timeline.signatures_before(parent_path, cutoff);
}
// 必须扫描完整父文件并逐行应用 cutoff,不能在首个未来时间戳处 break:
// rollout 写入顺序不承诺时间戳严格单调。
let mut events = Vec::new();
let mut max_timestamp: Option<DateTime<Utc>> = None;
let mut has_token_without_timestamp = false;
// 必须扫描完整父文件,不能在首个未来时间戳处 break:rollout 写入顺序
// 不承诺时间戳严格单调。缓存完整时间线后,不同 child cutoff 只需内存过滤。
for line in BufReader::new(file).lines() {
let Ok(line) = line else {
continue;
@@ -802,29 +927,31 @@ fn parent_signatures_before(
continue;
};
let Some(timestamp) = timestamp else {
return Err(format!(
"父 rollout {} 的 token_count 缺少有效 timestamp",
parent_path.display()
));
has_token_without_timestamp = true;
continue;
};
if timestamp <= cutoff {
signatures.push(signature);
}
events.push(TimestampedTokenSignature {
timestamp,
signature,
});
}
if max_timestamp.is_none_or(|timestamp| timestamp < cutoff) {
return Err(format!(
"父 rollout {} 尚未写到 child fork 时刻",
parent_path.display()
));
let timeline = Arc::new(ParentTokenTimeline {
events,
max_timestamp,
has_token_without_timestamp,
});
let result = timeline.signatures_before(parent_path, cutoff);
if let (Some(stamp), Ok(mut caches)) = (stamp, replay_caches().lock()) {
caches.parent_timelines.insert(
parent_path.to_path_buf(),
CachedParentTimeline {
stamp,
timeline: Arc::clone(&timeline),
},
);
}
if let Ok(mut caches) = replay_caches().lock() {
caches
.parent_signatures
.insert(cache_key, signatures.clone());
}
Ok(signatures)
result
}
fn resolve_parent_signatures(
@@ -993,10 +1120,14 @@ fn sync_single_codex_file(
),
));
};
let cache_key = (file_path.to_path_buf(), file_modified, file_size);
if let Ok(caches) = replay_caches().lock() {
if let Some(prefix) = caches.replay_prefixes.get(&cache_key) {
*prefix
if let Some(prefix) = caches
.replay_prefixes
.get(file_path)
.filter(|cached| cached.modified == file_modified && cached.size == file_size)
.map(|cached| cached.prefix)
{
prefix
} else {
drop(caches);
let parent_signatures =
@@ -1018,7 +1149,14 @@ fn sync_single_codex_file(
};
let prefix = matching_replay_prefix(&parsed.token_events, &parent_signatures);
if let Ok(mut caches) = replay_caches().lock() {
caches.replay_prefixes.insert(cache_key, prefix);
caches.replay_prefixes.insert(
file_path.to_path_buf(),
CachedReplayPrefix {
modified: file_modified,
size: file_size,
prefix,
},
);
}
prefix
}
@@ -1285,6 +1423,15 @@ mod tests {
token_count_at(input, cached, output, "2026-07-10T03:00:02Z")
}
fn token_count_without_timestamp(input: u64, cached: u64, output: u64) -> serde_json::Value {
let mut value = token_count(input, cached, output);
value
.as_object_mut()
.expect("token_count must be an object")
.remove("timestamp");
value
}
fn sync_test_file(
db: &Database,
file: &Path,
@@ -1407,6 +1554,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_thread_spawn_parent_strips_replay_and_keeps_live_usage() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1449,6 +1597,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_filtered_parent_events_use_subsequence_prefix_alignment() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1481,6 +1630,158 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_parent_rollout_is_cached_once_across_fork_cutoffs() -> Result<(), AppError> {
clear_codex_replay_caches();
let temp = tempdir().unwrap();
let parent = rollout_path(temp.path(), PARENT_ID);
write_jsonl(
&parent,
&[
session_meta(PARENT_ID),
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
token_count_at(200, 100, 20, "2026-07-10T03:00:10Z"),
turn_context_at("2026-07-10T03:00:20Z"),
],
);
let early = "2026-07-10T03:00:05Z".parse::<DateTime<Utc>>().unwrap();
let late = "2026-07-10T03:00:15Z".parse::<DateTime<Utc>>().unwrap();
assert_eq!(parent_signatures_before(&parent, early).unwrap().len(), 1);
let first_timeline =
Arc::clone(&replay_caches().lock().unwrap().parent_timelines[&parent].timeline);
assert_eq!(parent_signatures_before(&parent, late).unwrap().len(), 2);
let caches = replay_caches().lock().unwrap();
assert_eq!(caches.parent_timelines.len(), 1);
assert!(Arc::ptr_eq(
&first_timeline,
&caches.parent_timelines[&parent].timeline
));
Ok(())
}
#[test]
#[serial_test::serial]
fn test_parent_rollout_cache_invalidates_after_append() -> Result<(), AppError> {
clear_codex_replay_caches();
let temp = tempdir().unwrap();
let parent = rollout_path(temp.path(), PARENT_ID);
let cutoff = "2026-07-10T03:00:15Z".parse::<DateTime<Utc>>().unwrap();
write_jsonl(
&parent,
&[
session_meta(PARENT_ID),
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
turn_context_at("2026-07-10T03:00:20Z"),
],
);
assert_eq!(parent_signatures_before(&parent, cutoff).unwrap().len(), 1);
write_jsonl(
&parent,
&[
session_meta(PARENT_ID),
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
token_count_at(200, 100, 20, "2026-07-10T03:00:10Z"),
turn_context_at("2026-07-10T03:00:20Z"),
],
);
assert_eq!(parent_signatures_before(&parent, cutoff).unwrap().len(), 2);
let caches = replay_caches().lock().unwrap();
assert_eq!(caches.parent_timelines.len(), 1);
Ok(())
}
#[test]
#[serial_test::serial]
fn test_parent_rollout_content_error_cache_preserves_open_errors() {
clear_codex_replay_caches();
let temp = tempdir().unwrap();
let parent = rollout_path(temp.path(), PARENT_ID);
let cutoff = "2026-07-10T03:00:05Z".parse::<DateTime<Utc>>().unwrap();
write_jsonl(
&parent,
&[
session_meta(PARENT_ID),
token_count_without_timestamp(100, 50, 10),
turn_context_at("2026-07-10T03:00:20Z"),
],
);
let first_error = parent_signatures_before(&parent, cutoff).unwrap_err();
assert!(first_error.contains("token_count 缺少有效 timestamp"));
let cached_timeline =
|| Arc::clone(&replay_caches().lock().unwrap().parent_timelines[&parent].timeline);
let first_timeline = cached_timeline();
let second_error = parent_signatures_before(&parent, cutoff).unwrap_err();
assert_eq!(second_error, first_error);
assert!(Arc::ptr_eq(&first_timeline, &cached_timeline()));
fs::remove_file(&parent).unwrap();
let open_error = parent_signatures_before(&parent, cutoff).unwrap_err();
assert!(open_error.contains("无法打开父 rollout"));
}
#[test]
#[serial_test::serial]
fn test_parent_rollout_nanosecond_cutoffs_are_exact() {
clear_codex_replay_caches();
let temp = tempdir().unwrap();
let parent = rollout_path(temp.path(), PARENT_ID);
write_jsonl(
&parent,
&[
session_meta(PARENT_ID),
token_count_at(100, 50, 10, "2026-07-10T03:00:00.000000500Z"),
turn_context_at("2026-07-10T03:00:00.000000900Z"),
],
);
let before = "2026-07-10T03:00:00.000000300Z"
.parse::<DateTime<Utc>>()
.unwrap();
let after = "2026-07-10T03:00:00.000000700Z"
.parse::<DateTime<Utc>>()
.unwrap();
assert!(parent_signatures_before(&parent, before)
.unwrap()
.is_empty());
assert_eq!(parent_signatures_before(&parent, after).unwrap().len(), 1);
assert_eq!(replay_caches().lock().unwrap().parent_timelines.len(), 1);
}
#[cfg(any(unix, windows))]
#[test]
fn test_parent_file_stamp_distinguishes_same_size_same_mtime_files() {
let temp = tempdir().unwrap();
let parent = rollout_path(temp.path(), PARENT_ID);
let replacement = temp.path().join("replacement.jsonl");
let values = [session_meta(PARENT_ID), token_count(100, 50, 10)];
write_jsonl(&parent, &values);
write_jsonl(&replacement, &values);
let original_file = fs::File::open(&parent).unwrap();
let original_metadata = original_file.metadata().unwrap();
let replacement_file = fs::OpenOptions::new()
.write(true)
.open(&replacement)
.unwrap();
replacement_file
.set_times(fs::FileTimes::new().set_modified(original_metadata.modified().unwrap()))
.unwrap();
let original_stamp = ParentFileStamp::from_file(&original_file).unwrap();
let replacement_stamp = ParentFileStamp::from_file(&replacement_file).unwrap();
assert_eq!(
(original_stamp.size, original_stamp.modified_nanos),
(replacement_stamp.size, replacement_stamp.modified_nanos)
);
assert_ne!(original_stamp, replacement_stamp);
}
#[test]
#[serial_test::serial]
fn test_empty_fork_imports_no_parent_usage() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1526,6 +1827,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_conflicting_explicit_parents_are_deferred() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1551,6 +1853,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_parent_future_signature_cannot_extend_replay_prefix() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1582,6 +1885,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_missing_parent_is_deferred_and_recovered_without_child_change() -> Result<(), AppError>
{
clear_codex_replay_caches();
@@ -1615,6 +1919,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_billable_file_without_meta_is_deferred_without_cursor() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1641,6 +1946,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_non_billable_file_without_meta_advances_cursor() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;
@@ -1661,6 +1967,7 @@ mod tests {
}
#[test]
#[serial_test::serial]
fn test_subagents_use_filename_thread_ids() -> Result<(), AppError> {
clear_codex_replay_caches();
let db = Database::memory()?;