Skip to content
Merged
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
53 changes: 44 additions & 9 deletions src/transport/framing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,16 +12,21 @@ pub fn write_ndjson_stdout(payload: &str) -> io::Result<()> {
/// Generic NDJSON writer that works with any `std::io::Write` sink (useful for testing).
pub fn write_ndjson<W: Write>(writer: &mut W, payload: &str) -> io::Result<()> {
// If the JSON payload contains raw '\n' or '\r' bytes (not escaped inside strings),
// replace them with spaces or strip to guarantee exact NDJSON framing.
if payload.contains('\n') || payload.contains('\r') {
let sanitized: String = payload
.chars()
.map(|c| if c == '\n' || c == '\r' { ' ' } else { c })
.collect();
writer.write_all(sanitized.as_bytes())?;
} else {
writer.write_all(payload.as_bytes())?;
// replace them with spaces to guarantee exact NDJSON framing. Both are single-byte
// ASCII code points, so this writes existing byte slices straight to `writer`
// between them instead of collecting a sanitized copy of the whole payload onto
// the heap first (issue #50) — a no-op payload (the common case) is written in
// one `write_all` call, same as before.
let bytes = payload.as_bytes();
let mut start = 0;
for (i, &b) in bytes.iter().enumerate() {
if b == b'\n' || b == b'\r' {
writer.write_all(&bytes[start..i])?;
writer.write_all(b" ")?;
start = i + 1;
}
}
writer.write_all(&bytes[start..])?;
writer.write_all(b"\n")?;
writer.flush()?;
Ok(())
Expand Down Expand Up @@ -56,6 +61,36 @@ mod tests {
);
}

#[test]
fn test_write_ndjson_sanitizes_edge_positions_and_runs() {
// Regression for issue #50: the byte-slice rewrite must handle a
// newline as the very first/last byte and consecutive newlines
// (an empty slice between them) without panicking or dropping bytes.
let mut buffer = Vec::new();
let dirty = "\nleading\r\rmiddle\n\ntrailing\n";
write_ndjson(&mut buffer, dirty).unwrap();
let output = String::from_utf8(buffer).unwrap();
assert_eq!(output, " leading middle trailing \n");
}

#[test]
fn test_write_ndjson_sanitizes_large_payload() {
// A payload well past any small-buffer fast path, to exercise the
// byte-slice rewrite over a realistic multi-megabyte tabular result.
let mut payload = "x".repeat(2 * 1024 * 1024);
payload.push('\n');
payload.push_str(&"y".repeat(1024));

let mut buffer = Vec::new();
write_ndjson(&mut buffer, &payload).unwrap();
let output = String::from_utf8(buffer).unwrap();

assert_eq!(output.len(), payload.len() + 1);
assert_eq!(&output[..2 * 1024 * 1024], "x".repeat(2 * 1024 * 1024));
assert_eq!(output.as_bytes()[2 * 1024 * 1024], b' ');
assert!(output.ends_with(&format!("{}\n", "y".repeat(1024))));
}

#[test]
fn test_write_ndjson_error_payload() {
let mut buffer = Vec::new();
Expand Down
Loading