From 8553db3f276f0bf0fc526214b64c6382d478b079 Mon Sep 17 00:00:00 2001 From: jeremyhi Date: Wed, 30 Sep 2026 05:20:59 +0000 Subject: [PATCH] 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 * 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 * test(mito): rename the recovery state and the crash test in the WAL recovery tests Signed-off-by: jeremyhi * test(mito): drop a comment that restates the empty recovery state Signed-off-by: jeremyhi --------- Signed-off-by: jeremyhi --- src/log-store/src/object_store_wal.rs | 3 +- src/log-store/src/object_store_wal/batch.rs | 2 +- src/log-store/src/object_store_wal/store.rs | 203 ++- src/mito2/src/engine.rs | 2 + .../engine/object_store_wal_recovery_test.rs | 1150 +++++++++++++++++ src/mito2/src/error.rs | 16 +- src/mito2/src/flush.rs | 18 + src/mito2/src/wal.rs | 56 +- src/mito2/src/worker/handle_flush.rs | 3 + src/mito2/src/worker/handle_manifest.rs | 39 +- 10 files changed, 1465 insertions(+), 27 deletions(-) create mode 100644 src/mito2/src/engine/object_store_wal_recovery_test.rs diff --git a/src/log-store/src/object_store_wal.rs b/src/log-store/src/object_store_wal.rs index f36108f3bd5..2ec36b30463 100644 --- a/src/log-store/src/object_store_wal.rs +++ b/src/log-store/src/object_store_wal.rs @@ -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; diff --git a/src/log-store/src/object_store_wal/batch.rs b/src/log-store/src/object_store_wal/batch.rs index 078af921def..b1873c8a219 100644 --- a/src/log-store/src/object_store_wal/batch.rs +++ b/src/log-store/src/object_store_wal/batch.rs @@ -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 } diff --git a/src/log-store/src/object_store_wal/store.rs b/src/log-store/src/object_store_wal/store.rs index 2a0c75f1048..5a6454123dd 100644 --- a/src/log-store/src/object_store_wal/store.rs +++ b/src/log-store/src/object_store_wal/store.rs @@ -104,6 +104,12 @@ pub struct ObjectStoreLogStore { creates_held: watch::Sender, #[cfg(any(test, feature = "testing"))] creates_fail: Arc, + #[cfg(any(test, feature = "testing"))] + next_create_fails_after_write: Arc, + #[cfg(any(test, feature = "testing"))] + parked_creates: watch::Receiver, + #[cfg(any(test, feature = "testing"))] + durability_waits: watch::Receiver, } type ObsoleteEntryIds = Arc>>; @@ -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 { + 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>, }, + #[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, #[cfg(any(test, feature = "testing"))] creates_fail: Arc, + #[cfg(any(test, feature = "testing"))] + next_create_fails_after_write: Arc, + #[cfg(any(test, feature = "testing"))] + parked_creates: Arc>, + #[cfg(any(test, feature = "testing"))] + durability_waits: watch::Sender, } 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()); diff --git a/src/mito2/src/engine.rs b/src/mito2/src/engine.rs index abfb470c1be..47752adbe75 100644 --- a/src/mito2/src/engine.rs +++ b/src/mito2/src/engine.rs @@ -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; diff --git a/src/mito2/src/engine/object_store_wal_recovery_test.rs b/src/mito2/src/engine/object_store_wal_recovery_test.rs new file mode 100644 index 00000000000..6017f745d2e --- /dev/null +++ b/src/mito2/src/engine/object_store_wal_recovery_test.rs @@ -0,0 +1,1150 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Recovery tests for regions on a real [`ObjectStoreLogStore`] over an +//! in-memory object store. The store seals objects only through its testing +//! hooks, so every test controls which entries are durable. +//! +//! Every open of the store writes a start object above every present object, +//! so the first data object of a store opened on an empty prefix has +//! sequence 1. + +use std::collections::HashMap; +use std::sync::Arc; +use std::sync::atomic::Ordering; +use std::time::Duration; + +use api::v1::Rows; +use common_base::readable_size::ReadableSize; +use common_error::ext::{BoxedError, ErrorExt}; +use common_error::status_code::StatusCode; +use common_recordbatch::RecordBatches; +use common_wal::config::object_store::{AckMode, ObjectStoreWalConfig, STANDALONE_GENERATION}; +use common_wal::options::{ObjectStoreWalOptions, WAL_OPTIONS_KEY, WalOptions}; +use log_store::ObjectStoreLogStore; +use log_store::object_store_wal::entry_id; +use object_store::ObjectStore; +use object_store::services::Memory; +use rstest::rstest; +use store_api::logstore::LogStore; +use store_api::logstore::provider::Provider; +use store_api::mito_engine_options::SKIP_WAL_KEY; +use store_api::region_engine::{RegionEngine, RegionRole}; +use store_api::region_request::{ + PathType, RegionDropRequest, RegionFlushRequest, RegionOpenRequest, RegionPutRequest, + RegionRequest, RegionTruncateRequest, +}; +use store_api::storage::{RegionId, ScanRequest}; + +use crate::config::MitoConfig; +use crate::engine::MitoEngine; +use crate::region::MitoRegionRef; +use crate::test_util::{ + CreateRequestBuilder, TestEnv, build_rows, build_rows_for_key, flush_region, put_rows, + rows_schema, +}; + +/// The root the store derives its node prefix from. +const ROOT: &str = "cluster-a/wal"; +/// The node prefix of datanode 0 in the standalone generation under [`ROOT`]. +const PREFIX: &str = "cluster-a/wal/datanodes/0/epochs/0"; +/// The engine runs two workers and these regions map to different ones, so +/// their writes can be admitted into the same open batch. +const REGION_A: RegionId = RegionId::new(1, 1); +const REGION_B: RegionId = RegionId::new(1, 2); +const WAIT: Duration = Duration::from_secs(30); + +fn memory_store() -> ObjectStore { + ObjectStore::new(Memory::default()).unwrap() +} + +/// The configuration of a store with `ack_mode` that never seals a batch on +/// its own. +fn store_config(ack_mode: AckMode) -> ObjectStoreWalConfig { + ObjectStoreWalConfig { + storage_provider: String::new(), + prefix: ROOT.to_string(), + flush_interval: Duration::from_secs(3600), + max_batch_bytes: ReadableSize(u64::MAX), + ack_mode, + ..Default::default() + } +} + +/// Opens a store under [`PREFIX`] with `ack_mode` that never seals a batch +/// on its own. +async fn open_store(object_store: &ObjectStore, ack_mode: AckMode) -> Arc { + ObjectStoreLogStore::try_new( + object_store.clone(), + &store_config(ack_mode), + 0, + STANDALONE_GENERATION, + ) + .await + .unwrap() +} + +async fn new_engine(env: &mut TestEnv, store: Arc) -> MitoEngine { + let config = MitoConfig { + num_workers: 2, + ..Default::default() + }; + env.create_engine_with_log_store(config, store).await +} + +fn wal_options() -> HashMap { + let options = WalOptions::ObjectStore(ObjectStoreWalOptions::new(PREFIX.to_string())); + HashMap::from([( + WAL_OPTIONS_KEY.to_string(), + serde_json::to_string(&options).unwrap(), + )]) +} + +fn provider(region_id: RegionId) -> Provider { + Provider::object_store_provider(region_id, PREFIX.to_string()) +} + +/// Creates a region on the object store WAL with `extra_options` and returns +/// its table dir and row schema. +async fn create_region( + engine: &MitoEngine, + region_id: RegionId, + extra_options: &[(&str, &str)], +) -> (String, Vec) { + let mut builder = + CreateRequestBuilder::new().insert_option(WAL_OPTIONS_KEY, &wal_options()[WAL_OPTIONS_KEY]); + for (key, value) in extra_options { + builder = builder.insert_option(key, value); + } + let request = builder.build(); + let table_dir = request.table_dir.clone(); + let schema = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + (table_dir, schema) +} + +/// Opens a region whose options select the object store WAL under +/// [`PREFIX`] plus `extra_options`, and makes it writable. +async fn open_region( + engine: &MitoEngine, + region_id: RegionId, + table_dir: &str, + extra_options: &[(&str, &str)], +) -> std::result::Result<(), BoxedError> { + let mut options = wal_options(); + options.extend( + extra_options + .iter() + .map(|(key, value)| (key.to_string(), value.to_string())), + ); + engine + .handle_request( + region_id, + RegionRequest::Open(RegionOpenRequest { + engine: String::new(), + table_dir: table_dir.to_string(), + options, + skip_wal_replay: false, + path_type: PathType::Bare, + checkpoint: None, + requirements: Default::default(), + }), + ) + .await?; + engine.set_region_role(region_id, RegionRole::Leader) +} + +fn rows(schema: &[api::v1::ColumnSchema], start: usize, end: usize) -> Rows { + Rows { + schema: schema.to_vec(), + rows: build_rows(start, end), + } +} + +/// Returns every row of the region as a table, so two scans can be compared +/// row by row. +async fn scan_rows(engine: &MitoEngine, region_id: RegionId) -> String { + let stream = engine + .scan_to_stream(region_id, ScanRequest::default()) + .await + .unwrap(); + RecordBatches::try_collect(stream) + .await + .unwrap() + .pretty_print() + .unwrap() +} + +fn region(engine: &MitoEngine, region_id: RegionId) -> MitoRegionRef { + engine.get_region(region_id).unwrap() +} + +/// The entry ids a region tracks in memory and in its manifest, and the +/// rows its memtables hold. +#[derive(Debug, PartialEq, Eq)] +struct RegionRecoveryState { + flushed_entry_id: u64, + last_entry_id: u64, + topic_latest_entry_id: u64, + manifest_flushed_entry_id: u64, + memtable_rows: u64, +} + +async fn region_recovery_state(engine: &MitoEngine, region_id: RegionId) -> RegionRecoveryState { + let region = region(engine, region_id); + let current = region.version_control.current(); + let manifest = region.manifest_ctx.manifest().await; + RegionRecoveryState { + flushed_entry_id: current.version.flushed_entry_id, + last_entry_id: current.last_entry_id, + topic_latest_entry_id: region.topic_latest_entry_id.load(Ordering::Relaxed), + manifest_flushed_entry_id: manifest.flushed_entry_id, + memtable_rows: current.version.memtables.num_rows(), + } +} + +const EMPTY_RECOVERY_STATE: RegionRecoveryState = RegionRecoveryState { + flushed_entry_id: 0, + last_entry_id: 0, + topic_latest_entry_id: 0, + manifest_flushed_entry_id: 0, + memtable_rows: 0, +}; + +/// Writes rows through the engine and seals them into one WAL object. +/// +/// A put in the `durable` mode blocks until its entries are durable, so the +/// writes run in the background until the store has admitted all of them and +/// the batch is sealed by hand. +struct SealedWriter { + store: Arc, + admitted: usize, +} + +impl SealedWriter { + fn new(store: &Arc) -> Self { + Self { + store: store.clone(), + admitted: 0, + } + } + + async fn put_and_seal(&mut self, engine: &MitoEngine, writes: Vec<(RegionId, Rows)>) { + let handles = writes + .into_iter() + .map(|(region_id, rows)| { + let engine = engine.clone(); + tokio::spawn(async move { put_rows(&engine, region_id, rows).await }) + }) + .collect::>(); + self.admitted += handles.len(); + tokio::time::timeout(WAIT, self.store.wait_for_admitted_appends(self.admitted)) + .await + .expect("writes must reach the store") + .unwrap(); + self.store.seal_open_batch().await.unwrap(); + for handle in handles { + handle.await.unwrap(); + } + } +} + +fn latest(store: &ObjectStoreLogStore, region_id: RegionId) -> u64 { + store.latest_entry_id(&provider(region_id)).unwrap() +} + +/// Returns the sequences of the WAL objects under the prefix, in order. +async fn wal_object_seqs(object_store: &ObjectStore) -> Vec { + let mut seqs = object_store + .list(&format!("{PREFIX}/objects/")) + .await + .unwrap() + .into_iter() + .filter(|entry| entry.metadata().is_file()) + .map(|entry| { + let path = entry.path(); + path.rsplit('/') + .next() + .and_then(|name| name.strip_suffix(".wal")) + .and_then(|seq| seq.parse().ok()) + .unwrap_or_else(|| panic!("unexpected WAL object key {path}")) + }) + .collect::>(); + seqs.sort_unstable(); + seqs +} + +#[rstest] +#[case(AckMode::Durable)] +#[case(AckMode::Enqueued)] +#[tokio::test] +async fn test_reopen_after_partial_flush_replays_only_unflushed_regions(#[case] ack_mode: AckMode) { + let mut env = TestEnv::with_prefix("object-store-wal-partial-flush").await; + let object_store = memory_store(); + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir_a, schema_a) = create_region(&engine, REGION_A, &[]).await; + let (table_dir_b, schema_b) = create_region(&engine, REGION_B, &[]).await; + + // Objects 1 and 2 each hold one entry of both regions: the ids of a + // region are its position under the sequence of its object. + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal( + &engine, + vec![ + (REGION_A, rows(&schema_a, 0, 2)), + (REGION_B, rows(&schema_b, 0, 3)), + ], + ) + .await; + writer + .put_and_seal( + &engine, + vec![ + (REGION_B, rows(&schema_b, 3, 5)), + (REGION_A, rows(&schema_a, 2, 4)), + ], + ) + .await; + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + assert_eq!(entry_id(2, 1), latest(&store, REGION_B)); + + // Flushing empties the memtables, so the topic latest entry id follows + // the store. + flush_region(&engine, REGION_A, None).await; + let flushed_a = RegionRecoveryState { + flushed_entry_id: entry_id(2, 1), + last_entry_id: entry_id(2, 1), + topic_latest_entry_id: entry_id(2, 1), + manifest_flushed_entry_id: entry_id(2, 1), + memtable_rows: 0, + }; + assert_eq!(flushed_a, region_recovery_state(&engine, REGION_A).await); + let rows_a = scan_rows(&engine, REGION_A).await; + let rows_b = scan_rows(&engine, REGION_B).await; + + engine.stop().await.unwrap(); + drop(engine); + drop(writer); + drop(store); + + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir_a, &[]) + .await + .unwrap(); + open_region(&engine, REGION_B, &table_dir_b, &[]) + .await + .unwrap(); + + // Region A replays nothing: it flushed both entries, so the topic latest + // entry id comes from the store. Region B replays both of its entries. + assert_eq!(flushed_a, region_recovery_state(&engine, REGION_A).await); + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(2, 1), + memtable_rows: 5, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_B).await + ); + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + assert_eq!(entry_id(2, 1), latest(&store, REGION_B)); + assert_eq!(rows_a, scan_rows(&engine, REGION_A).await); + assert_eq!(rows_b, scan_rows(&engine, REGION_B).await); + assert_eq!(4, engine.get_region_statistic(REGION_A).unwrap().num_rows); + assert_eq!(5, engine.get_region_statistic(REGION_B).unwrap().num_rows); +} + +#[rstest] +#[case(AckMode::Durable)] +#[case(AckMode::Enqueued)] +#[tokio::test] +async fn test_reopen_after_crash_replays_durable_entries_once(#[case] ack_mode: AckMode) { + let mut env = TestEnv::with_prefix("object-store-wal-crash").await; + let object_store = memory_store(); + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + + // The entry is durable as an object but the manifest never learns of it. + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 0, 3))]) + .await; + assert_eq!(entry_id(1, 1), latest(&store, REGION_A)); + let rows_before = scan_rows(&engine, REGION_A).await; + + drop(writer); + crash(engine, store).await; + + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(1, 1), + memtable_rows: 3, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(entry_id(1, 1), latest(&store, REGION_A)); + assert_eq!(rows_before, scan_rows(&engine, REGION_A).await); + assert_eq!(3, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[tokio::test] +async fn test_open_region_rejects_mismatched_wal_prefix() { + let mut env = TestEnv::with_prefix("object-store-wal-prefix-mismatch").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Durable).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 0, 3))]) + .await; + + engine.stop().await.unwrap(); + drop(engine); + drop(writer); + drop(store); + + // The process now runs its store under the next generation while the + // region still persists the prefix it was created with. + let config = store_config(AckMode::Durable); + let other_prefix = config.node_prefix(0, STANDALONE_GENERATION + 1); + let store = + ObjectStoreLogStore::try_new(object_store.clone(), &config, 0, STANDALONE_GENERATION + 1) + .await + .unwrap(); + let engine = new_engine(&mut env, store).await; + let err = open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap_err(); + assert_eq!(StatusCode::InvalidArguments, err.status_code()); + let message = err.output_msg(); + assert!( + message.contains(PREFIX) && message.contains(&other_prefix), + "unexpected error: {message}" + ); + assert!(!engine.is_region_exists(REGION_A)); +} + +#[tokio::test] +async fn test_reopen_on_empty_prefix_without_durable_entries() { + let mut env = TestEnv::with_prefix("object-store-wal-empty-prefix").await; + let store = open_store(&memory_store(), AckMode::Durable).await; + let engine = new_engine(&mut env, store.clone()).await; + // Region A never writes; region B skips the WAL and flushes its rows. + let (table_dir_a, schema_a) = create_region(&engine, REGION_A, &[]).await; + let (table_dir_b, schema_b) = create_region(&engine, REGION_B, &[(SKIP_WAL_KEY, "true")]).await; + put_rows(&engine, REGION_B, rows(&schema_b, 0, 3)).await; + flush_region(&engine, REGION_B, None).await; + let rows_b = scan_rows(&engine, REGION_B).await; + + engine.stop().await.unwrap(); + drop(engine); + drop(store); + + // A fresh object store holds no object under the prefix. + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Durable).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir_a, &[]) + .await + .unwrap(); + open_region(&engine, REGION_B, &table_dir_b, &[(SKIP_WAL_KEY, "true")]) + .await + .unwrap(); + + assert_eq!(0, latest(&store, REGION_A)); + assert_eq!(0, latest(&store, REGION_B)); + assert_eq!( + EMPTY_RECOVERY_STATE, + region_recovery_state(&engine, REGION_A).await + ); + // Region B keeps the object store provider but wrote nothing to it; its + // rows come back from the SST alone. + assert_eq!(provider(REGION_B), region(&engine, REGION_B).provider); + assert_eq!( + 0, + region_recovery_state(&engine, REGION_B).await.memtable_rows + ); + assert_eq!(rows_b, scan_rows(&engine, REGION_B).await); + + // The region is usable and its first entry lands in the first object + // after the start object. + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema_a, 0, 2))]) + .await; + assert_eq!(vec![0, 1], wal_object_seqs(&object_store).await); + assert_eq!(entry_id(1, 1), latest(&store, REGION_A)); + assert_eq!( + entry_id(1, 1), + region_recovery_state(&engine, REGION_A).await.last_entry_id + ); + assert_eq!(2, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[rstest] +#[case(AckMode::Durable)] +#[case(AckMode::Enqueued)] +#[tokio::test] +async fn test_entry_ids_continue_across_two_restarts(#[case] ack_mode: AckMode) { + let mut env = TestEnv::with_prefix("object-store-wal-two-restarts").await; + let object_store = memory_store(); + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + + // The entry of object 1 is flushed, the entry of object 2 is not. + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 0, 2))]) + .await; + flush_region(&engine, REGION_A, None).await; + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 2, 4))]) + .await; + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + engine.stop().await.unwrap(); + drop(engine); + drop(writer); + drop(store); + + // The first restart writes start object 3 and replays the entry of + // object 2 only. + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + flushed_entry_id: entry_id(1, 1), + last_entry_id: entry_id(2, 1), + topic_latest_entry_id: entry_id(1, 1), + manifest_flushed_entry_id: entry_id(1, 1), + memtable_rows: 2, + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(4, engine.get_region_statistic(REGION_A).unwrap().num_rows); + + // The sequence continues at object 4, whose entry is flushed; the entry + // of object 5 is not. + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 4, 6))]) + .await; + assert_eq!(entry_id(4, 1), latest(&store, REGION_A)); + assert_eq!( + entry_id(4, 1), + region_recovery_state(&engine, REGION_A).await.last_entry_id + ); + flush_region(&engine, REGION_A, None).await; + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 6, 8))]) + .await; + assert_eq!(entry_id(5, 1), latest(&store, REGION_A)); + let rows_before = scan_rows(&engine, REGION_A).await; + engine.stop().await.unwrap(); + drop(engine); + drop(writer); + drop(store); + + // The second restart replays the entry of object 5 only. + let store = open_store(&object_store, ack_mode).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + flushed_entry_id: entry_id(4, 1), + last_entry_id: entry_id(5, 1), + topic_latest_entry_id: entry_id(4, 1), + manifest_flushed_entry_id: entry_id(4, 1), + memtable_rows: 2, + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(entry_id(5, 1), latest(&store, REGION_A)); + assert_eq!(rows_before, scan_rows(&engine, REGION_A).await); + assert_eq!(8, engine.get_region_statistic(REGION_A).unwrap().num_rows); + assert_eq!( + 8, + region(&engine, REGION_A) + .version_control + .committed_sequence() + ); +} + +/// Sends a put without asserting its result. +async fn try_put( + engine: &MitoEngine, + region_id: RegionId, + rows: Rows, +) -> std::result::Result<(), BoxedError> { + engine + .handle_request( + region_id, + RegionRequest::Put(RegionPutRequest { + skip_wal: false, + rows, + hint: None, + partition_expr_version: None, + }), + ) + .await + .map(|_| ()) +} + +#[tokio::test] +async fn test_reopen_after_a_failed_create_shows_only_the_acknowledged_write() { + let mut env = TestEnv::with_prefix("object-store-wal-failed-create").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Durable).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + let key_rows = |value| Rows { + schema: schema.clone(), + rows: build_rows_for_key("k", 0, 1, value), + }; + + // Object 1 is stored, but its create is reported as failed, so the write + // fails. + let failed = { + let engine = engine.clone(); + let rows = key_rows(1); + tokio::spawn(async move { try_put(&engine, REGION_A, rows).await }) + }; + tokio::time::timeout(WAIT, store.wait_for_admitted_appends(1)) + .await + .unwrap() + .unwrap(); + store.fail_next_create_after_write(); + store.seal_open_batch().await.unwrap_err(); + assert!(failed.await.unwrap().is_err()); + assert_eq!(vec![0, 1], wal_object_seqs(&object_store).await); + assert_eq!(0, latest(&store, REGION_A)); + + // The same key and timestamp is written again under the same row + // sequence, and acknowledged from object 2. + let mut writer = SealedWriter { + store: store.clone(), + admitted: 1, + }; + writer + .put_and_seal(&engine, vec![(REGION_A, key_rows(2))]) + .await; + assert_eq!(vec![0, 1, 2], wal_object_seqs(&object_store).await); + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + let region_a = region(&engine, REGION_A); + assert_eq!(1, region_a.version_control.committed_sequence()); + let expected = "\ ++-------+---------+---------------------+ +| tag_0 | field_0 | ts | ++-------+---------+---------------------+ +| k | 2.0 | 1970-01-01T00:00:00 | ++-------+---------+---------------------+"; + assert_eq!(expected, scan_rows(&engine, REGION_A).await); + + engine.stop().await.unwrap(); + drop(engine); + drop(writer); + drop(store); + + // Object 1 is off the chain, so only the acknowledged write replays. + let store = open_store(&object_store, AckMode::Durable).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(2, 1), + memtable_rows: 1, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(expected, scan_rows(&engine, REGION_A).await); +} + +/// Requests a flush of the region in the background and returns its handle; +/// the result is up to the caller. +fn spawn_flush( + engine: &MitoEngine, + region_id: RegionId, +) -> tokio::task::JoinHandle> { + let engine = engine.clone(); + tokio::spawn(async move { + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest { + row_group_size: None, + reason: None, + }), + ) + .await + .map(|_| ()) + }) +} + +/// Seals the open batch in the background; with creates held, the seal +/// returns only once they are released. +fn spawn_seal( + store: &Arc, +) -> tokio::task::JoinHandle> { + let store = store.clone(); + tokio::spawn(async move { store.seal_open_batch().await }) +} + +/// Waits until a testing hook of the store observes the state it waits for. +async fn wait_for(hook: impl Future>) { + tokio::time::timeout(WAIT, hook) + .await + .expect("the store must reach the state") + .unwrap(); +} + +/// Tears the engine and its store down as a crash of the process would: the +/// store writes nothing more and fails every caller still waiting for it, +/// then the engine stops. Returns once nothing holds the store any more, so a +/// store reopened on the prefix never runs beside the old one. +async fn crash(engine: MitoEngine, store: Arc) { + store.crash().await; + engine.stop().await.unwrap(); + drop(engine); + let released = Arc::downgrade(&store); + drop(store); + tokio::time::timeout(WAIT, async { + while released.strong_count() > 0 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("nothing may hold the crashed store"); +} + +#[tokio::test] +async fn test_flush_waits_until_the_wal_is_durable() { + let mut env = TestEnv::with_prefix("object-store-wal-flush-barrier").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (_, schema) = create_region(&engine, REGION_A, &[]).await; + + // The write is acknowledged on admission; its object is not created + // while creates are held. + store.hold_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(1, 1), + memtable_rows: 3, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(0, latest(&store, REGION_A)); + + // The flush writes its SST, then waits for the entry to be durable + // before it records the entry as flushed in the manifest. + let flush = spawn_flush(&engine, REGION_A); + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + wait_for(store.wait_for_durability_waits(1)).await; + assert!(!flush.is_finished()); + assert!(!seal.is_finished()); + assert_eq!( + 0, + region_recovery_state(&engine, REGION_A) + .await + .manifest_flushed_entry_id + ); + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + + store.release_creates(); + tokio::time::timeout(WAIT, seal) + .await + .unwrap() + .unwrap() + .unwrap(); + tokio::time::timeout(WAIT, flush) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!(entry_id(1, 1), latest(&store, REGION_A)); + assert_eq!( + RegionRecoveryState { + flushed_entry_id: entry_id(1, 1), + last_entry_id: entry_id(1, 1), + topic_latest_entry_id: entry_id(1, 1), + manifest_flushed_entry_id: entry_id(1, 1), + memtable_rows: 0, + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(3, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[tokio::test] +async fn test_enqueued_crash_before_the_object_exists_replays_durable_entries() { + let mut env = TestEnv::with_prefix("object-store-wal-enqueued-crash").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + + // The entry is acknowledged and a flush is attempted, but its object is + // never created: the flush waits and the manifest keeps watermark 0. + store.hold_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + assert_eq!( + entry_id(1, 1), + region_recovery_state(&engine, REGION_A).await.last_entry_id + ); + let flush = spawn_flush(&engine, REGION_A); + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + wait_for(store.wait_for_durability_waits(1)).await; + assert!(!flush.is_finished()); + assert_eq!( + 0, + region_recovery_state(&engine, REGION_A) + .await + .manifest_flushed_entry_id + ); + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + + // The process dies with the create still held: the waiting flush and + // the seal fail, and the object is never written. + store.crash().await; + assert!(flush.await.unwrap().is_err()); + assert!(seal.await.unwrap().is_err()); + crash(engine, store).await; + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + + // Nothing is durable, so nothing replays. The restart writes start + // object 1 and the region's next entry lands in object 2; the lost entry + // was inside the unpersisted backlog. + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + EMPTY_RECOVERY_STATE, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(0, latest(&store, REGION_A)); + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 3, 5))]) + .await; + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + assert_eq!( + entry_id(2, 1), + region_recovery_state(&engine, REGION_A).await.last_entry_id + ); + assert_eq!(vec![0, 1, 2], wal_object_seqs(&object_store).await); + let rows_before = scan_rows(&engine, REGION_A).await; + assert_eq!(2, engine.get_region_statistic(REGION_A).unwrap().num_rows); + drop(writer); + crash(engine, store).await; + + // The second restart replays the entry that became durable and skips + // nothing: the manifest never named an entry that was not durable. + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(2, 1), + memtable_rows: 2, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + assert_eq!(rows_before, scan_rows(&engine, REGION_A).await); + assert_eq!(2, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[tokio::test] +async fn test_drop_cancels_a_flush_waiting_for_wal_durability() { + let mut env = TestEnv::with_prefix("object-store-wal-flush-barrier-drop").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (_, schema) = create_region(&engine, REGION_A, &[]).await; + + // The flush reaches the durability barrier and waits for a create that + // is held. + store.hold_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + let flush = spawn_flush(&engine, REGION_A); + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + wait_for(store.wait_for_durability_waits(1)).await; + assert!(!flush.is_finished()); + + // The drop cancels the flush instead of waiting behind the upload. + tokio::time::timeout( + WAIT, + engine.handle_request( + REGION_A, + RegionRequest::Drop(RegionDropRequest { + fast_path: false, + force: false, + partial_drop: false, + }), + ), + ) + .await + .unwrap() + .unwrap(); + assert!(!engine.is_region_exists(REGION_A)); + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + assert!( + tokio::time::timeout(WAIT, flush) + .await + .unwrap() + .unwrap() + .is_err() + ); + + store.release_creates(); + tokio::time::timeout(WAIT, seal) + .await + .unwrap() + .unwrap() + .unwrap(); +} + +#[rstest] +#[case(RegionTruncateRequest::All)] +#[case(RegionTruncateRequest::Unflushed)] +#[tokio::test] +async fn test_enqueued_truncate_waits_until_the_wal_is_durable( + #[case] request: RegionTruncateRequest, +) { + let mut env = TestEnv::with_prefix("object-store-wal-truncate-barrier").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + + // The entry is acknowledged and a truncate is attempted, but its object + // is never created: the truncate waits and the manifest keeps frontier 0. + store.hold_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + assert_eq!( + entry_id(1, 1), + region_recovery_state(&engine, REGION_A).await.last_entry_id + ); + let truncate = { + let engine = engine.clone(); + tokio::spawn(async move { + engine + .handle_request(REGION_A, RegionRequest::Truncate(request)) + .await + .map(|_| ()) + }) + }; + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + wait_for(store.wait_for_durability_waits(1)).await; + assert!(!truncate.is_finished()); + let manifest = region(&engine, REGION_A).manifest_ctx.manifest().await; + assert_eq!(0, manifest.flushed_entry_id); + assert_eq!(None, manifest.truncated_entry_id); + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + + // The process dies with the create still held: the waiting truncate and + // the seal fail. + store.crash().await; + assert!(truncate.await.unwrap().is_err()); + assert!(seal.await.unwrap().is_err()); + crash(engine, store).await; + + // Nothing is durable and the manifest names no entry, so the region + // starts over after start object 1, and its next entry becomes durable + // in object 2. + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + EMPTY_RECOVERY_STATE, + region_recovery_state(&engine, REGION_A).await + ); + let mut writer = SealedWriter::new(&store); + writer + .put_and_seal(&engine, vec![(REGION_A, rows(&schema, 3, 5))]) + .await; + assert_eq!(entry_id(2, 1), latest(&store, REGION_A)); + let rows_before = scan_rows(&engine, REGION_A).await; + drop(writer); + crash(engine, store).await; + + // The second restart replays the durable entry instead of skipping it. + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + last_entry_id: entry_id(2, 1), + memtable_rows: 2, + ..EMPTY_RECOVERY_STATE + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(rows_before, scan_rows(&engine, REGION_A).await); + assert_eq!(2, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[rstest] +#[case(RegionTruncateRequest::All)] +#[case(RegionTruncateRequest::Unflushed)] +#[tokio::test] +async fn test_enqueued_truncate_completes_once_the_wal_is_durable( + #[case] request: RegionTruncateRequest, +) { + let mut env = TestEnv::with_prefix("object-store-wal-truncate-durable").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (table_dir, schema) = create_region(&engine, REGION_A, &[]).await; + let truncated_entry_id = + matches!(request, RegionTruncateRequest::All).then_some(entry_id(1, 1)); + + // The truncate waits for the acknowledged entry, then records it once + // its object is created. + store.hold_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + let truncate = { + let engine = engine.clone(); + tokio::spawn(async move { + engine + .handle_request(REGION_A, RegionRequest::Truncate(request)) + .await + .map(|_| ()) + }) + }; + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + wait_for(store.wait_for_durability_waits(1)).await; + assert!(!truncate.is_finished()); + store.release_creates(); + tokio::time::timeout(WAIT, seal) + .await + .unwrap() + .unwrap() + .unwrap(); + tokio::time::timeout(WAIT, truncate) + .await + .unwrap() + .unwrap() + .unwrap(); + let manifest = region(&engine, REGION_A).manifest_ctx.manifest().await; + assert_eq!(entry_id(1, 1), manifest.flushed_entry_id); + assert_eq!(truncated_entry_id, manifest.truncated_entry_id); + assert_eq!(0, engine.get_region_statistic(REGION_A).unwrap().num_rows); + engine.stop().await.unwrap(); + drop(engine); + drop(store); + + // The restart replays nothing: the discarded rows stay discarded. + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + open_region(&engine, REGION_A, &table_dir, &[]) + .await + .unwrap(); + assert_eq!( + RegionRecoveryState { + flushed_entry_id: entry_id(1, 1), + last_entry_id: entry_id(1, 1), + topic_latest_entry_id: entry_id(1, 1), + manifest_flushed_entry_id: entry_id(1, 1), + memtable_rows: 0, + }, + region_recovery_state(&engine, REGION_A).await + ); + assert_eq!(0, engine.get_region_statistic(REGION_A).unwrap().num_rows); +} + +#[tokio::test] +async fn test_flush_does_not_publish_a_frontier_for_a_lost_enqueued_backlog() { + let mut env = TestEnv::with_prefix("object-store-wal-lost-backlog").await; + let object_store = memory_store(); + let store = open_store(&object_store, AckMode::Enqueued).await; + let engine = new_engine(&mut env, store.clone()).await; + let (_, schema) = create_region(&engine, REGION_A, &[]).await; + + // The entry is acknowledged, then its create fails while stop has begun + // but the stop command is not handled yet: the backlog is lost. + store.hold_creates(); + store.fail_creates(); + put_rows(&engine, REGION_A, rows(&schema, 0, 3)).await; + let seal = spawn_seal(&store); + wait_for(store.wait_for_parked_creates(1)).await; + store.begin_stop(); + store.release_creates(); + assert!( + tokio::time::timeout(WAIT, seal) + .await + .unwrap() + .unwrap() + .is_err() + ); + assert_eq!(vec![0], wal_object_seqs(&object_store).await); + + // A flush in that window fails at the barrier instead of recording the + // lost entry as flushed. + let flush = spawn_flush(&engine, REGION_A); + assert!( + tokio::time::timeout(WAIT, flush) + .await + .unwrap() + .unwrap() + .is_err() + ); + assert_eq!( + 0, + region_recovery_state(&engine, REGION_A) + .await + .manifest_flushed_entry_id + ); + assert!(store.stop().await.is_err()); +} diff --git a/src/mito2/src/error.rs b/src/mito2/src/error.rs index 7d43548dd2b..f738044c9ef 100644 --- a/src/mito2/src/error.rs +++ b/src/mito2/src/error.rs @@ -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 }, @@ -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(), diff --git a/src/mito2/src/flush.rs b/src/mito2/src/flush.rs index a5d9a7d0a73..b5d8c4e56ea 100644 --- a/src/mito2/src/flush.rs +++ b/src/mito2/src/flush.rs @@ -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, + /// 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. diff --git a/src/mito2/src/wal.rs b/src/mito2/src/wal.rs index fa549622898..772aac336cf 100644 --- a/src/mito2/src/wal.rs +++ b/src/mito2/src/wal.rs @@ -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 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 Wal { } } + /// 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, diff --git a/src/mito2/src/worker/handle_flush.rs b/src/mito2/src/worker/handle_flush.rs index bfdfea5d1d9..1d4442d8d03 100644 --- a/src/mito2/src/worker/handle_flush.rs +++ b/src/mito2/src/worker/handle_flush.rs @@ -243,6 +243,9 @@ impl RegionWorkerLoop { 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, ®ion.provider), } } } diff --git a/src/mito2/src/worker/handle_manifest.rs b/src/mito2/src/worker/handle_manifest.rs index aa8475a49f1..b574e490ac7 100644 --- a/src/mito2/src/worker/handle_manifest.rs +++ b/src/mito2/src/worker/handle_manifest.rs @@ -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 RegionWorkerLoop { 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, ®ion.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 RegionWorkerLoop { 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, ®ion.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 RegionWorkerLoop { 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,