diff --git a/src/log-store/Cargo.toml b/src/log-store/Cargo.toml index 5d068ca344d..7e4da1557e1 100644 --- a/src/log-store/Cargo.toml +++ b/src/log-store/Cargo.toml @@ -11,6 +11,9 @@ protoc-rust.workspace = true [lints] workspace = true +[features] +testing = [] + [dependencies] async-stream.workspace = true async-trait.workspace = true diff --git a/src/log-store/src/error.rs b/src/log-store/src/error.rs index 8be0b8bd3da..85000dbea3c 100644 --- a/src/log-store/src/error.rs +++ b/src/log-store/src/error.rs @@ -350,6 +350,13 @@ pub enum Error { location: Location, }, + #[snafu(display("Incomplete multipart WAL entry of region {}", region_id))] + IncompleteWalEntry { + region_id: RegionId, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("WAL object sequence is exhausted, last sequence: {}", last_object_seq))] WalObjectSequenceExhausted { last_object_seq: u64, @@ -357,6 +364,16 @@ pub enum Error { location: Location, }, + #[snafu(display( + "WAL object sequence {} is not settled: the open batch has assigned entry ids under it", + object_seq + ))] + WalObjectSequenceUnsettled { + object_seq: u64, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display( "WAL entry positions of region {} in one object are exhausted", region_id @@ -374,6 +391,20 @@ pub enum Error { location: Location, }, + #[snafu(display( + "WAL object {} was written by the earlier epoch {}, this store writes epoch {}", + path, + existing_epoch, + epoch + ))] + StaleWalObject { + path: String, + existing_epoch: u64, + epoch: u64, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display( "WAL object {} that starts epoch {} was already present, so the open cannot tell whether it wrote it", path, @@ -417,16 +448,15 @@ pub enum Error { location: Location, }, - #[snafu(display("Object store WAL operation failed"))] - ObjectStoreWal { - source: Arc, + #[snafu(display("Object store WAL log store is stopped"))] + ObjectStoreWalStopped { #[snafu(implicit)] location: Location, }, - /// Appending to the object store WAL is unsupported. - #[snafu(display("Object store WAL operation is not supported"))] - UnsupportedObjectStoreWalOperation { + #[snafu(display("Object store WAL operation failed"))] + ObjectStoreWal { + source: Arc, #[snafu(implicit)] location: Location, }, @@ -464,6 +494,7 @@ impl ErrorExt for Error { | InvalidWalObjectStore { .. } | MismatchedWalPrefix { .. } | MismatchedWalRegion { .. } + | IncompleteWalEntry { .. } | InvalidWalEntryRange { .. } => StatusCode::InvalidArguments, StartWalTask { .. } | StopWalTask { .. } @@ -477,14 +508,15 @@ impl ErrorExt for Error { | OrderedBatchProducerStopped { .. } | WaitProduceResultReceiver { .. } | WaitDumpIndex { .. } - | MetaLengthExceededLimit { .. } => StatusCode::Internal, + | MetaLengthExceededLimit { .. } + | ObjectStoreWalStopped { .. } => StatusCode::Internal, CorruptedWalObject { .. } | WalObjectConflict { .. } | WalObjectSequenceExhausted { .. } | WalEntryPositionExhausted { .. } => StatusCode::Unexpected, + WalObjectSequenceUnsettled { .. } => StatusCode::IllegalState, - UnsupportedObjectStoreWalOperation { .. } => StatusCode::Unsupported, InvalidWalObject { source, .. } => source.status_code(), ObjectStoreWal { source, .. } => source.status_code(), @@ -493,6 +525,7 @@ impl ErrorExt for Error { | WriteIndex { .. } | ReadIndex { .. } | WalObjectStore { .. } + | StaleWalObject { .. } | UnconfirmedWalEpochStart { .. } | Io { .. } => StatusCode::StorageUnavailable, // Raft engine @@ -529,6 +562,8 @@ impl ErrorExt for Error { FetchEntry { .. } | RaftEngine { .. } | AddEntryLogBatch { .. } + | WalObjectSequenceUnsettled { .. } + | StaleWalObject { .. } | UnconfirmedWalEpochStart { .. } => RetryHint::Retryable, ProduceRecord { error, .. } => match error { rskafka::client::producer::Error::Client(error) => { diff --git a/src/log-store/src/object_store_wal/store.rs b/src/log-store/src/object_store_wal/store.rs index 7b024bb1a80..ee925ae0a38 100644 --- a/src/log-store/src/object_store_wal/store.rs +++ b/src/log-store/src/object_store_wal/store.rs @@ -12,9 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. -//! Object store WAL construction, recovery and region reads. +//! Object store WAL construction, recovery, writes and region reads. -use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; +use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque}; use std::fmt; use std::ops::Range; use std::sync::atomic::{AtomicBool, Ordering}; @@ -25,6 +25,8 @@ use async_stream::try_stream; use bytes::Bytes; use common_telemetry::info; use common_wal::config::object_store::ObjectStoreWalConfig; +use futures::future::BoxFuture; +use futures::stream::FuturesUnordered; use futures::{StreamExt, TryStreamExt}; use object_store::ObjectStore; use snafu::{IntoError, OptionExt, ResultExt, ensure}; @@ -32,25 +34,38 @@ use store_api::logstore::entry::{Entry, NaiveEntry}; use store_api::logstore::provider::{ObjectStoreProvider, Provider}; use store_api::logstore::{AppendBatchResponse, EntryId, LogStore, SendableEntryStream, WalIndex}; use store_api::storage::RegionId; +#[cfg(any(test, feature = "testing"))] +use tokio::sync::watch; use tokio::sync::{mpsc, oneshot}; +use tokio::time::MissedTickBehavior; use crate::error::{ - CorruptedWalObjectSnafu, Error, InvalidProviderSnafu, InvalidWalObjectSnafu, - InvalidWalObjectStoreSnafu, MismatchedWalPrefixSnafu, MismatchedWalRegionSnafu, - ObjectStoreWalSnafu, Result, UnconfirmedWalEpochStartSnafu, - UnsupportedObjectStoreWalOperationSnafu, WalObjectSequenceExhaustedSnafu, + CorruptedWalObjectSnafu, Error, IncompleteWalEntrySnafu, InvalidProviderSnafu, + InvalidWalObjectSnafu, InvalidWalObjectStoreSnafu, MismatchedWalPrefixSnafu, + MismatchedWalRegionSnafu, ObjectStoreWalSnafu, ObjectStoreWalStoppedSnafu, Result, + StaleWalObjectSnafu, UnconfirmedWalEpochStartSnafu, WalObjectSequenceExhaustedSnafu, + WalObjectSequenceUnsettledSnafu, }; -use crate::object_store_wal::batch::OBJECT_SEQ_LIMIT; +use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, OpenBatch, sequence_floor}; use crate::object_store_wal::catalog::ObjectCatalog; use crate::object_store_wal::format::{ - ChainLink, FixedTrailer, FooterEntry, HEADER_LEN, Header, MIN_OBJECT_LEN, TRAILER_LEN, - decode_footer, decode_header, decode_segment, decode_trailer, encode_object, footer_range, - verify_segment_ranges, + ChainLink, EncodedObject, FixedTrailer, FooterEntry, HEADER_LEN, Header, MIN_OBJECT_LEN, + Record, TRAILER_LEN, decode_footer, decode_header, decode_segment, decode_trailer, + encode_object, footer_range, verify_segment_ranges, }; use crate::object_store_wal::io::{ListedObject, ObjectStoreIo, PutResult}; const COMMAND_BUFFER: usize = 1024; +/// Appends queued for the actor; a caller beyond them waits to send. +const APPEND_BUFFER: usize = 16; const MIN_FLUSH_INTERVAL: Duration = Duration::from_millis(10); +const MAX_IN_FLIGHT_CREATES: usize = 4; +/// Number of sealed batches that may wait to be created or indexed. One +/// admission can seal two batches, the open batch before an append that would +/// exhaust its positions and the append's own batch, so the actor takes an +/// append from its channel only while two more fit: a stalled object store +/// holds callers back instead of growing the backlog. +const MAX_SEALED_BATCHES: usize = 2 * MAX_IN_FLIGHT_CREATES; /// Number of objects whose footers recovery fetches at a time. const RECOVERY_CONCURRENCY: usize = 8; /// Bytes recovery reads from the end of an object in one request. The window @@ -59,6 +74,16 @@ const RECOVERY_CONCURRENCY: usize = 8; const RECOVERY_TAIL_WINDOW: usize = 64 * 1024; /// A log store over the immutable WAL objects under one prefix. +/// +/// Appends are admitted into an open batch. A background actor seals the batch +/// when it reaches the size limit or the flush interval elapses and creates +/// the object under the next sequence while it keeps admitting entries into +/// the next batch; up to [`MAX_IN_FLIGHT_CREATES`] creates run at a time. +/// Objects are indexed in the catalog and their batches acknowledged in +/// sequence order, so an acknowledged entry never has a missing predecessor. +/// An append returns once its object is durable and indexed. While +/// [`MAX_SEALED_BATCHES`] batches wait to become durable no append is +/// admitted. pub(crate) struct ObjectStoreLogStore { prefix: String, io: Arc, @@ -67,6 +92,13 @@ pub(crate) struct ObjectStoreLogStore { terminal_error: TerminalError, stopped: Arc, command_tx: mpsc::Sender, + append_tx: mpsc::Sender, + #[cfg(any(test, feature = "testing"))] + admitted_appends: watch::Receiver, + #[cfg(any(test, feature = "testing"))] + creates_held: watch::Sender, + #[cfg(any(test, feature = "testing"))] + creates_fail: Arc, } type ObsoleteEntryIds = Arc>>; @@ -112,7 +144,7 @@ impl ObjectStoreLogStore { } ); - positive_bytes(config.max_batch_bytes.as_bytes(), "max batch bytes")?; + let max_batch_bytes = positive_bytes(config.max_batch_bytes.as_bytes(), "max batch bytes")?; let Recovered { mut catalog, next_object_seq, @@ -144,31 +176,62 @@ impl ObjectStoreLogStore { path: io.object_path(start.object_seq), })?; let catalog = Arc::new(RwLock::new(catalog)); + let obsolete_entry_ids = ObsoleteEntryIds::default(); let terminal_error = TerminalError::default(); let stopped = Arc::new(AtomicBool::new(false)); let (command_tx, command_rx) = mpsc::channel(COMMAND_BUFFER); + let (append_tx, append_rx) = mpsc::channel(APPEND_BUFFER); + #[cfg(any(test, feature = "testing"))] + let (admitted_appends_tx, admitted_appends_rx) = watch::channel(0); + #[cfg(any(test, feature = "testing"))] + let (creates_held_tx, creates_held_rx) = watch::channel(false); + #[cfg(any(test, feature = "testing"))] + let creates_fail = Arc::new(AtomicBool::new(false)); + let actor = Actor { + io: io.clone(), + catalog: catalog.clone(), + obsolete_entry_ids: obsolete_entry_ids.clone(), terminal_error: terminal_error.clone(), stopped: stopped.clone(), command_rx, + append_rx, + open_batch: OpenBatch::new(max_batch_bytes), issued_entry_ids: durable_entry_ids, + pending: Vec::new(), + sealed: VecDeque::new(), + creates: FuturesUnordered::new(), + stop: Vec::new(), next_object_seq: start .object_seq .checked_add(1) .filter(|next_object_seq| *next_object_seq < OBJECT_SEQ_LIMIT), epoch, - chain_tip: start, - stop: Vec::new(), + last_indexed: start, + flush_interval: config.flush_interval, + #[cfg(any(test, feature = "testing"))] + admitted_appends: admitted_appends_tx, + #[cfg(any(test, feature = "testing"))] + creates_held: creates_held_rx, + #[cfg(any(test, feature = "testing"))] + creates_fail: creates_fail.clone(), }; common_runtime::spawn_global(actor.run()); Ok(Arc::new(Self { prefix, io, catalog, - obsolete_entry_ids: ObsoleteEntryIds::default(), + obsolete_entry_ids, terminal_error, stopped, command_tx, + append_tx, + #[cfg(any(test, feature = "testing"))] + admitted_appends: admitted_appends_rx, + #[cfg(any(test, feature = "testing"))] + creates_held: creates_held_tx, + #[cfg(any(test, feature = "testing"))] + creates_fail, })) } @@ -236,10 +299,70 @@ fn positive_bytes(bytes: u64, name: &str) -> Result { }) } +#[cfg(any(test, feature = "testing"))] +impl ObjectStoreLogStore { + /// Waits until the actor has admitted at least `expected` append calls + /// since the store was built. + pub(crate) async fn wait_for_admitted_appends(&self, expected: usize) -> Result<()> { + self.admitted_appends + .clone() + .wait_for(|count| *count >= expected) + .await + .ok() + .map(|_| ()) + .context(ObjectStoreWalStoppedSnafu) + } + + /// Seals the open batch regardless of its size and age and returns once + /// its object is durable and indexed, or with the error that failed it. + pub(crate) async fn seal_open_batch(&self) -> Result<()> { + ensure!( + !self.stopped.load(Ordering::Acquire), + ObjectStoreWalStoppedSnafu + ); + let (response_tx, response_rx) = oneshot::channel(); + self.command_tx + .send(Command::Seal { + response: response_tx, + }) + .await + .ok() + .context(ObjectStoreWalStoppedSnafu)?; + response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)? + } + + /// Parks every conditional create that starts from now on until + /// [`release_creates`](Self::release_creates), so a test can observe + /// entries that are admitted but not durable. A create that is parked when + /// the store is dropped never runs. + pub(crate) fn hold_creates(&self) { + self.creates_held.send_replace(true); + } + + /// Lets the creates parked by [`hold_creates`](Self::hold_creates) run. + pub(crate) fn release_creates(&self) { + self.creates_held.send_replace(false); + } + + /// Makes every create that runs from now on fail with a transient object + /// store error instead of writing. + pub(crate) fn fail_creates(&self) { + self.creates_fail.store(true, Ordering::Release); + } + + /// 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(crate) fn begin_stop(&self) { + self.stopped.store(true, Ordering::Release); + } +} + #[async_trait::async_trait] impl LogStore for ObjectStoreLogStore { type Error = Error; + /// Stops the store. Creates in flight run to completion and acknowledge + /// their entries if they succeed. async fn stop(&self) -> Result<()> { self.stopped.store(true, Ordering::Release); let (response_tx, response_rx) = oneshot::channel(); @@ -256,8 +379,34 @@ impl LogStore for ObjectStoreLogStore { Ok(()) } - async fn append_batch(&self, _entries: Vec) -> Result { - UnsupportedObjectStoreWalOperationSnafu.fail() + async fn append_batch(&self, entries: Vec) -> Result { + ensure!( + !self.stopped.load(Ordering::Acquire), + ObjectStoreWalStoppedSnafu + ); + self.check_terminal()?; + if entries.is_empty() { + return Ok(AppendBatchResponse::default()); + } + for entry in &entries { + let region_id = self.region_of(entry.provider())?; + ensure!( + region_id == entry.region_id(), + MismatchedWalRegionSnafu { + region_id: entry.region_id(), + reason: format!("provider belongs to region {region_id}"), + } + ); + ensure!(entry.is_complete(), IncompleteWalEntrySnafu { region_id }); + } + + let (response_tx, response_rx) = oneshot::channel(); + self.append_tx + .send((entries, response_tx)) + .await + .ok() + .context(ObjectStoreWalStoppedSnafu)?; + response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)? } /// Reads the provider's region from `entry_id`, hiding obsolete entries. @@ -347,7 +496,20 @@ impl LogStore for ObjectStoreLogStore { .collect()) } - /// Records a memory-only watermark. The caller re-establishes it after a restart. + /// Moves the obsolete watermark of the region up to `entry_id` and makes + /// every id the region is assigned from now on greater than `entry_id`: + /// the sequence of the next object is raised above the object `entry_id` + /// names unless it is there already, so a watermark assigned under + /// another prefix, or one whose object this prefix no longer holds, is + /// never passed by a new id. Zero names no object and needs no floor. + /// + /// The watermark and the floor are applied together by the actor: a call + /// cancelled before its command is queued publishes neither, while a + /// command already queued still applies both. While the open batch holds + /// ids under the next sequence the floor cannot move, and the call fails + /// with [`Error::WalObjectSequenceUnsettled`] and records neither. The + /// watermark is kept in memory; the caller re-establishes it after a + /// restart. async fn obsolete( &self, provider: &Provider, @@ -355,8 +517,26 @@ impl LogStore for ObjectStoreLogStore { entry_id: EntryId, ) -> Result<()> { self.check_region(provider, region_id)?; - record_obsolete(&self.obsolete_entry_ids, region_id, entry_id); - Ok(()) + let (response_tx, response_rx) = oneshot::channel(); + let command = Command::Obsolete { + region_id, + entry_id, + response: response_tx, + }; + let answered = match self.command_tx.send(command).await { + Ok(()) => response_rx.await.ok(), + Err(_) => None, + }; + match answered { + Some(result) => result, + // The actor exited, before the command was queued or with it + // still queued behind the stop: nothing is assigned an id any + // more, so the watermark alone is consistent. + None => { + record_obsolete(&self.obsolete_entry_ids, region_id, entry_id); + Ok(()) + } + } } /// Hides every entry of the region in memory. The caller re-establishes @@ -390,40 +570,569 @@ impl LogStore for ObjectStoreLogStore { } } +type AppendResponse = oneshot::Sender>; + +type QueuedAppend = (Vec, AppendResponse); + enum Command { + /// Answered once the watermark is recorded and no id of the region at + /// or below `entry_id` can be assigned, or with the reason neither was + /// done. + Obsolete { + region_id: RegionId, + entry_id: EntryId, + response: oneshot::Sender>, + }, Stop { response: oneshot::Sender>, }, + #[cfg(any(test, feature = "testing"))] + Seal { + response: oneshot::Sender>, + }, } +/// An append waiting for the object that holds its entries. +struct PendingAppend { + last_entry_ids: HashMap, + response: AppendResponse, +} + +/// Where the conditional create of a sealed batch stands. +enum CreateState { + /// The create has not started because no slot was free. + Pending, + InFlight, + Created, +} + +/// A sealed and encoded batch that holds its object sequence and waits to be +/// created, indexed and acknowledged. +struct SealedBatch { + object_seq: u64, + bytes: Bytes, + footer: Vec, + waiters: Vec, + #[cfg(any(test, feature = "testing"))] + seal_waiters: Vec>>, + state: CreateState, +} + +impl SealedBatch { + fn fail(self, error: impl Fn() -> Error) { + for waiter in self.waiters { + let _ = waiter.response.send(Err(error())); + } + #[cfg(any(test, feature = "testing"))] + for waiter in self.seal_waiters { + let _ = waiter.send(Err(error())); + } + } +} + +type CreateOutcome = (u64, Result); + +/// The actor that owns the open batch and the sealed batches until they are +/// durable. +/// +/// Sealed batches form a pipeline in sequence order. A batch is created under +/// its sequence as soon as one of [`MAX_IN_FLIGHT_CREATES`] slots is free, but +/// it is indexed and acknowledged only once every earlier batch is, so the +/// acknowledged history of a region never has a missing predecessor. Every +/// batch links to its pipeline predecessor, the batch sealed before it, or to +/// the last indexed object when no earlier batch waits, so recovery replays a +/// batch only on the chain of a later object. The outcomes of a create are: +/// +/// | Situation | Sequence | Waiters | Store | +/// | --- | --- | --- | --- | +/// | created, or identical retry | advances | acknowledged in order | healthy | +/// | transient error | never reused: a create that reported an error may have written its object | every batch that is not indexed fails, created ones included; retries are assigned new ids | healthy: the next batch links to the last indexed object, so no chain reaches a failed batch | +/// | an object of an earlier epoch holds the sequence | never reused | as for a transient error | healthy: that object can never be on the chain | +/// | an object of the same or a later epoch holds the sequence, encoding or catalog error | unchanged | every batch that is not indexed fails | poisoned | +/// | created at the last representable sequence | cannot advance | acknowledged | poisoned: no later batch can be allocated a sequence | +/// +/// A create still in flight when its batch failed runs to completion and its +/// outcome is ignored: whatever it stored is off the chain. After `stop` began +/// nothing is admitted and no create starts; creates in flight run to +/// completion and acknowledge if they succeed. +/// +/// Appends arrive on their own channel, which the actor does not read while +/// [`MAX_SEALED_BATCHES`] batches wait, so commands such as `stop` are handled +/// while admission is held back. struct Actor { + io: Arc, + catalog: Arc>, + obsolete_entry_ids: ObsoleteEntryIds, terminal_error: TerminalError, stopped: Arc, command_rx: mpsc::Receiver, + append_rx: mpsc::Receiver, + open_batch: OpenBatch, issued_entry_ids: HashMap, - next_object_seq: Option, - epoch: u64, - /// The object a newly sealed batch extends. - chain_tip: ChainLink, + /// Waiters of the open batch. + pending: Vec, + /// Batches that are not durable yet, in sequence order. + sealed: VecDeque, + /// Creates in flight, including those of batches that already failed. + creates: FuturesUnordered>, + /// Callers of `stop`, answered once nothing is in flight. stop: Vec>>, + /// Sequence of the next sealed batch, `None` once the sequence is + /// exhausted. The open batch assigns its entry ids from it. It is above + /// every sealed batch and never moves back while the store runs. + next_object_seq: Option, + /// Epoch of this instance, carried by every object it writes. + epoch: u64, + /// The last indexed object, which is the start object until a batch is + /// indexed. A batch sealed while no earlier batch waits extends it. + last_indexed: ChainLink, + flush_interval: Duration, + #[cfg(any(test, feature = "testing"))] + admitted_appends: watch::Sender, + #[cfg(any(test, feature = "testing"))] + creates_held: watch::Receiver, + #[cfg(any(test, feature = "testing"))] + creates_fail: Arc, } impl Actor { async fn run(mut self) { - while let Some(Command::Stop { response }) = self.command_rx.recv().await { - self.handle_stop(response); + let mut interval = tokio::time::interval(self.flush_interval); + interval.set_missed_tick_behavior(MissedTickBehavior::Delay); + // The first tick completes immediately. + interval.tick().await; + + loop { + tokio::select! { + _ = interval.tick() => { + // Nothing starts after stop began. + if !self.is_stopped() { + self.flush_open_batch(); + } + } + Some((object_seq, result)) = self.creates.next(), if !self.creates.is_empty() => { + self.on_create_completed(object_seq, result); + } + Some((entries, response)) = self.append_rx.recv(), if self.sealed.len() + 2 <= MAX_SEALED_BATCHES => { + self.handle_append(entries, response); + } + command = self.command_rx.recv() => match command { + Some(Command::Obsolete { region_id, entry_id, response }) => { + self.handle_obsolete(region_id, entry_id, response); + } + Some(Command::Stop { response }) => { + self.handle_stop(response); + } + #[cfg(any(test, feature = "testing"))] + Some(Command::Seal { response }) => { + self.handle_seal(response); + } + // Every sender is gone: the store was dropped without + // `stop`. The creates in flight are dropped with the actor. + None => return, + }, + } if self.finish_stop() { return; } } } - fn handle_stop(&mut self, response: oneshot::Sender>) { - self.stop.push(response); + fn handle_append(&mut self, entries: Vec, response: AppendResponse) { + // The append was queued before `stop` set the flag; nothing that is + // not durable yet gets admitted once it is set. + if self.is_stopped() { + let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build())); + return; + } + if let Some(error) = terminal(&self.terminal_error) { + let _ = response.send(Err(shared(&error))); + return; + } + self.admit(entries, response); } + /// Admits `entries` into the open batch, which assigns their ids under + /// the next object sequence. The caller waits for the object. + fn admit(&mut self, entries: Vec, response: AppendResponse) { + // A region would run past the position range of the open batch: the + // batch is sealed and the entries open the next one. The size limit + // seals a batch long before a million entries of one region, so this + // is a theoretical bound. + if !self.open_batch.is_empty() && self.open_batch.would_exhaust_positions(&entries) { + self.flush_open_batch(); + if let Some(error) = terminal(&self.terminal_error) { + let _ = response.send(Err(shared(&error))); + return; + } + } + let Some(object_seq) = self.next_object_seq else { + let error = self.poison( + WalObjectSequenceExhaustedSnafu { + last_object_seq: OBJECT_SEQ_LIMIT - 1, + } + .build(), + ); + let _ = response.send(Err(shared(&error))); + return; + }; + let last_entry_ids = match self.open_batch.admit(object_seq, entries) { + Ok(last_entry_ids) => last_entry_ids, + // Nothing was admitted: the append alone runs past the position + // range, which no object can hold. + Err(error) => { + let _ = response.send(Err(error)); + return; + } + }; + self.pending.push(PendingAppend { + last_entry_ids, + response, + }); + #[cfg(any(test, feature = "testing"))] + self.admitted_appends.send_modify(|count| *count += 1); + if self.open_batch.should_seal() { + self.flush_open_batch(); + } + } + + /// Seals the open batch as the object `next_object_seq` and starts its + /// create when a slot is free. Returns whether a batch was sealed. + fn flush_open_batch(&mut self) -> bool { + if self.open_batch.is_empty() { + return false; + } + if let Some(error) = terminal(&self.terminal_error) { + self.fail_unacknowledged(|| shared(&error)); + return false; + } + let Some(object_seq) = self.next_object_seq else { + self.poison( + WalObjectSequenceExhaustedSnafu { + last_object_seq: OBJECT_SEQ_LIMIT - 1, + } + .build(), + ); + return false; + }; + + let (entries, _) = self.open_batch.seal(); + let header = Header { + object_seq, + epoch: self.epoch, + prev: Some( + self.sealed + .back() + .map_or(self.last_indexed, |batch| ChainLink { + object_seq: batch.object_seq, + epoch: self.epoch, + }), + ), + }; + let encoded = match encode_batch(header, entries) { + Ok(encoded) => encoded, + Err(error) => { + self.poison(error); + return false; + } + }; + self.next_object_seq = object_seq + .checked_add(1) + .filter(|next_object_seq| *next_object_seq < OBJECT_SEQ_LIMIT); + if self.next_object_seq.is_none() { + // The batch takes the last representable sequence: it is created + // and acknowledged, but no later batch can be allocated one. + set_terminal( + &self.terminal_error, + WalObjectSequenceExhaustedSnafu { + last_object_seq: object_seq, + } + .build(), + ); + } + self.sealed.push_back(SealedBatch { + object_seq, + bytes: encoded.bytes, + footer: encoded.footer, + waiters: std::mem::take(&mut self.pending), + #[cfg(any(test, feature = "testing"))] + seal_waiters: Vec::new(), + state: CreateState::Pending, + }); + self.start_creates(); + true + } + + /// Starts the creates of pending batches in sequence order while fewer + /// than [`MAX_IN_FLIGHT_CREATES`] are in flight, counting the creates of + /// batches that already failed. Nothing starts once stop began. + fn start_creates(&mut self) { + if self.is_stopped() { + return; + } + let mut in_flight = self.creates.len(); + for batch in self.sealed.iter_mut() { + if in_flight >= MAX_IN_FLIGHT_CREATES { + break; + } + if !matches!(batch.state, CreateState::Pending) { + continue; + } + batch.state = CreateState::InFlight; + in_flight += 1; + let io = self.io.clone(); + let object_seq = batch.object_seq; + let bytes = batch.bytes.clone(); + let epoch = self.epoch; + #[cfg(any(test, feature = "testing"))] + let mut creates_held = self.creates_held.clone(); + #[cfg(any(test, feature = "testing"))] + let creates_fail = self.creates_fail.clone(); + self.creates.push(Box::pin(async move { + // The store was dropped while the create was parked: it never + // runs. + #[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); + } + let result = match io.put_if_absent(object_seq, bytes).await { + Err(error @ Error::WalObjectConflict { .. }) => { + stale_conflict(io.as_ref(), object_seq, epoch, error).await + } + result => result, + }; + (object_seq, result) + })); + } + } + + fn on_create_completed(&mut self, object_seq: u64, result: Result) { + // A batch that already failed, or that the store gave up on when it + // poisoned itself: the object may exist, but it is off the chain. + let Some(index) = self + .sealed + .iter() + .position(|batch| batch.object_seq == object_seq) + else { + self.start_creates(); + return; + }; + match result { + Ok(_) => self.sealed[index].state = CreateState::Created, + // The object store did not confirm the object, which may still + // exist or land later, or an earlier epoch holds the sequence and + // can never be on the chain. The caller retries the append itself. + Err(error @ (Error::WalObjectStore { .. } | Error::StaleWalObject { .. })) => { + self.roll_back(Arc::new(error)) + } + Err(error) => { + self.poison(error); + return; + } + } + self.settle(); + } + + /// Indexes and acknowledges the sealed batches from the front as far as + /// they are created, and starts the creates that a free slot allows. + fn settle(&mut self) { + while let Some(front) = self.sealed.front() { + if !matches!(front.state, CreateState::Created) || !self.index_front() { + break; + } + } + self.start_creates(); + } + + /// Indexes the created object at the front and acknowledges its waiters. + /// Returns false when the catalog rejected it, which poisons the store. + fn index_front(&mut self) -> bool { + let Some(front) = self.sealed.front() else { + return false; + }; + let indexed = self + .catalog + .write() + .unwrap_or_else(PoisonError::into_inner) + .insert_object(front.object_seq, front.footer.clone()); + if let Err(error) = indexed { + self.poison(error); + return false; + } + let Some(batch) = self.sealed.pop_front() else { + return false; + }; + self.last_indexed = ChainLink { + object_seq: batch.object_seq, + epoch: self.epoch, + }; + for waiter in batch.waiters { + let _ = waiter.response.send(Ok(AppendBatchResponse { + last_entry_ids: waiter.last_entry_ids, + })); + } + #[cfg(any(test, feature = "testing"))] + for waiter in batch.seal_waiters { + let _ = waiter.send(Ok(())); + } + true + } + + /// Fails every batch that is not indexed after a create failed + /// transiently, including batches already created and batches whose + /// create is in flight. Their sequences are not reused: a create that + /// reported an error may still have written its object. The next batch + /// links to the last indexed object, so no later chain reaches a failed + /// batch. Waiters of a store that was stopped meanwhile learn that + /// instead of the I/O error, like every other entry that never became + /// durable. + fn roll_back(&mut self, error: Arc) { + let stopped = self.is_stopped(); + let failure = || { + if stopped { + ObjectStoreWalStoppedSnafu.build() + } else { + shared(&error) + } + }; + for batch in self.sealed.drain(..) { + batch.fail(failure); + } + self.reset_open_batch(); + self.fail_unacknowledged(failure); + } + + /// Records `error` as terminal and fails every waiter that is not + /// acknowledged with it, or with the stopped error if the store was + /// stopped meanwhile. Creates in flight run to completion, but their + /// outcome is ignored: the entries of an object they create were never + /// acknowledged, like those of a crash between creation and + /// acknowledgement. Returns the recorded error. + fn poison(&mut self, error: Error) -> Arc { + let error = set_terminal(&self.terminal_error, error); + let stopped = self.is_stopped(); + let failure = || { + if stopped { + ObjectStoreWalStoppedSnafu.build() + } else { + shared(&error) + } + }; + for batch in self.sealed.drain(..) { + batch.fail(failure); + } + self.reset_open_batch(); + self.fail_unacknowledged(failure); + error + } + + /// Drops the entries of the open batch; the next admission hands out the + /// same ids again under the sequence the batch is at. + fn reset_open_batch(&mut self) { + self.open_batch.reset(); + } + + /// Fails the waiters of the open batch. + fn fail_unacknowledged(&mut self, error: impl Fn() -> Error) { + for pending in self.pending.drain(..) { + let _ = pending.response.send(Err(error())); + } + } + + /// Makes sure no id of the region at or below `entry_id` is assigned + /// from now on, then records the obsolete watermark of the region; + /// neither is done when the sequence cannot be raised. + fn handle_obsolete( + &mut self, + region_id: RegionId, + entry_id: EntryId, + response: oneshot::Sender>, + ) { + let result = self.raise_sequence_floor(entry_id).map(|_| { + record_obsolete(&self.obsolete_entry_ids, region_id, entry_id); + }); + let _ = response.send(result); + } + + /// Moves the sequence of the next object above the object that holds + /// `entry_id` unless it is there already. The move fails while the open + /// batch has handed out ids under the next sequence. A floor that does + /// not fit an entry id poisons the store like an exhausted sequence. + fn raise_sequence_floor(&mut self, entry_id: EntryId) -> Result<()> { + // Nothing is assigned an id after stop began. + if self.is_stopped() { + return Ok(()); + } + if let Some(error) = terminal(&self.terminal_error) { + return Err(shared(&error)); + } + let sequence_floor = sequence_floor(entry_id); + let Some(next_object_seq) = self.next_object_seq else { + return Ok(()); + }; + if next_object_seq >= sequence_floor { + return Ok(()); + } + ensure!( + self.open_batch.is_empty(), + WalObjectSequenceUnsettledSnafu { + object_seq: next_object_seq, + } + ); + if sequence_floor < OBJECT_SEQ_LIMIT { + self.next_object_seq = Some(sequence_floor); + Ok(()) + } else { + self.next_object_seq = None; + let error = self.poison( + WalObjectSequenceExhaustedSnafu { + last_object_seq: OBJECT_SEQ_LIMIT - 1, + } + .build(), + ); + Err(shared(&error)) + } + } + + /// Begins stopping. Nothing is admitted from now on; the open batch and + /// the batches whose create has not started are dropped. Stop is answered by [`finish_stop`](Self::finish_stop) + /// once nothing is in flight. + fn handle_stop(&mut self, response: oneshot::Sender>) { + self.stop.push(response); + self.reset_open_batch(); + for pending in self.pending.drain(..) { + let _ = pending + .response + .send(Err(ObjectStoreWalStoppedSnafu.build())); + } + if let Some(index) = self + .sealed + .iter() + .position(|batch| matches!(batch.state, CreateState::Pending)) + { + for batch in self.sealed.drain(index..) { + batch.fail(|| ObjectStoreWalStoppedSnafu.build()); + } + } + } + + /// Answers the callers of `stop` once every sealed batch is settled and + /// every create has completed. Returns true when the actor is done. fn finish_stop(&mut self) -> bool { - if !self.is_stopped() || self.stop.is_empty() { + if self.stop.is_empty() || !self.sealed.is_empty() || !self.creates.is_empty() { return false; } for response in self.stop.drain(..) { @@ -432,6 +1141,22 @@ impl Actor { true } + #[cfg(any(test, feature = "testing"))] + fn handle_seal(&mut self, response: oneshot::Sender>) { + if self.is_stopped() { + let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build())); + return; + } + if self.flush_open_batch() + && let Some(batch) = self.sealed.back_mut() + { + batch.seal_waiters.push(response); + return; + } + let result = terminal(&self.terminal_error).map_or(Ok(()), |error| Err(shared(&error))); + let _ = response.send(result); + } + fn is_stopped(&self) -> bool { self.stopped.load(Ordering::Acquire) } @@ -459,6 +1184,18 @@ fn shared(error: &Arc) -> Error { ObjectStoreWalSnafu.into_error(error.clone()) } +fn encode_batch(header: Header, entries: Vec) -> Result { + let records = entries + .into_iter() + .map(|entry| Record { + region_id: entry.region_id(), + entry_id: entry.entry_id(), + payload: Bytes::from(entry.into_bytes()), + }) + .collect::>(); + encode_object(header, &records) +} + /// What recovery rebuilt from the objects under a prefix. #[derive(Debug)] struct Recovered { @@ -644,11 +1381,7 @@ async fn start_epoch( .fail(); } Err(error @ Error::WalObjectConflict { .. }) => { - let head = io.get_range(object_seq, 0, HEADER_LEN as u64).await?; - let existing = decode_header(&head).with_context(|_| InvalidWalObjectSnafu { - path: io.object_path(object_seq), - })?; - if existing.epoch >= epoch { + if epoch_of(io, object_seq).await? >= epoch { return Err(error); } object_seq += 1; @@ -658,6 +1391,36 @@ async fn start_epoch( } } +async fn epoch_of(io: &dyn WalObjectIo, object_seq: u64) -> Result { + let head = io.get_range(object_seq, 0, HEADER_LEN as u64).await?; + let header = decode_header(&head).with_context(|_| InvalidWalObjectSnafu { + path: io.object_path(object_seq), + })?; + Ok(header.epoch) +} + +/// Resolves a create of this store's `epoch` that met a different object: +/// an object of an earlier epoch can never be on the chain, so the create +/// fails like a transient error; an object of the same or a later epoch +/// belongs to another writer and `conflict` is returned. +async fn stale_conflict( + io: &dyn WalObjectIo, + object_seq: u64, + epoch: u64, + conflict: Error, +) -> Result { + let existing_epoch = epoch_of(io, object_seq).await?; + if existing_epoch >= epoch { + return Err(conflict); + } + StaleWalObjectSnafu { + path: io.object_path(object_seq), + existing_epoch, + epoch, + } + .fail() +} + /// Fetches and verifies the headers and footers of `objects`, up to /// `concurrency` objects at a time, and returns them ordered by object sequence /// whatever the order the fetches complete in. The first failure abandons the @@ -836,14 +1599,13 @@ mod tests { use common_base::readable_size::ReadableSize; use common_error::ext::{ErrorExt, RetryHint}; use object_store::services::Memory; + use store_api::logstore::entry::{MultiplePartEntry, MultiplePartHeader}; use tokio::time::timeout; use super::*; use crate::error::WalObjectStoreSnafu; - use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, entry_id}; - use crate::object_store_wal::format::{ - FOOTER_ENTRY_LEN, Header, Record, decode_object, encode_object, - }; + use crate::object_store_wal::batch::{POSITION_LIMIT, entry_id}; + use crate::object_store_wal::format::{FOOTER_ENTRY_LEN, decode_object}; const PREFIX: &str = "wal/datanodes/1/epochs/2"; const WAIT: Duration = Duration::from_secs(30); @@ -861,8 +1623,14 @@ mod tests { } } + /// Every append reaches the size limit, so it is persisted on its own. fn eager() -> ObjectStoreWalConfig { - config(Duration::from_millis(10), 1) + config(Duration::from_secs(3600), 1) + } + + /// Nothing is persisted until a test seals the open batch. + fn manual() -> ObjectStoreWalConfig { + config(Duration::from_secs(3600), u64::MAX) } async fn open( @@ -890,6 +1658,126 @@ mod tests { fn id(object_seq: u64, position: u64) -> EntryId { entry_id(object_seq, position) } + + fn entry(store: &ObjectStoreLogStore, region_id: RegionId, data: &str) -> Entry { + store + .entry(data.as_bytes().to_vec(), 0, region_id, &provider(region_id)) + .unwrap() + } + + async fn append( + store: &ObjectStoreLogStore, + region_id: RegionId, + data: &str, + ) -> Result { + store + .append_batch(vec![entry(store, region_id, data)]) + .await + } + + fn spawn_append_batch( + store: &Arc, + entries: Vec, + ) -> tokio::task::JoinHandle> { + let store = store.clone(); + tokio::spawn(async move { store.append_batch(entries).await }) + } + + /// Appends `count` single-entry batches, waiting for each to be admitted + /// before the next is sent, so they are admitted in order. + async fn spawn_appends( + store: &Arc, + region_id: RegionId, + count: usize, + ) -> Vec>> { + let admitted = *store.admitted_appends.borrow(); + let mut handles = Vec::with_capacity(count); + for index in 1..=count { + let entries = vec![entry(store, region_id, &format!("a{index}"))]; + handles.push(spawn_append_batch(store, entries)); + store + .wait_for_admitted_appends(admitted + index) + .await + .unwrap(); + } + handles + } + + async fn object_seqs(io: &dyn WalObjectIo) -> Vec { + io.list() + .await + .unwrap() + .into_iter() + .map(|object| object.object_seq) + .collect() + } + + fn unwrap_shared(error: &Error) -> &Error { + match error { + Error::ObjectStoreWal { source, .. } => source, + other => panic!("expected a shared error, actual {other:?}"), + } + } + + fn assert_stopped(error: &Error) { + assert!( + matches!(error, Error::ObjectStoreWalStopped { .. }), + "unexpected error: {error:?}" + ); + } + + /// Waits for the next create that `ParkedIo` parked. + async fn next_create( + parked: &mut mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + ) -> (u64, oneshot::Sender) { + timeout(WAIT, parked.recv()).await.unwrap().unwrap() + } + + /// Waits for `count` parked creates and returns their releases by + /// sequence. + async fn parked_creates( + parked: &mut mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + count: usize, + ) -> HashMap> { + let mut releases = HashMap::new(); + for _ in 0..count { + let (object_seq, release) = next_create(parked).await; + releases.insert(object_seq, release); + } + releases + } + + async fn open_over( + io: Arc, + config: &ObjectStoreWalConfig, + ) -> Arc { + ObjectStoreLogStore::open(io, config, PREFIX.to_string()) + .await + .unwrap() + } + + async fn open_parking_creates( + object_store: ObjectStore, + config: &ObjectStoreWalConfig, + ) -> ( + Arc, + Arc, + mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + ) { + let (io, parked) = ParkedIo::parking_creates(object_store); + let store = open_over(io.clone(), config).await; + // The start object of the open is created without parking. + io.creates_parked.store(true, Ordering::SeqCst); + (store, io, parked) + } + + /// Round-trips a command through the actor, so an assertion that + /// something did not happen runs after the actor handled every command + /// sent before. The open batch must be empty, as under the eager config. + async fn round_trip_actor(store: &ObjectStoreLogStore) { + store.seal_open_batch().await.unwrap(); + } + #[tokio::test] async fn test_store_rejects_invalid_config() { for config in [ @@ -1071,6 +1959,18 @@ mod tests { io.put_if_absent(object_seq, encoded.bytes).await.unwrap(); } + /// Puts an object of another writer with `epoch` under `object_seq`. A + /// store opened on an empty prefix writes epoch 1. + async fn put_foreign(object_store: &ObjectStore, object_seq: u64, epoch: u64) { + let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap(); + let header = Header { + object_seq, + epoch, + prev: None, + }; + put_object_with_header(&io, header, &[]).await; + } + fn object_path(object_store: &ObjectStore, object_seq: u64) -> String { ObjectStoreIo::new(object_store.clone(), PREFIX) .unwrap() @@ -1359,10 +2259,51 @@ mod tests { HashMap::from([(region(1), id(5, 7)), (region(2), 8)]), recovered.durable_entry_ids ); - let store = open(object_store, &eager()).await; + let store = open(object_store.clone(), &eager()).await; assert_eq!(id(5, 7), latest(&store, region(1))); assert_eq!(8, latest(&store, region(2))); assert_eq!(0, latest(&store, region(3))); + + // The sequence resumes above the object the largest id names, where + // the start object of the store takes object 6, so the first new id + // of every region is greater than every old one. + let response = store + .append_batch(vec![ + entry(&store, region(1), "a"), + entry(&store, region(2), "b"), + ]) + .await + .unwrap(); + assert_eq!( + HashMap::from([(region(1), id(7, 1)), (region(2), id(7, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![2, 4, 6, 7], object_seqs(store.io.as_ref()).await); + assert_eq!( + expected_entries( + region(1), + &[ + (1, "e1"), + (id(5, 7), &format!("e{}", id(5, 7))), + (id(7, 1), "a") + ] + ), + read_entries(&store, region(1), 0).await + ); + store.stop().await.unwrap(); + + // The sequence continues after the start object of the restart. + let store = open(object_store, &eager()).await; + assert_eq!(id(7, 1), latest(&store, region(1))); + let response = append(&store, region(2), "b2").await.unwrap(); + assert_eq!( + HashMap::from([(region(2), id(9, 1))]), + response.last_entry_ids + ); + assert_eq!( + expected_entries(region(2), &[(8, "e8"), (id(7, 1), "b"), (id(9, 1), "b2")]), + read_entries(&store, region(2), 0).await + ); store.stop().await.unwrap(); } @@ -1401,7 +2342,14 @@ mod tests { let other = Provider::object_store_provider(region_id, "other/prefix".to_string()); let raft = Provider::raft_engine_provider(region_id.as_u64()); for foreign in [&other, &raft] { + let entry = Entry::Naive(NaiveEntry { + provider: foreign.clone(), + region_id, + entry_id: 0, + data: b"a1".to_vec(), + }); let errors = [ + store.append_batch(vec![entry]).await.unwrap_err(), store.latest_entry_id(foreign).unwrap_err(), store.read(foreign, 0, None).await.err().unwrap(), store.create_namespace(foreign).await.unwrap_err(), @@ -1435,12 +2383,14 @@ mod tests { timeout(WAIT, store.command_tx.closed()).await.unwrap(); store.stop().await.unwrap(); assert!(store.stopped.load(Ordering::Acquire)); + // Even an empty append fails once stop began. + assert_stopped(&store.append_batch(Vec::new()).await.unwrap_err()); } - #[tokio::test] - async fn test_store_stop_with_actor_gone() { + /// Builds a store whose commands the test receives instead of an actor. + fn store_without_actor() -> (ObjectStoreLogStore, mpsc::Receiver) { let (command_tx, command_rx) = mpsc::channel(COMMAND_BUFFER); - drop(command_rx); + let (append_tx, _) = mpsc::channel(APPEND_BUFFER); let store = ObjectStoreLogStore { prefix: PREFIX.to_string(), io: Arc::new(ObjectStoreIo::new(memory_store(), PREFIX).unwrap()), @@ -1449,7 +2399,18 @@ mod tests { terminal_error: Arc::default(), stopped: Arc::new(AtomicBool::new(false)), command_tx, + append_tx, + admitted_appends: watch::channel(0).1, + creates_held: watch::channel(false).0, + creates_fail: Arc::default(), }; + (store, command_rx) + } + + #[tokio::test] + async fn test_store_stop_with_actor_gone() { + let (store, command_rx) = store_without_actor(); + drop(command_rx); store.stop().await.unwrap(); store.stop().await.unwrap(); } @@ -1485,6 +2446,7 @@ mod tests { let region_id = region(1); let provider = provider(region_id); let errors = [ + store.append_batch(Vec::new()).await.unwrap_err(), store.latest_entry_id(&provider).unwrap_err(), store.read(&provider, 0, None).await.err().unwrap(), store.create_namespace(&provider).await.unwrap_err(), @@ -1641,21 +2603,6 @@ mod tests { assert_invalid_object(&error, &object.path, "fewer bytes than the listed"); } - #[tokio::test] - async fn test_store_append_is_unsupported() { - let store = open(memory_store(), &eager()).await; - let error = store.append_batch(Vec::new()).await.unwrap_err(); - assert!(matches!( - error, - Error::UnsupportedObjectStoreWalOperation { .. } - )); - assert_eq!( - common_error::status_code::StatusCode::Unsupported, - error.status_code() - ); - store.stop().await.unwrap(); - } - async fn read_entries( store: &ObjectStoreLogStore, region_id: RegionId, @@ -1799,10 +2746,7 @@ mod tests { store.stop().await.unwrap(); let reopened = open(object_store, &eager()).await; assert_eq!(all, read_entries(&reopened, region(1), 0).await); - reopened - .obsolete(&p, region(1), EntryId::MAX) - .await - .unwrap(); + reopened.obsolete(&p, region(1), 20).await.unwrap(); assert_eq!( Vec::::new(), read_entries(&reopened, region(1), 0).await @@ -1909,6 +2853,1046 @@ mod tests { store.stop().await.unwrap(); } + #[tokio::test] + async fn test_store_assigns_ids_from_the_object_sequence() { + let store = open(memory_store(), &manual()).await; + let region_one = region(1); + let region_two = region(2); + + // Appends admitted into one batch share its object; the regions take + // their own positions in admission order, and every append learns + // the last id of each of its regions. + let first = spawn_append_batch( + &store, + vec![ + entry(&store, region_one, "a1"), + entry(&store, region_two, "b1"), + entry(&store, region_one, "a2"), + ], + ); + store.wait_for_admitted_appends(1).await.unwrap(); + let second = spawn_append_batch(&store, vec![entry(&store, region_two, "b2")]); + store.wait_for_admitted_appends(2).await.unwrap(); + assert!(!first.is_finished()); + assert_eq!(0, latest(&store, region_one)); + store.seal_open_batch().await.unwrap(); + assert_eq!( + HashMap::from([(region_one, id(1, 2)), (region_two, id(1, 1))]), + first.await.unwrap().unwrap().last_entry_ids + ); + assert_eq!( + HashMap::from([(region_two, id(1, 2))]), + second.await.unwrap().unwrap().last_entry_ids + ); + + // The next object starts every region at position one again. + let third = spawn_append_batch( + &store, + vec![ + entry(&store, region_two, "b3"), + entry(&store, region_one, "a3"), + ], + ); + store.wait_for_admitted_appends(3).await.unwrap(); + store.seal_open_batch().await.unwrap(); + assert_eq!( + HashMap::from([(region_one, id(2, 1)), (region_two, id(2, 1))]), + third.await.unwrap().unwrap().last_entry_ids + ); + assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await); + assert_eq!( + expected_entries( + region_one, + &[(id(1, 1), "a1"), (id(1, 2), "a2"), (id(2, 1), "a3")] + ), + read_entries(&store, region_one, 1).await + ); + assert_eq!( + expected_entries(region_two, &[(id(1, 2), "b2"), (id(2, 1), "b3")]), + read_entries(&store, region_two, id(1, 2)).await + ); + assert_eq!(id(2, 1), latest(&store, region_one)); + assert_eq!(id(2, 1), latest(&store, region_two)); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_seals_when_a_region_exhausts_its_positions() { + let store = open(memory_store(), &manual()).await; + let region_id = region(1); + let entries_of = |count: u64| { + (0..count) + .map(|_| entry(&store, region_id, "")) + .collect::>() + }; + + // The open batch holds every position of the region. + let full = spawn_append_batch(&store, entries_of(POSITION_LIMIT - 1)); + store.wait_for_admitted_appends(1).await.unwrap(); + assert_eq!(vec![0], object_seqs(store.io.as_ref()).await); + + // The next entry of the region seals the batch and opens the next + // object; another region would still have fit. + let next = spawn_append_batch(&store, vec![entry(&store, region_id, "next")]); + store.wait_for_admitted_appends(2).await.unwrap(); + let response = timeout(WAIT, full).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(1, POSITION_LIMIT - 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await); + + // An append that alone runs past the range fits no object: it seals + // the open batch like any append the batch cannot take, then it is + // refused without poisoning the store. + let error = store + .append_batch(entries_of(POSITION_LIMIT)) + .await + .unwrap_err(); + assert!( + matches!(error, Error::WalEntryPositionExhausted { .. }), + "unexpected error: {error:?}" + ); + let response = timeout(WAIT, next).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + let after = spawn_append_batch(&store, vec![entry(&store, region_id, "after")]); + store.wait_for_admitted_appends(3).await.unwrap(); + store.seal_open_batch().await.unwrap(); + let response = timeout(WAIT, after).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(3, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1, 2, 3], object_seqs(store.io.as_ref()).await); + assert_eq!( + expected_entries(region_id, &[(id(2, 1), "next"), (id(3, 1), "after")]), + read_entries(&store, region_id, id(2, 1)).await + ); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_flushes_at_the_minimum_interval() { + // An append is acknowledged once its object is durable, so a completed + // append proves the tick after it sealed the batch. The appends are + // spaced far apart, so every one of them lands in its own tick and + // the ticks in between find an empty batch. + let interval = MIN_FLUSH_INTERVAL; + let store = open(memory_store(), &config(interval, u64::MAX)).await; + let data = ["a1", "a2", "a3"]; + for (index, data) in data.iter().enumerate() { + tokio::time::sleep(interval * 5).await; + timeout(WAIT, append(&store, region(1), data)) + .await + .unwrap() + .unwrap(); + let expected_seqs = (0..=index as u64 + 1).collect::>(); + assert_eq!(expected_seqs, object_seqs(store.io.as_ref()).await); + } + + // Ticks with an empty open batch do not create objects. + tokio::time::sleep(interval * 5).await; + let seqs = object_seqs(store.io.as_ref()).await; + assert_eq!(vec![0, 1, 2, 3], seqs); + // Object 0 is the start object of the store. + for object_seq in seqs.into_iter().skip(1) { + let bytes = store.io.get(object_seq).await.unwrap(); + let decoded = decode_object(&bytes).unwrap(); + assert_eq!(1, decoded.records.len(), "object {object_seq} is empty"); + } + assert_eq!( + expected_entries( + region(1), + &[(id(1, 1), "a1"), (id(2, 1), "a2"), (id(3, 1), "a3")] + ), + read_entries(&store, region(1), 1).await + ); + } + + #[tokio::test] + async fn test_store_rejects_entries_of_another_region() { + let store = open(memory_store(), &manual()).await; + let region_id = region(1); + + let entry = Entry::Naive(NaiveEntry { + provider: provider(region_id), + region_id: region(2), + entry_id: 0, + data: Vec::new(), + }); + let error = store.append_batch(vec![entry]).await.unwrap_err(); + assert!( + matches!(error, Error::MismatchedWalRegion { .. }), + "unexpected error: {error:?}" + ); + let incomplete = Entry::MultiplePart(MultiplePartEntry { + provider: provider(region_id), + region_id, + entry_id: 0, + headers: vec![MultiplePartHeader::First], + parts: vec![b"a1".to_vec()], + }); + let error = store.append_batch(vec![incomplete]).await.unwrap_err(); + assert!( + matches!(error, Error::IncompleteWalEntry { region_id: actual, .. } if actual == region_id), + "unexpected error: {error:?}" + ); + assert_eq!( + common_error::status_code::StatusCode::InvalidArguments, + error.status_code() + ); + assert!( + store + .append_batch(Vec::new()) + .await + .unwrap() + .last_entry_ids + .is_empty() + ); + // Nothing was admitted. + store.seal_open_batch().await.unwrap(); + assert_eq!(vec![0], object_seqs(store.io.as_ref()).await); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_pipelines_creates_and_acknowledges_in_sequence_order() { + let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let region_id = region(1); + let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 2).await; + + // At most the limit of creates run at a time; the rest wait for a slot. + let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await; + assert_eq!( + (1..=MAX_IN_FLIGHT_CREATES as u64).collect::>(), + releases.keys().copied().collect::>() + ); + round_trip_actor(&store).await; + assert!(parked.try_recv().is_err()); + assert!(appends.iter().all(|append| !append.is_finished())); + + // Objects 3 and 2 become durable before object 1 and free a slot + // each, but nothing is acknowledged ahead of object 1. + for (released, expected_next) in [(3, 5), (2, 6)] { + releases.remove(&released).unwrap().send(true).unwrap(); + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(expected_next, object_seq); + releases.insert(object_seq, release); + } + round_trip_actor(&store).await; + assert_eq!(vec![0, 2, 3], object_seqs(io.as_ref()).await); + assert!(appends.iter().all(|append| !append.is_finished())); + assert_eq!(0, latest(&store, region_id)); + + // Object 1 releases the acknowledgements of objects 1 to 3. + releases.remove(&1).unwrap().send(true).unwrap(); + let mut appends = appends.into_iter(); + for object_seq in 1..4 { + let response = timeout(WAIT, appends.next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq, 1))]), + response.last_entry_ids + ); + } + assert_eq!(id(3, 1), latest(&store, region_id)); + + // Object 6 before object 5: the append of object 6 waits for it. + releases.remove(&4).unwrap().send(true).unwrap(); + releases.remove(&6).unwrap().send(true).unwrap(); + let fourth = timeout(WAIT, appends.next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(4, 1))]), + fourth.last_entry_ids + ); + let appends = appends.collect::>(); + round_trip_actor(&store).await; + assert!(appends.iter().all(|append| !append.is_finished())); + releases.remove(&5).unwrap().send(true).unwrap(); + for (append, object_seq) in appends.into_iter().zip(5..) { + let response = timeout(WAIT, append).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq, 1))]), + response.last_entry_ids + ); + } + assert_eq!( + MAX_IN_FLIGHT_CREATES, + io.max_in_flight.load(Ordering::SeqCst) + ); + assert_eq!(vec![0, 1, 2, 3, 4, 5, 6], object_seqs(io.as_ref()).await); + // Every batch extends the batch sealed before it, whichever was + // durable first. + for object_seq in 1..=6 { + let header = decode_header(&io.get(object_seq).await.unwrap()).unwrap(); + assert_eq!( + Some(object_seq - 1), + header.prev.map(|link| link.object_seq) + ); + } + assert_eq!( + expected_entries( + region_id, + &[ + (id(1, 1), "a1"), + (id(2, 1), "a2"), + (id(3, 1), "a3"), + (id(4, 1), "a4"), + (id(5, 1), "a5"), + (id(6, 1), "a6") + ] + ), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_transient_failure_fails_every_batch_that_is_not_indexed() { + let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let region_id = region(1); + // Every slot is taken; the last batch waits for one. + let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 1).await; + let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await; + round_trip_actor(&store).await; + assert!(parked.try_recv().is_err()); + + // Object 1 fails while objects 2 to 4 are in flight and object 5 waits + // for a slot: every batch fails at once. + releases.remove(&1).unwrap().send(false).unwrap(); + for append in appends { + let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + assert_eq!(RetryHint::Retryable, error.retry_hint()); + } + assert_eq!(0, latest(&store, region_id)); + + // The retried entries take the sequences after the last sealed one. + // The creates of the failed batches still hold their slots, so only + // one retry starts until they complete; what they store is off the + // chain. + let retries = spawn_appends(&store, region_id, 2).await; + let (object_seq, first) = next_create(&mut parked).await; + assert_eq!(6, object_seq); + round_trip_actor(&store).await; + assert!(parked.try_recv().is_err()); + for release in releases.into_values() { + release.send(true).unwrap(); + } + let (object_seq, second) = next_create(&mut parked).await; + assert_eq!(7, object_seq); + first.send(true).unwrap(); + second.send(true).unwrap(); + for (retry, object_seq) in retries.into_iter().zip(6..) { + let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq, 1))]), + response.last_entry_ids + ); + } + assert_eq!(vec![0, 2, 3, 4, 6, 7], object_seqs(io.as_ref()).await); + assert_eq!( + expected_entries(region_id, &[(id(6, 1), "a1"), (id(7, 1), "a2")]), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_transient_failure_before_a_created_object_keeps_it_off_the_chain() { + let object_store = memory_store(); + let (store, io, mut parked) = open_parking_creates(object_store.clone(), &eager()).await; + let region_id = region(1); + let appends = spawn_appends(&store, region_id, 2).await; + let mut releases = parked_creates(&mut parked, 2).await; + + // Object 2 is created while object 1 failed: both batches fail, and + // the store keeps serving. + releases.remove(&2).unwrap().send(true).unwrap(); + timeout(WAIT, async { + while object_seqs(io.as_ref()).await != [0, 2] { + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .unwrap(); + releases.remove(&1).unwrap().send(false).unwrap(); + for append in appends { + let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + } + assert_eq!(0, latest(&store, region_id)); + + // The retry extends the start object, not the failed batches. + let retry = spawn_appends(&store, region_id, 1).await; + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(3, object_seq); + release.send(true).unwrap(); + let response = timeout(WAIT, retry.into_iter().next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(3, 1))]), + response.last_entry_ids + ); + let bytes = io.get(3).await.unwrap(); + let header = decode_header(&bytes).unwrap(); + assert_eq!( + Some(ChainLink { + object_seq: 0, + epoch: 1, + }), + header.prev + ); + store.stop().await.unwrap(); + + // After a restart only the acknowledged entry replays: object 2 is + // off the chain although it exists. + let store = open(object_store, &eager()).await; + assert_eq!(id(3, 1), latest(&store, region_id)); + assert_eq!( + expected_entries(region_id, &[(id(3, 1), "a1")]), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_permanent_failure_before_a_durable_object_poisons() { + let object_store = memory_store(); + let (store, _, mut parked) = open_parking_creates(object_store.clone(), &eager()).await; + let region_id = region(1); + let foreign = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap(); + // Another writer of the same epoch took sequence 1. + put_foreign(&object_store, 1, 1).await; + let appends = spawn_appends(&store, region_id, 2).await; + let mut releases = parked_creates(&mut parked, 2).await; + + // Object 2 is created; object 1 conflicts with the foreign object. + releases.remove(&2).unwrap().send(true).unwrap(); + timeout(WAIT, async { + while object_seqs(&foreign).await != [0, 1, 2] { + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .unwrap(); + round_trip_actor(&store).await; + assert!(appends.iter().all(|append| !append.is_finished())); + releases.remove(&1).unwrap().send(true).unwrap(); + for append in appends { + let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }), + "unexpected error: {error:?}" + ); + } + let error = store.latest_entry_id(&provider(region_id)).unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }), + "unexpected error: {error:?}" + ); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_conflict_with_an_earlier_epoch_fails_the_batch_without_poisoning() { + let object_store = memory_store(); + let (store, io, mut parked) = open_parking_creates(object_store.clone(), &eager()).await; + let region_id = region(1); + // A late object of an earlier epoch lands under sequence 1. + put_foreign(&object_store, 1, 0).await; + let append = spawn_appends(&store, region_id, 1).await; + next_create(&mut parked).await.1.send(true).unwrap(); + let error = timeout(WAIT, append.into_iter().next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap_err(); + assert!( + matches!( + unwrap_shared(&error), + Error::StaleWalObject { + existing_epoch: 0, + epoch: 1, + .. + } + ), + "unexpected error: {error:?}" + ); + assert_eq!(RetryHint::Retryable, error.retry_hint()); + + // The store keeps serving: the retry takes the next sequence. + let retry = spawn_appends(&store, region_id, 1).await; + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(2, object_seq); + release.send(true).unwrap(); + let response = timeout(WAIT, retry.into_iter().next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1, 2], object_seqs(io.as_ref()).await); + store.stop().await.unwrap(); + + // The stale object is off the chain after a restart. + let store = open(object_store, &eager()).await; + assert_eq!( + expected_entries(region_id, &[(id(2, 1), "a1")]), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_obsolete_raises_the_sequence_floor() { + let store = open(memory_store(), &manual()).await; + let region_one = region(1); + let region_two = region(2); + let mut admitted = 0; + let mut spawn_append = |region_id, data: &str| { + admitted += 1; + ( + spawn_append_batch(&store, vec![entry(&store, region_id, data)]), + admitted, + ) + }; + let (first, count) = spawn_append(region_one, "a1"); + store.wait_for_admitted_appends(count).await.unwrap(); + store.seal_open_batch().await.unwrap(); + let response = timeout(WAIT, first).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_one, id(1, 1))]), + response.last_entry_ids + ); + let (second, count) = spawn_append(region_one, "a2"); + store.wait_for_admitted_appends(count).await.unwrap(); + + // A durable watermark names an object below the next sequence and + // changes nothing: the open batch goes on under sequence 2. + store + .obsolete(&provider(region_one), region_one, id(1, 1)) + .await + .unwrap(); + let (third, count) = spawn_append(region_one, "a3"); + store.wait_for_admitted_appends(count).await.unwrap(); + + // A watermark that names a later object, as one inherited from + // another prefix does, cannot move the sequence while the open batch + // handed out ids under it: neither the sequence nor the watermark + // moves. + let error = store + .obsolete(&provider(region_two), region_two, id(3, 7)) + .await + .unwrap_err(); + assert!( + matches!( + error, + Error::WalObjectSequenceUnsettled { object_seq: 2, .. } + ), + "unexpected error: {error:?}" + ); + assert_eq!(RetryHint::Retryable, error.retry_hint()); + assert!( + store + .obsolete_entry_ids + .lock() + .unwrap() + .get(®ion_two) + .is_none() + ); + + // Once the batch is durable the sequence moves above that object, so + // the region's next id is greater than the watermark. + store.seal_open_batch().await.unwrap(); + let response = timeout(WAIT, second).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_one, id(2, 1))]), + response.last_entry_ids + ); + let response = timeout(WAIT, third).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_one, id(2, 2))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await); + store + .obsolete(&provider(region_two), region_two, id(3, 7)) + .await + .unwrap(); + assert_eq!( + Some(&id(3, 7)), + store.obsolete_entry_ids.lock().unwrap().get(®ion_two) + ); + let (fourth, count) = spawn_append(region_two, "b1"); + store.wait_for_admitted_appends(count).await.unwrap(); + store.seal_open_batch().await.unwrap(); + let response = timeout(WAIT, fourth).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_two, id(4, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1, 2, 4], object_seqs(store.io.as_ref()).await); + // The watermark hides nothing of this prefix; the region has no + // entry at or below it. + assert_eq!( + expected_entries(region_two, &[(id(4, 1), "b1")]), + read_entries(&store, region_two, 1).await + ); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_does_not_reuse_the_sequence_of_a_failed_create() { + let (io, _) = RecordingIo::over(memory_store()); + let store = open_over(io.clone(), &manual()).await; + let region_one = region(1); + let region_two = region(2); + + // Two appends share object 1, which is written but reported as + // failed: it exists and is not indexed. + let first = spawn_append_batch(&store, vec![entry(&store, region_one, "a1")]); + store.wait_for_admitted_appends(1).await.unwrap(); + let second = spawn_append_batch(&store, vec![entry(&store, region_two, "b1")]); + store.wait_for_admitted_appends(2).await.unwrap(); + io.fail_after_next_put.store(true, Ordering::Relaxed); + store.seal_open_batch().await.unwrap_err(); + for append in [first, second] { + timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + } + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + assert_eq!(0, latest(&store, region_one)); + + // Retried apart, the entries land in later objects rather than + // conflicting with object 1. + for (region_id, data, object_seq) in [(region_one, "a1", 2), (region_two, "b1", 3)] { + let retry = spawn_append_batch(&store, vec![entry(&store, region_id, data)]); + store + .wait_for_admitted_appends(object_seq as usize + 1) + .await + .unwrap(); + store.seal_open_batch().await.unwrap(); + let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq, 1))]), + response.last_entry_ids + ); + } + assert_eq!(vec![0, 1, 2, 3], object_seqs(io.as_ref()).await); + store.stop().await.unwrap(); + + // Object 2 extends the start object, so recovery leaves object 1 off + // the chain: only the retries replay. + let store = open_over(io, &manual()).await; + assert_eq!( + expected_entries(region_one, &[(id(2, 1), "a1")]), + read_entries(&store, region_one, 0).await + ); + assert_eq!( + expected_entries(region_two, &[(id(3, 1), "b1")]), + read_entries(&store, region_two, 0).await + ); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_floor_applies_while_a_create_is_in_flight() { + let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let region_one = region(1); + let region_two = region(2); + let pending = spawn_append_batch(&store, vec![entry(&store, region_one, "a1")]); + let (_, release) = next_create(&mut parked).await; + + // Object 1 is in flight and the open batch is empty: the floor moves + // the next sequence, which a failure of object 1 does not move back. + store + .obsolete(&provider(region_two), region_two, id(5, 1)) + .await + .unwrap(); + release.send(false).unwrap(); + timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err(); + + let write = spawn_append_batch(&store, vec![entry(&store, region_two, "b1")]); + next_create(&mut parked).await.1.send(true).unwrap(); + let response = timeout(WAIT, write).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_two, id(6, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 6], object_seqs(io.as_ref()).await); + } + + #[tokio::test] + async fn test_store_obsolete_queued_behind_stop_records_the_watermark() { + let (store, mut command_rx) = store_without_actor(); + let region_id = region(1); + let provider_one = provider(region_id); + + // The actor exits with the command still queued and never answers + // it: nothing is assigned an id any more, so the watermark alone holds. + let (result, ()) = tokio::join!(store.obsolete(&provider_one, region_id, 1), async { + let command = timeout(WAIT, command_rx.recv()).await.unwrap(); + assert!(matches!( + command, + Some(Command::Obsolete { entry_id: 1, .. }) + )); + }); + result.unwrap(); + assert_eq!( + Some(&1), + store.obsolete_entry_ids.lock().unwrap().get(®ion_id) + ); + + // The same holds for a call that finds the actor gone. + drop(command_rx); + store + .obsolete(&provider(region(2)), region(2), id(1, 1)) + .await + .unwrap(); + assert_eq!( + Some(&id(1, 1)), + store.obsolete_entry_ids.lock().unwrap().get(®ion(2)) + ); + } + + #[tokio::test] + async fn test_store_hooks_hold_and_fail_creates() { + let store = open(memory_store(), &eager()).await; + let region_id = region(1); + store.hold_creates(); + let held = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]); + store.wait_for_admitted_appends(1).await.unwrap(); + round_trip_actor(&store).await; + assert!(!held.is_finished()); + assert_eq!(vec![0], object_seqs(store.io.as_ref()).await); + store.release_creates(); + timeout(WAIT, held).await.unwrap().unwrap().unwrap(); + assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await); + + store.fail_creates(); + let error = append(&store, region_id, "a2").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); + + store.begin_stop(); + assert_stopped(&append(&store, region_id, "a3").await.unwrap_err()); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_create_held_when_the_store_is_dropped_never_runs() { + let (io, _) = RecordingIo::over(memory_store()); + let store = open_over(io.clone(), &eager()).await; + store.hold_creates(); + let append = spawn_append_batch(&store, vec![entry(&store, region(1), "a1")]); + store.wait_for_admitted_appends(1).await.unwrap(); + append.abort(); + let _ = append.await; + // A command sender outlives the store, so the actor keeps running and + // the parked create completes on the closed hold channel instead of + // being dropped with the actor. + let command_tx = store.command_tx.clone(); + let terminal_error = store.terminal_error.clone(); + drop(store); + timeout(WAIT, async { + while terminal(&terminal_error).is_none() && object_seqs(io.as_ref()).await == [0] { + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .unwrap(); + assert_eq!(vec![0], object_seqs(io.as_ref()).await); + let error = terminal(&terminal_error).unwrap(); + assert!( + matches!(error.as_ref(), Error::ObjectStoreWalStopped { .. }), + "unexpected error: {error:?}" + ); + drop(command_tx); + } + + #[tokio::test] + async fn test_store_rollback_and_poison_drop_the_open_batch() { + let object_store = memory_store(); + let (store, io, mut parked) = open_parking_creates(object_store.clone(), &manual()).await; + let region_id = region(1); + let spawn_seal = || { + let store = store.clone(); + tokio::spawn(async move { store.seal_open_batch().await }) + }; + + // Object 1 is in flight and a2 waits in the open batch under + // sequence 2 when the create fails: both fail. + let sealed = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]); + store.wait_for_admitted_appends(1).await.unwrap(); + let seal = spawn_seal(); + let (_, release) = next_create(&mut parked).await; + let open = spawn_append_batch(&store, vec![entry(&store, region_id, "a2")]); + store.wait_for_admitted_appends(2).await.unwrap(); + release.send(false).unwrap(); + for append in [sealed, open] { + let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + } + timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err(); + + // The open batch had not taken its sequence, so the next entry takes + // the first id of sequence 2 again, alone. + let retry = spawn_append_batch(&store, vec![entry(&store, region_id, "b1")]); + store.wait_for_admitted_appends(3).await.unwrap(); + let seal = spawn_seal(); + next_create(&mut parked).await.1.send(true).unwrap(); + timeout(WAIT, seal).await.unwrap().unwrap().unwrap(); + let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + assert_eq!( + expected_entries(region_id, &[(id(2, 1), "b1")]), + read_entries(&store, region_id, 0).await + ); + + // Object 3 conflicts while c2 waits in the open batch: both fail. + let sealed = spawn_append_batch(&store, vec![entry(&store, region_id, "c1")]); + store.wait_for_admitted_appends(4).await.unwrap(); + let seal = spawn_seal(); + let (_, release) = next_create(&mut parked).await; + let open = spawn_append_batch(&store, vec![entry(&store, region_id, "c2")]); + store.wait_for_admitted_appends(5).await.unwrap(); + put_foreign(&object_store, 3, 1).await; + release.send(true).unwrap(); + for append in [sealed, open] { + let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }), + "unexpected error: {error:?}" + ); + } + timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err(); + assert_eq!(vec![0, 2, 3], object_seqs(io.as_ref()).await); + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_poisons_once_the_last_sequence_is_taken() { + let object_store = memory_store(); + let last = OBJECT_SEQ_LIMIT - 1; + // The start object of the store takes the sequence before the last. + put_object(&object_store, last - 2, region(1), &[id(last - 2, 1)]).await; + let store = open(object_store, &eager()).await; + let next = entry(&store, region(1), "next"); + + // The batch at the last sequence is created and acknowledged, and the + // store is poisoned as soon as it takes that sequence. + let response = append(&store, region(1), "last").await.unwrap(); + assert_eq!( + HashMap::from([(region(1), id(last, 1))]), + response.last_entry_ids + ); + for error in [ + store.latest_entry_id(&provider(region(1))).unwrap_err(), + store.append_batch(vec![next]).await.unwrap_err(), + ] { + assert!( + matches!( + unwrap_shared(&error), + Error::WalObjectSequenceExhausted { last_object_seq, .. } if *last_object_seq == last + ), + "unexpected error: {error:?}" + ); + } + store.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_store_holds_appends_back_while_the_sealed_batches_are_full() { + let (store, _, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let region_id = region(1); + let admitted = || *store.admitted_appends.borrow(); + + // Every create is parked: once no two more sealed batches fit, + // further appends wait in the channel and are not admitted. + let mut appends = spawn_appends(&store, region_id, MAX_SEALED_BATCHES - 1).await; + for index in 0..3 { + let entries = vec![entry(&store, region_id, &format!("w{index}"))]; + appends.push(spawn_append_batch(&store, entries)); + } + let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await; + round_trip_actor(&store).await; + assert_eq!(MAX_SEALED_BATCHES - 1, admitted()); + + // A durable object frees a place, and one more append is admitted. + releases.remove(&1).unwrap().send(true).unwrap(); + timeout(WAIT, store.wait_for_admitted_appends(MAX_SEALED_BATCHES)) + .await + .unwrap() + .unwrap(); + round_trip_actor(&store).await; + assert_eq!(MAX_SEALED_BATCHES, admitted()); + + // Stop is handled while appends are held back: the batches whose create + // has not started and the appends still waiting learn of the stop. + let stop = begin_spawned_stop(&store).await; + for release in releases.into_values() { + release.send(true).unwrap(); + } + next_create(&mut parked).await.1.send(true).unwrap(); + timeout(WAIT, stop).await.unwrap().unwrap().unwrap(); + let mut acknowledged = 0; + for append in appends { + match timeout(WAIT, append).await.unwrap().unwrap() { + Ok(_) => acknowledged += 1, + Err(error) => assert_stopped(&error), + } + } + assert_eq!(MAX_IN_FLIGHT_CREATES + 1, acknowledged); + } + + /// Starts `stop` and waits until it has set the stopped flag. + async fn begin_spawned_stop( + store: &Arc, + ) -> tokio::task::JoinHandle> { + let stop = { + let store = store.clone(); + tokio::spawn(async move { store.stop().await }) + }; + while !store.stopped.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + stop + } + + #[tokio::test] + async fn test_store_stop_during_failed_flush_reports_stopped() { + let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let pending = spawn_append_batch(&store, vec![entry(&store, region(1), "a1")]); + let (_, release) = next_create(&mut parked).await; + let stop = begin_spawned_stop(&store).await; + assert!(!stop.is_finished()); + + // The create fails after stop began: its waiter learns of the stop. + release.send(false).unwrap(); + timeout(WAIT, stop).await.unwrap().unwrap().unwrap(); + assert_stopped(&timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err()); + assert_eq!(vec![0], object_seqs(io.as_ref()).await); + assert_stopped(&append(&store, region(1), "a2").await.unwrap_err()); + } + + #[tokio::test] + async fn test_store_stop_with_several_creates_in_flight() { + let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await; + let region_id = region(1); + // Four creates are in flight, the fifth batch waits for a slot. + let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 1).await; + let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await; + assert!(parked.try_recv().is_err()); + + let stop = begin_spawned_stop(&store).await; + assert!(!stop.is_finished()); + // Nothing is assigned an id after stop began, so a watermark needs no + // floor even while creates are in flight. + store + .obsolete(&provider(region(2)), region(2), id(9, 1)) + .await + .unwrap(); + assert_eq!( + Some(&id(9, 1)), + store.obsolete_entry_ids.lock().unwrap().get(®ion(2)) + ); + + // The creates in flight run to completion in any order and are + // acknowledged; the batch that never started learns of the stop. + for object_seq in [3, 1, 4, 2] { + releases.remove(&object_seq).unwrap().send(true).unwrap(); + } + timeout(WAIT, stop).await.unwrap().unwrap().unwrap(); + let mut appends = appends.into_iter(); + for object_seq in 1..=MAX_IN_FLIGHT_CREATES as u64 { + let response = timeout(WAIT, appends.next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq, 1))]), + response.last_entry_ids + ); + } + let error = timeout(WAIT, appends.next().unwrap()) + .await + .unwrap() + .unwrap() + .unwrap_err(); + assert_stopped(&error); + // No create was started after stop began. + assert!(parked.try_recv().is_err()); + assert_eq!(vec![0, 1, 2, 3, 4], object_seqs(io.as_ref()).await); + assert_eq!(id(4, 1), latest(&store, region_id)); + timeout(WAIT, store.command_tx.closed()).await.unwrap(); + } + + #[tokio::test] + async fn test_store_stop_during_conflicting_flush_reports_stopped_and_poisons() { + let object_store = memory_store(); + let (store, _, mut parked) = open_parking_creates(object_store.clone(), &eager()).await; + let region_id = region(1); + let pending = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]); + let (_, release) = next_create(&mut parked).await; + let stop = begin_spawned_stop(&store).await; + // Another writer of the same epoch takes the sequence before the + // create runs. + put_foreign(&object_store, 1, 1).await; + + release.send(true).unwrap(); + timeout(WAIT, stop).await.unwrap().unwrap().unwrap(); + assert_stopped(&timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err()); + for error in [ + store.latest_entry_id(&provider(region_id)).unwrap_err(), + store + .read(&provider(region_id), 1, None) + .await + .err() + .unwrap(), + ] { + assert!( + matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }), + "unexpected error: {error:?}" + ); + } + } + + fn injected_failure(operation: &'static str, path: String) -> Result { + Err( + object_store::Error::new(object_store::ErrorKind::Unexpected, "injected failure") + .set_temporary(), + ) + .context(WalObjectStoreSnafu { operation, path }) + } + fn chain_header(object_seq: u64, epoch: u64, prev: Option<(u64, u64)>) -> Header { Header { object_seq, @@ -2248,15 +4232,6 @@ mod tests { assert_eq!(vec![0], object_seqs(&io).await); } - async fn object_seqs(io: &ObjectStoreIo) -> Vec { - io.list() - .await - .unwrap() - .into_iter() - .map(|object| object.object_seq) - .collect() - } - #[tokio::test] async fn test_store_open_fails_when_its_start_object_lands_with_an_unknown_outcome() { let object_store = memory_store(); @@ -2368,9 +4343,16 @@ mod tests { } } + /// Object access whose whole-object reads, or conditional creates, park + /// until the test releases them, counting how many are in flight. A + /// release of false fails the operation with a transient error before it + /// reaches the object store. struct ParkedIo { inner: ObjectStoreIo, parked: mpsc::UnboundedSender<(u64, oneshot::Sender)>, + park_creates: bool, + /// Whether creates park yet; set once the store opened. + creates_parked: AtomicBool, in_flight: AtomicUsize, max_in_flight: AtomicUsize, } @@ -2381,43 +4363,68 @@ mod tests { ) -> ( Arc, mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + ) { + Self::new(object_store, false) + } + + fn parking_creates( + object_store: ObjectStore, + ) -> ( + Arc, + mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + ) { + Self::new(object_store, true) + } + + fn new( + object_store: ObjectStore, + park_creates: bool, + ) -> ( + Arc, + mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, ) { let (parked, parked_rx) = mpsc::unbounded_channel(); ( Arc::new(Self { inner: ObjectStoreIo::new(object_store, PREFIX).unwrap(), parked, + park_creates, + creates_parked: AtomicBool::new(false), in_flight: AtomicUsize::new(0), max_in_flight: AtomicUsize::new(0), }), parked_rx, ) } + + /// Parks until the test releases `object_seq` and returns whether the + /// operation proceeds. + async fn park(&self, object_seq: u64) -> bool { + let in_flight = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; + self.max_in_flight.fetch_max(in_flight, Ordering::SeqCst); + let (release, released) = oneshot::channel(); + self.parked.send((object_seq, release)).unwrap(); + let proceed = released.await.unwrap(); + self.in_flight.fetch_sub(1, Ordering::SeqCst); + proceed + } } #[async_trait::async_trait] impl WalObjectIo for ParkedIo { async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result { + if self.park_creates + && self.creates_parked.load(Ordering::SeqCst) + && !self.park(object_seq).await + { + return injected_failure("write", self.object_path(object_seq)); + } self.inner.put_if_absent(object_seq, content).await } async fn get(&self, object_seq: u64) -> Result { - let in_flight = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; - self.max_in_flight.fetch_max(in_flight, Ordering::SeqCst); - let (release, released) = oneshot::channel(); - self.parked.send((object_seq, release)).unwrap(); - let success = released.await.unwrap(); - self.in_flight.fetch_sub(1, Ordering::SeqCst); - if !success { - return Err(object_store::Error::new( - object_store::ErrorKind::Unexpected, - "injected failure", - ) - .set_temporary()) - .context(WalObjectStoreSnafu { - operation: "read", - path: self.object_path(object_seq), - }); + if !self.park_creates && !self.park(object_seq).await { + return injected_failure("read", self.object_path(object_seq)); } self.inner.get(object_seq).await } @@ -2437,10 +4444,13 @@ mod tests { type RangeReads = Arc>>; - /// Object access that records every range read as (sequence, offset, length). + /// Object access that records every range read as (sequence, offset, + /// length) and, on request, reports the next conditional create as failed + /// after it wrote the object. struct RecordingIo { inner: ObjectStoreIo, reads: RangeReads, + fail_after_next_put: AtomicBool, } impl RecordingIo { @@ -2449,6 +4459,7 @@ mod tests { let io = Self { inner: ObjectStoreIo::new(object_store, PREFIX).unwrap(), reads: reads.clone(), + fail_after_next_put: AtomicBool::new(false), }; (Arc::new(io), reads) } @@ -2457,7 +4468,11 @@ mod tests { #[async_trait::async_trait] impl WalObjectIo for RecordingIo { async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result { - self.inner.put_if_absent(object_seq, content).await + let result = self.inner.put_if_absent(object_seq, content).await?; + if self.fail_after_next_put.swap(false, Ordering::Relaxed) { + return injected_failure("write", self.object_path(object_seq)); + } + Ok(result) } async fn get(&self, object_seq: u64) -> Result {