mirror of
https://github.com/openai/codex.git
synced 2026-09-29 16:57:06 +08:00
Attach compressed rollouts to diagnostic reports as JSONL (#44175)
## Why Diagnostic attachment reads assume the queued file path still exists and contains plain bytes. Compressed rollouts can therefore be omitted when only the logical `.jsonl` path is available, or attached as compressed data when a `.jsonl.zst` path is supplied. ## What changed - Read rollout attachments through a bounded decoder that resolves plain or compressed representations without materializing a durable JSONL file. - Use canonical `.jsonl` filenames for attachments and app-server report metadata, while preserving filename overrides. - Apply size limits to decoded bytes and preserve JSONL prefix truncation. ## Testing Add regression tests for compressed attachments, representation changes after queuing, plain-sibling preference, filename overrides, decoded size limits, truncation, nonregular files, and unrelated `.zst` attachments. GitOrigin-RevId: b30f7dd08a741b0c99283460a1ce8933d2920ddf
This commit is contained in:
Generated
+1
@@ -3423,6 +3423,7 @@ dependencies = [
|
||||
"codex-http-client",
|
||||
"codex-login",
|
||||
"codex-protocol",
|
||||
"codex-rollout",
|
||||
"flate2",
|
||||
"http 1.4.0",
|
||||
"httpdate",
|
||||
|
||||
@@ -206,7 +206,7 @@ impl FeedbackRequestProcessor {
|
||||
.await
|
||||
&& seen_attachment_paths.insert(rollout_path.clone())
|
||||
{
|
||||
thread.rollout_filename = rollout_path
|
||||
thread.rollout_filename = codex_rollout::plain_rollout_path(&rollout_path)
|
||||
.file_name()
|
||||
.map(|name| name.to_string_lossy().into_owned());
|
||||
attachment_paths.push(FeedbackAttachmentPath {
|
||||
|
||||
@@ -13,6 +13,7 @@ bytes = { workspace = true }
|
||||
codex-http-client = { workspace = true }
|
||||
codex-login = { workspace = true }
|
||||
codex-protocol = { workspace = true }
|
||||
codex-rollout = { workspace = true }
|
||||
flate2 = { workspace = true }
|
||||
http = { workspace = true }
|
||||
httpdate = { workspace = true }
|
||||
|
||||
@@ -395,6 +395,10 @@ enum AttachmentReadMode {
|
||||
Prefix,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "rollout_attachment_tests.rs"]
|
||||
mod rollout_attachment_tests;
|
||||
|
||||
impl FeedbackAttachmentPath {
|
||||
/// Read a whole regular file within the caller's size limit.
|
||||
pub fn read_attachment(&self, max_bytes: usize) -> io::Result<Option<FeedbackAttachment>> {
|
||||
@@ -406,18 +410,30 @@ impl FeedbackAttachmentPath {
|
||||
max_bytes: usize,
|
||||
mode: AttachmentReadMode,
|
||||
) -> io::Result<Option<FeedbackAttachment>> {
|
||||
let metadata = fs::metadata(&self.path)?;
|
||||
if !metadata.is_file()
|
||||
|| (metadata.len() > max_bytes as u64 && matches!(mode, AttachmentReadMode::Whole))
|
||||
{
|
||||
return Ok(None);
|
||||
}
|
||||
let mut buffer = Vec::new();
|
||||
// Keep one extra byte so the encoder can detect and label a truncated prefix,
|
||||
// including when the file grows after the metadata check.
|
||||
fs::File::open(&self.path)?
|
||||
.take(max_bytes as u64 + 1)
|
||||
.read_to_end(&mut buffer)?;
|
||||
let rollout_path = codex_rollout::rollout_id_from_path(&self.path)
|
||||
.map(|_| codex_rollout::plain_rollout_path(&self.path));
|
||||
let buffer = if rollout_path.is_some() {
|
||||
let Some(buffer) =
|
||||
codex_rollout::read_rollout_prefix(&self.path, max_bytes.saturating_add(1))?
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
buffer
|
||||
} else {
|
||||
let metadata = fs::metadata(&self.path)?;
|
||||
if !metadata.is_file()
|
||||
|| (metadata.len() > max_bytes as u64 && matches!(mode, AttachmentReadMode::Whole))
|
||||
{
|
||||
return Ok(None);
|
||||
}
|
||||
let mut buffer = Vec::new();
|
||||
// Keep one extra byte so the encoder can detect and label a truncated prefix,
|
||||
// including when the file grows after the metadata check.
|
||||
fs::File::open(&self.path)?
|
||||
.take(max_bytes as u64 + 1)
|
||||
.read_to_end(&mut buffer)?;
|
||||
buffer
|
||||
};
|
||||
if buffer.len() > max_bytes && matches!(mode, AttachmentReadMode::Whole) {
|
||||
return Ok(None);
|
||||
}
|
||||
@@ -425,7 +441,9 @@ impl FeedbackAttachmentPath {
|
||||
.attachment_filename_override
|
||||
.clone()
|
||||
.unwrap_or_else(|| {
|
||||
self.path
|
||||
rollout_path
|
||||
.as_ref()
|
||||
.unwrap_or(&self.path)
|
||||
.file_name()
|
||||
.map(|name| name.to_string_lossy().to_string())
|
||||
.unwrap_or_else(|| "extra-log.log".to_string())
|
||||
|
||||
@@ -0,0 +1,232 @@
|
||||
//! Exercises lazy feedback attachment reads across rollout representations.
|
||||
|
||||
use super::*;
|
||||
use pretty_assertions::assert_eq;
|
||||
|
||||
const JSONL: &[u8] = b"{\"message\":\"old diagnostic\"}\n{\"message\":\"more details\"}\n";
|
||||
// A fixed zstd frame for JSONL; the attachment reader does not need an encoder dependency.
|
||||
const COMPRESSED: &[u8] = &[
|
||||
40, 181, 47, 253, 0, 88, 165, 1, 0, 196, 2, 123, 34, 109, 101, 115, 115, 97, 103, 101, 34, 58,
|
||||
34, 111, 108, 100, 32, 100, 105, 97, 103, 110, 111, 115, 116, 105, 99, 34, 125, 10, 109, 111,
|
||||
114, 101, 32, 100, 101, 116, 97, 105, 108, 115, 34, 125, 10, 1, 0, 1, 78, 57, 1,
|
||||
];
|
||||
|
||||
struct Fixture {
|
||||
directory: PathBuf,
|
||||
plain: PathBuf,
|
||||
compressed: PathBuf,
|
||||
}
|
||||
|
||||
impl Fixture {
|
||||
fn new() -> Self {
|
||||
let thread_id = ThreadId::new();
|
||||
let directory = std::env::temp_dir().join(format!("codex-feedback-rollout-{thread_id}"));
|
||||
fs::create_dir(&directory).expect("create fixture directory");
|
||||
let plain = directory.join(format!("rollout-2026-09-09T12-00-00-{thread_id}.jsonl"));
|
||||
let compressed = plain.with_extension("jsonl.zst");
|
||||
fs::write(&compressed, COMPRESSED).expect("write compressed rollout");
|
||||
Self {
|
||||
directory,
|
||||
plain,
|
||||
compressed,
|
||||
}
|
||||
}
|
||||
|
||||
fn attachment(&self) -> FeedbackAttachmentPath {
|
||||
FeedbackAttachmentPath {
|
||||
path: self.plain.clone(),
|
||||
attachment_filename_override: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for Fixture {
|
||||
fn drop(&mut self) {
|
||||
let _ = fs::remove_dir_all(&self.directory);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unloaded_compressed_rollout_is_included_as_jsonl_attachment() {
|
||||
let fixture = Fixture::new();
|
||||
// The app-server supplies the DB's logical .jsonl path without loading this thread.
|
||||
let paths = [fixture.attachment()];
|
||||
let snapshot = CodexFeedback::new().snapshot(/*session_id*/ None);
|
||||
let attachments = snapshot
|
||||
.feedback_attachments(
|
||||
/*include_logs*/ false,
|
||||
&[],
|
||||
&paths,
|
||||
/*logs_override*/ None,
|
||||
)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
attachments
|
||||
.iter()
|
||||
.map(|attachment| (
|
||||
attachment.filename.as_str(),
|
||||
attachment.content_type.as_deref(),
|
||||
attachment.buffer.as_slice(),
|
||||
))
|
||||
.collect::<Vec<_>>(),
|
||||
vec![(
|
||||
fixture.plain.file_name().unwrap().to_str().unwrap(),
|
||||
Some("text/plain"),
|
||||
JSONL
|
||||
)]
|
||||
);
|
||||
assert!(
|
||||
!fixture.plain.exists(),
|
||||
"feedback must not materialize the durable rollout"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compressed_paths_keep_canonical_names_and_filename_overrides() {
|
||||
let mut fixture = Fixture::new();
|
||||
// Reverted rollouts carry a stable thread ID and a separate immutable rollout ID.
|
||||
let plain = fixture.directory.join(format!(
|
||||
"rollout-2026-09-09T12-00-00-{}_{}.jsonl",
|
||||
ThreadId::new(),
|
||||
ThreadId::new()
|
||||
));
|
||||
let compressed = plain.with_extension("jsonl.zst");
|
||||
fs::rename(&fixture.compressed, &compressed).unwrap();
|
||||
fixture.plain = plain;
|
||||
fixture.compressed = compressed;
|
||||
|
||||
for filename_override in [None, Some("reviewer-rollout.jsonl".to_string())] {
|
||||
let path = FeedbackAttachmentPath {
|
||||
path: fixture.compressed.clone(),
|
||||
attachment_filename_override: filename_override.clone(),
|
||||
};
|
||||
let attachment = path.read_attachment(JSONL.len()).unwrap().unwrap();
|
||||
let expected_filename = filename_override.unwrap_or_else(|| {
|
||||
fixture
|
||||
.plain
|
||||
.file_name()
|
||||
.unwrap()
|
||||
.to_str()
|
||||
.unwrap()
|
||||
.to_string()
|
||||
});
|
||||
assert_eq!(
|
||||
(
|
||||
attachment.filename,
|
||||
attachment.content_type,
|
||||
attachment.buffer
|
||||
),
|
||||
(
|
||||
expected_filename,
|
||||
Some("text/plain".to_string()),
|
||||
JSONL.to_vec()
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compressed_attachment_prefers_newer_plain_sibling() {
|
||||
let fixture = Fixture::new();
|
||||
let current = b"{\"message\":\"new diagnostic\"}\n";
|
||||
fs::write(&fixture.plain, current).unwrap();
|
||||
let attachment = FeedbackAttachmentPath {
|
||||
path: fixture.compressed.clone(),
|
||||
attachment_filename_override: None,
|
||||
}
|
||||
.read_attachment(/*max_bytes*/ 1024)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(attachment.buffer, current);
|
||||
assert_eq!(fs::read(&fixture.compressed).unwrap(), COMPRESSED);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn queued_attachments_follow_both_representation_transitions() {
|
||||
let fixture = Fixture::new();
|
||||
let compressed_path = FeedbackAttachmentPath {
|
||||
path: fixture.compressed.clone(),
|
||||
attachment_filename_override: None,
|
||||
};
|
||||
// The compressed file is materialized after the attachment path was queued.
|
||||
fs::write(&fixture.plain, JSONL).unwrap();
|
||||
fs::remove_file(&fixture.compressed).unwrap();
|
||||
assert_eq!(
|
||||
compressed_path
|
||||
.read_attachment(/*max_bytes*/ 1024)
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.buffer,
|
||||
JSONL
|
||||
);
|
||||
let plain_path = fixture.attachment();
|
||||
// The plain file is compressed after the attachment path was queued.
|
||||
fs::write(&fixture.compressed, COMPRESSED).unwrap();
|
||||
fs::remove_file(&fixture.plain).unwrap();
|
||||
assert_eq!(
|
||||
plain_path
|
||||
.read_attachment(/*max_bytes*/ 1024)
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.buffer,
|
||||
JSONL
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compressed_attachment_bounds_decoded_bytes_and_preserves_jsonl_truncation() {
|
||||
let fixture = Fixture::new();
|
||||
let path = fixture.attachment();
|
||||
assert!(path.read_attachment(/*max_bytes*/ 35).unwrap().is_none());
|
||||
assert_eq!(
|
||||
path.read_attachment(JSONL.len()).unwrap().unwrap().buffer,
|
||||
JSONL
|
||||
);
|
||||
// The valid first frame provides the needed prefix. Decoding the entire file would fail.
|
||||
fs::write(
|
||||
&fixture.compressed,
|
||||
[COMPRESSED, b"invalid trailing frame"].concat(),
|
||||
)
|
||||
.unwrap();
|
||||
let mut attachment = path
|
||||
.read_attachment_with_mode(/*max_bytes*/ 35, AttachmentReadMode::Prefix)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(attachment.buffer, JSONL[..36]);
|
||||
crate::attachment_truncation::truncate_attachment(
|
||||
&mut attachment.filename,
|
||||
&mut attachment.buffer,
|
||||
/*target_bytes*/ 35,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(attachment.buffer, b"{\"message\":\"old diagnostic\"}\n");
|
||||
assert!(attachment.filename.starts_with("truncated-rollout-"));
|
||||
assert!(attachment.filename.ends_with(".jsonl"));
|
||||
assert!(!fixture.plain.exists());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn nonregular_rollout_is_omitted_and_unrelated_zstd_attachment_stays_opaque() {
|
||||
let fixture = Fixture::new();
|
||||
fs::create_dir(&fixture.plain).unwrap();
|
||||
assert!(
|
||||
fixture
|
||||
.attachment()
|
||||
.read_attachment(/*max_bytes*/ 1024)
|
||||
.unwrap()
|
||||
.is_none()
|
||||
);
|
||||
let opaque_path = fixture.directory.join("diagnostics.zst");
|
||||
fs::rename(&fixture.compressed, &opaque_path).unwrap();
|
||||
let attachment = FeedbackAttachmentPath {
|
||||
path: opaque_path,
|
||||
attachment_filename_override: None,
|
||||
}
|
||||
.read_attachment(/*max_bytes*/ 1024)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
(attachment.filename.as_str(), attachment.buffer.as_slice()),
|
||||
("diagnostics.zst", COMPRESSED)
|
||||
);
|
||||
}
|
||||
@@ -99,6 +99,7 @@ pub use compression::open_rollout_line_reader;
|
||||
pub use compression::plain_rollout_path;
|
||||
pub use compression::spawn_rollout_compression_worker;
|
||||
pub use seekable_reader::open_rollout_seekable_reader;
|
||||
pub use seekable_reader::read_rollout_prefix;
|
||||
pub use seekable_reader::rollout_contains_prefix;
|
||||
|
||||
/// Materializes a compressed rollout as plain JSONL before another rollout references it.
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
//! Reads rollout byte positions independently of their plain or compressed representation.
|
||||
//! Reads rollout bytes and positions independently of their plain or compressed representation.
|
||||
//!
|
||||
//! Offsets always address the original JSONL bytes. Readers retain an open file, or an anonymous
|
||||
//! decoded snapshot, so a concurrent compression cannot invalidate an in-progress scan.
|
||||
@@ -39,6 +39,30 @@ impl RolloutReader {
|
||||
}
|
||||
}
|
||||
|
||||
/// Reads at most `max_bytes` decoded rollout bytes without materializing the durable file.
|
||||
///
|
||||
/// Returns `None` for nonregular files. This retains an open representation while reading, so
|
||||
/// compression or materialization cannot invalidate the read after resolution.
|
||||
pub fn read_rollout_prefix(path: &Path, max_bytes: usize) -> io::Result<Option<Vec<u8>>> {
|
||||
if std::fs::metadata(path).is_ok_and(|metadata| !metadata.is_file()) {
|
||||
return Ok(None);
|
||||
}
|
||||
let source = RolloutReader::open(path)?;
|
||||
let file = match &source {
|
||||
RolloutReader::Plain(file) | RolloutReader::Compressed(file) => file,
|
||||
};
|
||||
if !file.metadata()?.is_file() {
|
||||
return Ok(None);
|
||||
}
|
||||
let reader: Box<dyn Read> = match source {
|
||||
RolloutReader::Plain(file) => Box::new(file),
|
||||
RolloutReader::Compressed(file) => Box::new(zstd::stream::read::Decoder::new(file)?),
|
||||
};
|
||||
let mut bytes = Vec::new();
|
||||
reader.take(max_bytes as u64).read_to_end(&mut bytes)?;
|
||||
Ok(Some(bytes))
|
||||
}
|
||||
|
||||
/// Opens the original JSONL bytes for blocking offset reads without changing the rollout on disk.
|
||||
///
|
||||
/// Compressed files are decoded into an anonymous temporary file, keeping memory use bounded and
|
||||
|
||||
Reference in New Issue
Block a user