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,