From 5483f6970b7f109feb658a3b9adda2ed1a50702b Mon Sep 17 00:00:00 2001 From: "robert.blafford" Date: Mon, 10 Aug 2026 17:46:41 +0000 Subject: [PATCH 1/2] fix(disk buffer): coordinate published progress safely Publish reader-visible progress only after file I/O completes, gate reads on that boundary, serialize writer operations through an actor, and propagate shutdown into reader-progress waits. --- .../src/topology/channel/sender.rs | 307 +++++++++++++++--- .../src/variants/disk_v2/ledger.rs | 144 +++++--- .../src/variants/disk_v2/mod.rs | 5 + .../src/variants/disk_v2/reader.rs | 79 ++++- .../disk_v2/tests/acknowledgements.rs | 6 +- .../src/variants/disk_v2/tests/invariants.rs | 200 +++++++++++- .../src/variants/disk_v2/tests/mod.rs | 1 + .../src/variants/disk_v2/tests/writer.rs | 204 ++++++++++++ .../src/variants/disk_v2/writer.rs | 114 +++++-- 9 files changed, 940 insertions(+), 120 deletions(-) create mode 100644 lib/vector-buffers/src/variants/disk_v2/tests/writer.rs diff --git a/lib/vector-buffers/src/topology/channel/sender.rs b/lib/vector-buffers/src/topology/channel/sender.rs index 1ff0e4aeb5382..c96291352c967 100644 --- a/lib/vector-buffers/src/topology/channel/sender.rs +++ b/lib/vector-buffers/src/topology/channel/sender.rs @@ -1,11 +1,11 @@ // Derivative's Debug impl generates 'let _ = field.fmt(f)' which triggers this lint. #![allow(clippy::let_underscore_must_use)] -use std::{sync::Arc, time::Instant}; +use std::{fmt, io, sync::Arc, time::Instant}; use async_recursion::async_recursion; use derivative::Derivative; -use tokio::sync::Mutex; +use tokio::sync::{mpsc, oneshot, watch}; use tracing::Span; use vector_common::internal_event::{InternalEventHandle, Registered, register}; @@ -17,6 +17,129 @@ use crate::{ variants::disk_v2::{self, ProductionFilesystem, TryWriteOutcome}, }; +const DISK_V2_WRITER_QUEUE_CAPACITY: usize = 1; + +enum DiskV2WriterCommand { + Write { + item: T, + blocking: bool, + response: oneshot::Sender, disk_v2::WriterError>>, + }, + Flush { + response: oneshot::Sender>, + }, +} + +#[derive(Clone)] +pub struct DiskV2Sender { + commands: mpsc::Sender>, + _shutdown: watch::Sender<()>, +} + +impl DiskV2Sender { + fn new(mut writer: disk_v2::BufferWriter) -> Self { + let (commands, mut command_rx) = mpsc::channel(DISK_V2_WRITER_QUEUE_CAPACITY); + let (shutdown, shutdown_rx) = watch::channel(()); + writer.set_shutdown(shutdown_rx); + // The task is the cancellation boundary for disk I/O. While the topology is active, once + // the queue accepts a command, this owner finishes the write and its visibility flush even + // if the requesting future is dropped. When the last sender is dropped, only waits for + // reader progress are interrupted; any file I/O already in progress still finishes. + vector_common::spawn_in_current_span(async move { + while let Some(command) = command_rx.recv().await { + let (failed, shutting_down) = match command { + DiskV2WriterCommand::Write { + item, + blocking, + response, + } => { + let result = if blocking { + writer.write_record_outcome(item).await + } else { + writer.try_write_record(item).await + }; + let result = match result { + Ok(outcome) => writer + .flush() + .await + .map(|()| outcome) + .map_err(|source| disk_v2::WriterError::Io { source }), + other => other, + }; + let shutting_down = matches!(&result, Err(disk_v2::WriterError::Shutdown)); + let failed = result.is_err() && !shutting_down; + if failed { + writer.fail(); + } + let _ = response.send(result); + (failed, shutting_down) + } + DiskV2WriterCommand::Flush { response } => { + let result = writer.flush().await; + let failed = result.is_err(); + if failed { + writer.fail(); + } + let _ = response.send(result); + (failed, false) + } + }; + + if failed || shutting_down { + break; + } + } + }); + + Self { + commands, + _shutdown: shutdown, + } + } + + async fn write(&self, item: T, blocking: bool) -> crate::Result> { + let (response, response_rx) = oneshot::channel(); + self.commands + .send(DiskV2WriterCommand::Write { + item, + blocking, + response, + }) + .await + .map_err(|_| io::Error::other("disk buffer writer task stopped"))?; + + response_rx + .await + .map_err(|_| io::Error::other("disk buffer writer task stopped"))? + .map_err(|error| { + error!(%error, "Disk buffer writer encountered an unrecoverable error."); + error.into() + }) + } + + async fn flush(&self) -> crate::Result<()> { + let (response, response_rx) = oneshot::channel(); + self.commands + .send(DiskV2WriterCommand::Flush { response }) + .await + .map_err(|_| io::Error::other("disk buffer writer task stopped"))?; + + response_rx + .await + .map_err(|_| io::Error::other("disk buffer writer task stopped"))? + .map_err(|error| { + error!(%error, "Disk buffer writer encountered an unrecoverable error."); + error.into() + }) + } +} + +impl fmt::Debug for DiskV2Sender { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.debug_struct("DiskV2Sender").finish() + } +} + /// Adapter for papering over various sender backends. #[derive(Clone, Debug)] pub enum SenderAdapter { @@ -24,7 +147,7 @@ pub enum SenderAdapter { InMemory(LimitedSender), /// The disk v2 buffer. - DiskV2(Arc>>), + DiskV2(DiskV2Sender), } impl From> for SenderAdapter { @@ -35,7 +158,7 @@ impl From> for SenderAdapter { impl From> for SenderAdapter { fn from(v: disk_v2::BufferWriter) -> Self { - Self::DiskV2(Arc::new(Mutex::new(v))) + Self::DiskV2(DiskV2Sender::new(v)) } } @@ -50,20 +173,7 @@ where .await .map(|()| TryWriteOutcome::Written) .map_err(Into::into), - Self::DiskV2(writer) => { - let mut writer = writer.lock().await; - - writer.write_record_outcome(item).await.map_err(|e| { - // Record-level failures that can never succeed (a record too large to encode - // within the max record size) are handled inside the writer and surfaced as a - // dropped outcome. Anything that reaches this point -- I/O errors, - // serialization failures, an inconsistent writer state -- is genuinely - // unrecoverable. - error!("Disk buffer writer has encountered an unrecoverable error."); - - e.into() - }) - } + Self::DiskV2(writer) => writer.write(item, true).await, } } @@ -73,35 +183,14 @@ where .try_send(item) .map(|()| TryWriteOutcome::Written) .or_else(|e| Ok(TryWriteOutcome::Full(e.into_inner()))), - Self::DiskV2(writer) => { - let mut writer = writer.lock().await; - - writer.try_write_record(item).await.map_err(|e| { - // Record-level failures that can never succeed (a record too large to encode - // within the max record size) are handled inside the writer and surfaced as a - // dropped outcome. Anything that reaches this point -- I/O errors, - // serialization failures, an inconsistent writer state -- is genuinely - // unrecoverable. - error!("Disk buffer writer has encountered an unrecoverable error."); - - e.into() - }) - } + Self::DiskV2(writer) => writer.write(item, false).await, } } pub(crate) async fn flush(&mut self) -> crate::Result<()> { match self { Self::InMemory(_) => Ok(()), - Self::DiskV2(writer) => { - let mut writer = writer.lock().await; - writer.flush().await.map_err(|e| { - // Errors on the I/O path, which is all that flushing touches, are never recoverable. - error!("Disk buffer writer has encountered an unrecoverable error."); - - e.into() - }) - } + Self::DiskV2(writer) => writer.flush().await, } } @@ -301,3 +390,139 @@ impl BufferSender { Ok(()) } } + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use super::{DiskV2Sender, DiskV2WriterCommand, TryWriteOutcome, oneshot}; + use crate::{ + buffer_usage_data::BufferUsageHandle, + test::{SizedRecord, with_temp_dir}, + variants::disk_v2::{Buffer, DiskBufferConfigBuilder}, + }; + + #[tokio::test] + async fn disk_writer_finishes_accepted_write_after_request_is_cancelled() { + with_temp_dir(|data_dir| { + let data_dir = data_dir.to_path_buf(); + + async move { + let config = DiskBufferConfigBuilder::from_path(data_dir) + .build() + .expect("disk buffer config should be valid"); + let (writer, mut reader) = + Buffer::::from_config(config, BufferUsageHandle::noop()) + .await + .expect("disk buffer should initialize"); + let sender = DiskV2Sender::new(writer); + let first = SizedRecord::new(64); + let second = SizedRecord::new(96); + + // Enqueue the first write and then discard its response, exactly matching a caller + // future being cancelled after the writer task accepted the command. + let (response, response_rx) = oneshot::channel(); + sender + .commands + .send(DiskV2WriterCommand::Write { + item: first.clone(), + blocking: true, + response, + }) + .await + .expect("writer task should accept the command"); + drop(response_rx); + + // No later writer command is needed to complete or publish the cancelled request. + let first_read = tokio::time::timeout(Duration::from_secs(2), reader.next()) + .await + .expect("first read should not stall") + .expect("first read should succeed"); + assert_eq!(first_read, Some(first)); + + // Reusing the same writer after cancellation must assign a new ID rather than + // resubmitting the first record under its old ID. + assert_eq!( + sender + .write(second.clone(), true) + .await + .expect("second write should succeed"), + TryWriteOutcome::Written + ); + sender.flush().await.expect("writer flush should succeed"); + + let second_read = tokio::time::timeout(Duration::from_secs(2), reader.next()) + .await + .expect("second read should not stall") + .expect("second read should succeed"); + + assert_eq!(second_read, Some(second)); + } + }) + .await; + } + + #[tokio::test] + async fn disk_writer_stops_waiting_for_reader_when_last_sender_is_dropped() { + with_temp_dir(|data_dir| { + let data_dir = data_dir.to_path_buf(); + + async move { + let config = DiskBufferConfigBuilder::from_path(data_dir) + .max_buffer_size(4096) + .max_data_file_size(1024) + .max_record_size(1024) + .build() + .expect("disk buffer config should be valid"); + let (writer, _reader) = + Buffer::::from_config(config, BufferUsageHandle::noop()) + .await + .expect("disk buffer should initialize"); + let sender = DiskV2Sender::new(writer); + + let mut blocked_item = None; + for _ in 0..100 { + let item = SizedRecord::new(128); + match sender + .write(item, false) + .await + .expect("nonblocking write should succeed") + { + TryWriteOutcome::Written => {} + TryWriteOutcome::Full(item) => { + blocked_item = Some(item); + break; + } + TryWriteOutcome::Dropped => panic!("record should fit in the buffer"), + } + } + let blocked_item = blocked_item.expect("writes should fill the buffer"); + + let (response, response_rx) = oneshot::channel(); + sender + .commands + .send(DiskV2WriterCommand::Write { + item: blocked_item, + blocking: true, + response, + }) + .await + .expect("writer task should accept the blocking command"); + + // Closing the final sender is the topology shutdown signal. It must interrupt the + // reader-progress wait without cancelling any file I/O already in progress. + drop(sender); + + let result = tokio::time::timeout(Duration::from_secs(2), response_rx) + .await + .expect("writer task should observe shutdown") + .expect("writer task should return the command response"); + assert!(matches!( + result, + Err(crate::variants::disk_v2::WriterError::Shutdown) + )); + } + }) + .await; + } +} diff --git a/lib/vector-buffers/src/variants/disk_v2/ledger.rs b/lib/vector-buffers/src/variants/disk_v2/ledger.rs index b6d855cc9a37e..e9e6c38570a24 100644 --- a/lib/vector-buffers/src/variants/disk_v2/ledger.rs +++ b/lib/vector-buffers/src/variants/disk_v2/ledger.rs @@ -18,7 +18,7 @@ use snafu::{ResultExt, Snafu}; use tokio::{ fs, io::AsyncWriteExt, - sync::{Mutex, MutexGuard, Notify}, + sync::{Mutex, MutexGuard, watch}, }; use vector_common::finalizer::OrderedFinalizer; @@ -36,6 +36,47 @@ use crate::buffer_usage_data::BufferUsageHandle; pub const LEDGER_LEN: usize = align16(mem::size_of::()); +/// Latest writer position that is safe for the runtime reader to consume. +/// +/// This state is deliberately not persisted. Startup recovery derives it again from the durable +/// ledger and reconciled data files before the reader and writer are returned to the topology. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(super) struct WriterProgress { + revision: u64, + published_position: Option<(u16, u64)>, + next_record_id: u64, + done: bool, + failed: bool, +} + +impl WriterProgress { + fn initial(next_record_id: u64) -> Self { + Self { + revision: 0, + published_position: None, + next_record_id, + done: false, + failed: false, + } + } + + pub(super) fn published_position(self) -> Option<(u16, u64)> { + self.published_position + } + + pub(super) fn next_record_id(self) -> u64 { + self.next_record_id + } + + pub(super) fn is_done(self) -> bool { + self.done + } + + pub(super) fn is_failed(self) -> bool { + self.failed + } +} + /// Error that occurred during calls to [`Ledger`]. #[derive(Debug, Snafu)] pub enum LedgerLoadCreateError { @@ -241,10 +282,10 @@ where state: BackedArchive, // The total size, in bytes, of all unread records in the buffer. total_buffer_size: AtomicU64, - // Notifier for reader-related progress. - reader_notify: Notify, - // Notifier for writer-related progress. - writer_notify: Notify, + // Coalesced runtime reader progress observed by the writer. + reader_progress: watch::Sender, + // Coalesced runtime writer publication boundary observed by the reader. + writer_progress: watch::Sender, // Tracks when writer has fully shutdown. writer_done: AtomicBool, // Number of pending record acknowledgements that have yeet to be consumed by the reader. @@ -491,29 +532,19 @@ where Ok(deleted_files) } - /// Waits for a signal from the reader that progress has been made. - /// - /// This will only occur when a record is read, which may allow enough space (below the maximum - /// configured buffer size) for a write to occur, or similarly, when a data file is deleted. - #[cfg_attr(test, instrument(skip(self), level = "trace"))] - pub async fn wait_for_reader(&self) { - self.reader_notify.notified().await; + pub(super) fn subscribe_reader_progress(&self) -> watch::Receiver { + self.reader_progress.subscribe() } - /// Waits for a signal from the writer that progress has been made. - /// - /// This will occur when a record is written, or when a new data file is created. - /// - /// Writer progress is published before this notification is sent. - #[cfg_attr(test, instrument(skip(self), level = "trace"))] - pub async fn wait_for_writer(&self) { - self.writer_notify.notified().await; + pub(super) fn subscribe_writer_progress(&self) -> watch::Receiver { + self.writer_progress.subscribe() } /// Notifies all tasks waiting on progress by the reader. #[cfg_attr(test, instrument(skip(self), level = "trace"))] pub fn notify_reader_waiters(&self) { - self.reader_notify.notify_one(); + self.reader_progress + .send_modify(|revision| *revision = revision.wrapping_add(1)); } /// Notifies all tasks waiting on progress by the writer. @@ -521,19 +552,43 @@ where /// Callers must publish their shared state before notifying. #[cfg_attr(test, instrument(skip(self), level = "trace"))] pub fn notify_writer_waiters(&self) { - self.writer_notify.notify_one(); - } - - /// Publishes flushed writer progress and then wakes the reader. - pub fn publish_writer_progress(&self, event_count: u64, record_size: u64) -> u64 { - let next_record_id = self.state().increment_next_writer_record_id(event_count); + self.writer_progress + .send_modify(|progress| progress.revision = progress.revision.wrapping_add(1)); + } + + /// Publishes completed writer progress and then wakes the reader. + pub fn publish_writer_progress( + &self, + event_count: u64, + record_size: u64, + file_id: u16, + byte_offset: u64, + ) -> u64 { + // The next writer record ID is the reader's publication gate. Publish every other piece of + // shared state first so an acquire load that observes the new ID also observes the matching + // occupancy and usage accounting. self.increment_total_buffer_size(record_size); self.usage_handle .increment_received_event_count_and_byte_size(event_count, record_size); - self.notify_writer_waiters(); + let next_record_id = self.state().increment_next_writer_record_id(event_count); + self.writer_progress.send_modify(|progress| { + progress.revision = progress.revision.wrapping_add(1); + progress.published_position = Some((file_id, byte_offset)); + progress.next_record_id = next_record_id; + }); next_record_id } + /// Publishes a writer file/offset boundary without changing record or occupancy accounting. + pub(super) fn publish_writer_position(&self, position: Option<(u16, u64)>) { + let next_record_id = self.state().get_next_writer_record_id(); + self.writer_progress.send_modify(|progress| { + progress.revision = progress.revision.wrapping_add(1); + progress.published_position = position; + progress.next_record_id = next_record_id; + }); + } + /// Tracks the statistics of multiple successful reads. pub fn track_reads(&self, event_count: u64, total_record_size: u64) { self.decrement_total_buffer_size(total_record_size); @@ -555,14 +610,24 @@ where /// If the writer was not yet marked done, `false` is returned. Otherwise, `true` is returned, /// and the caller should handle any necessary logic for closing the writer. pub fn mark_writer_done(&self) -> bool { - self.writer_done - .compare_exchange_weak(false, true, Ordering::SeqCst, Ordering::SeqCst) - .is_ok() + let marked = self + .writer_done + .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) + .is_ok(); + if marked { + self.writer_progress.send_modify(|progress| { + progress.revision = progress.revision.wrapping_add(1); + progress.done = true; + }); + } + marked } - /// Returns `true` if the writer was marked as done. - pub fn is_writer_done(&self) -> bool { - self.writer_done.load(Ordering::Acquire) + pub(super) fn mark_writer_failed(&self) { + self.writer_progress.send_modify(|progress| { + progress.revision = progress.revision.wrapping_add(1); + progress.failed = true; + }); } /// Increments the pending acknowledgement counter by the given amount. @@ -796,14 +861,19 @@ where // Create the ledger object, and synchronize the buffer statistics with the buffer usage // handle. This handles making sure we account for the starting size of the buffer, and // what not. - let cleanup_reader_file_id = ledger_state.get_archive_ref().get_current_reader_file_id(); + let archived_state = ledger_state.get_archive_ref(); + let cleanup_reader_file_id = archived_state.get_current_reader_file_id(); + let initial_writer_progress = + WriterProgress::initial(archived_state.get_next_writer_record_id()); + let (reader_progress, _) = watch::channel(0); + let (writer_progress, _) = watch::channel(initial_writer_progress); let ledger = Ledger { config, lock, state: ledger_state, total_buffer_size: AtomicU64::new(0), - reader_notify: Notify::new(), - writer_notify: Notify::new(), + reader_progress, + writer_progress, writer_done: AtomicBool::new(false), pending_acks: AtomicU64::new(0), unacked_reader_file_id_offset: AtomicU16::new(0), diff --git a/lib/vector-buffers/src/variants/disk_v2/mod.rs b/lib/vector-buffers/src/variants/disk_v2/mod.rs index 970a416ba9e5a..1adbf810060e8 100644 --- a/lib/vector-buffers/src/variants/disk_v2/mod.rs +++ b/lib/vector-buffers/src/variants/disk_v2/mod.rs @@ -259,6 +259,10 @@ where .validate_last_write() .await .context(WriterSeekFailedSnafu)?; + // Validation may advance the runtime writer checkpoint beyond the value loaded from the + // ledger. Publish that reconciled state before constructing the reader so its persistent + // progress receiver cannot start from a stale record-ID boundary. + writer.publish_current_position(); let finalizer = Arc::clone(&ledger).spawn_finalizer(); @@ -281,6 +285,7 @@ where .seek_to_next_record() .await .context(ReaderSeekFailedSnafu)?; + writer.publish_current_position(); // Install the authoritative buffer size now that the reader is positioned at the first // unread record. Startup recovery treats the durable ledger checkpoint as the logical diff --git a/lib/vector-buffers/src/variants/disk_v2/reader.rs b/lib/vector-buffers/src/variants/disk_v2/reader.rs index c53c1b72f3e46..1b4310d09525a 100644 --- a/lib/vector-buffers/src/variants/disk_v2/reader.rs +++ b/lib/vector-buffers/src/variants/disk_v2/reader.rs @@ -9,13 +9,16 @@ use std::{ use crc32fast::Hasher; use rkyv::{AlignedVec, archived_root}; use snafu::{ResultExt, Snafu}; -use tokio::io::{AsyncBufReadExt, AsyncRead, BufReader}; +use tokio::{ + io::{AsyncBufReadExt, AsyncRead, BufReader}, + sync::watch, +}; use vector_common::{finalization::BatchNotifier, finalizer::OrderedFinalizer}; use super::{ Filesystem, common::create_crc32c_hasher, - ledger::Ledger, + ledger::{Ledger, WriterProgress}, record::{ArchivedRecord, Record, RecordStatus, validate_record_archive}, }; use crate::{ @@ -426,6 +429,7 @@ where FS: Filesystem, { ledger: Arc>, + writer_progress: watch::Receiver, reader: Option>, pending_read_token: Option, bytes_read: u64, @@ -450,9 +454,11 @@ where pub(crate) fn new(ledger: Arc>, finalizer: OrderedFinalizer) -> Self { let ledger_last_reader_record_id = ledger.state().get_last_reader_record_id(); let next_expected_record_id = ledger_last_reader_record_id + 1; + let writer_progress = ledger.subscribe_writer_progress(); Self { ledger, + writer_progress, reader: None, pending_read_token: None, bytes_read: 0, @@ -482,6 +488,17 @@ where self.data_file_start_record_id = None; } + fn is_record_published(&self, token: &ReadToken) -> bool { + token.record_id() < self.writer_progress.borrow().next_record_id() + } + + #[cfg_attr(test, instrument(skip(self), level = "trace"))] + async fn wait_for_writer(&mut self) { + match self.writer_progress.changed().await { + Ok(()) | Err(_) => {} + } + } + fn track_read(&mut self, record_id: u64, record_bytes: u64, event_count: NonZeroU64) { // We explicitly reduce the event count by one here in order to correctly calculate the // "last" record ID, which you can visualize as follows... @@ -819,7 +836,7 @@ where data_file_path = data_file_path.to_string_lossy().as_ref(), "Data file does not yet exist. Waiting for writer to create." ); - self.ledger.wait_for_writer().await; + self.wait_for_writer().await; } else { self.ledger.increment_acked_reader_file_id(); } @@ -1016,14 +1033,6 @@ where .context(IoSnafu)?; force_check_pending_data_files = false; - // Startup can read one valid frame beyond the durable reader checkpoint while using - // its ID to bound a preceding corrupt frame. The underlying file cursor has already - // advanced, but the frame itself remains in `RecordReader::aligned_buf`, so retain its - // token and deliver it through the normal read path once initialization is complete. - if let Some(token) = self.pending_read_token.take() { - break token; - } - // If the writer has marked themselves as done, and the buffer has been emptied, then // we're done and can return. We have to look at something besides simply the writer // being marked as done to know if we're actually done or not, and "buffer size" is better @@ -1034,17 +1043,53 @@ where // corrupted records, but hadn't yet had a "good" record that we could read, since the // "we skipped records due to corruption" logic requires performing valid read to // detect, and calculate a valid delta from. - if self.ledger.is_writer_done() { + let writer_progress = *self.writer_progress.borrow(); + if writer_progress.is_failed() { + return Err(ReaderError::Io { + source: io::Error::other("disk buffer writer failed"), + }); + } + if writer_progress.is_done() { let total_buffer_size = self.ledger.get_total_buffer_size(); if total_buffer_size == 0 { return Ok(None); } } + // Startup can read one valid frame beyond the durable reader checkpoint while using + // its ID to bound a preceding corrupt frame. Runtime reads can likewise observe a + // physically flushed record in the short interval before the writer publishes it to + // the ledger. The underlying file cursor has already advanced in either case, so keep + // the token until it is safe to deliver. An unrelated notification only causes this + // predicate to be checked again. + if let Some(token) = self.pending_read_token.take() { + if self.ready_to_read && !self.is_record_published(&token) { + self.pending_read_token = Some(token); + self.wait_for_writer().await; + continue; + } + break token; + } + self.ensure_ready_for_read().await.context(IoSnafu)?; let (reader_file_id, writer_file_id) = self.ledger.get_current_reader_writer_file_id(); + // The watch value is the latest completed writer boundary. Mark this revision as + // observed before checking it so a concurrent publication makes `changed()` return + // immediately. The buffered file reader may have prefetched later physical bytes, but + // record delivery remains bounded by both this byte position and the record-ID gate. + let writer_progress = *self.writer_progress.borrow_and_update(); + if self.ready_to_read + && let Some((published_file_id, published_byte_offset)) = + writer_progress.published_position() + && reader_file_id == published_file_id + && self.bytes_read >= published_byte_offset + { + self.wait_for_writer().await; + continue; + } + // Essentially: is the writer still writing to this data file or not, and are we // actually ready to read (aka initialized)? // @@ -1077,6 +1122,14 @@ where self.pending_read_token = Some(token); return Ok(None); } + // The file write becomes visible before the writer can synchronously publish its + // record and byte counters. Do not let a reader already draining the file consume + // and acknowledge that record against the old counters. + Ok(Some(token)) if !self.is_record_published(&token) => { + self.pending_read_token = Some(token); + self.wait_for_writer().await; + continue; + } // We got a valid record within the startup replay window, or a normal runtime // record, so keep the token. Ok(Some(token)) => break token, @@ -1149,7 +1202,7 @@ where continue; } - self.ledger.wait_for_writer().await; + self.wait_for_writer().await; } else { debug!( bytes_read = self.bytes_read, diff --git a/lib/vector-buffers/src/variants/disk_v2/tests/acknowledgements.rs b/lib/vector-buffers/src/variants/disk_v2/tests/acknowledgements.rs index a63ed70635a40..9eefbbec63734 100644 --- a/lib/vector-buffers/src/variants/disk_v2/tests/acknowledgements.rs +++ b/lib/vector-buffers/src/variants/disk_v2/tests/acknowledgements.rs @@ -68,7 +68,8 @@ async fn ack_wakes_reader() { let ledger = Arc::new(ledger); let finalizer = Arc::clone(&ledger).spawn_finalizer(); - let mut wait_for_writer = spawn(ledger.wait_for_writer()); + let mut writer_progress = ledger.subscribe_writer_progress(); + let mut wait_for_writer = spawn(writer_progress.changed()); assert_pending!(wait_for_writer.poll()); assert!(!wait_for_writer.is_woken()); @@ -78,7 +79,8 @@ async fn ack_wakes_reader() { acknowledge(batch).await; assert!(wait_for_writer.is_woken()); - assert_ready!(wait_for_writer.poll()); + assert_ready!(wait_for_writer.poll()) + .expect("writer progress channel should remain open"); } }) .await; diff --git a/lib/vector-buffers/src/variants/disk_v2/tests/invariants.rs b/lib/vector-buffers/src/variants/disk_v2/tests/invariants.rs index 29be00646a2cc..0d8f60f0b8b05 100644 --- a/lib/vector-buffers/src/variants/disk_v2/tests/invariants.rs +++ b/lib/vector-buffers/src/variants/disk_v2/tests/invariants.rs @@ -14,6 +14,7 @@ use crate::{ create_buffer_v2_with_write_buffer_size, create_default_buffer_v2, get_corrected_max_record_size, get_minimum_data_file_size_for_record_payload, }, + writer::RecordWriter, }, }; @@ -90,17 +91,27 @@ async fn publishing_writer_progress_updates_ledger_before_waking_reader() { let (_writer, _reader, ledger) = create_default_buffer_v2::<_, SizedRecord>(data_dir).await; - let mut wait_for_writer = spawn(ledger.wait_for_writer()); + let mut writer_progress = ledger.subscribe_writer_progress(); + let mut wait_for_writer = spawn(writer_progress.changed()); assert_pending!(wait_for_writer.poll()); assert!(!wait_for_writer.is_woken()); - let next_record_id = ledger.publish_writer_progress(1, 128); + let next_record_id = + ledger.publish_writer_progress(1, 128, ledger.get_current_writer_file_id(), 128); assert_eq!(next_record_id, 2); assert_eq!(ledger.state().get_next_writer_record_id(), 2); assert_eq!(ledger.get_total_buffer_size(), 128); assert!(wait_for_writer.is_woken()); - assert_ready!(wait_for_writer.poll()); + assert_ready!(wait_for_writer.poll()) + .expect("writer progress channel should remain open"); + drop(wait_for_writer); + let published = *writer_progress.borrow(); + assert_eq!( + published.published_position(), + Some((ledger.get_current_writer_file_id(), 128)) + ); + assert_eq!(published.next_record_id(), 2); } }); @@ -108,6 +119,189 @@ async fn publishing_writer_progress_updates_ledger_before_waking_reader() { fut.instrument(parent.or_current()).await; } +#[tokio::test] +async fn writer_progress_published_before_wait_is_not_missed() { + let _a = install_tracing_helpers(); + + with_temp_dir(|dir| { + let data_dir = dir.to_path_buf(); + + async move { + let (_writer, _reader, ledger) = + create_default_buffer_v2::<_, SizedRecord>(data_dir).await; + let mut writer_progress = ledger.subscribe_writer_progress(); + + ledger.publish_writer_progress(1, 128, ledger.get_current_writer_file_id(), 128); + + writer_progress + .changed() + .await + .expect("writer progress channel should remain open"); + let published = *writer_progress.borrow_and_update(); + assert_eq!( + published.published_position(), + Some((ledger.get_current_writer_file_id(), 128)) + ); + assert_eq!(published.next_record_id(), 2); + } + }) + .await; +} + +#[tokio::test] +async fn reader_progress_published_before_wait_is_not_missed() { + let _a = install_tracing_helpers(); + + with_temp_dir(|dir| { + let data_dir = dir.to_path_buf(); + + async move { + let (_writer, _reader, ledger) = + create_default_buffer_v2::<_, SizedRecord>(data_dir).await; + let mut reader_progress = ledger.subscribe_reader_progress(); + + ledger.notify_reader_waiters(); + + reader_progress + .changed() + .await + .expect("reader progress channel should remain open"); + } + }) + .await; +} + +#[tokio::test] +async fn reader_uses_record_gate_when_writer_position_is_unavailable() { + let _a = install_tracing_helpers(); + + with_temp_dir(|dir| { + let data_dir = dir.to_path_buf(); + + async move { + let (mut writer, mut reader, ledger) = + create_default_buffer_v2::<_, SizedRecord>(data_dir).await; + let record = SizedRecord::new(64); + + writer + .write_record(record.clone()) + .await + .expect("write should not fail"); + writer.flush().await.expect("flush should not fail"); + + // Recovery can temporarily have no open writer file while retaining readable, + // checkpointed records. In that state, the record-ID boundary remains authoritative. + ledger.publish_writer_position(None); + + assert_eq!(read_next(&mut reader).await, Some(record)); + } + }) + .await; +} + +#[tokio::test] +async fn writer_failure_terminates_reader_wait() { + let _a = install_tracing_helpers(); + + with_temp_dir(|dir| { + let data_dir = dir.to_path_buf(); + + async move { + let (_writer, mut reader, ledger) = + create_default_buffer_v2::<_, SizedRecord>(data_dir).await; + ledger.mark_writer_failed(); + + let error = reader + .next() + .await + .expect_err("writer failure should terminate the reader"); + assert!(matches!( + error, + crate::variants::disk_v2::ReaderError::Io { .. } + )); + } + }) + .await; +} + +#[tokio::test] +async fn reader_waits_for_visible_record_to_be_published() { + let assertion_registry = install_tracing_helpers(); + + let fut = with_temp_dir(|dir| { + let data_dir = dir.to_path_buf(); + + async move { + let (_writer, mut reader, ledger) = + create_default_buffer_v2::<_, SizedRecord>(data_dir).await; + let config = ledger.config().clone(); + let data_file = OpenOptions::new() + .append(true) + .open(ledger.get_current_writer_data_file_path()) + .await + .expect("data file should open"); + let mut physical_writer = RecordWriter::new( + data_file, + 0, + config.write_buffer_size, + config.max_data_file_size, + config.max_record_size, + ); + let record = SizedRecord::new(64); + let (record_bytes, _) = physical_writer + .write_record(1, record.clone()) + .await + .expect("physical write should succeed"); + physical_writer + .flush() + .await + .expect("physical flush should succeed"); + + assert_eq!(ledger.state().get_next_writer_record_id(), 1); + assert_buffer_is_empty!(ledger); + + let waiting_for_publication = assertion_registry + .build() + .with_name("wait_for_writer") + .with_parent_name("reader_waits_for_visible_record_to_be_published") + .was_entered() + .finalize(); + let mut blocked_read = spawn(read_next(&mut reader)); + while !waiting_for_publication.try_assert() { + assert_pending!(blocked_read.poll()); + } + assert_pending!(blocked_read.poll()); + assert!(!blocked_read.is_woken()); + + ledger.notify_writer_waiters(); + + assert!(blocked_read.is_woken()); + assert_pending!(blocked_read.poll()); + + ledger.publish_writer_progress( + 1, + record_bytes as u64, + ledger.get_current_writer_file_id(), + record_bytes as u64, + ); + + assert!(blocked_read.is_woken()); + let mut read_result = None; + for _ in 0..100 { + if let std::task::Poll::Ready(result) = blocked_read.poll() { + read_result = Some(result); + break; + } + tokio::task::yield_now().await; + } + assert_eq!(read_result, Some(Some(record))); + } + }); + + let parent = trace_span!("reader_waits_for_visible_record_to_be_published"); + fut.instrument(parent.or_current()).await; +} + #[tokio::test] async fn last_record_is_valid_during_load_when_buffer_correctly_flushed_and_stopped() { let assertion_registry = install_tracing_helpers(); diff --git a/lib/vector-buffers/src/variants/disk_v2/tests/mod.rs b/lib/vector-buffers/src/variants/disk_v2/tests/mod.rs index 1f3a8f2e2d89e..c30469f8bb331 100644 --- a/lib/vector-buffers/src/variants/disk_v2/tests/mod.rs +++ b/lib/vector-buffers/src/variants/disk_v2/tests/mod.rs @@ -30,6 +30,7 @@ mod known_errors; mod model; mod record; mod size_limits; +mod writer; impl AsyncFile for DuplexStream { async fn metadata(&self) -> io::Result { diff --git a/lib/vector-buffers/src/variants/disk_v2/tests/writer.rs b/lib/vector-buffers/src/variants/disk_v2/tests/writer.rs new file mode 100644 index 0000000000000..98fb553fdce97 --- /dev/null +++ b/lib/vector-buffers/src/variants/disk_v2/tests/writer.rs @@ -0,0 +1,204 @@ +use std::{ + io, + pin::Pin, + sync::{Arc, Mutex}, + task::{Context, Poll, Waker}, +}; + +use tokio::io::{AsyncRead, AsyncWrite, ReadBuf}; +use tokio_test::{assert_pending, assert_ready, task::spawn}; + +use crate::{ + test::SizedRecord, + variants::disk_v2::{ + io::{AsyncFile, Metadata}, + writer::{FlushResult, RecordWriter}, + }, +}; + +const MAX_RECORD_SIZE: usize = 1024 * 1024; +const MAX_DATA_FILE_SIZE: u64 = MAX_RECORD_SIZE as u64; + +#[derive(Debug, Default)] +struct FlushState { + pending: Vec, + visible: Vec, + flush_allowed: bool, + flush_waker: Option, +} + +#[derive(Clone, Debug, Default)] +struct FlushControl { + state: Arc>, +} + +impl FlushControl { + fn allow_flush(&self) { + let waker = { + let mut state = self + .state + .lock() + .expect("flush state should not be poisoned"); + state.flush_allowed = true; + state.flush_waker.take() + }; + if let Some(waker) = waker { + waker.wake(); + } + } + + fn pending_len(&self) -> usize { + self.state + .lock() + .expect("flush state should not be poisoned") + .pending + .len() + } + + fn visible_len(&self) -> usize { + self.state + .lock() + .expect("flush state should not be poisoned") + .visible + .len() + } +} + +#[derive(Clone, Debug, Default)] +struct FlushGatedFile { + control: FlushControl, +} + +impl AsyncRead for FlushGatedFile { + fn poll_read( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + _buf: &mut ReadBuf<'_>, + ) -> Poll> { + Poll::Ready(Ok(())) + } +} + +impl AsyncWrite for FlushGatedFile { + fn poll_write( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + let mut state = self + .control + .state + .lock() + .expect("flush state should not be poisoned"); + state.pending.extend_from_slice(buf); + Poll::Ready(Ok(buf.len())) + } + + fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let mut state = self + .control + .state + .lock() + .expect("flush state should not be poisoned"); + if !state.flush_allowed { + state.flush_waker = Some(cx.waker().clone()); + return Poll::Pending; + } + + let pending = std::mem::take(&mut state.pending); + state.visible.extend_from_slice(&pending); + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + self.poll_flush(cx) + } +} + +impl AsyncFile for FlushGatedFile { + async fn metadata(&self) -> io::Result { + let state = self + .control + .state + .lock() + .expect("flush state should not be poisoned"); + Ok(Metadata { + len: (state.pending.len() + state.visible.len()) as u64, + }) + } + + async fn truncate(&self, _size: u64) -> io::Result<()> { + Ok(()) + } + + async fn sync_all(&self) -> io::Result<()> { + Ok(()) + } +} + +fn record_writer( + write_buffer_size: usize, +) -> (RecordWriter, FlushControl) { + let inner = FlushGatedFile::default(); + let control = inner.control.clone(); + ( + RecordWriter::new( + inner, + 0, + write_buffer_size, + MAX_DATA_FILE_SIZE, + MAX_RECORD_SIZE, + ), + control, + ) +} + +#[tokio::test] +async fn direct_write_waits_for_visibility_before_reporting_progress() { + let (mut writer, control) = record_writer(1); + + let mut write = spawn(writer.write_record(1, SizedRecord::new(64))); + assert_pending!(write.poll()); + assert!(control.pending_len() > 0); + assert_eq!(control.visible_len(), 0); + + control.allow_flush(); + + assert!(write.is_woken()); + let (record_bytes, result) = assert_ready!(write.poll()).expect("write should succeed"); + assert_eq!( + result, + Some(FlushResult { + events_flushed: 1, + bytes_flushed: record_bytes as u64, + }) + ); + assert_eq!(control.visible_len(), record_bytes); +} + +#[tokio::test] +async fn buffered_write_waits_for_visibility_before_reporting_flush_progress() { + let (mut writer, control) = record_writer(MAX_RECORD_SIZE); + let (record_bytes, result) = writer + .write_record(1, SizedRecord::new(64)) + .await + .expect("write should succeed"); + assert_eq!(result, None); + + let mut flush = spawn(writer.flush()); + assert_pending!(flush.poll()); + assert!(control.pending_len() > 0); + assert_eq!(control.visible_len(), 0); + + control.allow_flush(); + + assert!(flush.is_woken()); + assert_eq!( + assert_ready!(flush.poll()).expect("flush should succeed"), + Some(FlushResult { + events_flushed: 1, + bytes_flushed: record_bytes as u64, + }) + ); + assert_eq!(control.visible_len(), record_bytes); +} diff --git a/lib/vector-buffers/src/variants/disk_v2/writer.rs b/lib/vector-buffers/src/variants/disk_v2/writer.rs index 2c02d044300a9..5f04d589fdb87 100644 --- a/lib/vector-buffers/src/variants/disk_v2/writer.rs +++ b/lib/vector-buffers/src/variants/disk_v2/writer.rs @@ -22,7 +22,10 @@ use rkyv::{ }, }; use snafu::{ResultExt, Snafu}; -use tokio::io::{AsyncWrite, AsyncWriteExt}; +use tokio::{ + io::{AsyncWrite, AsyncWriteExt}, + sync::watch, +}; use super::{ common::{DiskBufferConfig, create_crc32c_hasher}, @@ -48,6 +51,10 @@ pub enum WriterError where T: Bufferable, { + /// The writer was asked to shut down while waiting for reader progress. + #[snafu(display("writer shutting down"))] + Shutdown, + /// A general I/O error occurred. /// /// Different methods will capture specific I/O errors depending on the situation, as some @@ -375,6 +382,7 @@ impl TrackingBufWriter { // If the given buffer is too large to be buffered at all, then bypass the internal buffer. if buf.len() >= self.buf.capacity() { self.inner.write_all(buf).await?; + self.inner.flush().await?; let flush_result = flush_result.get_or_insert(FlushResult::default()); flush_result.events_flushed += event_count as u64; @@ -409,6 +417,10 @@ impl TrackingBufWriter { let bytes_flushed = self.buf.len() as u64; let result = self.inner.write_all(&self.buf[..]).await; + let result = match result { + Ok(()) => self.inner.flush().await, + Err(error) => Err(error), + }; self.unflushed_events = 0; self.buf.clear(); @@ -876,6 +888,8 @@ where FS::File: Unpin, { ledger: Arc>, + reader_progress: watch::Receiver, + shutdown: Option>, config: DiskBufferConfig, writer: Option>, next_record_id: u64, @@ -898,8 +912,11 @@ where pub(crate) fn new(ledger: Arc>) -> Self { let config = ledger.config().clone(); let next_record_id = ledger.state().get_next_writer_record_id(); + let reader_progress = ledger.subscribe_reader_progress(); BufferWriter { ledger, + reader_progress, + shutdown: None, config, writer: None, data_file_size: 0, @@ -935,9 +952,48 @@ where self.unflushed_events -= flushed_events; self.unflushed_bytes -= flushed_bytes; - self.next_record_id = self - .ledger - .publish_writer_progress(flushed_events, flushed_bytes); + let published_byte_offset = self.data_file_size - self.unflushed_bytes; + self.next_record_id = self.ledger.publish_writer_progress( + flushed_events, + flushed_bytes, + self.ledger.get_current_writer_file_id(), + published_byte_offset, + ); + } + + pub(super) fn publish_current_position(&self) { + let position = self.writer.as_ref().map(|_| { + ( + self.ledger.get_current_writer_file_id(), + self.data_file_size - self.unflushed_bytes, + ) + }); + self.ledger.publish_writer_position(position); + } + + #[cfg_attr(test, instrument(skip(self), level = "trace"))] + async fn wait_for_reader(&mut self) -> bool { + if let Some(shutdown) = self.shutdown.as_mut() { + tokio::select! { + biased; + _ = shutdown.changed() => true, + result = self.reader_progress.changed() => { + match result { + Ok(()) | Err(_) => {} + } + false + } + } + } else { + match self.reader_progress.changed().await { + Ok(()) | Err(_) => {} + } + false + } + } + + pub(crate) fn set_shutdown(&mut self, shutdown: watch::Receiver<()>) { + self.shutdown = Some(shutdown); } fn can_write(&self) -> bool { @@ -1050,7 +1106,7 @@ where let previous_writer_file_id = self.ledger.get_current_writer_file_id(); self.reset(); self.ledger.state().increment_writer_file_id(); - self.ensure_ready_for_write().await.context(IoSnafu)?; + self.ensure_ready_for_write().await?; self.reconcile_current_data_file_with_checkpoint().await?; self.ledger.flush().context(IoSnafu)?; @@ -1279,7 +1335,7 @@ where current_writer_data_file = ?self.ledger.get_current_writer_data_file_path(), "Validating last written record in current data file." ); - self.ensure_ready_for_write().await.context(IoSnafu)?; + self.ensure_ready_for_write().await?; // If our current file is empty, there's no sense doing this check. if self.data_file_size == 0 { @@ -1443,7 +1499,7 @@ where self.reset(); self.mark_for_skip(); - self.ensure_ready_for_write().await.context(IoSnafu)?; + self.ensure_ready_for_write().await?; self.ledger.flush().context(IoSnafu)?; debug!( @@ -1465,7 +1521,7 @@ where // The inline antithesis assertion block pushes this over the line limit. Its // source lines count even when the feature is off, so the allow is unconditional. #[allow(clippy::too_many_lines)] - async fn ensure_ready_for_write(&mut self) -> io::Result<()> { + async fn ensure_ready_for_write(&mut self) -> Result<(), WriterError> { // Check the overall size of the buffer and figure out if we can write. loop { // If we haven't yet exceeded the maximum buffer size, then we can proceed. Likewise, if @@ -1498,7 +1554,9 @@ where ); } - self.ledger.wait_for_reader().await; + if self.wait_for_reader().await { + return Err(WriterError::Shutdown); + } } // If we already have an open writer, and we have no more space in the data file to write, @@ -1520,7 +1578,7 @@ where // // We still flush ourselves to disk, etc, to make sure all of the data is there. should_open_next = true; - self.flush_inner(true).await?; + self.flush_inner(true).await.context(IoSnafu)?; self.reset(); } @@ -1577,8 +1635,9 @@ where .ledger .filesystem() .open_file_writable(&data_file_path) - .await?; - let metadata = data_file.metadata().await?; + .await + .context(IoSnafu)?; + let metadata = data_file.metadata().await.context(IoSnafu)?; let file_len = metadata.len(); if file_len == 0 || !should_open_next { // The file is either empty, which means we created it and "own it" now, @@ -1595,7 +1654,7 @@ where } } // Legitimate I/O error with the operation, bubble this up. - _ => return Err(e), + _ => return Err(WriterError::Io { source: e }), }, }; @@ -1608,7 +1667,7 @@ where ); // Make sure the file is flushed to disk, especially if we just created it. - data_file.sync_all().await?; + data_file.sync_all().await.context(IoSnafu)?; self.writer = Some(RecordWriter::new( data_file, @@ -1623,7 +1682,7 @@ where // file ID now to signal that the writer has moved on. if should_open_next { self.ledger.state().increment_writer_file_id(); - self.ledger.notify_writer_waiters(); + self.publish_current_position(); // The writer just rolled to a fresh data file, the boundary the // crash, partial-write, and file-id-rollover faults act on. @@ -1655,7 +1714,9 @@ where // Wait until the reader signals progress and try again. debug!("Target data file is still present and not yet processed. Waiting for reader."); - self.ledger.wait_for_reader().await; + if self.wait_for_reader().await { + return Err(WriterError::Shutdown); + } } } @@ -1712,7 +1773,7 @@ where // Make sure we have an open data file to write to, which might also be us opening the // next data file because our first attempt at writing had to finalize a data file that // was already full. - self.ensure_ready_for_write().await.context(IoSnafu)?; + self.ensure_ready_for_write().await?; let writer = self .writer @@ -1800,7 +1861,8 @@ where // error like `DataFileFull`, as that gets checked during archiving. The guard fires // Errored automatically on `?` exit. let result = writer.flush_record(token).await?; - // Record is durable on disk; disarm so finalizers resolve as Delivered. + // The record has been accepted by the file writer; disarm so finalizers resolve as + // Delivered. The flush result below independently controls when progress is published. record_finalizers.disarm(); result } else { @@ -1876,7 +1938,9 @@ where Ok(bytes_written) => return Ok(bytes_written), Err(old_record) => { record = old_record; - self.ledger.wait_for_reader().await; + if self.wait_for_reader().await { + return Err(WriterError::Shutdown); + } } } } @@ -1901,10 +1965,9 @@ where #[instrument(skip(self), level = "debug")] async fn flush_inner(&mut self, force_full_flush: bool) -> io::Result<()> { - // We always flush the `BufWriter` when this is called, but we don't always flush to disk or - // flush the ledger. This is enough for readers on Linux since the file ends up in the page - // cache, as we don't do any O_DIRECT fanciness, and the new contents can be immediately - // read. + // We always flush the `BufWriter` and wait for the underlying asynchronous file write when + // this is called, but we don't always synchronize the file or flush the ledger. Completion + // is enough to make the new contents visible to readers through the page cache on Linux. // // TODO: Windows has a page cache as well, and macOS _should_, but we should verify this // behavior works on those platforms as well. @@ -1967,9 +2030,12 @@ where pub fn close(&mut self) { if self.ledger.mark_writer_done() { debug!("Writer marked as closed."); - self.ledger.notify_writer_waiters(); } } + + pub(crate) fn fail(&self) { + self.ledger.mark_writer_failed(); + } } impl Drop for BufferWriter From 2065990ca8bb95dc4e1ee7bab7722faa734801af Mon Sep 17 00:00:00 2001 From: "robert.blafford" Date: Mon, 10 Aug 2026 17:56:27 +0000 Subject: [PATCH 2/2] Add changelog fragment --- changelog.d/disk_v2_reader_writer_coordination.fix.md | 5 +++++ 1 file changed, 5 insertions(+) create mode 100644 changelog.d/disk_v2_reader_writer_coordination.fix.md diff --git a/changelog.d/disk_v2_reader_writer_coordination.fix.md b/changelog.d/disk_v2_reader_writer_coordination.fix.md new file mode 100644 index 0000000000000..db08fc70e5615 --- /dev/null +++ b/changelog.d/disk_v2_reader_writer_coordination.fix.md @@ -0,0 +1,5 @@ +Prevent disk buffers from stalling when reader or writer progress occurs immediately before the +other side begins waiting. Newly written records are also held until the writer publishes the +corresponding buffer accounting. + +authors: graphcareful