From 194bc2fb3c48ea1c55cff8fdd5886daa8b60d8ad Mon Sep 17 00:00:00 2001 From: jeremyhi Date: Thu, 24 Sep 2026 11:07:29 +0000 Subject: [PATCH] feat(log-store): add the object store WAL durable write path (#9320) * feat(log-store): add the object store WAL durable write path ObjectStoreLogStore can now write. append_batch admits entries into one open batch and assigns object-sequence-major entry ids at admission. The batch is sealed by size, by the flush interval, or before a region would run past the position range, and sealed batches are uploaded with at most four conditional creates in flight, started in sequence order. Created objects are indexed in sequence order, and an append is acknowledged only once its object is durable and indexed. A transient create failure rolls back the failed batch and every later batch unless a later object is already durable, in which case the store poisons itself with a history-gap error. A conflicting object, an encoding or catalog error, or taking the last representable sequence poisons the store. obsolete now goes through the actor and raises the sequence floor together with the watermark, refusing with a retryable error while the next sequence is not settled. stop drops the open batch and the batches whose create has not started, and lets creates in flight finish. Only the durable acknowledgement mode exists. A testing feature exposes hooks to wait for admissions, seal the open batch, and hold or fail creates. Signed-off-by: jeremyhi * test(log-store): synchronize object store WAL write tests with the actor Tests that assert nothing happened round-trip a command through the actor instead of yielding the test task, the conflict test waits for the object it expects, and the obsolete-behind-stop test holds the command channel itself so the unanswered command is deterministic. Signed-off-by: jeremyhi * fix(log-store): bound admitted WAL appends and never reuse a failed sequence A create that reports an error may still have written its object, so the sequences of failed batches are no longer handed out again: the next batch keeps the sequence after the last sealed one and retries are assigned new ids. A retry batched differently can no longer conflict with that object. As the next sequence never moves back, the sequence floor of obsolete only waits for an open batch that has handed out ids. Appends now arrive on their own bounded channel, which the actor stops reading while MAX_SEALED_BATCHES batches wait to become durable, so a stalled object store holds callers back instead of growing the backlog, while stop and obsolete are still handled. Signed-off-by: jeremyhi * fix(log-store): reserve room for two seals before admitting a WAL append One admission can seal the open batch before an append that would exhaust its positions and then the append's own batch, so the actor takes an append only while two more sealed batches fit under MAX_SEALED_BATCHES. Signed-off-by: jeremyhi * feat(log-store): chain object store WAL objects and recover the latest chain A conditional create can fail with an unknown outcome while its object is stored, or still lands later, and Mito reuses the row sequences of a failed append. Replaying such an object next to a later acknowledged one can let the unacknowledged rows win after a restart. Format version 2 gives every object header its writer's epoch, a link to the object it extends (sequence and writer instance) and a header CRC32, and allows objects without segments. Recovery replays only the chain ending at the complete object with the largest epoch and sequence: a link holds when its predecessor is present with the recorded writer instance, or is missing below every present object. Objects off the chain are orphans that are never replayed but keep their sequences. Each open writes an empty object that starts an epoch above every present object, linked to the recovered tip, before it accepts writes, so a late object of an earlier instance never ends the chain. A start object that meets an object of an earlier epoch moves to the next sequence; one of an equal or later epoch fails the open. Signed-off-by: jeremyhi * fix(log-store): break the WAL chain at every missing predecessor Accepting a missing predecessor below every present object lets a late object that lands below the chain change which links hold. Nothing collects objects yet, so a missing predecessor now always breaks the link, and recovery fails when objects are present but none completes a chain. Signed-off-by: jeremyhi * fix(log-store): keep the object store WAL format at version 1 The object store WAL has not been enabled anywhere, so no object in the previous layout exists and the chained header can stay version 1. Signed-off-by: jeremyhi * refactor(log-store): link object store WAL objects by epoch Only one store instance writes under an epoch, since every open starts an epoch above every present object and a start object that meets the same or a later epoch fails the open. The epoch therefore identifies the instance, and the random writer instance id is dropped from the header. A link now records the sequence and the epoch of the object it extends, and holds when the predecessor carries that epoch. The header shrinks to 46 bytes, and the store logs its epoch when it opens. Signed-off-by: jeremyhi * fix(log-store): fail an open whose start object is already present Without a random writer instance, two opens that recover the same objects encode byte-identical start objects, and a conditional create treats the same bytes as its own retry. A start object that is already present therefore fails the open with a retryable error instead of letting both opens claim the epoch; the next open counts the object and starts a later epoch. Signed-off-by: jeremyhi * fix(log-store): derive the WAL epoch from the claimed start sequence Two opens can recover different listings when a late object lands across a sequence gap between them, pick the same largest epoch plus one, and both create their start objects under different sequences. The epoch of an instance is now one above the sequence its start object claims, so a successful create decides the epoch and no two instances share one. It stays above every epoch recovery listed, and an object that carries an epoch above the next sequence fails the open. Signed-off-by: jeremyhi * fix(log-store): never run a held WAL create after the store is dropped The create test hook ignored the closed hold channel, so a create parked when the store was dropped could still run. It now returns without creating. Drop the per-admission bookkeeping of issued entry ids, which nothing reads, and move the parked I/O documentation to its helper. Signed-off-by: jeremyhi * test(log-store): cover a WAL start object stored with an unknown outcome Add a fault-injection test in which the create of the start object stores the object but reports an error: the open fails without moving to another sequence, and the next open counts the stored object and claims a later epoch. Rename the test helper that writes a whole object from a given header to put_object_with_header, and drop a needless clone. Signed-off-by: jeremyhi * test(log-store): make the held-create drop test deterministic Keep the actor running while the store drops the hold sender, so the parked create always completes on the closed channel instead of racing the actor's exit. Drop a comment that restates epoch_of. Signed-off-by: jeremyhi --------- Signed-off-by: jeremyhi --- src/log-store/Cargo.toml | 3 + src/log-store/src/error.rs | 51 +- src/log-store/src/object_store_wal/store.rs | 2191 ++++++++++++++++++- 3 files changed, 2149 insertions(+), 96 deletions(-) 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 {