diff --git a/Cargo.lock b/Cargo.lock index b2a035483..704031afe 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2229,6 +2229,7 @@ dependencies = [ "proptest", "prost 0.14.3", "reqwest 0.12.28", + "rustix 1.1.4", "ryu", "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index bf9254b96..947a2758d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -128,6 +128,7 @@ xxhash-rust = { version = "0.8.15", features = ["xxh32", "xxh64"] } rustls = { version = "0.23", default-features = false, features = [ "aws_lc_rs", ] } +rustix = { version = "1", features = ["fs"] } rcgen = "0.14" [workspace.lints.clippy] diff --git a/crates/ffwd-output/Cargo.toml b/crates/ffwd-output/Cargo.toml index 06bdd483a..6f9fef44d 100644 --- a/crates/ffwd-output/Cargo.toml +++ b/crates/ffwd-output/Cargo.toml @@ -16,7 +16,6 @@ ffwd-core = { version = "0.1.0", path = "../ffwd-core" } ffwd-otap-proto = { version = "0.1.0", path = "../ffwd-otap-proto" } ffwd-types = { version = "0.1.0", path = "../ffwd-types" } flate2 = "1" -libc = { workspace = true } httpdate = "1" itoa = "1" memchr = "2" @@ -45,6 +44,12 @@ tracing = { workspace = true } ureq = { version = "3", default-features = false, features = ["rustls"] } zstd = "0.13" +[target.'cfg(target_os = "linux")'.dependencies] +rustix = { workspace = true } + +[target.'cfg(target_os = "macos")'.dependencies] +libc = { workspace = true } + [dev-dependencies] arrow = { workspace = true, features = ["test_utils"] } insta = "1" diff --git a/crates/ffwd-output/src/pipelined.rs b/crates/ffwd-output/src/pipelined.rs index 7aedb4ef6..91034a093 100644 --- a/crates/ffwd-output/src/pipelined.rs +++ b/crates/ffwd-output/src/pipelined.rs @@ -1,7 +1,7 @@ //! Pipelined sink: separates serialization (CPU) from transport (I/O). //! -//! The pipeline uses a dedicated OS thread for I/O, overlapping serialization -//! of batch N+1 with writing of batch N. A triple-buffered pool avoids +//! The pipeline uses a dedicated OS thread for I/O and a small buffer pool to +//! keep blocking file writes off the caller's serialization path without //! allocation churn. //! //! # Architecture @@ -26,7 +26,7 @@ use std::future::Future; use std::io; use std::pin::Pin; use std::sync::Arc; -use std::sync::mpsc::{self, SyncSender}; +use std::sync::mpsc::{self as std_mpsc, SyncSender}; use std::thread::{self, JoinHandle}; use arrow::record_batch::RecordBatch; @@ -110,16 +110,16 @@ impl Default for PipelineConfig { enum WriterMsg { /// A filled buffer ready to be written. - Data(Vec), + Data { + buf: Vec, + ack: tokio::sync::oneshot::Sender>, + }, /// Flush request with a one-shot ack channel. Flush(SyncSender>), /// Shutdown request — writer should flush, sync, then exit. Shutdown(SyncSender>), } -/// Sent from writer thread to signal a write failure. -struct WriteError(io::Error); - // --------------------------------------------------------------------------- // PipelinedSink // --------------------------------------------------------------------------- @@ -134,11 +134,10 @@ pub struct PipelinedSink { filled_tx: Option>, /// Channel to receive empty (recycled) buffers back from the writer. - empty_rx: mpsc::Receiver>, + empty_rx: std_mpsc::Receiver>, - /// Channel to receive write errors from the writer thread. - /// Non-blocking: checked before each send_batch to surface I/O failures. - error_rx: mpsc::Receiver, + /// Caller-owned spare buffer for serialization failures. + spare_buf: Option>, /// Handle to the writer thread (for join on shutdown). writer_handle: Option>, @@ -162,9 +161,8 @@ impl PipelinedSink { )); } - let (filled_tx, filled_rx) = mpsc::sync_channel::(config.num_buffers); - let (empty_tx, empty_rx) = mpsc::sync_channel::>(config.num_buffers); - let (error_tx, error_rx) = mpsc::sync_channel::(config.num_buffers); + let (filled_tx, filled_rx) = std_mpsc::sync_channel::(config.num_buffers); + let (empty_tx, empty_rx) = std_mpsc::sync_channel::>(config.num_buffers); // Seed the empty buffer pool. for _ in 0..config.num_buffers { @@ -174,7 +172,7 @@ impl PipelinedSink { let writer_handle = thread::Builder::new() .name(format!("ffwd-writer-{}", serializer.name())) .spawn(move || { - writer_thread_loop(writer, filled_rx, empty_tx, error_tx); + writer_thread_loop(writer, filled_rx, empty_tx); }) .map_err(io::Error::other)?; @@ -183,48 +181,50 @@ impl PipelinedSink { stats, filled_tx: Some(filled_tx), empty_rx, - error_rx, + spare_buf: None, writer_handle: Some(writer_handle), }) } - /// Check for write errors from the writer thread (non-blocking). - /// - /// Returns the first queued error, if any. This surfaces I/O failures - /// from previous batches — a write error on batch N is reported during - /// `send_batch` for batch N+1. - fn check_writer_error(&self) -> io::Result<()> { - match self.error_rx.try_recv() { - Ok(WriteError(e)) => Err(e), - Err(_) => Ok(()), + /// Try to get an empty buffer from the pool without blocking the async + /// worker thread. + fn acquire_buffer(&mut self) -> io::Result> { + if let Some(buf) = self.spare_buf.take() { + return Ok(buf); } + + match self.empty_rx.try_recv() { + Ok(buf) => Ok(buf), + Err(std_mpsc::TryRecvError::Empty) => Err(io::Error::new( + io::ErrorKind::WouldBlock, + "pipelined writer has no available buffer", + )), + Err(std_mpsc::TryRecvError::Disconnected) => Err(io::Error::new( + io::ErrorKind::BrokenPipe, + "writer thread exited unexpectedly", + )), + } + } + + /// Clone the writer-thread sender before awaiting an acknowledgement. + fn filled_sender(&self) -> io::Result> { + self.filled_tx + .clone() + .ok_or_else(|| io::Error::new(io::ErrorKind::NotConnected, "sink already shut down")) } - /// Get an empty buffer from the pool, blocking until one is available. - fn acquire_buffer(&self) -> io::Result> { - self.empty_rx.recv().map_err(|_recv| { + /// Send a filled buffer to the writer thread and wait for it to be written. + async fn write_filled(tx: SyncSender, buf: Vec) -> io::Result<()> { + let (ack, ack_rx) = tokio::sync::oneshot::channel(); + tx.send(WriterMsg::Data { buf, ack }).map_err(|_send| { io::Error::new( io::ErrorKind::BrokenPipe, "writer thread exited unexpectedly", ) - }) - } - - /// Send a filled buffer to the writer thread. - fn send_filled(&self, buf: Vec) -> io::Result<()> { - if let Some(ref tx) = self.filled_tx { - tx.send(WriterMsg::Data(buf)).map_err(|_send| { - io::Error::new( - io::ErrorKind::BrokenPipe, - "writer thread exited unexpectedly", - ) - }) - } else { - Err(io::Error::new( - io::ErrorKind::NotConnected, - "sink already shut down", - )) - } + })?; + ack_rx.await.map_err(|_recv| { + io::Error::new(io::ErrorKind::BrokenPipe, "writer thread exited before ack") + })? } /// Send a control message and wait for the ack. @@ -232,7 +232,7 @@ impl PipelinedSink { where F: FnOnce(SyncSender>) -> WriterMsg, { - let (ack_tx, ack_rx) = mpsc::sync_channel(1); + let (ack_tx, ack_rx) = std_mpsc::sync_channel(1); if let Some(ref tx) = self.filled_tx { tx.send(make_msg(ack_tx)).map_err(|_send| { io::Error::new( @@ -259,12 +259,10 @@ impl Sink for PipelinedSink { metadata: &'a BatchMetadata, ) -> Pin + Send + 'a>> { Box::pin(async move { - // Check for write errors from previous batches before accepting new work. - if let Err(e) = self.check_writer_error() { - return SendResult::from_io_error(e); - } - - // Acquire an empty buffer from the pool (blocks if writer is behind). + // Acquire an empty buffer from the pool without blocking the Tokio + // worker thread. If all buffers are still owned by in-flight + // timed-out writes, surface WouldBlock so the worker pool applies + // its normal transient-error backoff. let mut buf = match self.acquire_buffer() { Ok(b) => b, Err(e) => return SendResult::from_io_error(e), @@ -275,19 +273,23 @@ impl Sink for PipelinedSink { let rows = match self.serializer.serialize(batch, metadata, &mut buf) { Ok(n) => n, Err(e) => { - // Return buffer to pool to avoid exhaustion on repeated errors. - // Send it through the writer thread as empty data — the writer - // will write zero bytes (no-op) and recycle the buffer. + // Keep the caller-owned buffer available without depending + // on the writer thread, which may already be unhealthy. buf.clear(); - let _ = self.send_filled(buf); + self.spare_buf = Some(buf); return SendResult::from_io_error(e); } }; let bytes = buf.len() as u64; - // Hand off the filled buffer to the writer thread. - if let Err(e) = self.send_filled(buf) { + // Hand off the filled buffer and wait for the writer to report the + // OS write result before acknowledging delivery to the worker pool. + let filled_tx = match self.filled_sender() { + Ok(tx) => tx, + Err(e) => return SendResult::from_io_error(e), + }; + if let Err(e) = Self::write_filled(filled_tx, buf).await { return SendResult::from_io_error(e); } @@ -339,23 +341,23 @@ impl Drop for PipelinedSink { fn writer_thread_loop( mut writer: impl BatchWriter, - filled_rx: mpsc::Receiver, + filled_rx: std_mpsc::Receiver, empty_tx: SyncSender>, - error_tx: SyncSender, ) { while let Ok(msg) = filled_rx.recv() { match msg { - WriterMsg::Data(buf) => { - if let Err(e) = writer.write(&buf) { + WriterMsg::Data { buf, ack } => { + let result = writer.write(&buf); + if let Err(e) = &result { tracing::error!(error = %e, "pipelined writer: write failed"); - // Notify the caller about the write failure. - let _ = error_tx.try_send(WriteError(e)); } // Recycle the buffer back to the pool. let mut recycled = buf; recycled.clear(); - // If send fails, the sink is being dropped — just let the buffer go. - let _ = empty_tx.try_send(recycled); + // Preserve the fixed-size pool. If send fails, the sink is + // being dropped, so letting the buffer go is fine. + let _ = empty_tx.send(recycled); + let _ = ack.send(result); } WriterMsg::Flush(ack) => { let result = writer.flush(); @@ -378,10 +380,11 @@ fn writer_thread_loop( /// Pre-allocation chunk size: 64 MB. /// -/// When the file is opened (or needs more space), we ask the OS to reserve -/// this much contiguous space. This avoids per-write metadata updates on -/// ext4/xfs/APFS that slow down sequential extending writes. -const PREALLOC_BYTES: i64 = 64 * 1024 * 1024; +/// On Linux/macOS, when the file is opened, we ask the OS to reserve this much +/// space up front. This reduces per-write metadata updates on filesystems such +/// as ext4, XFS, and APFS during sequential extending writes. +#[cfg(any(target_os = "linux", target_os = "macos"))] +const PREALLOC_BYTES: u64 = 64 * 1024 * 1024; /// A [`BatchWriter`] that writes to a file using blocking std::fs I/O. /// @@ -389,7 +392,7 @@ const PREALLOC_BYTES: i64 = 64 * 1024 * 1024; /// /// | Platform | Optimization | /// |----------|-------------| -/// | Linux | `fallocate` pre-allocation, `POSIX_FADV_DONTNEED` after writes | +/// | Linux | `fallocate` pre-allocation, `POSIX_FADV_DONTNEED` after writes via `rustix` | /// | macOS | `F_PREALLOCATE`, `F_NOCACHE` (bypass UBC for written pages) | /// | Windows | (standard `write_all` — no extra APIs yet) | pub struct FileWriter { @@ -407,7 +410,7 @@ impl FileWriter { // Initialize bytes_written to the current file position so that // FADV_DONTNEED offsets are correct when appending to existing files. #[cfg(target_os = "linux")] - let bytes_written = file.metadata().map_or(0, |m| m.len()); + let bytes_written = file.metadata().map_or_else(|_| 0, |m| m.len()); Self { file, @@ -422,63 +425,52 @@ impl FileWriter { Self::apply_linux_hints(file); #[cfg(target_os = "macos")] Self::apply_macos_hints(file); + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + let _ = file; } /// Linux: pre-allocate space with `fallocate` and set sequential advice. #[cfg(target_os = "linux")] fn apply_linux_hints(file: &std::fs::File) { - use std::os::unix::io::AsRawFd; - let fd = file.as_raw_fd(); + use rustix::fs::{Advice, FallocateFlags, fadvise, fallocate}; // Pre-allocate disk space to avoid metadata updates on every // extending write. KEEP_SIZE means the visible file size is not // changed — only the underlying blocks are reserved. - // SAFETY: fd is valid (owned by `file`), flags and offsets are correct. - // fallocate failure is benign (returns -1, we ignore). - unsafe { - let _ = libc::fallocate(fd, libc::FALLOC_FL_KEEP_SIZE, 0, PREALLOC_BYTES); - } + let _ = fallocate(file, FallocateFlags::KEEP_SIZE, 0, PREALLOC_BYTES); // Tell the kernel we're doing sequential writes so it can optimize // read-ahead and page cache management. - // SAFETY: fd is valid (owned by `file`), POSIX_FADV_SEQUENTIAL is an - // advisory hint that cannot cause UB regardless of return value. - unsafe { - let _ = libc::posix_fadvise(fd, 0, 0, libc::POSIX_FADV_SEQUENTIAL); - } + let _ = fadvise(file, 0, None, Advice::Sequential); } /// macOS: set `F_NOCACHE` and pre-allocate with `F_PREALLOCATE`. #[cfg(target_os = "macos")] fn apply_macos_hints(file: &std::fs::File) { use std::os::unix::io::AsRawFd; + let fd = file.as_raw_fd(); - // F_NOCACHE: bypass the Unified Buffer Cache for this fd. - // Written pages won't linger in memory — good for a forwarder that - // produces write-only data consumers won't re-read via this fd. - // SAFETY: fd is valid (owned by `file`), F_NOCACHE with arg=1 is a - // hint that cannot cause UB regardless of return value. + // F_NOCACHE bypasses the Unified Buffer Cache for this descriptor so + // write-only log data does not linger in process-visible cache. + // SAFETY: `fd` is owned by `file`, and F_NOCACHE with arg=1 is an + // advisory descriptor hint that is safe to ignore on failure. unsafe { let _ = libc::fcntl(fd, libc::F_NOCACHE, 1i32); } - // F_PREALLOCATE: reserve contiguous disk space. - // This is macOS's equivalent of Linux fallocate — it avoids - // fragmentation and repeated metadata updates during sequential - // extending writes. - // SAFETY: fd is valid (owned by `file`), fstore_t is zero-initialized - // then fully populated. fcntl(F_PREALLOCATE) reads the struct by - // pointer and never writes through it. + // F_PREALLOCATE is macOS's fallocate equivalent. It reserves space + // without changing the visible file length. + // SAFETY: `fd` is owned by `file`; `store` is fully initialized before + // the kernel reads it, and fcntl does not retain the pointer. unsafe { let mut store: libc::fstore_t = std::mem::zeroed(); store.fst_flags = libc::F_ALLOCATECONTIG; store.fst_posmode = libc::F_PEOFPOSMODE; store.fst_offset = 0; - store.fst_length = PREALLOC_BYTES; + store.fst_length = PREALLOC_BYTES as libc::off_t; let ret = libc::fcntl(fd, libc::F_PREALLOCATE, &store); if ret == -1 { - // Contiguous allocation failed — retry allowing fragments. store.fst_flags = libc::F_ALLOCATEALL; let _ = libc::fcntl(fd, libc::F_PREALLOCATE, &store); } @@ -491,18 +483,14 @@ impl FileWriter { /// write-only log data that will never be re-read by this process. #[cfg(target_os = "linux")] fn advise_dontneed(&self, offset: u64, len: u64) { - use std::os::unix::io::AsRawFd; - let fd = self.file.as_raw_fd(); - // SAFETY: fd is valid (owned by `self.file`), offset and len describe - // already-written data, POSIX_FADV_DONTNEED is an advisory hint. - unsafe { - let _ = libc::posix_fadvise( - fd, - offset as libc::off_t, - len as libc::off_t, - libc::POSIX_FADV_DONTNEED, - ); - } + use std::num::NonZeroU64; + + use rustix::fs::{Advice, fadvise}; + + let Some(len) = NonZeroU64::new(len) else { + return; + }; + let _ = fadvise(&self.file, offset, Some(len), Advice::DontNeed); } } @@ -578,3 +566,475 @@ impl BatchSerializer for JsonBatchSerializer { self.inner.name() } } + +#[cfg(test)] +mod tests { + use std::collections::VecDeque; + use std::sync::Mutex; + use std::sync::atomic::Ordering; + use std::time::Duration; + + use arrow::datatypes::Schema; + use proptest::prelude::*; + use proptest::test_runner::TestCaseError; + + use super::*; + + struct FixedSerializer { + rows: u64, + payload: &'static [u8], + } + + impl BatchSerializer for FixedSerializer { + fn serialize( + &mut self, + _batch: &RecordBatch, + _metadata: &BatchMetadata, + buf: &mut Vec, + ) -> io::Result { + buf.extend_from_slice(self.payload); + Ok(self.rows) + } + + fn name(&self) -> &str { + "fixed" + } + } + + struct FailsOnceSerializer { + failed_once: bool, + rows: u64, + payload: &'static [u8], + } + + impl BatchSerializer for FailsOnceSerializer { + fn serialize( + &mut self, + _batch: &RecordBatch, + _metadata: &BatchMetadata, + buf: &mut Vec, + ) -> io::Result { + if self.failed_once { + buf.extend_from_slice(self.payload); + Ok(self.rows) + } else { + self.failed_once = true; + Err(io::Error::other("encode failed")) + } + } + + fn name(&self) -> &str { + "fails-once" + } + } + + #[derive(Clone, Debug)] + struct ScriptStep { + serialize_ok: bool, + write_ok: bool, + rows: u64, + payload: Vec, + } + + #[derive(Clone, Debug)] + enum ScriptOp { + Send(ScriptStep), + Flush, + } + + type SharedWritePlan = Arc>>; + + struct ScriptedSerializer { + steps: VecDeque, + write_plan: SharedWritePlan, + } + + impl BatchSerializer for ScriptedSerializer { + fn serialize( + &mut self, + _batch: &RecordBatch, + _metadata: &BatchMetadata, + buf: &mut Vec, + ) -> io::Result { + let step = self + .steps + .pop_front() + .ok_or_else(|| io::Error::other("missing scripted serialize step"))?; + if !step.serialize_ok { + return Err(io::Error::other("scripted encode failure")); + } + buf.extend_from_slice(&step.payload); + self.write_plan + .lock() + .expect("scripted write plan mutex") + .push_back(step.clone()); + Ok(step.rows) + } + + fn name(&self) -> &str { + "scripted" + } + } + + struct RecordingWriter { + writes: Arc>>>, + } + + impl BatchWriter for RecordingWriter { + fn write(&mut self, buf: &[u8]) -> io::Result<()> { + self.writes + .lock() + .expect("recording writer mutex") + .push(buf.to_vec()); + Ok(()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + + fn shutdown(&mut self) -> io::Result<()> { + Ok(()) + } + } + + struct ScriptedWriter { + write_plan: SharedWritePlan, + writes: Arc>>>, + } + + impl BatchWriter for ScriptedWriter { + fn write(&mut self, buf: &[u8]) -> io::Result<()> { + let step = self + .write_plan + .lock() + .expect("scripted write plan mutex") + .pop_front() + .ok_or_else(|| io::Error::other("missing scripted write step"))?; + if buf != step.payload.as_slice() { + return Err(io::Error::other("scripted payload mismatch")); + } + if step.write_ok { + self.writes + .lock() + .expect("scripted writes mutex") + .push(buf.to_vec()); + Ok(()) + } else { + Err(io::Error::other("scripted write failure")) + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + + fn shutdown(&mut self) -> io::Result<()> { + Ok(()) + } + } + + struct FailingWriter; + + impl BatchWriter for FailingWriter { + fn write(&mut self, _buf: &[u8]) -> io::Result<()> { + Err(io::Error::other("disk full")) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + + fn shutdown(&mut self) -> io::Result<()> { + Ok(()) + } + } + + struct DelayedThenFailWriter { + release_first: std_mpsc::Receiver<()>, + writes: usize, + } + + impl BatchWriter for DelayedThenFailWriter { + fn write(&mut self, _buf: &[u8]) -> io::Result<()> { + self.writes += 1; + if self.writes == 1 { + self.release_first + .recv_timeout(Duration::from_secs(5)) + .map_err(io::Error::other)?; + Ok(()) + } else { + Err(io::Error::other("second write failed")) + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + + fn shutdown(&mut self) -> io::Result<()> { + Ok(()) + } + } + + fn empty_batch() -> RecordBatch { + RecordBatch::new_empty(Arc::new(Schema::empty())) + } + + fn metadata() -> BatchMetadata { + BatchMetadata { + resource_attrs: Arc::default(), + observed_time_ns: 0, + } + } + + #[tokio::test] + async fn send_batch_returns_only_after_writer_success() { + let stats = Arc::new(ComponentStats::new()); + let writes = Arc::new(Mutex::new(Vec::new())); + let mut sink = PipelinedSink::new( + FixedSerializer { + rows: 2, + payload: b"ok\n", + }, + RecordingWriter { + writes: Arc::clone(&writes), + }, + Arc::clone(&stats), + PipelineConfig { + num_buffers: 1, + buf_capacity: 16, + }, + ) + .expect("sink"); + + let batch = empty_batch(); + let result = sink.send_batch(&batch, &metadata()).await; + + assert!(result.is_ok(), "result: {result:?}"); + assert_eq!( + writes.lock().expect("recording writer mutex").as_slice(), + &[b"ok\n".to_vec()] + ); + assert_eq!(stats.lines_total.load(Ordering::Relaxed), 2); + assert_eq!(stats.bytes_total.load(Ordering::Relaxed), 3); + } + + #[tokio::test] + async fn send_batch_reports_writer_error_without_success_counters() { + let stats = Arc::new(ComponentStats::new()); + let mut sink = PipelinedSink::new( + FixedSerializer { + rows: 2, + payload: b"lost\n", + }, + FailingWriter, + Arc::clone(&stats), + PipelineConfig { + num_buffers: 1, + buf_capacity: 16, + }, + ) + .expect("sink"); + + let batch = empty_batch(); + let result = sink.send_batch(&batch, &metadata()).await; + + assert!( + matches!(result, SendResult::IoError(_)), + "result: {result:?}" + ); + assert_eq!(stats.lines_total.load(Ordering::Relaxed), 0); + assert_eq!(stats.bytes_total.load(Ordering::Relaxed), 0); + } + + #[tokio::test] + async fn serialization_error_keeps_buffer_available_for_retry() { + let stats = Arc::new(ComponentStats::new()); + let writes = Arc::new(Mutex::new(Vec::new())); + let mut sink = PipelinedSink::new( + FailsOnceSerializer { + failed_once: false, + rows: 1, + payload: b"retry\n", + }, + RecordingWriter { + writes: Arc::clone(&writes), + }, + Arc::clone(&stats), + PipelineConfig { + num_buffers: 1, + buf_capacity: 16, + }, + ) + .expect("sink"); + + let batch = empty_batch(); + let first = sink.send_batch(&batch, &metadata()).await; + assert!(matches!(first, SendResult::IoError(_)), "result: {first:?}"); + + let second = sink.send_batch(&batch, &metadata()).await; + assert!(second.is_ok(), "result: {second:?}"); + assert_eq!( + writes.lock().expect("recording writer mutex").as_slice(), + &[b"retry\n".to_vec()] + ); + assert_eq!(stats.lines_total.load(Ordering::Relaxed), 1); + assert_eq!(stats.bytes_total.load(Ordering::Relaxed), 6); + } + + #[tokio::test] + async fn timed_out_send_does_not_satisfy_next_batch_ack() { + let stats = Arc::new(ComponentStats::new()); + let (release_first_tx, release_first_rx) = std_mpsc::channel(); + let mut sink = PipelinedSink::new( + FixedSerializer { + rows: 1, + payload: b"batch\n", + }, + DelayedThenFailWriter { + release_first: release_first_rx, + writes: 0, + }, + Arc::clone(&stats), + PipelineConfig { + num_buffers: 1, + buf_capacity: 16, + }, + ) + .expect("sink"); + + let batch = empty_batch(); + let first = tokio::time::timeout( + Duration::from_millis(10), + sink.send_batch(&batch, &metadata()), + ) + .await; + assert!(first.is_err(), "first send should time out"); + + let second = + tokio::time::timeout(Duration::from_secs(1), sink.send_batch(&batch, &metadata())) + .await + .expect("buffer acquisition should not block the Tokio worker"); + assert!( + matches!(second, SendResult::IoError(ref error) if error.kind() == io::ErrorKind::WouldBlock), + "second send should report buffer backpressure, got {second:?}" + ); + + release_first_tx.send(()).expect("release first write"); + let third = tokio::time::timeout(Duration::from_secs(1), async { + loop { + let result = sink.send_batch(&batch, &metadata()).await; + if !matches!(result, SendResult::IoError(ref error) if error.kind() == io::ErrorKind::WouldBlock) + { + break result; + } + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .expect("writer should recycle the first timed-out buffer"); + assert!( + matches!(third, SendResult::IoError(_)), + "third send must receive its own write failure, got {third:?}" + ); + tokio::time::timeout(Duration::from_secs(1), sink.shutdown()) + .await + .expect("shutdown should not hang") + .expect("shutdown"); + assert_eq!(stats.lines_total.load(Ordering::Relaxed), 0); + assert_eq!(stats.bytes_total.load(Ordering::Relaxed), 0); + } + + fn script_step_strategy() -> impl Strategy { + ( + any::(), + any::(), + 0u64..8, + prop::collection::vec(any::(), 0..64), + ) + .prop_map(|(serialize_ok, write_ok, rows, payload)| ScriptStep { + serialize_ok, + write_ok, + rows, + payload, + }) + } + + fn script_op_strategy() -> impl Strategy { + prop_oneof![ + script_step_strategy().prop_map(ScriptOp::Send), + Just(ScriptOp::Flush), + ] + } + + proptest! { + #[test] + fn arbitrary_send_and_flush_sequences_count_only_acked_writes( + ops in prop::collection::vec(script_op_strategy(), 1..32) + ) { + let send_steps = ops + .iter() + .filter_map(|op| match op { + ScriptOp::Send(step) => Some(step.clone()), + ScriptOp::Flush => None, + }) + .collect(); + let write_plan = Arc::new(Mutex::new(VecDeque::new())); + let writes = Arc::new(Mutex::new(Vec::new())); + let stats = Arc::new(ComponentStats::new()); + let mut sink = PipelinedSink::new( + ScriptedSerializer { + steps: send_steps, + write_plan: Arc::clone(&write_plan), + }, + ScriptedWriter { + write_plan, + writes: Arc::clone(&writes), + }, + Arc::clone(&stats), + PipelineConfig { + num_buffers: 1, + buf_capacity: 64, + }, + ) + .expect("sink"); + let batch = empty_batch(); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio runtime"); + let (expected_writes, expected_rows, expected_bytes) = runtime.block_on(async { + let mut expected_writes = Vec::new(); + let mut expected_rows = 0u64; + let mut expected_bytes = 0u64; + + for op in &ops { + match op { + ScriptOp::Send(step) => { + let result = sink.send_batch(&batch, &metadata()).await; + let should_succeed = step.serialize_ok && step.write_ok; + prop_assert_eq!(result.is_ok(), should_succeed); + if should_succeed { + expected_rows += step.rows; + expected_bytes += step.payload.len() as u64; + expected_writes.push(step.payload.clone()); + } + } + ScriptOp::Flush => { + sink.flush().await.expect("flush"); + } + } + } + + sink.shutdown().await.expect("shutdown"); + Ok::<_, TestCaseError>((expected_writes, expected_rows, expected_bytes)) + })?; + + let actual_writes = writes.lock().expect("scripted writes mutex").clone(); + prop_assert_eq!(actual_writes, expected_writes); + prop_assert_eq!(stats.lines_total.load(Ordering::Relaxed), expected_rows); + prop_assert_eq!(stats.bytes_total.load(Ordering::Relaxed), expected_bytes); + } + } +} diff --git a/crates/ffwd-runtime/src/worker_pool/worker.rs b/crates/ffwd-runtime/src/worker_pool/worker.rs index 15d1425bd..7d59e8780 100644 --- a/crates/ffwd-runtime/src/worker_pool/worker.rs +++ b/crates/ffwd-runtime/src/worker_pool/worker.rs @@ -462,6 +462,7 @@ pub(super) async fn process_item( mod tests { use std::collections::BTreeSet; use std::future::Future; + use std::io; use std::pin::Pin; use std::sync::Arc; use std::time::Duration; @@ -525,7 +526,7 @@ mod tests { Box::pin(async { SendResult::Ok }) } - fn flush(&mut self) -> Pin> + Send + '_>> { + fn flush(&mut self) -> Pin> + Send + '_>> { Box::pin(async { Ok(()) }) } @@ -533,7 +534,7 @@ mod tests { "ok-sink" } - fn shutdown(&mut self) -> Pin> + Send + '_>> { + fn shutdown(&mut self) -> Pin> + Send + '_>> { Box::pin(async { Ok(()) }) } } @@ -549,7 +550,7 @@ mod tests { Box::pin(std::future::pending()) } - fn flush(&mut self) -> Pin> + Send + '_>> { + fn flush(&mut self) -> Pin> + Send + '_>> { Box::pin(async { Ok(()) }) } @@ -557,7 +558,43 @@ mod tests { "hanging-sink" } - fn shutdown(&mut self) -> Pin> + Send + '_>> { + fn shutdown(&mut self) -> Pin> + Send + '_>> { + Box::pin(async { Ok(()) }) + } + } + + struct WouldBlockThenOkSink { + calls: usize, + } + + impl Sink for WouldBlockThenOkSink { + fn send_batch<'a>( + &'a mut self, + _batch: &'a RecordBatch, + _metadata: &'a BatchMetadata, + ) -> Pin + Send + 'a>> { + let result = if self.calls == 0 { + self.calls += 1; + SendResult::IoError(io::Error::new( + io::ErrorKind::WouldBlock, + "buffer pool exhausted", + )) + } else { + self.calls += 1; + SendResult::Ok + }; + Box::pin(async move { result }) + } + + fn flush(&mut self) -> Pin> + Send + '_>> { + Box::pin(async { Ok(()) }) + } + + fn name(&self) -> &str { + "would-block-then-ok" + } + + fn shutdown(&mut self) -> Pin> + Send + '_>> { Box::pin(async { Ok(()) }) } } @@ -647,6 +684,36 @@ mod tests { assert_eq!(retries, 0); } + #[tokio::test] + async fn process_item_treats_would_block_as_transient_backpressure() { + let cancel = CancellationToken::new(); + let mut sink = WouldBlockThenOkSink { calls: 0 }; + let output_health = Arc::new(OutputHealthTracker::new(vec![])); + let metadata = BatchMetadata { + resource_attrs: Arc::default(), + observed_time_ns: 0, + }; + + let (outcome, _send_latency_ns, retries) = process_item( + ProcessItemContext { + worker_id: 0, + sink: &mut sink, + output_health: &output_health, + metadata: &metadata, + max_retry_delay: Duration::from_millis(1), + cancel: &cancel, + #[cfg(feature = "turmoil")] + batch_id: 0, // test only + }, + make_batch(), + ) + .await; + + assert_eq!(outcome, DeliveryOutcome::Delivered); + assert_eq!(retries, 1); + assert_eq!(sink.calls, 2); + } + #[test] fn terminalization_reducer_is_idempotent_for_any_two_step_schedule() { let actions = [ diff --git a/justfile b/justfile index 832f5985c..2ea5a5566 100644 --- a/justfile +++ b/justfile @@ -689,11 +689,11 @@ bench-source-metadata-fast *ARGS: # Profile OTLP decode/encode CPU with the normal allocator (flamegraph, per-mode timings). profile-otlp-io *ARGS: - cargo run -p ffwd-bench --release --features bench-tools --bin otlp_io_profile -- {{ARGS}} + cargo run -p ffwd-bench --release --features bench-tools,io-bench --bin otlp_io_profile -- {{ARGS}} # Profile OTLP decode/encode allocation counts with stats_alloc instrumentation. profile-otlp-io-alloc *ARGS: - cargo run -p ffwd-bench --release --features bench-tools,otlp-profile-alloc --bin otlp_io_profile -- {{ARGS}} + cargo run -p ffwd-bench --release --features bench-tools,io-bench,otlp-profile-alloc --bin otlp_io_profile -- {{ARGS}} # Generate microbenchmark report (markdown) bench-report: @@ -701,11 +701,11 @@ bench-report: # Profile FramedInput / format processing overhead and print a markdown report. bench-framed-input *ARGS: - cargo run -p ffwd-bench --release --features bench-tools --bin framed_input_profile -- {{ARGS}} + cargo run -p ffwd-bench --release --features bench-tools,io-bench --bin framed_input_profile -- {{ARGS}} # Allocation-focused FramedInput profiling (dhat-backed, slower; no throughput numbers). bench-framed-input-alloc *ARGS: - cargo run -p ffwd-bench --release --features bench-tools,dhat-heap --bin framed_input_profile -- --alloc-only {{ARGS}} + cargo run -p ffwd-bench --release --features bench-tools,io-bench,dhat-heap --bin framed_input_profile -- --alloc-only {{ARGS}} # Run sustained-load memory profiler (generator → SQL → null, default 5 minutes). # Use --quick (30s) for CI or --medium (120s) for quick checks. diff --git a/scripts/verify_tla_coverage.py b/scripts/verify_tla_coverage.py index b7b632974..dbeaa19a5 100755 --- a/scripts/verify_tla_coverage.py +++ b/scripts/verify_tla_coverage.py @@ -10,8 +10,10 @@ from __future__ import annotations import argparse +import os import pathlib import re +import shutil import subprocess import tempfile from typing import Sequence @@ -72,6 +74,36 @@ def parse_invariants(body_lines: Sequence[str]) -> list[str]: return invariants +def find_java() -> str: + """Find a Java runtime, matching the repo's justfile fallback behavior.""" + java_home = os.environ.get("JAVA_HOME") + candidates = [ + os.environ.get("JAVA_BIN"), + os.path.join(java_home, "bin", "java") if java_home else None, + "/opt/homebrew/opt/openjdk/bin/java", + "/usr/local/opt/openjdk/bin/java", + "java", + ] + for candidate in candidates: + if not candidate: + continue + java_bin = shutil.which(candidate) if os.path.basename(candidate) == candidate else candidate + if not java_bin: + continue + if not os.path.exists(java_bin): + continue + proc = subprocess.run( + [java_bin, "-version"], + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + if proc.returncode == 0: + return java_bin + raise RuntimeError( + "Java runtime not found. Install OpenJDK (e.g. 'brew install openjdk') or set JAVA_BIN." + ) + + def run_one(jar: pathlib.Path, tla_file: pathlib.Path, cfg_lines: Sequence[str], invariant: str) -> None: before, _body, after = split_cfg_sections(cfg_lines) rendered = before + ["INVARIANTS", f" {invariant}"] + after @@ -81,8 +113,9 @@ def run_one(jar: pathlib.Path, tla_file: pathlib.Path, cfg_lines: Sequence[str], tmp_cfg = pathlib.Path(fh.name) try: + java_bin = find_java() cmd = [ - "java", + java_bin, "-cp", str(jar), "tlc2.TLC", diff --git a/tla/DeliveryRetry.tla b/tla/DeliveryRetry.tla index 523415afa..fe54e6f46 100644 --- a/tla/DeliveryRetry.tla +++ b/tla/DeliveryRetry.tla @@ -13,7 +13,7 @@ * 1. Send batch to sink * 2. On Ok: terminal success (Delivered) * 3. On Rejected: terminal permanent failure (Rejected) - * 4. On IoError/RetryAfter/Timeout: exponential backoff, retry forever + * 4. On IoError/RetryAfter/Timeout: transient retry, retry forever * 5. On cancel (shutdown): terminal (PoolClosed) * * Key insight: the code deliberately retries FOREVER for transient @@ -162,10 +162,13 @@ ReceiveRejected == /\ outcome' = "Rejected" /\ UNCHANGED <> -\* ReceiveTransient: transient failure, compute next backoff and wait. +\* ReceiveTransient: transient failure, compute an abstract retry delay and wait. \* Maps to Ok(SendResult::IoError), Ok(SendResult::RetryAfter), and \* Err(_elapsed) (timeout) in process_item. All three follow the same -\* pattern: increment retry count, compute backoff, sleep, loop. +\* protocol pattern: increment retry count, sleep, loop. Timing details are +\* intentionally abstracted here: IoError/timeout use exponential backoff, +\* while RetryAfter uses the server-directed delay. Turmoil trace validators +\* cover that concrete timing distinction. ReceiveTransient == /\ state = "Sending" /\ sinkState = "Transient"