feat(mito): wait for WAL durability before publishing a manifest watermark (#9410)

* feat(mito): wait for WAL durability before publishing a manifest watermark

A log store that acknowledges appends before their entries are durable can
hand out entry ids that a crash loses. Before a flush, a full truncate or a
discard of unflushed data records an entry id in the manifest, Mito now waits
on `LogStore::wait_durable` for that id. The flush waits after its SSTs are
written and observes its cancellation, so a drop or truncate queued behind it
is not held up by the upload.

Add engine tests of the object store WAL across restarts and of the barrier,
and a log store testing hook that stores an object and then reports its create
as failed.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(mito): synchronize the WAL barrier tests on store signals

The barrier tests waited with fixed sleeps, which do not establish that a
flush or truncate has reached its durability wait, and the lost backlog test
could stop the store before the seal was handled. The object store WAL gets
testing hooks that observe parked creates and durability waits, and one that
ends the actor as a crash would; the tests wait on them, crash the store and
the engine before a reopen until nothing holds the old store, and cover a
truncate or discard of unflushed data that completes once its entry is
durable.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(mito): rename the recovery state and the crash test in the WAL recovery tests

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(mito): drop a comment that restates the empty recovery state

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
This commit is contained in:
jeremyhi
2026-09-30 05:20:59 +00:00
committed by GitHub
parent 78a7b93292
commit 8553db3f27
10 changed files with 1465 additions and 27 deletions
+1 -2
View File
@@ -62,6 +62,5 @@ mod format;
mod io;
mod store;
#[allow(unused_imports)]
pub(crate) use batch::entry_id;
pub use batch::entry_id;
pub use store::ObjectStoreLogStore;
+1 -1
View File
@@ -45,7 +45,7 @@ pub(crate) const OBJECT_SEQ_LIMIT: u64 = 1 << (u64::BITS - POSITION_BITS);
/// at one so that id zero, the watermark of a region without entries, is never
/// assigned. A region's ids increase with the object sequence and have gaps
/// wherever other regions or other positions took the sequence.
pub(crate) fn entry_id(object_seq: u64, position: u64) -> EntryId {
pub fn entry_id(object_seq: u64, position: u64) -> EntryId {
debug_assert!(object_seq < OBJECT_SEQ_LIMIT && (1..POSITION_LIMIT).contains(&position));
(object_seq << POSITION_BITS) | position
}
+193 -10
View File
@@ -104,6 +104,12 @@ pub struct ObjectStoreLogStore {
creates_held: watch::Sender<bool>,
#[cfg(any(test, feature = "testing"))]
creates_fail: Arc<AtomicBool>,
#[cfg(any(test, feature = "testing"))]
next_create_fails_after_write: Arc<AtomicBool>,
#[cfg(any(test, feature = "testing"))]
parked_creates: watch::Receiver<usize>,
#[cfg(any(test, feature = "testing"))]
durability_waits: watch::Receiver<usize>,
}
type ObsoleteEntryIds = Arc<Mutex<HashMap<RegionId, EntryId>>>;
@@ -203,6 +209,12 @@ impl ObjectStoreLogStore {
let (creates_held_tx, creates_held_rx) = watch::channel(false);
#[cfg(any(test, feature = "testing"))]
let creates_fail = Arc::new(AtomicBool::new(false));
#[cfg(any(test, feature = "testing"))]
let next_create_fails_after_write = Arc::new(AtomicBool::new(false));
#[cfg(any(test, feature = "testing"))]
let (parked_creates_tx, parked_creates_rx) = watch::channel(0);
#[cfg(any(test, feature = "testing"))]
let (durability_waits_tx, durability_waits_rx) = watch::channel(0);
let actor = Actor {
io: io.clone(),
@@ -237,6 +249,12 @@ impl ObjectStoreLogStore {
creates_held: creates_held_rx,
#[cfg(any(test, feature = "testing"))]
creates_fail: creates_fail.clone(),
#[cfg(any(test, feature = "testing"))]
next_create_fails_after_write: next_create_fails_after_write.clone(),
#[cfg(any(test, feature = "testing"))]
parked_creates: Arc::new(parked_creates_tx),
#[cfg(any(test, feature = "testing"))]
durability_waits: durability_waits_tx,
};
common_runtime::spawn_global(actor.run());
Ok(Arc::new(Self {
@@ -255,6 +273,12 @@ impl ObjectStoreLogStore {
creates_held: creates_held_tx,
#[cfg(any(test, feature = "testing"))]
creates_fail,
#[cfg(any(test, feature = "testing"))]
next_create_fails_after_write,
#[cfg(any(test, feature = "testing"))]
parked_creates: parked_creates_rx,
#[cfg(any(test, feature = "testing"))]
durability_waits: durability_waits_rx,
}))
}
@@ -313,6 +337,20 @@ impl ObjectStoreLogStore {
}
}
/// The transient object store error of a create that a testing hook fails.
#[cfg(any(test, feature = "testing"))]
fn injected_create_failure(io: &dyn WalObjectIo, object_seq: u64) -> Result<PutResult> {
let error = object_store::Error::new(
object_store::ErrorKind::Unexpected,
"injected create failure",
)
.set_temporary();
Err(error).context(crate::error::WalObjectStoreSnafu {
operation: "write",
path: io.object_path(object_seq),
})
}
/// Records the obsolete watermark of `region_id`, which never moves down.
fn record_obsolete(obsolete_entry_ids: &ObsoleteEntryIds, region_id: RegionId, entry_id: EntryId) {
obsolete_entry_ids
@@ -383,6 +421,57 @@ impl ObjectStoreLogStore {
self.creates_fail.store(true, Ordering::Release);
}
/// Makes the next create that writes its object report a transient object
/// store error afterwards, so the object exists although its create
/// failed.
pub fn fail_next_create_after_write(&self) {
self.next_create_fails_after_write
.store(true, Ordering::Release);
}
/// Waits until at least `expected` creates have been parked by
/// [`hold_creates`](Self::hold_creates) since the store was built.
pub async fn wait_for_parked_creates(&self, expected: usize) -> Result<()> {
self.parked_creates
.clone()
.wait_for(|count| *count >= expected)
.await
.ok()
.map(|_| ())
.context(ObjectStoreWalStoppedSnafu)
}
/// Waits until at least `expected` calls of
/// [`wait_durable`](LogStore::wait_durable) have had to wait for an entry
/// that is not durable since the store was built.
pub async fn wait_for_durability_waits(&self, expected: usize) -> Result<()> {
self.durability_waits
.clone()
.wait_for(|count| *count >= expected)
.await
.ok()
.map(|_| ())
.context(ObjectStoreWalStoppedSnafu)
}
/// Ends the actor the way a crash of the process would: nothing more is
/// written, creates that have not completed are dropped, and every caller
/// still waiting for the store fails. Returns once the actor has exited.
pub async fn crash(&self) {
self.stopped.store(true, Ordering::Release);
let (response_tx, response_rx) = oneshot::channel();
if self
.command_tx
.send(Command::Crash {
response: response_tx,
})
.await
.is_ok()
{
let _ = response_rx.await;
}
}
/// Sets the stopped flag without sending the stop command, which is the
/// state a store is in between the two steps of [`stop`](LogStore::stop).
pub fn begin_stop(&self) {
@@ -662,6 +751,8 @@ enum Command {
Seal {
response: oneshot::Sender<Result<()>>,
},
#[cfg(any(test, feature = "testing"))]
Crash { response: oneshot::Sender<()> },
}
/// An append waiting for the object that holds its entries.
@@ -797,6 +888,12 @@ struct Actor {
creates_held: watch::Receiver<bool>,
#[cfg(any(test, feature = "testing"))]
creates_fail: Arc<AtomicBool>,
#[cfg(any(test, feature = "testing"))]
next_create_fails_after_write: Arc<AtomicBool>,
#[cfg(any(test, feature = "testing"))]
parked_creates: Arc<watch::Sender<usize>>,
#[cfg(any(test, feature = "testing"))]
durability_waits: watch::Sender<usize>,
}
impl Actor {
@@ -835,6 +932,12 @@ impl Actor {
Some(Command::Seal { response }) => {
self.handle_seal(response);
}
#[cfg(any(test, feature = "testing"))]
Some(Command::Crash { response }) => {
self.handle_crash();
let _ = response.send(());
return;
}
// Every sender is gone: the store was dropped without
// `stop`. The creates in flight are dropped with the actor.
None => return,
@@ -1076,6 +1179,10 @@ impl Actor {
let mut creates_held = self.creates_held.clone();
#[cfg(any(test, feature = "testing"))]
let creates_fail = self.creates_fail.clone();
#[cfg(any(test, feature = "testing"))]
let next_create_fails_after_write = self.next_create_fails_after_write.clone();
#[cfg(any(test, feature = "testing"))]
let parked_creates = self.parked_creates.clone();
self.creates.push(Box::pin(async move {
if !delay.is_zero() {
tokio::time::sleep(delay).await;
@@ -1083,21 +1190,16 @@ impl Actor {
// The store was dropped while the create was parked: it never
// runs.
#[cfg(any(test, feature = "testing"))]
if *creates_held.borrow() {
parked_creates.send_modify(|count| *count += 1);
}
#[cfg(any(test, feature = "testing"))]
if creates_held.wait_for(|held| !*held).await.is_err() {
return (object_seq, Err(ObjectStoreWalStoppedSnafu.build()));
}
#[cfg(any(test, feature = "testing"))]
if creates_fail.load(Ordering::Acquire) {
let error = object_store::Error::new(
object_store::ErrorKind::Unexpected,
"injected create failure",
)
.set_temporary();
let result = Err(error).context(crate::error::WalObjectStoreSnafu {
operation: "write",
path: io.object_path(object_seq),
});
return (object_seq, result);
return (object_seq, injected_create_failure(io.as_ref(), object_seq));
}
// Any conflict poisons an `enqueued` store, so only the
// `durable` mode reads the epoch of the existing object.
@@ -1109,6 +1211,10 @@ impl Actor {
}
result => result,
};
#[cfg(any(test, feature = "testing"))]
if result.is_ok() && next_create_fails_after_write.swap(false, Ordering::AcqRel) {
return (object_seq, injected_create_failure(io.as_ref(), object_seq));
}
(object_seq, result)
}));
}
@@ -1354,6 +1460,8 @@ impl Actor {
entry_id,
response,
});
#[cfg(any(test, feature = "testing"))]
self.durability_waits.send_modify(|count| *count += 1);
}
/// Makes sure no id of the region at or below `entry_id` is assigned
@@ -1489,6 +1597,21 @@ impl Actor {
let _ = response.send(result);
}
/// Fails every caller that waits for the store; the creates are
/// dropped with the actor.
#[cfg(any(test, feature = "testing"))]
fn handle_crash(&mut self) {
let stopped = || ObjectStoreWalStoppedSnafu.build();
for batch in self.sealed.drain(..) {
batch.fail(stopped);
}
self.reset_open_batch();
self.fail_unacknowledged(stopped);
for response in self.stop.drain(..) {
let _ = response.send(Err(stopped()));
}
}
fn is_stopped(&self) -> bool {
self.stopped.load(Ordering::Acquire)
}
@@ -2759,6 +2882,9 @@ mod tests {
admitted_appends: watch::channel(0).1,
creates_held: watch::channel(false).0,
creates_fail: Arc::default(),
next_create_fails_after_write: Arc::default(),
parked_creates: watch::channel(0).1,
durability_waits: watch::channel(0).1,
};
(store, command_rx)
}
@@ -3956,6 +4082,63 @@ mod tests {
store.stop().await.unwrap();
}
#[tokio::test]
async fn test_store_hook_fails_the_next_create_after_it_writes() {
let store = open(memory_store(), &eager()).await;
let region_id = region(1);
// Object 1 is stored, but its create reports a failure, so the
// append fails and the region has no durable entry.
store.fail_next_create_after_write();
let error = append(&store, region_id, "a1").await.unwrap_err();
assert!(
matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
"unexpected error: {error:?}"
);
assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
assert_eq!(0, latest(&store, region_id));
// Only that create failed.
let response = append(&store, region_id, "a2").await.unwrap();
assert_eq!(
HashMap::from([(region_id, id(2, 1))]),
response.last_entry_ids
);
assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await);
store.stop().await.unwrap();
}
#[tokio::test]
async fn test_store_hooks_observe_parked_creates_waits_and_crash() {
let store = open(memory_store(), &enqueued(eager())).await;
let region_id = region(1);
// The append is acknowledged and its create parks; a wait for its
// entry has to wait.
store.hold_creates();
append(&store, region_id, "a1").await.unwrap();
timeout(WAIT, store.wait_for_parked_creates(1))
.await
.unwrap()
.unwrap();
let wait = {
let store = store.clone();
tokio::spawn(async move { store.wait_durable(&provider(region_id), id(1, 1)).await })
};
timeout(WAIT, store.wait_for_durability_waits(1))
.await
.unwrap()
.unwrap();
assert!(!wait.is_finished());
// The crash fails the waiter and the store, and the parked create
// never writes its object.
store.crash().await;
assert_stopped(&timeout(WAIT, wait).await.unwrap().unwrap().unwrap_err());
assert_stopped(&append(&store, region_id, "a2").await.unwrap_err());
assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
}
#[tokio::test]
async fn test_store_create_held_when_the_store_is_dropped_never_runs() {
let (io, _) = RecordingIo::over(memory_store());
+2
View File
@@ -51,6 +51,8 @@ pub mod listener;
#[cfg(test)]
mod merge_mode_test;
#[cfg(test)]
mod object_store_wal_recovery_test;
#[cfg(test)]
mod object_store_wal_test;
#[cfg(test)]
mod open_test;
File diff suppressed because it is too large Load Diff
+13 -3
View File
@@ -453,6 +453,14 @@ pub enum Error {
source: BoxedError,
},
#[snafu(display("Failed to wait for the WAL to be durable, region_id: {}", region_id))]
WaitWalDurable {
region_id: RegionId,
#[snafu(implicit)]
location: Location,
source: BoxedError,
},
// Shared error for each writer in the write group.
#[snafu(display("Failed to write region"))]
WriteGroup { source: Arc<Error> },
@@ -1439,9 +1447,10 @@ impl ErrorExt for Error {
OpenDal { .. } | ManifestDeltaNotFound { .. } | ReadParquet { .. } => {
StatusCode::StorageUnavailable
}
WriteWal { source, .. } | ReadWal { source, .. } | DeleteWal { source, .. } => {
source.status_code()
}
WriteWal { source, .. }
| ReadWal { source, .. }
| DeleteWal { source, .. }
| WaitWalDurable { source, .. } => source.status_code(),
CompressObject { .. }
| DecompressObject { .. }
| SerdeJson { .. }
@@ -1651,6 +1660,7 @@ impl ErrorExt for Error {
WriteWal { source, .. }
| ReadWal { source, .. }
| DeleteWal { source, .. }
| WaitWalDurable { source, .. }
| FetchManifests { source, .. }
| External { source, .. } => source.retry_hint(),
+18
View File
@@ -69,6 +69,7 @@ use crate::sst::parquet::{
DEFAULT_READ_BATCH_SIZE, DEFAULT_ROW_GROUP_SIZE, SstInfo, WriteOptions, flat_format,
};
use crate::sst::{FlatSchemaOptions, FormatType, to_flat_sst_arrow_schema};
use crate::wal::DurabilityBarrier;
use crate::worker::WorkerListener;
/// Global write buffer (memtable) manager.
@@ -282,6 +283,8 @@ pub(crate) struct RegionFlushTask {
///
/// This is used to generate the file meta.
pub(crate) partition_expr: Option<String>,
/// Waits until the WAL is durable through the entry id the flush records.
pub(crate) durability_barrier: DurabilityBarrier,
}
struct FlushTaskWaiters {
@@ -509,6 +512,17 @@ impl RegionFlushTask {
.await;
}
// The manifest may name only a durable entry as flushed. A log store
// that acknowledges appends before their entries are durable can have
// handed out `last_entry_id` for an entry that is still in its backlog.
// A DDL that cancels the flush must not wait behind that upload.
let durable = self.durability_barrier.wait(version_data.last_entry_id);
tokio::pin!(durable);
match CancellableFuture::new(durable.as_mut(), state.cancel_handle()).await {
Ok(result) => result?,
Err(_) => return FlushCancelledSnafu.fail(),
}
let edit = RegionEdit {
files_to_add: file_metas,
files_to_remove: Vec::new(),
@@ -1703,6 +1717,7 @@ mod tests {
flush_semaphore: Arc::new(Semaphore::new(2)),
is_staging: false,
partition_expr: None,
durability_barrier: DurabilityBarrier::noop(),
}
}
@@ -1853,6 +1868,7 @@ mod tests {
flush_semaphore: Arc::new(Semaphore::new(2)),
is_staging: false,
partition_expr: None,
durability_barrier: DurabilityBarrier::noop(),
};
task.push_sender(OptionOutputTx::from(output_tx));
scheduler
@@ -2143,6 +2159,7 @@ mod tests {
flush_semaphore: Arc::new(Semaphore::new(2)),
is_staging: false,
partition_expr: None,
durability_barrier: DurabilityBarrier::noop(),
})
.collect();
// Schedule first task.
@@ -2447,6 +2464,7 @@ mod tests {
flush_semaphore: Arc::new(Semaphore::new(2)),
is_staging: false,
partition_expr: None,
durability_barrier: DurabilityBarrier::noop(),
})
.collect();
// Schedule first task.
+55 -1
View File
@@ -36,7 +36,7 @@ use store_api::logstore::provider::Provider;
use store_api::logstore::{AppendBatchResponse, LogStore, WalIndex};
use store_api::storage::RegionId;
use crate::error::{BuildEntrySnafu, DeleteWalSnafu, Result, WriteWalSnafu};
use crate::error::{BuildEntrySnafu, DeleteWalSnafu, Result, WaitWalDurableSnafu, WriteWalSnafu};
use crate::wal::entry_reader::{LogStoreEntryReader, WalEntryReader};
use crate::wal::raw_entry_reader::{LogStoreRawEntryReader, RegionRawEntryReader};
@@ -45,6 +45,36 @@ pub type EntryId = store_api::logstore::entry::Id;
/// A stream that yields tuple of WAL entry id and corresponding entry.
pub type WalEntryStream<'a> = BoxStream<'a, Result<(EntryId, WalEntry)>>;
/// Waits until the WAL of one region is durable through an entry id.
///
/// A flush calls it before it records an entry id as flushed in the manifest:
/// a log store that acknowledges an append before its entries are durable
/// hands out entry ids that a crash can lose, and a manifest watermark that
/// names such an id would skip the entries assigned again after a restart.
#[derive(Clone)]
pub(crate) struct DurabilityBarrier(
Arc<dyn Fn(EntryId) -> BoxFuture<'static, Result<()>> + Send + Sync>,
);
impl DurabilityBarrier {
/// Returns once the region is durable through `entry_id`.
pub(crate) async fn wait(&self, entry_id: EntryId) -> Result<()> {
(self.0)(entry_id).await
}
/// A barrier that never waits, for tests without a log store.
#[cfg(test)]
pub(crate) fn noop() -> Self {
Self(Arc::new(|_| Box::pin(async { Ok(()) })))
}
}
impl std::fmt::Debug for DurabilityBarrier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("DurabilityBarrier")
}
}
/// Write ahead log.
///
/// All regions in the engine shares the same WAL instance.
@@ -104,6 +134,30 @@ impl<S: LogStore> Wal<S> {
}
}
/// Returns the [DurabilityBarrier] of the region written through `provider`.
pub(crate) fn durability_barrier(
&self,
region_id: RegionId,
provider: &Provider,
) -> DurabilityBarrier {
let store = self.store.clone();
let provider = provider.clone();
DurabilityBarrier(Arc::new(move |entry_id| {
let store = store.clone();
let provider = provider.clone();
Box::pin(async move {
if let Provider::Noop = provider {
return Ok(());
}
store
.wait_durable(&provider, entry_id)
.await
.map_err(BoxedError::new)
.context(WaitWalDurableSnafu { region_id })
})
}))
}
/// Returns a [WalEntryReader]
pub(crate) fn wal_entry_reader(
&self,
+3
View File
@@ -243,6 +243,9 @@ impl<S: LogStore> RegionWorkerLoop<S> {
flush_semaphore: self.flush_semaphore.clone(),
is_staging: region.is_staging(),
partition_expr: region.maybe_staging_partition_expr_str(),
durability_barrier: self
.wal
.durability_barrier(region.region_id, &region.provider),
}
}
}
+29 -10
View File
@@ -32,7 +32,7 @@ use crate::cache::file_cache::{FileType, IndexKey};
use crate::config::IndexBuildMode;
use crate::error::{EditRegionSnafu, RegionBusySnafu, RegionNotFoundSnafu, Result};
use crate::manifest::action::{
RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionTruncate,
RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionTruncate, TruncateKind,
};
use crate::memtable::MemtableBuilderProvider;
use crate::metrics::WRITE_CACHE_INFLIGHT_DOWNLOAD;
@@ -484,17 +484,31 @@ impl<S: LogStore> RegionWorkerLoop<S> {
let request_sender = self.sender.clone();
let manifest_ctx = region.manifest_ctx.clone();
let is_staging = region.is_staging();
let durability_barrier = self
.wal
.durability_barrier(region.region_id, &region.provider);
// Updates manifest in background.
common_runtime::spawn_global(async move {
// The truncated entry id becomes the replay frontier, so the WAL
// must be durable through it before the manifest names it.
let durable = match &truncate.kind {
TruncateKind::All {
truncated_entry_id, ..
} => durability_barrier.wait(*truncated_entry_id).await,
_ => Ok(()),
};
// Write region truncated to manifest.
let action_list =
RegionMetaActionList::with_action(RegionMetaAction::Truncate(truncate.clone()));
let result = manifest_ctx
.update_manifest(RegionLeaderState::Truncating, action_list, is_staging)
.await
.map(|_| ());
let result = match durable {
Ok(()) => manifest_ctx
.update_manifest(RegionLeaderState::Truncating, action_list, is_staging)
.await
.map(|_| ()),
Err(e) => Err(e),
};
// Sends the result back to the request sender.
let truncate_result = TruncateResult {
@@ -531,10 +545,12 @@ impl<S: LogStore> RegionWorkerLoop<S> {
let region_id = region.region_id;
let request_sender = self.sender.clone();
let manifest_ctx = region.manifest_ctx.clone();
let durability_barrier = self.wal.durability_barrier(region_id, &region.provider);
common_runtime::spawn_global(async move {
// The frontier moves to the last written entry and sequence, so replaying the
// WAL after a restart skips everything the memtables held.
// WAL after a restart skips everything the memtables held. The WAL must be
// durable through that entry before the manifest names it.
let edit = RegionEdit {
files_to_add: Vec::new(),
files_to_remove: Vec::new(),
@@ -545,10 +561,13 @@ impl<S: LogStore> RegionWorkerLoop<S> {
committed_sequence: None,
};
let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
let result = manifest_ctx
.update_manifest(RegionLeaderState::Truncating, action_list, false)
.await
.map(|_| ());
let result = match durability_barrier.wait(discarded_entry_id).await {
Ok(()) => manifest_ctx
.update_manifest(RegionLeaderState::Truncating, action_list, false)
.await
.map(|_| ()),
Err(e) => Err(e),
};
let result = DiscardUnflushedResult {
region_id,