Skip to main content

log_store/object_store_wal/
store.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Object store WAL construction, recovery, writes and region reads.
16
17use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque};
18use std::fmt;
19use std::ops::Range;
20use std::sync::atomic::{AtomicBool, Ordering};
21use std::sync::{Arc, Mutex, PoisonError, RwLock};
22use std::time::{Duration, Instant};
23
24use async_stream::try_stream;
25use bytes::Bytes;
26use common_telemetry::info;
27use common_wal::config::object_store::{AckMode, ObjectStoreWalConfig};
28use futures::future::BoxFuture;
29use futures::stream::FuturesUnordered;
30use futures::{StreamExt, TryStreamExt};
31use object_store::ObjectStore;
32use snafu::{IntoError, OptionExt, ResultExt, ensure};
33use store_api::logstore::entry::{Entry, NaiveEntry};
34use store_api::logstore::provider::{ObjectStoreProvider, Provider};
35use store_api::logstore::{AppendBatchResponse, EntryId, LogStore, SendableEntryStream, WalIndex};
36use store_api::storage::RegionId;
37#[cfg(any(test, feature = "testing"))]
38use tokio::sync::watch;
39use tokio::sync::{mpsc, oneshot};
40use tokio::time::MissedTickBehavior;
41
42use crate::error::{
43    CorruptedWalObjectSnafu, Error, IncompleteWalEntrySnafu, InvalidProviderSnafu,
44    InvalidWalObjectSnafu, InvalidWalObjectStoreSnafu, MismatchedWalPrefixSnafu,
45    MismatchedWalRegionSnafu, ObjectStoreWalSnafu, ObjectStoreWalStoppedSnafu, Result,
46    StaleWalObjectSnafu, UnconfirmedWalEpochStartSnafu, WalObjectSequenceExhaustedSnafu,
47    WalObjectSequenceUnsettledSnafu,
48};
49use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, OpenBatch, entry_id, sequence_floor};
50use crate::object_store_wal::catalog::ObjectCatalog;
51use crate::object_store_wal::format::{
52    ChainLink, EncodedObject, FixedTrailer, FooterEntry, HEADER_LEN, Header, MIN_OBJECT_LEN,
53    Record, TRAILER_LEN, decode_footer, decode_header, decode_segment, decode_trailer,
54    encode_object, footer_range, verify_segment_ranges,
55};
56use crate::object_store_wal::io::{ListedObject, ObjectStoreIo, PutResult};
57
58const COMMAND_BUFFER: usize = 1024;
59/// Appends queued for the actor; a caller beyond them waits to send.
60const APPEND_BUFFER: usize = 16;
61const MIN_FLUSH_INTERVAL: Duration = Duration::from_millis(10);
62const MAX_IN_FLIGHT_CREATES: usize = 4;
63/// Delay before a create that failed transiently is attempted again in the
64/// `enqueued` acknowledgement mode, where no caller is left to retry it.
65const CREATE_RETRY_DELAY: Duration = Duration::from_millis(100);
66/// Number of sealed batches that may wait to be created or indexed. One
67/// admission can seal two batches, the open batch before an append that would
68/// exhaust its positions and the append's own batch, so the actor takes an
69/// append from its channel only while two more fit: a stalled object store
70/// holds callers back instead of growing the backlog.
71const MAX_SEALED_BATCHES: usize = 2 * MAX_IN_FLIGHT_CREATES;
72/// Number of objects whose footers recovery fetches at a time.
73const RECOVERY_CONCURRENCY: usize = 8;
74/// Bytes recovery reads from the end of an object in one request. The window
75/// holds the trailer and the footer of an object with up to 1364 regions of
76/// 48 bytes each, so a second request for the footer is rare.
77const RECOVERY_TAIL_WINDOW: usize = 64 * 1024;
78
79/// A log store over the immutable WAL objects under one prefix.
80///
81/// Appends are admitted into an open batch. A background actor seals the batch
82/// when it reaches the size limit or the flush interval elapses and creates
83/// the object under the next sequence while it keeps admitting entries into
84/// the next batch; up to [`MAX_IN_FLIGHT_CREATES`] creates run at a time.
85/// Objects are indexed in the catalog in sequence order. In the `durable`
86/// acknowledgement mode an append returns once its object is durable and
87/// indexed, so an acknowledged entry never has a missing predecessor; in the
88/// `enqueued` mode it returns on admission and the object is created in the
89/// background. While [`MAX_SEALED_BATCHES`] batches
90/// wait to become durable no append is admitted.
91pub struct ObjectStoreLogStore {
92    prefix: String,
93    ack_mode: AckMode,
94    io: Arc<dyn WalObjectIo>,
95    catalog: Arc<RwLock<ObjectCatalog>>,
96    obsolete_entry_ids: ObsoleteEntryIds,
97    terminal_error: TerminalError,
98    stopped: Arc<AtomicBool>,
99    command_tx: mpsc::Sender<Command>,
100    append_tx: mpsc::Sender<QueuedAppend>,
101    #[cfg(any(test, feature = "testing"))]
102    admitted_appends: watch::Receiver<usize>,
103    #[cfg(any(test, feature = "testing"))]
104    creates_held: watch::Sender<bool>,
105    #[cfg(any(test, feature = "testing"))]
106    creates_fail: Arc<AtomicBool>,
107    #[cfg(any(test, feature = "testing"))]
108    next_create_fails_after_write: Arc<AtomicBool>,
109    #[cfg(any(test, feature = "testing"))]
110    parked_creates: watch::Receiver<usize>,
111    #[cfg(any(test, feature = "testing"))]
112    durability_waits: watch::Receiver<usize>,
113}
114
115type ObsoleteEntryIds = Arc<Mutex<HashMap<RegionId, EntryId>>>;
116
117type TerminalError = Arc<Mutex<Option<Arc<Error>>>>;
118
119impl fmt::Debug for ObjectStoreLogStore {
120    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
121        f.debug_struct("ObjectStoreLogStore")
122            .field("prefix", &self.prefix)
123            .field("ack_mode", &self.ack_mode)
124            .finish_non_exhaustive()
125    }
126}
127
128impl ObjectStoreLogStore {
129    /// Builds the store under the node and generation prefix derived from `config`,
130    /// recovering the catalog from the objects that already exist and writing
131    /// the object that starts the epoch of this instance. Recovery fails on the
132    /// first corrupted or conflicting object.
133    pub async fn try_new(
134        object_store: ObjectStore,
135        config: &ObjectStoreWalConfig,
136        node_id: u64,
137        generation: u64,
138    ) -> Result<Arc<Self>> {
139        let prefix = config.node_prefix(node_id, generation);
140        let io = ObjectStoreIo::new(object_store, &prefix)?;
141        Self::open(Arc::new(io), config, prefix).await
142    }
143
144    async fn open(
145        io: Arc<dyn WalObjectIo>,
146        config: &ObjectStoreWalConfig,
147        prefix: String,
148    ) -> Result<Arc<Self>> {
149        ensure!(
150            config.flush_interval >= MIN_FLUSH_INTERVAL,
151            InvalidWalObjectStoreSnafu {
152                reason: format!(
153                    "flush interval {:?} is shorter than {MIN_FLUSH_INTERVAL:?}",
154                    config.flush_interval
155                ),
156            }
157        );
158
159        let max_batch_bytes = positive_bytes(config.max_batch_bytes.as_bytes(), "max batch bytes")?;
160        let max_unpersisted_bytes = positive_bytes(
161            config.max_unpersisted_bytes.as_bytes(),
162            "max unpersisted bytes",
163        )?;
164        ensure!(
165            config.max_unpersisted_age > Duration::ZERO,
166            InvalidWalObjectStoreSnafu {
167                reason: "max unpersisted age is zero",
168            }
169        );
170        let Recovered {
171            mut catalog,
172            next_object_seq,
173            durable_entry_ids,
174            tip,
175            max_epoch,
176        } = recover(io.as_ref()).await?;
177        // Every epoch is one above the sequence of the start object its
178        // instance created, and every object is at or above its start object,
179        // so an epoch above the next sequence names no instance that ran.
180        ensure!(
181            max_epoch <= next_object_seq,
182            CorruptedWalObjectSnafu {
183                reason: format!(
184                    "an object carries epoch {max_epoch}, above the next sequence {next_object_seq}"
185                ),
186            }
187        );
188        let start = start_epoch(io.as_ref(), next_object_seq, tip).await?;
189        let epoch = start.epoch;
190        // The epoch identifies this instance in every object it writes.
191        info!(
192            "Opened object store WAL under {prefix} at epoch {epoch}, start object {}",
193            start.object_seq
194        );
195        catalog
196            .insert_object(start.object_seq, Vec::new())
197            .with_context(|_| InvalidWalObjectSnafu {
198                path: io.object_path(start.object_seq),
199            })?;
200        let catalog = Arc::new(RwLock::new(catalog));
201        let obsolete_entry_ids = ObsoleteEntryIds::default();
202        let terminal_error = TerminalError::default();
203        let stopped = Arc::new(AtomicBool::new(false));
204        let (command_tx, command_rx) = mpsc::channel(COMMAND_BUFFER);
205        let (append_tx, append_rx) = mpsc::channel(APPEND_BUFFER);
206        #[cfg(any(test, feature = "testing"))]
207        let (admitted_appends_tx, admitted_appends_rx) = watch::channel(0);
208        #[cfg(any(test, feature = "testing"))]
209        let (creates_held_tx, creates_held_rx) = watch::channel(false);
210        #[cfg(any(test, feature = "testing"))]
211        let creates_fail = Arc::new(AtomicBool::new(false));
212        #[cfg(any(test, feature = "testing"))]
213        let next_create_fails_after_write = Arc::new(AtomicBool::new(false));
214        #[cfg(any(test, feature = "testing"))]
215        let (parked_creates_tx, parked_creates_rx) = watch::channel(0);
216        #[cfg(any(test, feature = "testing"))]
217        let (durability_waits_tx, durability_waits_rx) = watch::channel(0);
218
219        let actor = Actor {
220            io: io.clone(),
221            catalog: catalog.clone(),
222            obsolete_entry_ids: obsolete_entry_ids.clone(),
223            terminal_error: terminal_error.clone(),
224            stopped: stopped.clone(),
225            command_rx,
226            append_rx,
227            ack_mode: config.ack_mode,
228            max_unpersisted_bytes,
229            max_unpersisted_age: config.max_unpersisted_age,
230            open_batch: OpenBatch::new(max_batch_bytes),
231            issued_entry_ids: durable_entry_ids,
232            pending: Vec::new(),
233            sealed: VecDeque::new(),
234            creates: FuturesUnordered::new(),
235            stalled: None,
236            durable_waiters: Vec::new(),
237            stop: Vec::new(),
238            stop_error: None,
239            next_object_seq: start
240                .object_seq
241                .checked_add(1)
242                .filter(|next_object_seq| *next_object_seq < OBJECT_SEQ_LIMIT),
243            epoch,
244            last_indexed: start,
245            flush_interval: config.flush_interval,
246            #[cfg(any(test, feature = "testing"))]
247            admitted_appends: admitted_appends_tx,
248            #[cfg(any(test, feature = "testing"))]
249            creates_held: creates_held_rx,
250            #[cfg(any(test, feature = "testing"))]
251            creates_fail: creates_fail.clone(),
252            #[cfg(any(test, feature = "testing"))]
253            next_create_fails_after_write: next_create_fails_after_write.clone(),
254            #[cfg(any(test, feature = "testing"))]
255            parked_creates: Arc::new(parked_creates_tx),
256            #[cfg(any(test, feature = "testing"))]
257            durability_waits: durability_waits_tx,
258        };
259        common_runtime::spawn_global(actor.run());
260        Ok(Arc::new(Self {
261            prefix,
262            ack_mode: config.ack_mode,
263            io,
264            catalog,
265            obsolete_entry_ids,
266            terminal_error,
267            stopped,
268            command_tx,
269            append_tx,
270            #[cfg(any(test, feature = "testing"))]
271            admitted_appends: admitted_appends_rx,
272            #[cfg(any(test, feature = "testing"))]
273            creates_held: creates_held_tx,
274            #[cfg(any(test, feature = "testing"))]
275            creates_fail,
276            #[cfg(any(test, feature = "testing"))]
277            next_create_fails_after_write,
278            #[cfg(any(test, feature = "testing"))]
279            parked_creates: parked_creates_rx,
280            #[cfg(any(test, feature = "testing"))]
281            durability_waits: durability_waits_rx,
282        }))
283    }
284
285    /// Returns the largest entry id of the provider's region whose object is
286    /// durable and indexed, or zero for a region without such entries. In the
287    /// `enqueued` acknowledgement mode an entry id that an append returned
288    /// stays above this value until the object holding it is created.
289    pub(crate) fn durable_entry_id(&self, provider: &Provider) -> Result<EntryId> {
290        self.check_terminal()?;
291        let region_id = self.region_of(provider)?;
292        Ok(self.durable_entry_id_of(region_id))
293    }
294
295    fn durable_entry_id_of(&self, region_id: RegionId) -> EntryId {
296        let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner);
297        catalog.region_max_entry_id(region_id).unwrap_or(0)
298    }
299
300    /// Returns the region of `provider`, which must select this store's prefix.
301    fn region_of(&self, provider: &Provider) -> Result<RegionId> {
302        let provider =
303            provider
304                .as_object_store_provider()
305                .with_context(|| InvalidProviderSnafu {
306                    expected: ObjectStoreProvider::type_name(),
307                    actual: provider.type_name(),
308                })?;
309        ensure!(
310            provider.prefix == self.prefix,
311            MismatchedWalPrefixSnafu {
312                expected: self.prefix.clone(),
313                actual: provider.prefix.clone(),
314            }
315        );
316        Ok(provider.region_id)
317    }
318
319    fn check_region(&self, provider: &Provider, region_id: RegionId) -> Result<()> {
320        self.check_terminal()?;
321        let provider_region = self.region_of(provider)?;
322        ensure!(
323            provider_region == region_id,
324            MismatchedWalRegionSnafu {
325                region_id,
326                reason: format!("provider belongs to region {provider_region}"),
327            }
328        );
329        Ok(())
330    }
331
332    fn check_terminal(&self) -> Result<()> {
333        match terminal(&self.terminal_error) {
334            Some(error) => Err(shared(&error)),
335            None => Ok(()),
336        }
337    }
338}
339
340/// The transient object store error of a create that a testing hook fails.
341#[cfg(any(test, feature = "testing"))]
342fn injected_create_failure(io: &dyn WalObjectIo, object_seq: u64) -> Result<PutResult> {
343    let error = object_store::Error::new(
344        object_store::ErrorKind::Unexpected,
345        "injected create failure",
346    )
347    .set_temporary();
348    Err(error).context(crate::error::WalObjectStoreSnafu {
349        operation: "write",
350        path: io.object_path(object_seq),
351    })
352}
353
354/// Records the obsolete watermark of `region_id`, which never moves down.
355fn record_obsolete(obsolete_entry_ids: &ObsoleteEntryIds, region_id: RegionId, entry_id: EntryId) {
356    obsolete_entry_ids
357        .lock()
358        .unwrap_or_else(PoisonError::into_inner)
359        .entry(region_id)
360        .and_modify(|current| *current = (*current).max(entry_id))
361        .or_insert(entry_id);
362}
363
364fn positive_bytes(bytes: u64, name: &str) -> Result<usize> {
365    usize::try_from(bytes)
366        .ok()
367        .filter(|bytes| *bytes > 0)
368        .with_context(|| InvalidWalObjectStoreSnafu {
369            reason: format!("{name} {bytes} is zero or too large"),
370        })
371}
372
373#[cfg(any(test, feature = "testing"))]
374impl ObjectStoreLogStore {
375    /// Waits until the actor has admitted at least `expected` append calls
376    /// since the store was built.
377    pub async fn wait_for_admitted_appends(&self, expected: usize) -> Result<()> {
378        self.admitted_appends
379            .clone()
380            .wait_for(|count| *count >= expected)
381            .await
382            .ok()
383            .map(|_| ())
384            .context(ObjectStoreWalStoppedSnafu)
385    }
386
387    /// Seals the open batch regardless of its size and age and returns once
388    /// its object is durable and indexed, or with the error that failed it.
389    pub async fn seal_open_batch(&self) -> Result<()> {
390        ensure!(
391            !self.stopped.load(Ordering::Acquire),
392            ObjectStoreWalStoppedSnafu
393        );
394        let (response_tx, response_rx) = oneshot::channel();
395        self.command_tx
396            .send(Command::Seal {
397                response: response_tx,
398            })
399            .await
400            .ok()
401            .context(ObjectStoreWalStoppedSnafu)?;
402        response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)?
403    }
404
405    /// Parks every conditional create that starts from now on until
406    /// [`release_creates`](Self::release_creates), so a test can observe
407    /// entries that are admitted but not durable. A create that is parked when
408    /// the store is dropped never runs.
409    pub fn hold_creates(&self) {
410        self.creates_held.send_replace(true);
411    }
412
413    /// Lets the creates parked by [`hold_creates`](Self::hold_creates) run.
414    pub fn release_creates(&self) {
415        self.creates_held.send_replace(false);
416    }
417
418    /// Makes every create that runs from now on fail with a transient object
419    /// store error instead of writing.
420    pub fn fail_creates(&self) {
421        self.creates_fail.store(true, Ordering::Release);
422    }
423
424    /// Makes the next create that writes its object report a transient object
425    /// store error afterwards, so the object exists although its create
426    /// failed.
427    pub fn fail_next_create_after_write(&self) {
428        self.next_create_fails_after_write
429            .store(true, Ordering::Release);
430    }
431
432    /// Waits until at least `expected` creates have been parked by
433    /// [`hold_creates`](Self::hold_creates) since the store was built.
434    pub async fn wait_for_parked_creates(&self, expected: usize) -> Result<()> {
435        self.parked_creates
436            .clone()
437            .wait_for(|count| *count >= expected)
438            .await
439            .ok()
440            .map(|_| ())
441            .context(ObjectStoreWalStoppedSnafu)
442    }
443
444    /// Waits until at least `expected` calls of
445    /// [`wait_durable`](LogStore::wait_durable) have had to wait for an entry
446    /// that is not durable since the store was built.
447    pub async fn wait_for_durability_waits(&self, expected: usize) -> Result<()> {
448        self.durability_waits
449            .clone()
450            .wait_for(|count| *count >= expected)
451            .await
452            .ok()
453            .map(|_| ())
454            .context(ObjectStoreWalStoppedSnafu)
455    }
456
457    /// Ends the actor the way a crash of the process would: nothing more is
458    /// written, creates that have not completed are dropped, and every caller
459    /// still waiting for the store fails. Returns once the actor has exited.
460    pub async fn crash(&self) {
461        self.stopped.store(true, Ordering::Release);
462        let (response_tx, response_rx) = oneshot::channel();
463        if self
464            .command_tx
465            .send(Command::Crash {
466                response: response_tx,
467            })
468            .await
469            .is_ok()
470        {
471            let _ = response_rx.await;
472        }
473    }
474
475    /// Sets the stopped flag without sending the stop command, which is the
476    /// state a store is in between the two steps of [`stop`](LogStore::stop).
477    pub fn begin_stop(&self) {
478        self.stopped.store(true, Ordering::Release);
479    }
480}
481
482#[async_trait::async_trait]
483impl LogStore for ObjectStoreLogStore {
484    type Error = Error;
485
486    /// Stops the store. Creates in flight run to completion and acknowledge
487    /// their entries if they succeed; in the `enqueued` mode the remaining
488    /// backlog is uploaded first, and a failure of that upload is returned.
489    async fn stop(&self) -> Result<()> {
490        self.stopped.store(true, Ordering::Release);
491        let (response_tx, response_rx) = oneshot::channel();
492        let sent = self
493            .command_tx
494            .send(Command::Stop {
495                response: response_tx,
496            })
497            .await;
498        // A closed channel means the actor already exited.
499        if sent.is_ok() {
500            return response_rx.await.unwrap_or(Ok(()));
501        }
502        Ok(())
503    }
504
505    async fn append_batch(&self, entries: Vec<Entry>) -> Result<AppendBatchResponse> {
506        ensure!(
507            !self.stopped.load(Ordering::Acquire),
508            ObjectStoreWalStoppedSnafu
509        );
510        self.check_terminal()?;
511        if entries.is_empty() {
512            return Ok(AppendBatchResponse::default());
513        }
514        for entry in &entries {
515            let region_id = self.region_of(entry.provider())?;
516            ensure!(
517                region_id == entry.region_id(),
518                MismatchedWalRegionSnafu {
519                    region_id: entry.region_id(),
520                    reason: format!("provider belongs to region {region_id}"),
521                }
522            );
523            ensure!(entry.is_complete(), IncompleteWalEntrySnafu { region_id });
524        }
525
526        let (response_tx, response_rx) = oneshot::channel();
527        self.append_tx
528            .send((entries, response_tx))
529            .await
530            .ok()
531            .context(ObjectStoreWalStoppedSnafu)?;
532        response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)?
533    }
534
535    /// Reads the provider's region from `entry_id`, hiding obsolete entries.
536    /// The catalog locates each segment, so a caller-supplied index is unnecessary.
537    async fn read(
538        &self,
539        provider: &Provider,
540        entry_id: EntryId,
541        _index: Option<WalIndex>,
542    ) -> Result<SendableEntryStream<'static, Entry, Error>> {
543        self.check_terminal()?;
544        let region_id = self.region_of(provider)?;
545        let obsolete = self
546            .obsolete_entry_ids
547            .lock()
548            .unwrap_or_else(PoisonError::into_inner)
549            .get(&region_id)
550            .copied();
551        let start_entry_id =
552            entry_id.max(obsolete.map_or(0, |obsolete| obsolete.saturating_add(1)));
553        let objects = {
554            let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner);
555            match catalog.region_max_entry_id(region_id) {
556                Some(max_entry_id)
557                    if start_entry_id <= max_entry_id && obsolete != Some(EntryId::MAX) =>
558                {
559                    catalog
560                        .objects_for_entry_range(region_id, start_entry_id, max_entry_id)?
561                        .into_iter()
562                        .map(|(object_seq, entry)| (object_seq, entry.clone()))
563                        .collect::<Vec<_>>()
564                }
565                _ => Vec::new(),
566            }
567        };
568
569        let io = self.io.clone();
570        let provider = provider.clone();
571        Ok(Box::pin(try_stream! {
572            for (object_seq, footer_entry) in objects {
573                let bytes = io
574                    .get_range(object_seq, footer_entry.segment_offset, footer_entry.segment_len)
575                    .await?;
576                let records = decode_segment(&bytes, &footer_entry)
577                    .with_context(|_| InvalidWalObjectSnafu {
578                        path: io.object_path(object_seq),
579                    })?;
580                let entries = records
581                    .into_iter()
582                    .filter(|record| record.entry_id >= start_entry_id)
583                    .map(|record| {
584                        Entry::Naive(NaiveEntry {
585                            provider: provider.clone(),
586                            region_id,
587                            entry_id: record.entry_id,
588                            data: record.payload.into(),
589                        })
590                    })
591                    .collect::<Vec<_>>();
592                if !entries.is_empty() {
593                    yield entries;
594                }
595            }
596        }))
597    }
598
599    async fn create_namespace(&self, ns: &Provider) -> Result<()> {
600        self.check_terminal()?;
601        self.region_of(ns).map(|_| ())
602    }
603
604    async fn delete_namespace(&self, ns: &Provider) -> Result<()> {
605        self.check_terminal()?;
606        self.region_of(ns).map(|_| ())
607    }
608
609    async fn list_namespaces(&self) -> Result<Vec<Provider>> {
610        self.check_terminal()?;
611        let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner);
612        let regions = catalog
613            .objects_in_order()
614            .flat_map(|(_, footer)| footer.iter().map(|entry| entry.region_id))
615            .collect::<BTreeSet<_>>();
616        Ok(regions
617            .into_iter()
618            .map(|region_id| Provider::object_store_provider(region_id, self.prefix.clone()))
619            .collect())
620    }
621
622    /// Moves the obsolete watermark of the region up to `entry_id`. In the
623    /// `enqueued` acknowledgement mode the watermark never passes the durable
624    /// entry id: an entry that is not durable yet is replayed after a crash,
625    /// and hiding it would skip it.
626    ///
627    /// Every id the region is assigned from now on is greater than
628    /// `entry_id`, whatever watermark is recorded: the sequence of the next
629    /// object is raised above the object `entry_id` names unless it is there
630    /// already, so a watermark assigned under another prefix, or one whose
631    /// object this prefix no longer holds, is never passed by a new id. Zero
632    /// names no object and needs no floor.
633    ///
634    /// The watermark and the floor are applied together by the actor: a call
635    /// cancelled before its command is queued publishes neither, while a
636    /// command already queued still applies both. While the open batch holds
637    /// ids under the next sequence the floor cannot move, and the call fails
638    /// with [`Error::WalObjectSequenceUnsettled`] and records neither. The
639    /// watermark is kept in memory; the caller re-establishes it after a
640    /// restart. In the `enqueued` mode an `entry_id` this store handed out
641    /// needs no floor: it is never handed out again while the store runs.
642    async fn obsolete(
643        &self,
644        provider: &Provider,
645        region_id: RegionId,
646        entry_id: EntryId,
647    ) -> Result<()> {
648        self.check_region(provider, region_id)?;
649        let watermark = match self.ack_mode {
650            AckMode::Durable => entry_id,
651            AckMode::Enqueued => entry_id.min(self.durable_entry_id_of(region_id)),
652        };
653        let (response_tx, response_rx) = oneshot::channel();
654        let command = Command::Obsolete {
655            region_id,
656            entry_id,
657            watermark,
658            response: response_tx,
659        };
660        let answered = match self.command_tx.send(command).await {
661            Ok(()) => response_rx.await.ok(),
662            Err(_) => None,
663        };
664        match answered {
665            Some(result) => result,
666            // The actor exited, before the command was queued or with it
667            // still queued behind the stop: nothing is assigned an id any
668            // more, so the watermark alone is consistent.
669            None => {
670                record_obsolete(&self.obsolete_entry_ids, region_id, watermark);
671                Ok(())
672            }
673        }
674    }
675
676    /// Hides every entry of the region in memory. The caller re-establishes
677    /// the watermark after a restart.
678    async fn obsolete_all(&self, provider: &Provider, region_id: RegionId) -> Result<()> {
679        self.check_region(provider, region_id)?;
680        record_obsolete(&self.obsolete_entry_ids, region_id, EntryId::MAX);
681        Ok(())
682    }
683
684    fn entry(
685        &self,
686        data: Vec<u8>,
687        entry_id: EntryId,
688        region_id: RegionId,
689        provider: &Provider,
690    ) -> Result<Entry> {
691        self.check_region(provider, region_id)?;
692        Ok(Entry::Naive(NaiveEntry {
693            provider: provider.clone(),
694            region_id,
695            entry_id,
696            data,
697        }))
698    }
699
700    fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId> {
701        self.durable_entry_id(provider)
702    }
703
704    /// Waits until every entry of the provider's region with an id at or
705    /// below `entry_id` is durable and indexed. Returns at once in the
706    /// `durable` acknowledgement mode, where a caller only holds durable ids.
707    async fn wait_durable(&self, provider: &Provider, entry_id: EntryId) -> Result<()> {
708        self.check_terminal()?;
709        let region_id = self.region_of(provider)?;
710        if entry_id <= self.durable_entry_id_of(region_id) {
711            return Ok(());
712        }
713        let (response_tx, response_rx) = oneshot::channel();
714        self.command_tx
715            .send(Command::WaitDurable {
716                region_id,
717                entry_id,
718                response: response_tx,
719            })
720            .await
721            .ok()
722            .context(ObjectStoreWalStoppedSnafu)?;
723        response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)?
724    }
725}
726
727type AppendResponse = oneshot::Sender<Result<AppendBatchResponse>>;
728
729type QueuedAppend = (Vec<Entry>, AppendResponse);
730
731enum Command {
732    /// Answered once the region is durable through `entry_id`.
733    WaitDurable {
734        region_id: RegionId,
735        entry_id: EntryId,
736        response: oneshot::Sender<Result<()>>,
737    },
738    /// Answered once the watermark is recorded and no id of the region at
739    /// or below `entry_id` can be assigned, or with the reason neither was
740    /// done.
741    Obsolete {
742        region_id: RegionId,
743        entry_id: EntryId,
744        watermark: EntryId,
745        response: oneshot::Sender<Result<()>>,
746    },
747    Stop {
748        response: oneshot::Sender<Result<()>>,
749    },
750    #[cfg(any(test, feature = "testing"))]
751    Seal {
752        response: oneshot::Sender<Result<()>>,
753    },
754    #[cfg(any(test, feature = "testing"))]
755    Crash { response: oneshot::Sender<()> },
756}
757
758/// An append waiting for the object that holds its entries.
759struct PendingAppend {
760    last_entry_ids: HashMap<RegionId, EntryId>,
761    response: AppendResponse,
762}
763
764/// A caller waiting until a region is durable through an entry id.
765struct DurableWaiter {
766    region_id: RegionId,
767    entry_id: EntryId,
768    response: oneshot::Sender<Result<()>>,
769}
770
771/// Where the conditional create of a sealed batch stands.
772enum CreateState {
773    /// The create has not started: no slot was free, or in the `enqueued`
774    /// mode an earlier attempt failed transiently and is repeated.
775    Pending,
776    InFlight,
777    Created,
778}
779
780/// A sealed and encoded batch that holds its object sequence and waits to be
781/// created, indexed and acknowledged.
782struct SealedBatch {
783    object_seq: u64,
784    bytes: Bytes,
785    footer: Vec<FooterEntry>,
786    first_admitted_at: Instant,
787    /// Number of creates that were attempted for the batch.
788    attempts: u32,
789    waiters: Vec<PendingAppend>,
790    #[cfg(any(test, feature = "testing"))]
791    seal_waiters: Vec<oneshot::Sender<Result<()>>>,
792    state: CreateState,
793}
794
795impl SealedBatch {
796    fn is_in_flight(&self) -> bool {
797        matches!(self.state, CreateState::InFlight)
798    }
799
800    fn fail(self, error: impl Fn() -> Error) {
801        for waiter in self.waiters {
802            let _ = waiter.response.send(Err(error()));
803        }
804        #[cfg(any(test, feature = "testing"))]
805        for waiter in self.seal_waiters {
806            let _ = waiter.send(Err(error()));
807        }
808    }
809}
810
811type CreateOutcome = (u64, Result<PutResult>);
812
813/// The actor that owns the open batch and the sealed batches until they are
814/// durable.
815///
816/// Sealed batches form a pipeline in sequence order. A batch is created under
817/// its sequence as soon as one of [`MAX_IN_FLIGHT_CREATES`] slots is free, but
818/// it is indexed and acknowledged only once every earlier batch is, so the
819/// acknowledged history of a region never has a missing predecessor. Every
820/// batch links to its pipeline predecessor, the batch sealed before it, or to
821/// the last indexed object when no earlier batch waits, so recovery replays a
822/// batch only on the chain of a later object. The outcomes of a create are:
823///
824/// | Situation | Sequence | Waiters | Store |
825/// | --- | --- | --- | --- |
826/// | created, or identical retry | advances | acknowledged in order | healthy |
827/// | transient error, `durable` mode | 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 |
828/// | transient error, `enqueued` mode | unchanged | already acknowledged | healthy: the create is repeated with the same bytes after [`CREATE_RETRY_DELAY`], which an identical retry accepts |
829/// | transient error, `enqueued` mode after `stop` began | never reused | already acknowledged | the backlog is dropped and `stop` reports the error |
830/// | an object of an earlier epoch holds the sequence, `durable` mode | never reused | as for a transient error | healthy: that object can never be on the chain |
831/// | an object of the same or a later epoch holds the sequence, any conflicting object in the `enqueued` mode, encoding or catalog error | unchanged | every batch that is not indexed fails | poisoned |
832/// | created at the last representable sequence | cannot advance | acknowledged | poisoned: no later batch can be allocated a sequence |
833///
834/// A create still in flight when its batch failed runs to completion and its
835/// outcome is ignored: whatever it stored is off the chain. After `stop` began
836/// nothing is admitted and no create starts, except that the `enqueued` mode
837/// uploads its backlog; creates in flight run to completion and acknowledge
838/// if they succeed.
839///
840/// Appends arrive on their own channel, which the actor does not read while
841/// [`MAX_SEALED_BATCHES`] batches wait or an append is held back at a backlog
842/// threshold of the `enqueued` mode, so commands such as `stop` are handled
843/// while admission is held back.
844struct Actor {
845    io: Arc<dyn WalObjectIo>,
846    catalog: Arc<RwLock<ObjectCatalog>>,
847    obsolete_entry_ids: ObsoleteEntryIds,
848    terminal_error: TerminalError,
849    stopped: Arc<AtomicBool>,
850    command_rx: mpsc::Receiver<Command>,
851    append_rx: mpsc::Receiver<QueuedAppend>,
852    ack_mode: AckMode,
853    max_unpersisted_bytes: usize,
854    max_unpersisted_age: Duration,
855    open_batch: OpenBatch,
856    /// Largest entry id ever handed out per region, whether it became
857    /// durable or failed. Unlike the accepted ids of the open batch it never
858    /// moves down.
859    issued_entry_ids: HashMap<RegionId, EntryId>,
860    /// Waiters of the open batch in the `durable` mode.
861    pending: Vec<PendingAppend>,
862    /// Batches that are not durable yet, in sequence order.
863    sealed: VecDeque<SealedBatch>,
864    /// Creates in flight, including those of batches that already failed.
865    creates: FuturesUnordered<BoxFuture<'static, CreateOutcome>>,
866    /// The append held back in the `enqueued` mode while the unpersisted
867    /// backlog is at a threshold. Later appends wait in their channel.
868    stalled: Option<QueuedAppend>,
869    durable_waiters: Vec<DurableWaiter>,
870    /// Callers of `stop`, answered once nothing is in flight.
871    stop: Vec<oneshot::Sender<Result<()>>>,
872    /// The failure `stop` reports in the `enqueued` mode once an
873    /// acknowledged backlog was dropped, recorded when it happens.
874    stop_error: Option<Arc<Error>>,
875    /// Sequence of the next sealed batch, `None` once the sequence is
876    /// exhausted. The open batch assigns its entry ids from it. It is above
877    /// every sealed batch and never moves back while the store runs.
878    next_object_seq: Option<u64>,
879    /// Epoch of this instance, carried by every object it writes.
880    epoch: u64,
881    /// The last indexed object, which is the start object until a batch is
882    /// indexed. A batch sealed while no earlier batch waits extends it.
883    last_indexed: ChainLink,
884    flush_interval: Duration,
885    #[cfg(any(test, feature = "testing"))]
886    admitted_appends: watch::Sender<usize>,
887    #[cfg(any(test, feature = "testing"))]
888    creates_held: watch::Receiver<bool>,
889    #[cfg(any(test, feature = "testing"))]
890    creates_fail: Arc<AtomicBool>,
891    #[cfg(any(test, feature = "testing"))]
892    next_create_fails_after_write: Arc<AtomicBool>,
893    #[cfg(any(test, feature = "testing"))]
894    parked_creates: Arc<watch::Sender<usize>>,
895    #[cfg(any(test, feature = "testing"))]
896    durability_waits: watch::Sender<usize>,
897}
898
899impl Actor {
900    async fn run(mut self) {
901        let mut interval = tokio::time::interval(self.flush_interval);
902        interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
903        // The first tick completes immediately.
904        interval.tick().await;
905
906        loop {
907            tokio::select! {
908                _ = interval.tick() => {
909                    // Nothing starts after stop began; the `enqueued` mode
910                    // sealed its backlog when stop was requested.
911                    if !self.is_stopped() {
912                        self.flush_open_batch();
913                    }
914                }
915                Some((object_seq, result)) = self.creates.next(), if !self.creates.is_empty() => {
916                    self.on_create_completed(object_seq, result);
917                }
918                Some((entries, response)) = self.append_rx.recv(), if self.stalled.is_none() && self.sealed.len() + 2 <= MAX_SEALED_BATCHES => {
919                    self.handle_append(entries, response);
920                }
921                command = self.command_rx.recv() => match command {
922                    Some(Command::WaitDurable { region_id, entry_id, response }) => {
923                        self.handle_wait_durable(region_id, entry_id, response);
924                    }
925                    Some(Command::Obsolete { region_id, entry_id, watermark, response }) => {
926                        self.handle_obsolete(region_id, entry_id, watermark, response);
927                    }
928                    Some(Command::Stop { response }) => {
929                        self.handle_stop(response);
930                    }
931                    #[cfg(any(test, feature = "testing"))]
932                    Some(Command::Seal { response }) => {
933                        self.handle_seal(response);
934                    }
935                    #[cfg(any(test, feature = "testing"))]
936                    Some(Command::Crash { response }) => {
937                        self.handle_crash();
938                        let _ = response.send(());
939                        return;
940                    }
941                    // Every sender is gone: the store was dropped without
942                    // `stop`. The creates in flight are dropped with the actor.
943                    None => return,
944                },
945            }
946            if self.finish_stop() {
947                return;
948            }
949        }
950    }
951
952    fn handle_append(&mut self, entries: Vec<Entry>, response: AppendResponse) {
953        // The append was queued before `stop` set the flag; nothing that is
954        // not durable yet gets admitted once it is set.
955        if self.is_stopped() {
956            let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build()));
957            return;
958        }
959        if let Some(error) = terminal(&self.terminal_error) {
960            let _ = response.send(Err(shared(&error)));
961            return;
962        }
963        if self.ack_mode == AckMode::Enqueued && self.backlog_at_threshold() {
964            self.stalled = Some((entries, response));
965            self.ensure_create_in_flight();
966            return;
967        }
968        self.admit(entries, response);
969    }
970
971    /// Admits `entries` into the open batch, which assigns their ids under
972    /// the next object sequence. In the `durable` mode the caller waits for
973    /// the object, in the `enqueued` mode it is answered now.
974    fn admit(&mut self, entries: Vec<Entry>, response: AppendResponse) {
975        // A region would run past the position range of the open batch: the
976        // batch is sealed and the entries open the next one. The size limit
977        // seals a batch long before a million entries of one region, so this
978        // is a theoretical bound.
979        if !self.open_batch.is_empty() && self.open_batch.would_exhaust_positions(&entries) {
980            self.flush_open_batch();
981            if let Some(error) = terminal(&self.terminal_error) {
982                let _ = response.send(Err(shared(&error)));
983                return;
984            }
985        }
986        let Some(object_seq) = self.next_object_seq else {
987            let error = self.poison(
988                WalObjectSequenceExhaustedSnafu {
989                    last_object_seq: OBJECT_SEQ_LIMIT - 1,
990                }
991                .build(),
992            );
993            let _ = response.send(Err(shared(&error)));
994            return;
995        };
996        let last_entry_ids = match self.open_batch.admit(object_seq, entries) {
997            Ok(last_entry_ids) => last_entry_ids,
998            // Nothing was admitted: the append alone runs past the position
999            // range, which no object can hold.
1000            Err(error) => {
1001                let _ = response.send(Err(error));
1002                return;
1003            }
1004        };
1005        for (region_id, entry_id) in &last_entry_ids {
1006            self.issued_entry_ids
1007                .entry(*region_id)
1008                .and_modify(|issued| *issued = (*issued).max(*entry_id))
1009                .or_insert(*entry_id);
1010        }
1011        match self.ack_mode {
1012            AckMode::Durable => self.pending.push(PendingAppend {
1013                last_entry_ids,
1014                response,
1015            }),
1016            AckMode::Enqueued => {
1017                let _ = response.send(Ok(AppendBatchResponse { last_entry_ids }));
1018            }
1019        }
1020        #[cfg(any(test, feature = "testing"))]
1021        self.admitted_appends.send_modify(|count| *count += 1);
1022        if self.open_batch.should_seal() {
1023            self.flush_open_batch();
1024        }
1025    }
1026
1027    /// Returns true once the unpersisted backlog, the open batch and every
1028    /// sealed batch that is not durable, reaches the size or the age threshold.
1029    fn backlog_at_threshold(&self) -> bool {
1030        let bytes = self.open_batch.estimated_bytes()
1031            + self
1032                .sealed
1033                .iter()
1034                .map(|batch| batch.bytes.len())
1035                .sum::<usize>();
1036        if bytes >= self.max_unpersisted_bytes {
1037            return true;
1038        }
1039        let oldest = self
1040            .sealed
1041            .front()
1042            .map(|batch| batch.first_admitted_at)
1043            .or_else(|| self.open_batch.first_admitted_at());
1044        oldest.is_some_and(|admitted_at| admitted_at.elapsed() >= self.max_unpersisted_age)
1045    }
1046
1047    /// Seals the open batch when no create is in flight, so that a stalled
1048    /// append has an upload to wait for.
1049    fn ensure_create_in_flight(&mut self) {
1050        if !self.sealed.iter().any(SealedBatch::is_in_flight) {
1051            self.flush_open_batch();
1052        }
1053    }
1054
1055    /// Admits the stalled append once the backlog is below the thresholds and
1056    /// the sealed batches leave room for the two an admission can seal.
1057    fn release_stalled(&mut self) {
1058        if self.stalled.is_none() {
1059            return;
1060        }
1061        if self.is_stopped() {
1062            if let Some((_, response)) = self.stalled.take() {
1063                let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build()));
1064            }
1065            return;
1066        }
1067        if self.backlog_at_threshold() {
1068            // The append needs an upload to complete.
1069            self.ensure_create_in_flight();
1070            return;
1071        }
1072        if self.sealed.len() + 2 > MAX_SEALED_BATCHES {
1073            return;
1074        }
1075        if let Some((entries, response)) = self.stalled.take() {
1076            self.admit(entries, response);
1077        }
1078    }
1079
1080    /// Seals the open batch as the object `next_object_seq` and starts its
1081    /// create when a slot is free. Returns whether a batch was sealed.
1082    fn flush_open_batch(&mut self) -> bool {
1083        if self.open_batch.is_empty() {
1084            return false;
1085        }
1086        if let Some(error) = terminal(&self.terminal_error) {
1087            self.fail_unacknowledged(|| shared(&error));
1088            return false;
1089        }
1090        let Some(object_seq) = self.next_object_seq else {
1091            self.poison(
1092                WalObjectSequenceExhaustedSnafu {
1093                    last_object_seq: OBJECT_SEQ_LIMIT - 1,
1094                }
1095                .build(),
1096            );
1097            return false;
1098        };
1099
1100        let (entries, first_admitted_at) = self.open_batch.seal();
1101        let header = Header {
1102            object_seq,
1103            epoch: self.epoch,
1104            prev: Some(
1105                self.sealed
1106                    .back()
1107                    .map_or(self.last_indexed, |batch| ChainLink {
1108                        object_seq: batch.object_seq,
1109                        epoch: self.epoch,
1110                    }),
1111            ),
1112        };
1113        let encoded = match encode_batch(header, entries) {
1114            Ok(encoded) => encoded,
1115            Err(error) => {
1116                self.poison(error);
1117                return false;
1118            }
1119        };
1120        self.next_object_seq = object_seq
1121            .checked_add(1)
1122            .filter(|next_object_seq| *next_object_seq < OBJECT_SEQ_LIMIT);
1123        if self.next_object_seq.is_none() {
1124            // The batch takes the last representable sequence: it is created
1125            // and acknowledged, but no later batch can be allocated one.
1126            set_terminal(
1127                &self.terminal_error,
1128                WalObjectSequenceExhaustedSnafu {
1129                    last_object_seq: object_seq,
1130                }
1131                .build(),
1132            );
1133        }
1134        self.sealed.push_back(SealedBatch {
1135            object_seq,
1136            bytes: encoded.bytes,
1137            footer: encoded.footer,
1138            first_admitted_at,
1139            attempts: 0,
1140            waiters: std::mem::take(&mut self.pending),
1141            #[cfg(any(test, feature = "testing"))]
1142            seal_waiters: Vec::new(),
1143            state: CreateState::Pending,
1144        });
1145        self.start_creates();
1146        true
1147    }
1148
1149    /// Starts the creates of pending batches in sequence order while fewer
1150    /// than [`MAX_IN_FLIGHT_CREATES`] are in flight, counting the creates of
1151    /// batches that already failed. Nothing starts once stop began, except
1152    /// the backlog of the `enqueued` mode.
1153    fn start_creates(&mut self) {
1154        if self.is_stopped() && self.ack_mode == AckMode::Durable {
1155            return;
1156        }
1157        let mut in_flight = self.creates.len();
1158        for batch in self.sealed.iter_mut() {
1159            if in_flight >= MAX_IN_FLIGHT_CREATES {
1160                break;
1161            }
1162            if !matches!(batch.state, CreateState::Pending) {
1163                continue;
1164            }
1165            batch.state = CreateState::InFlight;
1166            in_flight += 1;
1167            let delay = if batch.attempts == 0 {
1168                Duration::ZERO
1169            } else {
1170                CREATE_RETRY_DELAY
1171            };
1172            batch.attempts += 1;
1173            let io = self.io.clone();
1174            let object_seq = batch.object_seq;
1175            let bytes = batch.bytes.clone();
1176            let epoch = self.epoch;
1177            let ack_mode = self.ack_mode;
1178            #[cfg(any(test, feature = "testing"))]
1179            let mut creates_held = self.creates_held.clone();
1180            #[cfg(any(test, feature = "testing"))]
1181            let creates_fail = self.creates_fail.clone();
1182            #[cfg(any(test, feature = "testing"))]
1183            let next_create_fails_after_write = self.next_create_fails_after_write.clone();
1184            #[cfg(any(test, feature = "testing"))]
1185            let parked_creates = self.parked_creates.clone();
1186            self.creates.push(Box::pin(async move {
1187                if !delay.is_zero() {
1188                    tokio::time::sleep(delay).await;
1189                }
1190                // The store was dropped while the create was parked: it never
1191                // runs.
1192                #[cfg(any(test, feature = "testing"))]
1193                if *creates_held.borrow() {
1194                    parked_creates.send_modify(|count| *count += 1);
1195                }
1196                #[cfg(any(test, feature = "testing"))]
1197                if creates_held.wait_for(|held| !*held).await.is_err() {
1198                    return (object_seq, Err(ObjectStoreWalStoppedSnafu.build()));
1199                }
1200                #[cfg(any(test, feature = "testing"))]
1201                if creates_fail.load(Ordering::Acquire) {
1202                    return (object_seq, injected_create_failure(io.as_ref(), object_seq));
1203                }
1204                // Any conflict poisons an `enqueued` store, so only the
1205                // `durable` mode reads the epoch of the existing object.
1206                let result = match io.put_if_absent(object_seq, bytes).await {
1207                    Err(error @ Error::WalObjectConflict { .. })
1208                        if ack_mode == AckMode::Durable =>
1209                    {
1210                        stale_conflict(io.as_ref(), object_seq, epoch, error).await
1211                    }
1212                    result => result,
1213                };
1214                #[cfg(any(test, feature = "testing"))]
1215                if result.is_ok() && next_create_fails_after_write.swap(false, Ordering::AcqRel) {
1216                    return (object_seq, injected_create_failure(io.as_ref(), object_seq));
1217                }
1218                (object_seq, result)
1219            }));
1220        }
1221    }
1222
1223    fn on_create_completed(&mut self, object_seq: u64, result: Result<PutResult>) {
1224        // A batch that already failed, or that the store gave up on when it
1225        // poisoned itself: the object may exist, but it is off the chain.
1226        let Some(index) = self
1227            .sealed
1228            .iter()
1229            .position(|batch| batch.object_seq == object_seq)
1230        else {
1231            self.start_creates();
1232            return;
1233        };
1234        match result {
1235            Ok(_) => self.sealed[index].state = CreateState::Created,
1236            // Nobody is left to retry in the `enqueued` mode, so the store
1237            // repeats a create that failed transiently under the same sequence
1238            // with the same bytes, which an identical retry accepts.
1239            Err(ref error)
1240                if self.ack_mode == AckMode::Enqueued
1241                    && !self.is_stopped()
1242                    && is_transient(error) =>
1243            {
1244                self.sealed[index].state = CreateState::Pending;
1245            }
1246            // The object store did not confirm the object, which may still
1247            // exist or land later, or an earlier epoch holds the sequence and
1248            // can never be on the chain. A caller of the `durable` mode
1249            // retries the append itself; after stop began a transient failure
1250            // drops the `enqueued` backlog.
1251            Err(error @ Error::WalObjectStore { .. })
1252                if self.ack_mode == AckMode::Durable
1253                    || (self.is_stopped() && is_transient(&error)) =>
1254            {
1255                self.roll_back(Arc::new(error))
1256            }
1257            Err(error @ Error::StaleWalObject { .. }) if self.ack_mode == AckMode::Durable => {
1258                self.roll_back(Arc::new(error))
1259            }
1260            Err(error) => {
1261                self.poison(error);
1262                return;
1263            }
1264        }
1265        self.settle();
1266    }
1267
1268    /// Indexes and acknowledges the sealed batches from the front as far as
1269    /// they are created, and starts the creates that a free slot allows.
1270    fn settle(&mut self) {
1271        while let Some(front) = self.sealed.front() {
1272            if !matches!(front.state, CreateState::Created) || !self.index_front() {
1273                break;
1274            }
1275        }
1276        self.start_creates();
1277        self.release_stalled();
1278    }
1279
1280    /// Indexes the created object at the front and acknowledges its waiters.
1281    /// Returns false when the catalog rejected it, which poisons the store.
1282    fn index_front(&mut self) -> bool {
1283        let Some(front) = self.sealed.front() else {
1284            return false;
1285        };
1286        let indexed = self
1287            .catalog
1288            .write()
1289            .unwrap_or_else(PoisonError::into_inner)
1290            .insert_object(front.object_seq, front.footer.clone());
1291        if let Err(error) = indexed {
1292            self.poison(error);
1293            return false;
1294        }
1295        let Some(batch) = self.sealed.pop_front() else {
1296            return false;
1297        };
1298        self.last_indexed = ChainLink {
1299            object_seq: batch.object_seq,
1300            epoch: self.epoch,
1301        };
1302        for waiter in batch.waiters {
1303            let _ = waiter.response.send(Ok(AppendBatchResponse {
1304                last_entry_ids: waiter.last_entry_ids,
1305            }));
1306        }
1307        #[cfg(any(test, feature = "testing"))]
1308        for waiter in batch.seal_waiters {
1309            let _ = waiter.send(Ok(()));
1310        }
1311        self.resolve_durable_waiters();
1312        true
1313    }
1314
1315    fn resolve_durable_waiters(&mut self) {
1316        let waiters = std::mem::take(&mut self.durable_waiters);
1317        for waiter in waiters {
1318            if self.is_durable_through(waiter.region_id, waiter.entry_id) {
1319                let _ = waiter.response.send(Ok(()));
1320            } else {
1321                self.durable_waiters.push(waiter);
1322            }
1323        }
1324    }
1325
1326    /// Returns true when no id of the region at or below `entry_id` waits to
1327    /// become durable. An id of a batch that failed will never be durable and
1328    /// no longer waits.
1329    fn is_durable_through(&self, region_id: RegionId, entry_id: EntryId) -> bool {
1330        self.lowest_pending_entry_id(region_id)
1331            .is_none_or(|pending| pending > entry_id)
1332    }
1333
1334    /// Returns the lowest id of the region that was handed out and is not
1335    /// durable yet. The sealed batches are in sequence order and the open
1336    /// batch is above all of them.
1337    fn lowest_pending_entry_id(&self, region_id: RegionId) -> Option<EntryId> {
1338        self.sealed
1339            .iter()
1340            .flat_map(|batch| batch.footer.iter())
1341            .find(|entry| entry.region_id == region_id)
1342            .map(|entry| entry.min_entry_id)
1343            .or_else(|| {
1344                self.next_object_seq
1345                    .filter(|_| self.open_batch.holds_region(region_id))
1346                    .map(|object_seq| entry_id(object_seq, 1))
1347            })
1348    }
1349
1350    /// Fails every batch that is not indexed after a create failed
1351    /// transiently, including batches already created and batches whose
1352    /// create is in flight. Their sequences are not reused: a create that
1353    /// reported an error may still have written its object. The next batch
1354    /// links to the last indexed object, so no later chain reaches a failed
1355    /// batch. Waiters of a store that was stopped meanwhile learn that
1356    /// instead of the I/O error, like every other entry that never became
1357    /// durable.
1358    fn roll_back(&mut self, error: Arc<Error>) {
1359        let stopped = self.is_stopped();
1360        let failure = || {
1361            if stopped {
1362                ObjectStoreWalStoppedSnafu.build()
1363            } else {
1364                shared(&error)
1365            }
1366        };
1367        for batch in self.sealed.drain(..) {
1368            batch.fail(failure);
1369        }
1370        self.reset_open_batch();
1371        self.fail_unacknowledged(failure);
1372        // An acknowledged backlog was dropped: `stop` reports it, whether
1373        // its caller has arrived yet or not.
1374        if self.ack_mode == AckMode::Enqueued {
1375            self.stop_error.get_or_insert(error);
1376        }
1377    }
1378
1379    /// Records `error` as terminal and fails every waiter that is not
1380    /// acknowledged with it, or with the stopped error if the store was
1381    /// stopped meanwhile. Creates in flight run to completion, but their
1382    /// outcome is ignored: in the `durable` mode the entries of an object they
1383    /// create were never acknowledged, like those of a crash between creation
1384    /// and acknowledgement; in the `enqueued` mode the acknowledged backlog is
1385    /// discarded and `stop` reports it. Returns the recorded error.
1386    fn poison(&mut self, error: Error) -> Arc<Error> {
1387        let error = set_terminal(&self.terminal_error, error);
1388        let stopped = self.is_stopped();
1389        let failure = || {
1390            if stopped {
1391                ObjectStoreWalStoppedSnafu.build()
1392            } else {
1393                shared(&error)
1394            }
1395        };
1396        for batch in self.sealed.drain(..) {
1397            batch.fail(failure);
1398        }
1399        self.reset_open_batch();
1400        self.fail_unacknowledged(failure);
1401        if self.ack_mode == AckMode::Enqueued {
1402            self.stop_error.get_or_insert(error.clone());
1403        }
1404        error
1405    }
1406
1407    /// Drops the entries of the open batch; the next admission hands out the
1408    /// same ids again under the sequence the batch is at.
1409    fn reset_open_batch(&mut self) {
1410        self.open_batch.reset();
1411    }
1412
1413    /// Fails the waiters of the open batch, the stalled append and the
1414    /// durability waiters.
1415    fn fail_unacknowledged(&mut self, error: impl Fn() -> Error) {
1416        for pending in self.pending.drain(..) {
1417            let _ = pending.response.send(Err(error()));
1418        }
1419        if let Some((_, response)) = self.stalled.take() {
1420            let _ = response.send(Err(error()));
1421        }
1422        for waiter in self.durable_waiters.drain(..) {
1423            let _ = waiter.response.send(Err(error()));
1424        }
1425    }
1426
1427    fn handle_wait_durable(
1428        &mut self,
1429        region_id: RegionId,
1430        entry_id: EntryId,
1431        response: oneshot::Sender<Result<()>>,
1432    ) {
1433        if let Some(error) = terminal(&self.terminal_error) {
1434            let _ = response.send(Err(shared(&error)));
1435            return;
1436        }
1437        let durable = {
1438            let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner);
1439            catalog.region_max_entry_id(region_id).unwrap_or(0)
1440        };
1441        if entry_id <= durable {
1442            let _ = response.send(Ok(()));
1443            return;
1444        }
1445        // An acknowledged backlog was dropped: no entry that is not durable
1446        // can be certified any more, whether it was in that backlog or not.
1447        if let Some(error) = &self.stop_error {
1448            let _ = response.send(Err(shared(error)));
1449            return;
1450        }
1451        if self.is_durable_through(region_id, entry_id) {
1452            let _ = response.send(Ok(()));
1453            return;
1454        }
1455        // A caller that stopped waiting leaves a closed response behind.
1456        self.durable_waiters
1457            .retain(|waiter| !waiter.response.is_closed());
1458        self.durable_waiters.push(DurableWaiter {
1459            region_id,
1460            entry_id,
1461            response,
1462        });
1463        #[cfg(any(test, feature = "testing"))]
1464        self.durability_waits.send_modify(|count| *count += 1);
1465    }
1466
1467    /// Makes sure no id of the region at or below `entry_id` is assigned
1468    /// from now on, then records the obsolete watermark of the region;
1469    /// neither is done when the sequence cannot be raised.
1470    fn handle_obsolete(
1471        &mut self,
1472        region_id: RegionId,
1473        entry_id: EntryId,
1474        watermark: EntryId,
1475        response: oneshot::Sender<Result<()>>,
1476    ) {
1477        let result = self.raise_sequence_floor(region_id, entry_id).map(|_| {
1478            record_obsolete(&self.obsolete_entry_ids, region_id, watermark);
1479        });
1480        let _ = response.send(result);
1481    }
1482
1483    /// Moves the sequence of the next object above the object that holds
1484    /// `entry_id` unless it is there already. The move fails while the open
1485    /// batch has handed out ids under the next sequence. A floor that does
1486    /// not fit an entry id poisons the store like an exhausted sequence.
1487    fn raise_sequence_floor(&mut self, region_id: RegionId, entry_id: EntryId) -> Result<()> {
1488        // Nothing is assigned an id after stop began.
1489        if self.is_stopped() {
1490            return Ok(());
1491        }
1492        if let Some(error) = terminal(&self.terminal_error) {
1493            return Err(shared(&error));
1494        }
1495        // In the `enqueued` mode an id the store handed out is never handed
1496        // out again while it runs: a create that fails transiently is
1497        // repeated under its sequence and a permanent failure poisons the
1498        // store. Every later id of the region is greater, so the id needs no
1499        // floor even before it is durable.
1500        if self.ack_mode == AckMode::Enqueued
1501            && entry_id <= self.issued_entry_ids.get(&region_id).copied().unwrap_or(0)
1502        {
1503            return Ok(());
1504        }
1505        let sequence_floor = sequence_floor(entry_id);
1506        let Some(next_object_seq) = self.next_object_seq else {
1507            return Ok(());
1508        };
1509        if next_object_seq >= sequence_floor {
1510            return Ok(());
1511        }
1512        ensure!(
1513            self.open_batch.is_empty(),
1514            WalObjectSequenceUnsettledSnafu {
1515                object_seq: next_object_seq,
1516            }
1517        );
1518        if sequence_floor < OBJECT_SEQ_LIMIT {
1519            self.next_object_seq = Some(sequence_floor);
1520            Ok(())
1521        } else {
1522            self.next_object_seq = None;
1523            let error = self.poison(
1524                WalObjectSequenceExhaustedSnafu {
1525                    last_object_seq: OBJECT_SEQ_LIMIT - 1,
1526                }
1527                .build(),
1528            );
1529            Err(shared(&error))
1530        }
1531    }
1532
1533    /// Begins stopping. Nothing is admitted from now on; the `durable` mode
1534    /// drops the open batch and the batches whose create has not started,
1535    /// the `enqueued` mode seals its backlog so it is uploaded. Stop is
1536    /// answered by [`finish_stop`](Self::finish_stop) once nothing is in flight.
1537    fn handle_stop(&mut self, response: oneshot::Sender<Result<()>>) {
1538        self.stop.push(response);
1539        match self.ack_mode {
1540            AckMode::Durable => {
1541                self.reset_open_batch();
1542                for pending in self.pending.drain(..) {
1543                    let _ = pending
1544                        .response
1545                        .send(Err(ObjectStoreWalStoppedSnafu.build()));
1546                }
1547                if let Some(index) = self
1548                    .sealed
1549                    .iter()
1550                    .position(|batch| matches!(batch.state, CreateState::Pending))
1551                {
1552                    for batch in self.sealed.drain(index..) {
1553                        batch.fail(|| ObjectStoreWalStoppedSnafu.build());
1554                    }
1555                }
1556            }
1557            AckMode::Enqueued => {
1558                if let Some((_, response)) = self.stalled.take() {
1559                    let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build()));
1560                }
1561                self.flush_open_batch();
1562            }
1563        }
1564    }
1565
1566    /// Answers the callers of `stop` once every sealed batch is settled and
1567    /// every create has completed. Returns true when the actor is done.
1568    fn finish_stop(&mut self) -> bool {
1569        if self.stop.is_empty() || !self.sealed.is_empty() || !self.creates.is_empty() {
1570            return false;
1571        }
1572        for waiter in self.durable_waiters.drain(..) {
1573            let _ = waiter
1574                .response
1575                .send(Err(ObjectStoreWalStoppedSnafu.build()));
1576        }
1577        let error = self.stop_error.take();
1578        for response in self.stop.drain(..) {
1579            let _ = response.send(error.as_ref().map_or(Ok(()), |error| Err(shared(error))));
1580        }
1581        true
1582    }
1583
1584    #[cfg(any(test, feature = "testing"))]
1585    fn handle_seal(&mut self, response: oneshot::Sender<Result<()>>) {
1586        if self.is_stopped() {
1587            let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build()));
1588            return;
1589        }
1590        if self.flush_open_batch()
1591            && let Some(batch) = self.sealed.back_mut()
1592        {
1593            batch.seal_waiters.push(response);
1594            return;
1595        }
1596        let result = terminal(&self.terminal_error).map_or(Ok(()), |error| Err(shared(&error)));
1597        let _ = response.send(result);
1598    }
1599
1600    /// Fails every caller that waits for the store; the creates are
1601    /// dropped with the actor.
1602    #[cfg(any(test, feature = "testing"))]
1603    fn handle_crash(&mut self) {
1604        let stopped = || ObjectStoreWalStoppedSnafu.build();
1605        for batch in self.sealed.drain(..) {
1606            batch.fail(stopped);
1607        }
1608        self.reset_open_batch();
1609        self.fail_unacknowledged(stopped);
1610        for response in self.stop.drain(..) {
1611            let _ = response.send(Err(stopped()));
1612        }
1613    }
1614
1615    fn is_stopped(&self) -> bool {
1616        self.stopped.load(Ordering::Acquire)
1617    }
1618}
1619
1620fn terminal(terminal_error: &TerminalError) -> Option<Arc<Error>> {
1621    terminal_error
1622        .lock()
1623        .unwrap_or_else(PoisonError::into_inner)
1624        .clone()
1625}
1626
1627/// Records `error` as the terminal error unless one is already recorded, and
1628/// returns the recorded one.
1629fn set_terminal(terminal_error: &TerminalError, error: Error) -> Arc<Error> {
1630    terminal_error
1631        .lock()
1632        .unwrap_or_else(PoisonError::into_inner)
1633        .get_or_insert_with(|| Arc::new(error))
1634        .clone()
1635}
1636
1637/// Returns true for a storage error that a later attempt may not meet: one
1638/// the object store reports as temporary, or as persistent, which is how its
1639/// retry layer reports a temporary error that outlasted its retries.
1640fn is_transient(error: &Error) -> bool {
1641    matches!(error, Error::WalObjectStore { error, .. } if !error.is_permanent())
1642}
1643
1644/// Wraps an error that several callers receive.
1645fn shared(error: &Arc<Error>) -> Error {
1646    ObjectStoreWalSnafu.into_error(error.clone())
1647}
1648
1649fn encode_batch(header: Header, entries: Vec<Entry>) -> Result<EncodedObject> {
1650    let records = entries
1651        .into_iter()
1652        .map(|entry| Record {
1653            region_id: entry.region_id(),
1654            entry_id: entry.entry_id(),
1655            payload: Bytes::from(entry.into_bytes()),
1656        })
1657        .collect::<Vec<_>>();
1658    encode_object(header, &records)
1659}
1660
1661/// What recovery rebuilt from the objects under a prefix.
1662#[derive(Debug)]
1663struct Recovered {
1664    catalog: ObjectCatalog,
1665    /// Sequence above every present object and every id the catalog holds.
1666    next_object_seq: u64,
1667    durable_entry_ids: HashMap<RegionId, EntryId>,
1668    /// The object the next object extends, `None` on an empty prefix.
1669    tip: Option<ChainLink>,
1670    /// Largest epoch any present object carries, zero on an empty prefix.
1671    max_epoch: u64,
1672}
1673
1674/// An object as recovery fetched it: its key, header and footer.
1675struct FetchedObject {
1676    object: ListedObject,
1677    header: Header,
1678    footer: Vec<FooterEntry>,
1679}
1680
1681/// Rebuilds the catalog from object headers and footers,
1682/// so recovery costs a few small reads per object however large the objects
1683/// are. Segments are not read; a segment checksum is verified by the read that
1684/// decodes it. Footers are fetched for up to [`RECOVERY_CONCURRENCY`] objects
1685/// at a time and indexed in sequence order, so the catalog checks the entry
1686/// ranges of every object against its predecessors like a sequential replay.
1687/// Only the objects on the chain [`select_chain`] picks are indexed.
1688async fn recover(io: &dyn WalObjectIo) -> Result<Recovered> {
1689    let objects = io.list().await?;
1690    finish_recovery(fetch_footers(io, objects, RECOVERY_CONCURRENCY).await?)
1691}
1692
1693/// Indexes the chain of `objects`, which are ordered by sequence. Objects off
1694/// the chain are orphans: they are not indexed, but no later object takes
1695/// their sequences.
1696fn finish_recovery(objects: Vec<FetchedObject>) -> Result<Recovered> {
1697    let headers = objects
1698        .iter()
1699        .map(|fetched| (fetched.object.object_seq, &fetched.header))
1700        .collect::<BTreeMap<_, _>>();
1701    let chain = select_chain(&headers);
1702    // Every open writes an object that starts a chain or extends a complete
1703    // one, and nothing removes objects, so present objects without any
1704    // complete chain mean objects of the chain are gone.
1705    ensure!(
1706        headers.is_empty() || !chain.is_empty(),
1707        CorruptedWalObjectSnafu {
1708            reason: format!(
1709                "no object among {} present objects completes a chain",
1710                headers.len()
1711            ),
1712        }
1713    );
1714    let tip = chain.last().map(|object_seq| ChainLink {
1715        object_seq: *object_seq,
1716        epoch: headers[object_seq].epoch,
1717    });
1718    let max_epoch = headers
1719        .values()
1720        .map(|header| header.epoch)
1721        .max()
1722        .unwrap_or(0);
1723    let after_listed = match headers.last_key_value() {
1724        None => 0,
1725        Some((&last_object_seq, _)) => last_object_seq
1726            .checked_add(1)
1727            .context(WalObjectSequenceExhaustedSnafu { last_object_seq })?,
1728    };
1729    let chain = chain.into_iter().collect::<HashSet<_>>();
1730
1731    let mut catalog = ObjectCatalog::default();
1732    for FetchedObject { object, footer, .. } in objects {
1733        if chain.contains(&object.object_seq) {
1734            catalog
1735                .insert_object(object.object_seq, footer)
1736                .with_context(|_| InvalidWalObjectSnafu { path: object.path })?;
1737        }
1738    }
1739    let next_object_seq = catalog.next_object_seq()?.max(after_listed);
1740    ensure!(
1741        next_object_seq < OBJECT_SEQ_LIMIT,
1742        WalObjectSequenceExhaustedSnafu {
1743            last_object_seq: next_object_seq - 1,
1744        }
1745    );
1746    let durable_entry_ids = durable_entry_ids(&catalog);
1747    Ok(Recovered {
1748        catalog,
1749        next_object_seq,
1750        durable_entry_ids,
1751        tip,
1752        max_epoch,
1753    })
1754}
1755
1756/// Returns, in sequence order, the objects on the chain that ends at the tip:
1757/// the complete object with the largest epoch, then the largest sequence.
1758///
1759/// An object is complete when every link on its chain holds: the chain starts
1760/// at an object without a predecessor, and every other link names a present
1761/// object that carries the epoch the link records. Only one instance writes
1762/// under an epoch, so the epoch tells an object of the linking instance from
1763/// another object under the same sequence. A create that
1764/// was reported as failed may still leave its object, but every object
1765/// written after that failure links past it, and every instance writes under
1766/// an epoch above every object present when it opened, so neither such an
1767/// object nor a late object of an earlier instance ends the chosen chain.
1768fn select_chain(headers: &BTreeMap<u64, &Header>) -> Vec<u64> {
1769    // A predecessor precedes its successor, so one pass in sequence order
1770    // settles every object.
1771    let mut complete = HashMap::with_capacity(headers.len());
1772    for (&object_seq, header) in headers {
1773        let holds = match header.prev {
1774            None => true,
1775            Some(link) if link.object_seq >= object_seq => false,
1776            Some(link) => match headers.get(&link.object_seq) {
1777                Some(prev) => prev.epoch == link.epoch && complete[&link.object_seq],
1778                None => false,
1779            },
1780        };
1781        complete.insert(object_seq, holds);
1782    }
1783    let Some(tip) = headers
1784        .iter()
1785        .filter(|(object_seq, _)| complete[*object_seq])
1786        .max_by_key(|(object_seq, header)| (header.epoch, **object_seq))
1787        .map(|(object_seq, _)| *object_seq)
1788    else {
1789        return Vec::new();
1790    };
1791    let mut chain = vec![tip];
1792    while let Some(link) = headers[chain.last().expect("the chain holds the tip")].prev
1793        && headers.contains_key(&link.object_seq)
1794    {
1795        chain.push(link.object_seq);
1796    }
1797    chain.reverse();
1798    chain
1799}
1800
1801/// Writes the empty object that starts the epoch of this instance at
1802/// `object_seq`, linked to the recovered `tip`, and returns it as the tip
1803/// later objects extend. The epoch is one above the sequence the create
1804/// claims, so no two instances share one however their opens interleave, and
1805/// it is above the epoch of every object recovery listed.
1806///
1807/// An object an earlier instance left at that sequence after recovery listed
1808/// the prefix has a lower epoch and never ends a chain, so the start object
1809/// moves to the next sequence and epoch. An object of an equal or later epoch
1810/// belongs to another writer of the prefix, and the conflict fails the open,
1811/// as does a create whose outcome is unknown: the next open counts the object
1812/// in either case. So does a create that finds the same bytes present:
1813/// another open that recovered the same objects writes an identical start
1814/// object, so this open cannot claim the epoch even if the object is its own.
1815async fn start_epoch(
1816    io: &dyn WalObjectIo,
1817    mut object_seq: u64,
1818    tip: Option<ChainLink>,
1819) -> Result<ChainLink> {
1820    loop {
1821        ensure!(
1822            object_seq < OBJECT_SEQ_LIMIT,
1823            WalObjectSequenceExhaustedSnafu {
1824                last_object_seq: OBJECT_SEQ_LIMIT - 1,
1825            }
1826        );
1827        let epoch = object_seq + 1;
1828        let header = Header {
1829            object_seq,
1830            epoch,
1831            prev: tip,
1832        };
1833        match io
1834            .put_if_absent(object_seq, encode_object(header, &[])?.bytes)
1835            .await
1836        {
1837            Ok(PutResult::Created) => return Ok(ChainLink { object_seq, epoch }),
1838            Ok(PutResult::AlreadyPresent) => {
1839                return UnconfirmedWalEpochStartSnafu {
1840                    path: io.object_path(object_seq),
1841                    epoch,
1842                }
1843                .fail();
1844            }
1845            Err(error @ Error::WalObjectConflict { .. }) => {
1846                if epoch_of(io, object_seq).await? >= epoch {
1847                    return Err(error);
1848                }
1849                object_seq += 1;
1850            }
1851            Err(error) => return Err(error),
1852        }
1853    }
1854}
1855
1856async fn epoch_of(io: &dyn WalObjectIo, object_seq: u64) -> Result<u64> {
1857    let head = io.get_range(object_seq, 0, HEADER_LEN as u64).await?;
1858    let header = decode_header(&head).with_context(|_| InvalidWalObjectSnafu {
1859        path: io.object_path(object_seq),
1860    })?;
1861    Ok(header.epoch)
1862}
1863
1864/// Resolves a create of this store's `epoch` that met a different object:
1865/// an object of an earlier epoch can never be on the chain, so the create
1866/// fails like a transient error; an object of the same or a later epoch
1867/// belongs to another writer and `conflict` is returned.
1868async fn stale_conflict(
1869    io: &dyn WalObjectIo,
1870    object_seq: u64,
1871    epoch: u64,
1872    conflict: Error,
1873) -> Result<PutResult> {
1874    let existing_epoch = epoch_of(io, object_seq).await?;
1875    if existing_epoch >= epoch {
1876        return Err(conflict);
1877    }
1878    StaleWalObjectSnafu {
1879        path: io.object_path(object_seq),
1880        existing_epoch,
1881        epoch,
1882    }
1883    .fail()
1884}
1885
1886/// Fetches and verifies the headers and footers of `objects`, up to
1887/// `concurrency` objects at a time, and returns them ordered by object sequence
1888/// whatever the order the fetches complete in. The first failure abandons the
1889/// remaining fetches.
1890async fn fetch_footers(
1891    io: &dyn WalObjectIo,
1892    objects: Vec<ListedObject>,
1893    concurrency: usize,
1894) -> Result<Vec<FetchedObject>> {
1895    let mut fetched = futures::stream::iter(objects)
1896        .map(|object| async move {
1897            let (header, footer) = fetch_footer(io, &object).await?;
1898            Ok(FetchedObject {
1899                object,
1900                header,
1901                footer,
1902            })
1903        })
1904        .buffer_unordered(concurrency)
1905        .try_collect::<Vec<_>>()
1906        .await?;
1907    fetched.sort_unstable_by_key(|fetched| fetched.object.object_seq);
1908    Ok(fetched)
1909}
1910
1911/// Reads the header, trailer and footer of `object` and verifies them: the
1912/// header must carry the sequence of the key, the trailer must be well formed,
1913/// the footer must match the checksum the trailer holds and its segments must
1914/// tile the object body.
1915///
1916/// A short object is read whole. Otherwise the header and a window at the end
1917/// of the object are read concurrently, and the footer is read separately only
1918/// when it starts before the window.
1919async fn fetch_footer(
1920    io: &dyn WalObjectIo,
1921    object: &ListedObject,
1922) -> Result<(Header, Vec<FooterEntry>)> {
1923    let ListedObject {
1924        object_seq, size, ..
1925    } = *object;
1926    let invalid = |source: Error| {
1927        InvalidWalObjectSnafu {
1928            path: object.path.clone(),
1929        }
1930        .into_error(source)
1931    };
1932    let object_len = usize::try_from(size)
1933        .ok()
1934        .filter(|len| *len >= MIN_OBJECT_LEN)
1935        .with_context(|| CorruptedWalObjectSnafu {
1936            reason: format!(
1937                "truncated object, expected at least {MIN_OBJECT_LEN} bytes, actual {size}"
1938            ),
1939        })
1940        .map_err(invalid)?;
1941
1942    let window = object_len.min(RECOVERY_TAIL_WINDOW);
1943    let tail_start = object_len - window;
1944    let (head, tail) = if tail_start == 0 {
1945        let bytes = io.get(object_seq).await?;
1946        (bytes.clone(), bytes)
1947    } else {
1948        futures::try_join!(
1949            io.get_range(object_seq, 0, HEADER_LEN as u64),
1950            io.get_range(object_seq, tail_start as u64, window as u64)
1951        )?
1952    };
1953    if head.len() < HEADER_LEN || tail.len() != window {
1954        return Err(invalid(
1955            CorruptedWalObjectSnafu {
1956                reason: format!(
1957                    "object holds fewer bytes than the listed {size}, head {} bytes, tail {} bytes",
1958                    head.len(),
1959                    tail.len()
1960                ),
1961            }
1962            .build(),
1963        ));
1964    }
1965
1966    let (header, trailer, footer_range) =
1967        locate_footer(object_seq, object_len, &head, &tail).map_err(invalid)?;
1968    let footer = if footer_range.start >= tail_start {
1969        tail.slice(footer_range.start - tail_start..footer_range.end - tail_start)
1970    } else {
1971        io.get_range(
1972            object_seq,
1973            footer_range.start as u64,
1974            footer_range.len() as u64,
1975        )
1976        .await?
1977    };
1978    let footer = decode_footer(&footer, trailer).map_err(invalid)?;
1979    verify_segment_ranges(&footer, footer_range.start).map_err(invalid)?;
1980    Ok((header, footer))
1981}
1982
1983/// Verifies the header and trailer of the object `object_seq` of `object_len`
1984/// bytes from its first bytes `head` and its last bytes `tail`, and returns
1985/// the header and the trailer with the range the footer occupies in the object.
1986fn locate_footer(
1987    object_seq: u64,
1988    object_len: usize,
1989    head: &[u8],
1990    tail: &[u8],
1991) -> Result<(Header, FixedTrailer, Range<usize>)> {
1992    let header = decode_header(head)?;
1993    ensure!(
1994        header.object_seq == object_seq,
1995        CorruptedWalObjectSnafu {
1996            reason: format!(
1997                "header sequence {} does not match key sequence {object_seq}",
1998                header.object_seq
1999            ),
2000        }
2001    );
2002    let trailer = decode_trailer(&tail[tail.len() - TRAILER_LEN..])?;
2003    let footer_range = footer_range(trailer, object_len)?;
2004    Ok((header, trailer, footer_range))
2005}
2006
2007fn durable_entry_ids(catalog: &ObjectCatalog) -> HashMap<RegionId, EntryId> {
2008    let mut entry_ids = HashMap::new();
2009    for (_, footer) in catalog.objects_in_order() {
2010        for entry in footer {
2011            entry_ids
2012                .entry(entry.region_id)
2013                .and_modify(|current: &mut EntryId| *current = (*current).max(entry.max_entry_id))
2014                .or_insert(entry.max_entry_id);
2015        }
2016    }
2017    entry_ids
2018}
2019
2020/// Object access of the store, so tests can inject failures.
2021#[async_trait::async_trait]
2022pub(crate) trait WalObjectIo: Send + Sync {
2023    async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult>;
2024
2025    async fn get(&self, object_seq: u64) -> Result<Bytes>;
2026
2027    async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes>;
2028
2029    async fn list(&self) -> Result<Vec<ListedObject>>;
2030
2031    fn object_path(&self, object_seq: u64) -> String;
2032}
2033
2034#[async_trait::async_trait]
2035impl WalObjectIo for ObjectStoreIo {
2036    async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult> {
2037        ObjectStoreIo::put_if_absent(self, object_seq, content).await
2038    }
2039
2040    async fn get(&self, object_seq: u64) -> Result<Bytes> {
2041        ObjectStoreIo::get(self, object_seq).await
2042    }
2043
2044    async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes> {
2045        ObjectStoreIo::get_range(self, object_seq, offset, len).await
2046    }
2047
2048    async fn list(&self) -> Result<Vec<ListedObject>> {
2049        ObjectStoreIo::list(self).await
2050    }
2051
2052    fn object_path(&self, object_seq: u64) -> String {
2053        ObjectStoreIo::object_path(self, object_seq)
2054    }
2055}
2056
2057#[cfg(test)]
2058mod tests {
2059    use std::sync::atomic::AtomicUsize;
2060
2061    use common_base::readable_size::ReadableSize;
2062    use common_error::ext::{ErrorExt, RetryHint};
2063    use object_store::services::Memory;
2064    use store_api::logstore::entry::{MultiplePartEntry, MultiplePartHeader};
2065    use tokio::time::timeout;
2066
2067    use super::*;
2068    use crate::error::WalObjectStoreSnafu;
2069    use crate::object_store_wal::batch::{POSITION_LIMIT, entry_id};
2070    use crate::object_store_wal::format::{FOOTER_ENTRY_LEN, decode_object};
2071
2072    const PREFIX: &str = "wal/datanodes/1/epochs/2";
2073    const WAIT: Duration = Duration::from_secs(30);
2074
2075    fn memory_store() -> ObjectStore {
2076        ObjectStore::new(Memory::default()).unwrap()
2077    }
2078
2079    fn config(flush_interval: Duration, max_batch_bytes: u64) -> ObjectStoreWalConfig {
2080        ObjectStoreWalConfig {
2081            storage_provider: String::new(),
2082            flush_interval,
2083            max_batch_bytes: ReadableSize(max_batch_bytes),
2084            ..Default::default()
2085        }
2086    }
2087
2088    /// The same batching as `config`, acknowledging appends on admission.
2089    fn enqueued(config: ObjectStoreWalConfig) -> ObjectStoreWalConfig {
2090        ObjectStoreWalConfig {
2091            ack_mode: AckMode::Enqueued,
2092            ..config
2093        }
2094    }
2095
2096    /// Every append reaches the size limit, so it is persisted on its own.
2097    fn eager() -> ObjectStoreWalConfig {
2098        config(Duration::from_secs(3600), 1)
2099    }
2100
2101    /// Nothing is persisted until a test seals the open batch.
2102    fn manual() -> ObjectStoreWalConfig {
2103        config(Duration::from_secs(3600), u64::MAX)
2104    }
2105
2106    async fn open(
2107        object_store: ObjectStore,
2108        config: &ObjectStoreWalConfig,
2109    ) -> Arc<ObjectStoreLogStore> {
2110        ObjectStoreLogStore::try_new(object_store, config, 1, 2)
2111            .await
2112            .unwrap()
2113    }
2114
2115    fn region(number: u32) -> RegionId {
2116        RegionId::new(1, number)
2117    }
2118
2119    fn provider(region_id: RegionId) -> Provider {
2120        Provider::object_store_provider(region_id, PREFIX.to_string())
2121    }
2122
2123    fn latest(store: &ObjectStoreLogStore, region_id: RegionId) -> EntryId {
2124        store.latest_entry_id(&provider(region_id)).unwrap()
2125    }
2126
2127    /// The id of the entry at `position` of its region in object `object_seq`.
2128    fn id(object_seq: u64, position: u64) -> EntryId {
2129        entry_id(object_seq, position)
2130    }
2131
2132    fn entry(store: &ObjectStoreLogStore, region_id: RegionId, data: &str) -> Entry {
2133        store
2134            .entry(data.as_bytes().to_vec(), 0, region_id, &provider(region_id))
2135            .unwrap()
2136    }
2137
2138    async fn append(
2139        store: &ObjectStoreLogStore,
2140        region_id: RegionId,
2141        data: &str,
2142    ) -> Result<AppendBatchResponse> {
2143        store
2144            .append_batch(vec![entry(store, region_id, data)])
2145            .await
2146    }
2147
2148    fn spawn_append_batch(
2149        store: &Arc<ObjectStoreLogStore>,
2150        entries: Vec<Entry>,
2151    ) -> tokio::task::JoinHandle<Result<AppendBatchResponse>> {
2152        let store = store.clone();
2153        tokio::spawn(async move { store.append_batch(entries).await })
2154    }
2155
2156    /// Appends `count` single-entry batches, waiting for each to be admitted
2157    /// before the next is sent, so they are admitted in order.
2158    async fn spawn_appends(
2159        store: &Arc<ObjectStoreLogStore>,
2160        region_id: RegionId,
2161        count: usize,
2162    ) -> Vec<tokio::task::JoinHandle<Result<AppendBatchResponse>>> {
2163        let admitted = *store.admitted_appends.borrow();
2164        let mut handles = Vec::with_capacity(count);
2165        for index in 1..=count {
2166            let entries = vec![entry(store, region_id, &format!("a{index}"))];
2167            handles.push(spawn_append_batch(store, entries));
2168            store
2169                .wait_for_admitted_appends(admitted + index)
2170                .await
2171                .unwrap();
2172        }
2173        handles
2174    }
2175
2176    async fn object_seqs(io: &dyn WalObjectIo) -> Vec<u64> {
2177        io.list()
2178            .await
2179            .unwrap()
2180            .into_iter()
2181            .map(|object| object.object_seq)
2182            .collect()
2183    }
2184
2185    fn unwrap_shared(error: &Error) -> &Error {
2186        match error {
2187            Error::ObjectStoreWal { source, .. } => source,
2188            other => panic!("expected a shared error, actual {other:?}"),
2189        }
2190    }
2191
2192    fn assert_stopped(error: &Error) {
2193        assert!(
2194            matches!(error, Error::ObjectStoreWalStopped { .. }),
2195            "unexpected error: {error:?}"
2196        );
2197    }
2198
2199    /// Waits for the next create that `ParkedIo` parked.
2200    async fn next_create(
2201        parked: &mut mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
2202    ) -> (u64, oneshot::Sender<bool>) {
2203        timeout(WAIT, parked.recv()).await.unwrap().unwrap()
2204    }
2205
2206    /// Waits for `count` parked creates and returns their releases by
2207    /// sequence.
2208    async fn parked_creates(
2209        parked: &mut mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
2210        count: usize,
2211    ) -> HashMap<u64, oneshot::Sender<bool>> {
2212        let mut releases = HashMap::new();
2213        for _ in 0..count {
2214            let (object_seq, release) = next_create(parked).await;
2215            releases.insert(object_seq, release);
2216        }
2217        releases
2218    }
2219
2220    async fn open_over(
2221        io: Arc<dyn WalObjectIo>,
2222        config: &ObjectStoreWalConfig,
2223    ) -> Arc<ObjectStoreLogStore> {
2224        ObjectStoreLogStore::open(io, config, PREFIX.to_string())
2225            .await
2226            .unwrap()
2227    }
2228
2229    async fn open_parking_creates(
2230        object_store: ObjectStore,
2231        config: &ObjectStoreWalConfig,
2232    ) -> (
2233        Arc<ObjectStoreLogStore>,
2234        Arc<ParkedIo>,
2235        mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
2236    ) {
2237        let (io, parked) = ParkedIo::parking_creates(object_store);
2238        let store = open_over(io.clone(), config).await;
2239        // The start object of the open is created without parking.
2240        io.creates_parked.store(true, Ordering::SeqCst);
2241        (store, io, parked)
2242    }
2243
2244    /// Round-trips a command through the actor, so an assertion that
2245    /// something did not happen runs after the actor handled every command
2246    /// sent before. The open batch must be empty, as under the eager config.
2247    async fn round_trip_actor(store: &ObjectStoreLogStore) {
2248        store.seal_open_batch().await.unwrap();
2249    }
2250
2251    #[tokio::test]
2252    async fn test_store_rejects_invalid_config() {
2253        for config in [
2254            config(Duration::from_millis(9), 1),
2255            config(Duration::from_secs(1), 0),
2256            ObjectStoreWalConfig {
2257                prefix: "/absolute".to_string(),
2258                ..config(Duration::from_secs(1), 1)
2259            },
2260            ObjectStoreWalConfig {
2261                max_unpersisted_bytes: ReadableSize(0),
2262                ..config(Duration::from_secs(1), 1)
2263            },
2264            ObjectStoreWalConfig {
2265                max_unpersisted_age: Duration::ZERO,
2266                ..config(Duration::from_secs(1), 1)
2267            },
2268        ] {
2269            let error = ObjectStoreLogStore::try_new(memory_store(), &config, 1, 2)
2270                .await
2271                .err()
2272                .unwrap();
2273            assert!(
2274                matches!(error, Error::InvalidWalObjectStore { .. }),
2275                "unexpected error for {config:?}: {error:?}"
2276            );
2277        }
2278    }
2279
2280    #[tokio::test]
2281    async fn test_store_recovers_only_its_node_and_generation() {
2282        let object_store = memory_store();
2283        let config = ObjectStoreWalConfig::default();
2284        let identities = [(1, 2), (3, 2), (1, 4)];
2285        for (index, (node_id, generation)) in identities.iter().enumerate() {
2286            let io = ObjectStoreIo::new(
2287                object_store.clone(),
2288                config.node_prefix(*node_id, *generation),
2289            )
2290            .unwrap();
2291            put_fixture(
2292                &io,
2293                0,
2294                &[Record {
2295                    region_id: region(1),
2296                    entry_id: index as u64 + 1,
2297                    payload: Bytes::from_static(b"entry"),
2298                }],
2299            )
2300            .await;
2301        }
2302        // A root-level object must not be part of any node's recovery.
2303        let root_io = ObjectStoreIo::new(object_store.clone(), &config.prefix).unwrap();
2304        root_io
2305            .put_if_absent(0, Bytes::from_static(b"invalid"))
2306            .await
2307            .unwrap();
2308        for (index, (node_id, generation)) in identities.iter().enumerate() {
2309            let store =
2310                ObjectStoreLogStore::try_new(object_store.clone(), &config, *node_id, *generation)
2311                    .await
2312                    .unwrap();
2313            for (other_index, (other_node, other_generation)) in identities.iter().enumerate() {
2314                let provider = Provider::object_store_provider(
2315                    region(1),
2316                    config.node_prefix(*other_node, *other_generation),
2317                );
2318                if index == other_index {
2319                    assert_eq!(store.latest_entry_id(&provider).unwrap(), index as u64 + 1);
2320                } else {
2321                    assert!(matches!(
2322                        store.latest_entry_id(&provider),
2323                        Err(Error::MismatchedWalPrefix { .. })
2324                    ));
2325                }
2326            }
2327            let root_provider = Provider::object_store_provider(region(1), config.prefix.clone());
2328            assert!(matches!(
2329                store.latest_entry_id(&root_provider),
2330                Err(Error::MismatchedWalPrefix { .. })
2331            ));
2332            store.stop().await.unwrap();
2333        }
2334    }
2335
2336    /// Rebuilds the catalog by decoding whole objects as a recovery oracle.
2337    async fn recover_by_decoding(io: &dyn WalObjectIo) -> Result<Recovered> {
2338        let mut objects = Vec::new();
2339        for object in io.list().await? {
2340            let bytes = io.get(object.object_seq).await?;
2341            let decoded = decode_object(&bytes)
2342                .and_then(|decoded| {
2343                    ensure!(
2344                        decoded.header.object_seq == object.object_seq,
2345                        CorruptedWalObjectSnafu {
2346                            reason: format!(
2347                                "header sequence {} does not match key sequence {}",
2348                                decoded.header.object_seq, object.object_seq
2349                            ),
2350                        }
2351                    );
2352                    Ok(decoded)
2353                })
2354                .with_context(|_| InvalidWalObjectSnafu {
2355                    path: object.path.clone(),
2356                })?;
2357            objects.push(FetchedObject {
2358                object,
2359                header: decoded.header,
2360                footer: decoded.footer,
2361            });
2362        }
2363        finish_recovery(objects)
2364    }
2365
2366    fn assert_same_recovery(expected: &Recovered, actual: &Recovered) {
2367        assert_eq!(
2368            catalog_contents(&expected.catalog),
2369            catalog_contents(&actual.catalog)
2370        );
2371        assert_eq!(expected.next_object_seq, actual.next_object_seq);
2372        assert_eq!(expected.durable_entry_ids, actual.durable_entry_ids);
2373        assert_eq!(expected.tip, actual.tip);
2374        assert_eq!(expected.max_epoch, actual.max_epoch);
2375    }
2376
2377    fn catalog_contents(catalog: &ObjectCatalog) -> Vec<(u64, Vec<FooterEntry>)> {
2378        catalog
2379            .objects_in_order()
2380            .map(|(object_seq, footer)| (object_seq, footer.to_vec()))
2381            .collect()
2382    }
2383
2384    async fn populate(object_store: &ObjectStore, objects: usize, regions: u32) {
2385        for object in 0..objects {
2386            let records = (1..=regions)
2387                .filter(|number| !(object + *number as usize).is_multiple_of(3))
2388                .map(|number| Record {
2389                    region_id: region(number),
2390                    entry_id: id(object as u64, 1),
2391                    payload: Bytes::from(format!("o{object}-r{number}")),
2392                })
2393                .collect::<Vec<_>>();
2394            put_records(object_store, object as u64, &records).await;
2395        }
2396    }
2397
2398    async fn put_records(object_store: &ObjectStore, object_seq: u64, records: &[Record]) {
2399        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
2400        put_fixture(&io, object_seq, records).await;
2401    }
2402
2403    /// The header of a fixture object at `object_seq`, which extends the
2404    /// present object right below it under the epoch of that object, so the
2405    /// fixtures a test writes in sequence order form one chain.
2406    async fn fixture_header(io: &ObjectStoreIo, object_seq: u64) -> Header {
2407        let below = io
2408            .list()
2409            .await
2410            .unwrap()
2411            .into_iter()
2412            .map(|object| object.object_seq)
2413            .filter(|seq| *seq < object_seq)
2414            .max();
2415        let prev = match below {
2416            Some(seq) => Some(decode_header(&io.get(seq).await.unwrap()).unwrap()),
2417            None => None,
2418        };
2419        Header {
2420            object_seq,
2421            epoch: prev.as_ref().map_or(1, |prev| prev.epoch),
2422            prev: prev.map(|prev| ChainLink {
2423                object_seq: prev.object_seq,
2424                epoch: prev.epoch,
2425            }),
2426        }
2427    }
2428
2429    async fn put_fixture(io: &ObjectStoreIo, object_seq: u64, records: &[Record]) {
2430        let header = fixture_header(io, object_seq).await;
2431        put_object_with_header(io, header, records).await;
2432    }
2433
2434    async fn put_object_with_header(io: &ObjectStoreIo, header: Header, records: &[Record]) {
2435        let object_seq = header.object_seq;
2436        let encoded = encode_object(header, records).unwrap();
2437        io.put_if_absent(object_seq, encoded.bytes).await.unwrap();
2438    }
2439
2440    /// Puts an object of another writer with `epoch` under `object_seq`. A
2441    /// store opened on an empty prefix writes epoch 1.
2442    async fn put_foreign(object_store: &ObjectStore, object_seq: u64, epoch: u64) {
2443        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
2444        let header = Header {
2445            object_seq,
2446            epoch,
2447            prev: None,
2448        };
2449        put_object_with_header(&io, header, &[]).await;
2450    }
2451
2452    fn object_path(object_store: &ObjectStore, object_seq: u64) -> String {
2453        ObjectStoreIo::new(object_store.clone(), PREFIX)
2454            .unwrap()
2455            .object_path(object_seq)
2456    }
2457
2458    async fn corrupt_object(
2459        object_store: &ObjectStore,
2460        path: &str,
2461        corrupt: impl FnOnce(&mut Vec<u8>),
2462    ) {
2463        let mut bytes = object_store.read(path).await.unwrap().to_vec();
2464        corrupt(&mut bytes);
2465        object_store.write(path, bytes).await.unwrap();
2466    }
2467
2468    fn footer_of(bytes: &[u8]) -> (FixedTrailer, Vec<FooterEntry>) {
2469        let trailer = decode_trailer(&bytes[bytes.len() - TRAILER_LEN..]).unwrap();
2470        let footer =
2471            decode_footer(&bytes[footer_range(trailer, bytes.len()).unwrap()], trailer).unwrap();
2472        (trailer, footer)
2473    }
2474
2475    fn assert_invalid_object(error: &Error, path: &str, reason: &str) {
2476        match error {
2477            Error::InvalidWalObject {
2478                path: actual,
2479                source,
2480                ..
2481            } => {
2482                assert_eq!(path, actual);
2483                match &**source {
2484                    Error::CorruptedWalObject { reason: actual, .. } => assert!(
2485                        actual.contains(reason),
2486                        "expected reason to contain {reason:?}, actual {actual:?}"
2487                    ),
2488                    other => panic!("expected a corrupted object error, actual {other:?}"),
2489                }
2490            }
2491            other => panic!("expected an invalid object error, actual {other:?}"),
2492        }
2493    }
2494
2495    #[tokio::test]
2496    async fn test_store_footer_recovery_matches_full_decode() {
2497        let object_store = memory_store();
2498        populate(&object_store, 40, 5).await;
2499        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
2500
2501        let recovered = recover(&io).await.unwrap();
2502        let expected = recover_by_decoding(&io).await.unwrap();
2503
2504        assert_eq!(40, catalog_contents(&recovered.catalog).len());
2505        assert_same_recovery(&expected, &recovered);
2506        assert_eq!(40, recovered.next_object_seq);
2507        let durable = recovered.durable_entry_ids;
2508        assert_eq!(5, durable.len());
2509
2510        let store = open(object_store, &eager()).await;
2511        for number in 1..=5 {
2512            let region_id = region(number);
2513            // The objects `populate` gave the region one entry each.
2514            let objects = (0..40u64)
2515                .filter(|object| !(object + number as u64).is_multiple_of(3))
2516                .collect::<Vec<_>>();
2517            assert_eq!(id(*objects.last().unwrap(), 1), durable[&region_id]);
2518            assert_eq!(durable[&region_id], latest(&store, region_id));
2519        }
2520    }
2521
2522    #[tokio::test]
2523    async fn test_store_recovery_rejects_corrupted_trailer_version_and_footer() {
2524        type Corrupt = fn(&mut Vec<u8>);
2525        let cases: &[(&str, Corrupt)] = &[
2526            ("truncated object", |bytes| {
2527                bytes.truncate(MIN_OBJECT_LEN - 1)
2528            }),
2529            ("invalid header magic", |bytes| bytes[0] ^= 1),
2530            ("invalid footer range", |bytes| {
2531                let start = bytes.len() - TRAILER_LEN;
2532                bytes[start..start + 8].copy_from_slice(&0u64.to_be_bytes());
2533            }),
2534            ("overflows the object", |bytes| {
2535                let start = bytes.len() - TRAILER_LEN;
2536                bytes[start..start + 8].copy_from_slice(&u64::MAX.to_be_bytes());
2537            }),
2538            ("invalid trailer magic", |bytes| {
2539                let last = bytes.len() - 1;
2540                bytes[last] ^= 1;
2541            }),
2542            ("unsupported format version 2", |bytes| {
2543                bytes[8..10].copy_from_slice(&2u16.to_be_bytes());
2544            }),
2545            ("footer checksum mismatch", |bytes| {
2546                let (trailer, _) = footer_of(bytes);
2547                bytes[trailer.footer_offset as usize] ^= 1;
2548            }),
2549        ];
2550        for (reason, corrupt) in cases {
2551            let object_store = memory_store();
2552            populate(&object_store, 3, 2).await;
2553            let path = object_path(&object_store, 1);
2554            corrupt_object(&object_store, &path, *corrupt).await;
2555
2556            let error = ObjectStoreLogStore::try_new(object_store, &eager(), 1, 2)
2557                .await
2558                .unwrap_err();
2559            assert_invalid_object(&error, &path, reason);
2560        }
2561    }
2562
2563    #[tokio::test]
2564    async fn test_store_recovery_rejects_header_sequence_mismatch() {
2565        let object_store = memory_store();
2566        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
2567        let encoded = encode_object(
2568            fixture_header(&io, 0).await,
2569            &[Record {
2570                region_id: region(1),
2571                entry_id: 1,
2572                payload: Bytes::from_static(b"a1"),
2573            }],
2574        )
2575        .unwrap();
2576        io.put_if_absent(5, encoded.bytes).await.unwrap();
2577
2578        let error = ObjectStoreLogStore::try_new(object_store, &eager(), 1, 2)
2579            .await
2580            .unwrap_err();
2581        assert_invalid_object(
2582            &error,
2583            &io.object_path(5),
2584            "header sequence 0 does not match key sequence 5",
2585        );
2586    }
2587
2588    /// Overwrites the byte range of footer entry `index` and refreshes the
2589    /// footer checksum, so the footer is intact but describes the wrong bytes.
2590    fn rewrite_segment_range(bytes: &mut [u8], index: usize, offset: u64, len: u64) {
2591        let trailer_start = bytes.len() - TRAILER_LEN;
2592        let (trailer, _) = footer_of(bytes);
2593        let footer = footer_range(trailer, bytes.len()).unwrap();
2594        let entry = footer.start + 4 + index * FOOTER_ENTRY_LEN;
2595        bytes[entry + 28..entry + 36].copy_from_slice(&offset.to_be_bytes());
2596        bytes[entry + 36..entry + 44].copy_from_slice(&len.to_be_bytes());
2597        let checksum = crc32fast::hash(&bytes[footer]);
2598        bytes[trailer_start + 16..trailer_start + 20].copy_from_slice(&checksum.to_be_bytes());
2599    }
2600
2601    #[tokio::test]
2602    async fn test_store_recovery_rejects_segment_ranges_that_do_not_tile_the_object() {
2603        type Corrupt = fn(&mut Vec<u8>, &[FooterEntry]);
2604        let cases: [(&str, Corrupt); 5] = [
2605            ("overflows the object", |bytes, footer| {
2606                rewrite_segment_range(bytes, 0, u64::MAX, footer[0].segment_len);
2607            }),
2608            ("invalid segment range", |bytes, footer| {
2609                let second = &footer[1];
2610                rewrite_segment_range(bytes, 1, second.segment_offset + 1, second.segment_len);
2611            }),
2612            ("invalid segment range", |bytes, footer| {
2613                let second = &footer[1];
2614                rewrite_segment_range(bytes, 1, second.segment_offset - 1, second.segment_len);
2615            }),
2616            ("invalid segment range", |bytes, footer| {
2617                let second = &footer[1];
2618                rewrite_segment_range(bytes, 1, second.segment_offset, second.segment_len + 1);
2619            }),
2620            ("segments end at", |bytes, footer| {
2621                let second = &footer[1];
2622                rewrite_segment_range(bytes, 1, second.segment_offset, second.segment_len - 1);
2623            }),
2624        ];
2625        for (reason, corrupt) in cases {
2626            let object_store = memory_store();
2627            populate(&object_store, 2, 3).await;
2628            let path = object_path(&object_store, 1);
2629            corrupt_object(&object_store, &path, |bytes| {
2630                let (_, footer) = footer_of(bytes);
2631                assert_eq!(2, footer.len());
2632                corrupt(bytes, &footer);
2633                // The footer itself still verifies.
2634                footer_of(bytes);
2635            })
2636            .await;
2637
2638            let error = ObjectStoreLogStore::try_new(object_store.clone(), &eager(), 1, 2)
2639                .await
2640                .unwrap_err();
2641            assert_invalid_object(&error, &path, reason);
2642            let io = ObjectStoreIo::new(object_store, PREFIX).unwrap();
2643            let error = recover_by_decoding(&io).await.unwrap_err();
2644            assert_invalid_object(&error, &path, reason);
2645        }
2646    }
2647    #[tokio::test]
2648    async fn test_store_recovery_fetches_a_footer_longer_than_the_tail_window() {
2649        let object_store = memory_store();
2650        let regions = (RECOVERY_TAIL_WINDOW / FOOTER_ENTRY_LEN + 100) as u32;
2651        let records = (1..=regions)
2652            .map(|number| Record {
2653                region_id: region(number),
2654                entry_id: 1,
2655                payload: Bytes::from_static(b"wide"),
2656            })
2657            .collect::<Vec<_>>();
2658        put_records(&object_store, 0, &records).await;
2659        put_object(&object_store, 1, region(1), &[id(1, 1)]).await;
2660
2661        let (io, reads) = RecordingIo::over(object_store.clone());
2662        let objects = io.list().await.unwrap();
2663        let wide = &objects[0];
2664        let (trailer, _) = footer_of(&object_store.read(&wide.path).await.unwrap().to_vec());
2665        assert!(trailer.footer_len > RECOVERY_TAIL_WINDOW as u64);
2666
2667        let recovered = recover(io.as_ref()).await.unwrap();
2668        let expected = recover_by_decoding(io.as_ref()).await.unwrap();
2669        assert_same_recovery(&expected, &recovered);
2670        assert_eq!(regions as usize, recovered.durable_entry_ids.len());
2671
2672        // The wide object took the header, the tail window and the footer;
2673        // the narrow one was read whole.
2674        let mut wide_reads = reads
2675            .lock()
2676            .unwrap()
2677            .iter()
2678            .filter(|(object_seq, _, _)| *object_seq == wide.object_seq)
2679            .map(|(_, offset, len)| (*offset, *len))
2680            .collect::<Vec<_>>();
2681        wide_reads.sort_unstable();
2682        assert_eq!(
2683            vec![
2684                (0, HEADER_LEN as u64),
2685                (trailer.footer_offset, trailer.footer_len),
2686                (
2687                    wide.size - RECOVERY_TAIL_WINDOW as u64,
2688                    RECOVERY_TAIL_WINDOW as u64
2689                ),
2690            ],
2691            wide_reads
2692        );
2693        assert!(
2694            reads
2695                .lock()
2696                .unwrap()
2697                .iter()
2698                .all(|(object_seq, _, _)| *object_seq == wide.object_seq)
2699        );
2700
2701        let store = open(object_store, &eager()).await;
2702        assert_eq!(id(1, 1), latest(&store, region(1)));
2703        assert_eq!(1, latest(&store, region(regions)));
2704        store.stop().await.unwrap();
2705    }
2706
2707    async fn put_object(
2708        object_store: &ObjectStore,
2709        object_seq: u64,
2710        region_id: RegionId,
2711        entry_ids: &[EntryId],
2712    ) {
2713        let records = entry_ids
2714            .iter()
2715            .map(|entry_id| Record {
2716                region_id,
2717                entry_id: *entry_id,
2718                payload: Bytes::from(format!("e{entry_id}")),
2719            })
2720            .collect::<Vec<_>>();
2721        put_records(object_store, object_seq, &records).await;
2722    }
2723
2724    #[tokio::test]
2725    async fn test_store_resumes_sequence_and_durable_ids() {
2726        let object_store = memory_store();
2727        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
2728        let recovered = recover(&io).await.unwrap();
2729        assert_eq!(0, recovered.next_object_seq);
2730        assert!(recovered.durable_entry_ids.is_empty());
2731
2732        put_object(&object_store, 2, region(1), &[1, id(5, 7)]).await;
2733        put_object(&object_store, 4, region(2), &[8]).await;
2734        let recovered = recover(&io).await.unwrap();
2735        assert_eq!(6, recovered.next_object_seq);
2736        assert_eq!(
2737            HashMap::from([(region(1), id(5, 7)), (region(2), 8)]),
2738            recovered.durable_entry_ids
2739        );
2740        let store = open(object_store.clone(), &eager()).await;
2741        assert_eq!(id(5, 7), latest(&store, region(1)));
2742        assert_eq!(8, latest(&store, region(2)));
2743        assert_eq!(0, latest(&store, region(3)));
2744
2745        // The sequence resumes above the object the largest id names, where
2746        // the start object of the store takes object 6, so the first new id
2747        // of every region is greater than every old one.
2748        let response = store
2749            .append_batch(vec![
2750                entry(&store, region(1), "a"),
2751                entry(&store, region(2), "b"),
2752            ])
2753            .await
2754            .unwrap();
2755        assert_eq!(
2756            HashMap::from([(region(1), id(7, 1)), (region(2), id(7, 1))]),
2757            response.last_entry_ids
2758        );
2759        assert_eq!(vec![2, 4, 6, 7], object_seqs(store.io.as_ref()).await);
2760        assert_eq!(
2761            expected_entries(
2762                region(1),
2763                &[
2764                    (1, "e1"),
2765                    (id(5, 7), &format!("e{}", id(5, 7))),
2766                    (id(7, 1), "a")
2767                ]
2768            ),
2769            read_entries(&store, region(1), 0).await
2770        );
2771        store.stop().await.unwrap();
2772
2773        // The sequence continues after the start object of the restart.
2774        let store = open(object_store, &eager()).await;
2775        assert_eq!(id(7, 1), latest(&store, region(1)));
2776        let response = append(&store, region(2), "b2").await.unwrap();
2777        assert_eq!(
2778            HashMap::from([(region(2), id(9, 1))]),
2779            response.last_entry_ids
2780        );
2781        assert_eq!(
2782            expected_entries(region(2), &[(8, "e8"), (id(7, 1), "b"), (id(9, 1), "b2")]),
2783            read_entries(&store, region(2), 0).await
2784        );
2785        store.stop().await.unwrap();
2786    }
2787
2788    #[tokio::test]
2789    async fn test_store_rejects_exhausted_sequence() {
2790        for (object_seq, entry_id) in [(u64::MAX, 1), (OBJECT_SEQ_LIMIT - 1, 1), (0, u64::MAX)] {
2791            let object_store = memory_store();
2792            put_object(&object_store, object_seq, region(1), &[entry_id]).await;
2793            let error = ObjectStoreLogStore::try_new(object_store, &eager(), 1, 2)
2794                .await
2795                .unwrap_err();
2796            assert!(
2797                matches!(error, Error::WalObjectSequenceExhausted { .. }),
2798                "{error:?}"
2799            );
2800        }
2801    }
2802
2803    #[tokio::test]
2804    async fn test_store_recovery_rejects_conflicting_entry_ranges() {
2805        for entries in [[1, 2], [2, 3]] {
2806            let object_store = memory_store();
2807            put_object(&object_store, 0, region(1), &[1, 2]).await;
2808            put_object(&object_store, 1, region(1), &entries).await;
2809            let error = ObjectStoreLogStore::try_new(object_store.clone(), &eager(), 1, 2)
2810                .await
2811                .unwrap_err();
2812            assert_invalid_object(&error, &object_path(&object_store, 1), "entry");
2813        }
2814    }
2815
2816    #[tokio::test]
2817    async fn test_store_rejects_foreign_providers() {
2818        let store = open(memory_store(), &eager()).await;
2819        let region_id = region(1);
2820        let other = Provider::object_store_provider(region_id, "other/prefix".to_string());
2821        let raft = Provider::raft_engine_provider(region_id.as_u64());
2822        for foreign in [&other, &raft] {
2823            let entry = Entry::Naive(NaiveEntry {
2824                provider: foreign.clone(),
2825                region_id,
2826                entry_id: 0,
2827                data: b"a1".to_vec(),
2828            });
2829            let errors = [
2830                store.append_batch(vec![entry]).await.unwrap_err(),
2831                store.latest_entry_id(foreign).unwrap_err(),
2832                store.read(foreign, 0, None).await.err().unwrap(),
2833                store.create_namespace(foreign).await.unwrap_err(),
2834                store.delete_namespace(foreign).await.unwrap_err(),
2835                store.obsolete(foreign, region_id, 1).await.unwrap_err(),
2836                store.obsolete_all(foreign, region_id).await.unwrap_err(),
2837                store.entry(Vec::new(), 1, region_id, foreign).unwrap_err(),
2838            ];
2839            for error in errors {
2840                if foreign == &raft {
2841                    assert!(matches!(error, Error::InvalidProvider { .. }), "{error:?}");
2842                } else {
2843                    assert!(
2844                        matches!(&error,
2845                        Error::MismatchedWalPrefix { expected, actual, .. }
2846                        if expected == PREFIX && actual == "other/prefix"),
2847                        "{error:?}"
2848                    );
2849                }
2850            }
2851        }
2852        store.stop().await.unwrap();
2853    }
2854
2855    #[tokio::test]
2856    async fn test_store_stop_is_idempotent() {
2857        let store = open(memory_store(), &eager()).await;
2858        let (first, second) = tokio::join!(store.stop(), store.stop());
2859        first.unwrap();
2860        second.unwrap();
2861        timeout(WAIT, store.command_tx.closed()).await.unwrap();
2862        store.stop().await.unwrap();
2863        assert!(store.stopped.load(Ordering::Acquire));
2864        // Even an empty append fails once stop began.
2865        assert_stopped(&store.append_batch(Vec::new()).await.unwrap_err());
2866    }
2867
2868    /// Builds a store whose commands the test receives instead of an actor.
2869    fn store_without_actor() -> (ObjectStoreLogStore, mpsc::Receiver<Command>) {
2870        let (command_tx, command_rx) = mpsc::channel(COMMAND_BUFFER);
2871        let (append_tx, _) = mpsc::channel(APPEND_BUFFER);
2872        let store = ObjectStoreLogStore {
2873            prefix: PREFIX.to_string(),
2874            ack_mode: AckMode::Durable,
2875            io: Arc::new(ObjectStoreIo::new(memory_store(), PREFIX).unwrap()),
2876            catalog: Arc::default(),
2877            obsolete_entry_ids: ObsoleteEntryIds::default(),
2878            terminal_error: Arc::default(),
2879            stopped: Arc::new(AtomicBool::new(false)),
2880            command_tx,
2881            append_tx,
2882            admitted_appends: watch::channel(0).1,
2883            creates_held: watch::channel(false).0,
2884            creates_fail: Arc::default(),
2885            next_create_fails_after_write: Arc::default(),
2886            parked_creates: watch::channel(0).1,
2887            durability_waits: watch::channel(0).1,
2888        };
2889        (store, command_rx)
2890    }
2891
2892    #[tokio::test]
2893    async fn test_store_stop_with_actor_gone() {
2894        let (store, command_rx) = store_without_actor();
2895        drop(command_rx);
2896        store.stop().await.unwrap();
2897        store.stop().await.unwrap();
2898    }
2899
2900    #[tokio::test]
2901    async fn test_store_actor_exits_when_the_store_is_dropped() {
2902        let store = open(memory_store(), &eager()).await;
2903        let io = Arc::downgrade(&store.io);
2904        let stopped = Arc::downgrade(&store.stopped);
2905        drop(store);
2906        timeout(WAIT, async {
2907            while stopped.strong_count() != 0 {
2908                tokio::task::yield_now().await;
2909            }
2910        })
2911        .await
2912        .unwrap();
2913        assert!(io.upgrade().is_none());
2914    }
2915
2916    #[tokio::test]
2917    async fn test_store_retains_first_terminal_error() {
2918        let store = open(memory_store(), &eager()).await;
2919        let first = set_terminal(
2920            &store.terminal_error,
2921            CorruptedWalObjectSnafu { reason: "first" }.build(),
2922        );
2923        let second = set_terminal(
2924            &store.terminal_error,
2925            CorruptedWalObjectSnafu { reason: "second" }.build(),
2926        );
2927        assert!(Arc::ptr_eq(&first, &second));
2928        let region_id = region(1);
2929        let provider = provider(region_id);
2930        let errors = [
2931            store.append_batch(Vec::new()).await.unwrap_err(),
2932            store.latest_entry_id(&provider).unwrap_err(),
2933            store.read(&provider, 0, None).await.err().unwrap(),
2934            store.create_namespace(&provider).await.unwrap_err(),
2935            store.delete_namespace(&provider).await.unwrap_err(),
2936            store.list_namespaces().await.unwrap_err(),
2937            store.obsolete(&provider, region_id, 1).await.unwrap_err(),
2938            store.obsolete_all(&provider, region_id).await.unwrap_err(),
2939            store
2940                .entry(Vec::new(), 1, region_id, &provider)
2941                .unwrap_err(),
2942        ];
2943        for error in errors {
2944            match error {
2945                Error::ObjectStoreWal { source, .. } => assert!(Arc::ptr_eq(&first, &source)),
2946                error => panic!("{error:?}"),
2947            }
2948        }
2949        assert!(store.obsolete_entry_ids.lock().unwrap().is_empty());
2950        store.stop().await.unwrap();
2951    }
2952
2953    #[tokio::test]
2954    async fn test_store_recovery_bounds_concurrency_and_orders_by_sequence() {
2955        let object_store = memory_store();
2956        populate(&object_store, 16, 3).await;
2957        let (io, mut parked) = ParkedIo::over(object_store);
2958        let objects = io.list().await.unwrap();
2959
2960        let mut fetch = {
2961            let io = io.clone();
2962            tokio::spawn(async move {
2963                fetch_footers(io.as_ref(), objects, RECOVERY_CONCURRENCY)
2964                    .await
2965                    .unwrap()
2966                    .into_iter()
2967                    .map(|fetched| fetched.object.object_seq)
2968                    .collect::<Vec<_>>()
2969            })
2970        };
2971        let mut wave = Vec::new();
2972        for _ in 0..RECOVERY_CONCURRENCY {
2973            wave.push(timeout(WAIT, parked.recv()).await.unwrap().unwrap());
2974        }
2975        assert_eq!(
2976            (0..8).collect::<Vec<_>>(),
2977            wave.iter()
2978                .map(|(object_seq, _)| *object_seq)
2979                .collect::<Vec<_>>()
2980        );
2981        // No ninth fetch can start while eight are parked.
2982        for _ in 0..16 {
2983            tokio::task::yield_now().await;
2984        }
2985        assert!(parked.try_recv().is_err());
2986
2987        // Complete the initial wave in reverse order, parking each replacement.
2988        let mut replacements = Vec::new();
2989        for (expected_next, (_, release)) in (8..16).zip(wave.into_iter().rev()) {
2990            release.send(true).unwrap();
2991            let (object_seq, release) = timeout(WAIT, parked.recv()).await.unwrap().unwrap();
2992            assert_eq!(expected_next, object_seq);
2993            assert!(!fetch.is_finished());
2994            replacements.push(release);
2995        }
2996        for release in replacements {
2997            release.send(true).unwrap();
2998        }
2999        assert_eq!(
3000            (0..16).collect::<Vec<u64>>(),
3001            timeout(WAIT, &mut fetch).await.unwrap().unwrap()
3002        );
3003        assert_eq!(8, io.max_in_flight.load(Ordering::SeqCst));
3004    }
3005
3006    #[tokio::test]
3007    async fn test_store_recovery_abandons_pending_fetches_on_failure() {
3008        let object_store = memory_store();
3009        populate(&object_store, 16, 3).await;
3010        let (io, mut parked) = ParkedIo::over(object_store.clone());
3011        let fetch = {
3012            let io = io.clone();
3013            tokio::spawn(async move {
3014                ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string()).await
3015            })
3016        };
3017        let mut wave = Vec::new();
3018        for _ in 0..RECOVERY_CONCURRENCY {
3019            wave.push(timeout(WAIT, parked.recv()).await.unwrap().unwrap());
3020        }
3021        let (failed_seq, release) = wave.pop().unwrap();
3022        release.send(false).unwrap();
3023        let error = timeout(WAIT, fetch).await.unwrap().unwrap().unwrap_err();
3024        assert!(
3025            matches!(&error, Error::WalObjectStore { operation: "read", path, .. }
3026            if path == &io.object_path(failed_seq))
3027        );
3028        assert_eq!(RetryHint::Retryable, error.retry_hint());
3029        for (_, release) in wave {
3030            assert!(release.is_closed());
3031        }
3032        assert!(parked.try_recv().is_err());
3033        let store = open(object_store, &eager()).await;
3034        assert_eq!(id(15, 1), latest(&store, region(1)));
3035        // Sequence 16 holds the object that started the epoch of the store.
3036        assert_eq!(17, store.catalog.read().unwrap().next_object_seq().unwrap());
3037        store.stop().await.unwrap();
3038    }
3039
3040    #[tokio::test]
3041    async fn test_store_recovery_reads_only_header_and_tail_of_large_object() {
3042        let object_store = memory_store();
3043        put_records(
3044            &object_store,
3045            0,
3046            &[Record {
3047                region_id: region(1),
3048                entry_id: 1,
3049                payload: Bytes::from(vec![0; RECOVERY_TAIL_WINDOW * 2]),
3050            }],
3051        )
3052        .await;
3053        let path = object_path(&object_store, 0);
3054        corrupt_object(&object_store, &path, |bytes| bytes[HEADER_LEN + 20] ^= 1).await;
3055        let (io, reads) = RecordingIo::over(object_store);
3056        let object = io.list().await.unwrap().remove(0);
3057        let store = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string())
3058            .await
3059            .unwrap();
3060        assert_eq!(1, latest(&store, region(1)));
3061        let mut reads = reads.lock().unwrap().clone();
3062        reads.sort_unstable();
3063        assert_eq!(
3064            vec![
3065                (0, 0, HEADER_LEN as u64),
3066                (
3067                    0,
3068                    object.size - RECOVERY_TAIL_WINDOW as u64,
3069                    RECOVERY_TAIL_WINDOW as u64
3070                )
3071            ],
3072            reads
3073        );
3074        store.stop().await.unwrap();
3075    }
3076
3077    #[tokio::test]
3078    async fn test_store_recovery_rejects_short_read() {
3079        let object_store = memory_store();
3080        put_object(&object_store, 0, region(1), &[1]).await;
3081        let io = ObjectStoreIo::new(object_store, PREFIX).unwrap();
3082        let mut object = io.list().await.unwrap().remove(0);
3083        object.size += 1;
3084        let error = fetch_footer(&io, &object).await.unwrap_err();
3085        assert_invalid_object(&error, &object.path, "fewer bytes than the listed");
3086    }
3087
3088    async fn read_entries(
3089        store: &ObjectStoreLogStore,
3090        region_id: RegionId,
3091        start: EntryId,
3092    ) -> Vec<Entry> {
3093        store
3094            .read(&provider(region_id), start, None)
3095            .await
3096            .unwrap()
3097            .try_collect::<Vec<_>>()
3098            .await
3099            .unwrap()
3100            .into_iter()
3101            .flatten()
3102            .collect()
3103    }
3104
3105    fn expected_entries(region_id: RegionId, entries: &[(EntryId, &str)]) -> Vec<Entry> {
3106        entries
3107            .iter()
3108            .map(|(entry_id, data)| {
3109                Entry::Naive(NaiveEntry {
3110                    provider: provider(region_id),
3111                    region_id,
3112                    entry_id: *entry_id,
3113                    data: data.as_bytes().to_vec(),
3114                })
3115            })
3116            .collect()
3117    }
3118
3119    #[tokio::test]
3120    async fn test_store_reads_only_region_segments_in_entry_order() {
3121        let object_store = memory_store();
3122        let mut expected_ranges = HashMap::<_, Vec<_>>::new();
3123        for (seq, first) in [(1, 10), (3, 30), (7, 70)] {
3124            let records = [region(2), region(1)]
3125                .into_iter()
3126                .flat_map(|region_id| {
3127                    (first..first + 3).map(move |entry_id| Record {
3128                        region_id,
3129                        entry_id,
3130                        payload: Bytes::from(format!("r{}-e{entry_id}", region_id.region_number())),
3131                    })
3132                })
3133                .collect::<Vec<_>>();
3134            put_records(&object_store, seq, &records).await;
3135            let bytes = object_store
3136                .read(&object_path(&object_store, seq))
3137                .await
3138                .unwrap()
3139                .to_vec();
3140            for entry in footer_of(&bytes).1 {
3141                expected_ranges.entry(entry.region_id).or_default().push((
3142                    seq,
3143                    entry.segment_offset,
3144                    entry.segment_len,
3145                ));
3146            }
3147        }
3148        let (io, reads) = RecordingIo::over(object_store);
3149        let store = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string())
3150            .await
3151            .unwrap();
3152        for number in [1, 2] {
3153            reads.lock().unwrap().clear();
3154            let entries = read_entries(&store, region(number), 11).await;
3155            let ids = [11, 12, 30, 31, 32, 70, 71, 72];
3156            let expected = ids
3157                .into_iter()
3158                .map(|entry_id| {
3159                    Entry::Naive(NaiveEntry {
3160                        provider: provider(region(number)),
3161                        region_id: region(number),
3162                        entry_id,
3163                        data: format!("r{number}-e{entry_id}").into_bytes(),
3164                    })
3165                })
3166                .collect::<Vec<_>>();
3167            assert_eq!(expected, entries);
3168            assert_eq!(expected_ranges[&region(number)], *reads.lock().unwrap());
3169        }
3170        reads.lock().unwrap().clear();
3171        assert_eq!(
3172            Vec::<Entry>::new(),
3173            read_entries(&store, region(1), 73).await
3174        );
3175        assert_eq!(
3176            Vec::<Entry>::new(),
3177            read_entries(&store, region(3), 0).await
3178        );
3179        assert!(reads.lock().unwrap().is_empty());
3180        assert_eq!(
3181            vec![provider(region(1)), provider(region(2))],
3182            store.list_namespaces().await.unwrap()
3183        );
3184        store.create_namespace(&provider(region(3))).await.unwrap();
3185        store.delete_namespace(&provider(region(1))).await.unwrap();
3186        assert_eq!(
3187            vec![provider(region(1)), provider(region(2))],
3188            store.list_namespaces().await.unwrap()
3189        );
3190        store.stop().await.unwrap();
3191    }
3192
3193    #[tokio::test]
3194    async fn test_store_obsolete_hides_entries_from_read_only() {
3195        let object_store = memory_store();
3196        put_object(&object_store, 0, region(1), &[10, 11, 12]).await;
3197        put_object(&object_store, 1, region(1), &[20]).await;
3198        put_object(&object_store, 2, region(2), &[10]).await;
3199        let store = open(object_store.clone(), &eager()).await;
3200        let p = provider(region(1));
3201        let all = expected_entries(
3202            region(1),
3203            &[(10, "e10"), (11, "e11"), (12, "e12"), (20, "e20")],
3204        );
3205        store.obsolete(&p, region(1), 9).await.unwrap();
3206        assert_eq!(all, read_entries(&store, region(1), 0).await);
3207        store.obsolete(&p, region(1), 11).await.unwrap();
3208        let remaining = expected_entries(region(1), &[(12, "e12"), (20, "e20")]);
3209        assert_eq!(remaining, read_entries(&store, region(1), 0).await);
3210        assert_eq!(
3211            expected_entries(region(1), &[(20, "e20")]),
3212            read_entries(&store, region(1), 20).await
3213        );
3214        assert_eq!(20, latest(&store, region(1)));
3215        store.obsolete(&p, region(1), 10).await.unwrap();
3216        assert_eq!(remaining, read_entries(&store, region(1), 0).await);
3217        store.obsolete_all(&p, region(1)).await.unwrap();
3218        store.obsolete(&p, region(1), 0).await.unwrap();
3219        assert_eq!(
3220            Vec::<Entry>::new(),
3221            read_entries(&store, region(1), 0).await
3222        );
3223        assert_eq!(20, latest(&store, region(1)));
3224        assert_eq!(
3225            expected_entries(region(2), &[(10, "e10")]),
3226            read_entries(&store, region(2), 0).await
3227        );
3228        store.stop().await.unwrap();
3229        let reopened = open(object_store, &eager()).await;
3230        assert_eq!(all, read_entries(&reopened, region(1), 0).await);
3231        reopened.obsolete(&p, region(1), 20).await.unwrap();
3232        assert_eq!(
3233            Vec::<Entry>::new(),
3234            read_entries(&reopened, region(1), 0).await
3235        );
3236        assert_eq!(20, latest(&reopened, region(1)));
3237        reopened.stop().await.unwrap();
3238    }
3239
3240    #[tokio::test]
3241    async fn test_store_entry_and_watermarks_reject_another_region() {
3242        let store = open(memory_store(), &eager()).await;
3243        let p = provider(region(2));
3244        let errors = [
3245            store
3246                .entry(b"payload".to_vec(), 7, region(1), &p)
3247                .unwrap_err(),
3248            store.obsolete(&p, region(1), 7).await.unwrap_err(),
3249            store.obsolete_all(&p, region(1)).await.unwrap_err(),
3250        ];
3251        for error in errors {
3252            assert!(
3253                matches!(&error, Error::MismatchedWalRegion { region_id, reason, .. }
3254                if *region_id == region(1) && reason == &format!("provider belongs to region {}", region(2)))
3255            );
3256            assert_eq!(
3257                common_error::status_code::StatusCode::InvalidArguments,
3258                error.status_code()
3259            );
3260        }
3261        assert!(store.obsolete_entry_ids.lock().unwrap().is_empty());
3262        assert_eq!(
3263            expected_entries(region(2), &[(7, "payload")]),
3264            vec![store.entry(b"payload".to_vec(), 7, region(2), &p).unwrap()]
3265        );
3266        assert!(store.list_namespaces().await.unwrap().is_empty());
3267        store.stop().await.unwrap();
3268    }
3269
3270    #[tokio::test]
3271    async fn test_store_corrupted_segment_fails_the_read_that_decodes_it() {
3272        let object_store = memory_store();
3273        put_records(
3274            &object_store,
3275            0,
3276            &[
3277                Record {
3278                    region_id: region(1),
3279                    entry_id: 1,
3280                    payload: Bytes::from_static(b"a1"),
3281                },
3282                Record {
3283                    region_id: region(2),
3284                    entry_id: 1,
3285                    payload: Bytes::from_static(b"b1"),
3286                },
3287            ],
3288        )
3289        .await;
3290        put_object(&object_store, 1, region(2), &[10]).await;
3291        let path = object_path(&object_store, 0);
3292        let bytes = object_store.read(&path).await.unwrap().to_vec();
3293        let footer = footer_of(&bytes).1;
3294        let corrupt = &footer[1];
3295        corrupt_object(&object_store, &path, |bytes| {
3296            bytes[(corrupt.segment_offset + corrupt.segment_len - 1) as usize] ^= 1;
3297        })
3298        .await;
3299        let (io, reads) = RecordingIo::over(object_store);
3300        let store = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string())
3301            .await
3302            .unwrap();
3303        assert_eq!(1, latest(&store, region(1)));
3304        assert_eq!(10, latest(&store, region(2)));
3305        reads.lock().unwrap().clear();
3306        assert_eq!(
3307            expected_entries(region(1), &[(1, "a1")]),
3308            read_entries(&store, region(1), 0).await
3309        );
3310        assert_eq!(
3311            vec![(0, footer[0].segment_offset, footer[0].segment_len)],
3312            *reads.lock().unwrap()
3313        );
3314        assert_eq!(
3315            expected_entries(region(2), &[(10, "e10")]),
3316            read_entries(&store, region(2), 10).await
3317        );
3318        reads.lock().unwrap().clear();
3319        let error = store
3320            .read(&provider(region(2)), 0, None)
3321            .await
3322            .unwrap()
3323            .try_collect::<Vec<_>>()
3324            .await
3325            .unwrap_err();
3326        assert_invalid_object(
3327            &error,
3328            &path,
3329            &format!("segment of region {} checksum mismatch", region(2)),
3330        );
3331        assert_eq!(
3332            vec![(0, corrupt.segment_offset, corrupt.segment_len)],
3333            *reads.lock().unwrap()
3334        );
3335        store.stop().await.unwrap();
3336    }
3337
3338    #[tokio::test]
3339    async fn test_store_assigns_ids_from_the_object_sequence() {
3340        let store = open(memory_store(), &manual()).await;
3341        let region_one = region(1);
3342        let region_two = region(2);
3343
3344        // Appends admitted into one batch share its object; the regions take
3345        // their own positions in admission order, and every append learns
3346        // the last id of each of its regions.
3347        let first = spawn_append_batch(
3348            &store,
3349            vec![
3350                entry(&store, region_one, "a1"),
3351                entry(&store, region_two, "b1"),
3352                entry(&store, region_one, "a2"),
3353            ],
3354        );
3355        store.wait_for_admitted_appends(1).await.unwrap();
3356        let second = spawn_append_batch(&store, vec![entry(&store, region_two, "b2")]);
3357        store.wait_for_admitted_appends(2).await.unwrap();
3358        assert!(!first.is_finished());
3359        assert_eq!(0, latest(&store, region_one));
3360        store.seal_open_batch().await.unwrap();
3361        assert_eq!(
3362            HashMap::from([(region_one, id(1, 2)), (region_two, id(1, 1))]),
3363            first.await.unwrap().unwrap().last_entry_ids
3364        );
3365        assert_eq!(
3366            HashMap::from([(region_two, id(1, 2))]),
3367            second.await.unwrap().unwrap().last_entry_ids
3368        );
3369
3370        // The next object starts every region at position one again.
3371        let third = spawn_append_batch(
3372            &store,
3373            vec![
3374                entry(&store, region_two, "b3"),
3375                entry(&store, region_one, "a3"),
3376            ],
3377        );
3378        store.wait_for_admitted_appends(3).await.unwrap();
3379        store.seal_open_batch().await.unwrap();
3380        assert_eq!(
3381            HashMap::from([(region_one, id(2, 1)), (region_two, id(2, 1))]),
3382            third.await.unwrap().unwrap().last_entry_ids
3383        );
3384        assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await);
3385        assert_eq!(
3386            expected_entries(
3387                region_one,
3388                &[(id(1, 1), "a1"), (id(1, 2), "a2"), (id(2, 1), "a3")]
3389            ),
3390            read_entries(&store, region_one, 1).await
3391        );
3392        assert_eq!(
3393            expected_entries(region_two, &[(id(1, 2), "b2"), (id(2, 1), "b3")]),
3394            read_entries(&store, region_two, id(1, 2)).await
3395        );
3396        assert_eq!(id(2, 1), latest(&store, region_one));
3397        assert_eq!(id(2, 1), latest(&store, region_two));
3398        store.stop().await.unwrap();
3399    }
3400
3401    #[tokio::test]
3402    async fn test_store_seals_when_a_region_exhausts_its_positions() {
3403        let store = open(memory_store(), &manual()).await;
3404        let region_id = region(1);
3405        let entries_of = |count: u64| {
3406            (0..count)
3407                .map(|_| entry(&store, region_id, ""))
3408                .collect::<Vec<_>>()
3409        };
3410
3411        // The open batch holds every position of the region.
3412        let full = spawn_append_batch(&store, entries_of(POSITION_LIMIT - 1));
3413        store.wait_for_admitted_appends(1).await.unwrap();
3414        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
3415
3416        // The next entry of the region seals the batch and opens the next
3417        // object; another region would still have fit.
3418        let next = spawn_append_batch(&store, vec![entry(&store, region_id, "next")]);
3419        store.wait_for_admitted_appends(2).await.unwrap();
3420        let response = timeout(WAIT, full).await.unwrap().unwrap().unwrap();
3421        assert_eq!(
3422            HashMap::from([(region_id, id(1, POSITION_LIMIT - 1))]),
3423            response.last_entry_ids
3424        );
3425        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
3426
3427        // An append that alone runs past the range fits no object: it seals
3428        // the open batch like any append the batch cannot take, then it is
3429        // refused without poisoning the store.
3430        let error = store
3431            .append_batch(entries_of(POSITION_LIMIT))
3432            .await
3433            .unwrap_err();
3434        assert!(
3435            matches!(error, Error::WalEntryPositionExhausted { .. }),
3436            "unexpected error: {error:?}"
3437        );
3438        let response = timeout(WAIT, next).await.unwrap().unwrap().unwrap();
3439        assert_eq!(
3440            HashMap::from([(region_id, id(2, 1))]),
3441            response.last_entry_ids
3442        );
3443        let after = spawn_append_batch(&store, vec![entry(&store, region_id, "after")]);
3444        store.wait_for_admitted_appends(3).await.unwrap();
3445        store.seal_open_batch().await.unwrap();
3446        let response = timeout(WAIT, after).await.unwrap().unwrap().unwrap();
3447        assert_eq!(
3448            HashMap::from([(region_id, id(3, 1))]),
3449            response.last_entry_ids
3450        );
3451        assert_eq!(vec![0, 1, 2, 3], object_seqs(store.io.as_ref()).await);
3452        assert_eq!(
3453            expected_entries(region_id, &[(id(2, 1), "next"), (id(3, 1), "after")]),
3454            read_entries(&store, region_id, id(2, 1)).await
3455        );
3456        store.stop().await.unwrap();
3457    }
3458
3459    #[tokio::test]
3460    async fn test_store_flushes_at_the_minimum_interval() {
3461        // An append is acknowledged once its object is durable, so a completed
3462        // append proves the tick after it sealed the batch. The appends are
3463        // spaced far apart, so every one of them lands in its own tick and
3464        // the ticks in between find an empty batch.
3465        let interval = MIN_FLUSH_INTERVAL;
3466        let store = open(memory_store(), &config(interval, u64::MAX)).await;
3467        let data = ["a1", "a2", "a3"];
3468        for (index, data) in data.iter().enumerate() {
3469            tokio::time::sleep(interval * 5).await;
3470            timeout(WAIT, append(&store, region(1), data))
3471                .await
3472                .unwrap()
3473                .unwrap();
3474            let expected_seqs = (0..=index as u64 + 1).collect::<Vec<_>>();
3475            assert_eq!(expected_seqs, object_seqs(store.io.as_ref()).await);
3476        }
3477
3478        // Ticks with an empty open batch do not create objects.
3479        tokio::time::sleep(interval * 5).await;
3480        let seqs = object_seqs(store.io.as_ref()).await;
3481        assert_eq!(vec![0, 1, 2, 3], seqs);
3482        // Object 0 is the start object of the store.
3483        for object_seq in seqs.into_iter().skip(1) {
3484            let bytes = store.io.get(object_seq).await.unwrap();
3485            let decoded = decode_object(&bytes).unwrap();
3486            assert_eq!(1, decoded.records.len(), "object {object_seq} is empty");
3487        }
3488        assert_eq!(
3489            expected_entries(
3490                region(1),
3491                &[(id(1, 1), "a1"), (id(2, 1), "a2"), (id(3, 1), "a3")]
3492            ),
3493            read_entries(&store, region(1), 1).await
3494        );
3495    }
3496
3497    #[tokio::test]
3498    async fn test_store_rejects_entries_of_another_region() {
3499        let store = open(memory_store(), &manual()).await;
3500        let region_id = region(1);
3501
3502        let entry = Entry::Naive(NaiveEntry {
3503            provider: provider(region_id),
3504            region_id: region(2),
3505            entry_id: 0,
3506            data: Vec::new(),
3507        });
3508        let error = store.append_batch(vec![entry]).await.unwrap_err();
3509        assert!(
3510            matches!(error, Error::MismatchedWalRegion { .. }),
3511            "unexpected error: {error:?}"
3512        );
3513        let incomplete = Entry::MultiplePart(MultiplePartEntry {
3514            provider: provider(region_id),
3515            region_id,
3516            entry_id: 0,
3517            headers: vec![MultiplePartHeader::First],
3518            parts: vec![b"a1".to_vec()],
3519        });
3520        let error = store.append_batch(vec![incomplete]).await.unwrap_err();
3521        assert!(
3522            matches!(error, Error::IncompleteWalEntry { region_id: actual, .. } if actual == region_id),
3523            "unexpected error: {error:?}"
3524        );
3525        assert_eq!(
3526            common_error::status_code::StatusCode::InvalidArguments,
3527            error.status_code()
3528        );
3529        assert!(
3530            store
3531                .append_batch(Vec::new())
3532                .await
3533                .unwrap()
3534                .last_entry_ids
3535                .is_empty()
3536        );
3537        // Nothing was admitted.
3538        store.seal_open_batch().await.unwrap();
3539        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
3540        store.stop().await.unwrap();
3541    }
3542
3543    #[tokio::test]
3544    async fn test_store_pipelines_creates_and_acknowledges_in_sequence_order() {
3545        let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await;
3546        let region_id = region(1);
3547        let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 2).await;
3548
3549        // At most the limit of creates run at a time; the rest wait for a slot.
3550        let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await;
3551        assert_eq!(
3552            (1..=MAX_IN_FLIGHT_CREATES as u64).collect::<BTreeSet<_>>(),
3553            releases.keys().copied().collect::<BTreeSet<_>>()
3554        );
3555        round_trip_actor(&store).await;
3556        assert!(parked.try_recv().is_err());
3557        assert!(appends.iter().all(|append| !append.is_finished()));
3558
3559        // Objects 3 and 2 become durable before object 1 and free a slot
3560        // each, but nothing is acknowledged ahead of object 1.
3561        for (released, expected_next) in [(3, 5), (2, 6)] {
3562            releases.remove(&released).unwrap().send(true).unwrap();
3563            let (object_seq, release) = next_create(&mut parked).await;
3564            assert_eq!(expected_next, object_seq);
3565            releases.insert(object_seq, release);
3566        }
3567        round_trip_actor(&store).await;
3568        assert_eq!(vec![0, 2, 3], object_seqs(io.as_ref()).await);
3569        assert!(appends.iter().all(|append| !append.is_finished()));
3570        assert_eq!(0, latest(&store, region_id));
3571
3572        // Object 1 releases the acknowledgements of objects 1 to 3.
3573        releases.remove(&1).unwrap().send(true).unwrap();
3574        let mut appends = appends.into_iter();
3575        for object_seq in 1..4 {
3576            let response = timeout(WAIT, appends.next().unwrap())
3577                .await
3578                .unwrap()
3579                .unwrap()
3580                .unwrap();
3581            assert_eq!(
3582                HashMap::from([(region_id, id(object_seq, 1))]),
3583                response.last_entry_ids
3584            );
3585        }
3586        assert_eq!(id(3, 1), latest(&store, region_id));
3587
3588        // Object 6 before object 5: the append of object 6 waits for it.
3589        releases.remove(&4).unwrap().send(true).unwrap();
3590        releases.remove(&6).unwrap().send(true).unwrap();
3591        let fourth = timeout(WAIT, appends.next().unwrap())
3592            .await
3593            .unwrap()
3594            .unwrap()
3595            .unwrap();
3596        assert_eq!(
3597            HashMap::from([(region_id, id(4, 1))]),
3598            fourth.last_entry_ids
3599        );
3600        let appends = appends.collect::<Vec<_>>();
3601        round_trip_actor(&store).await;
3602        assert!(appends.iter().all(|append| !append.is_finished()));
3603        releases.remove(&5).unwrap().send(true).unwrap();
3604        for (append, object_seq) in appends.into_iter().zip(5..) {
3605            let response = timeout(WAIT, append).await.unwrap().unwrap().unwrap();
3606            assert_eq!(
3607                HashMap::from([(region_id, id(object_seq, 1))]),
3608                response.last_entry_ids
3609            );
3610        }
3611        assert_eq!(
3612            MAX_IN_FLIGHT_CREATES,
3613            io.max_in_flight.load(Ordering::SeqCst)
3614        );
3615        assert_eq!(vec![0, 1, 2, 3, 4, 5, 6], object_seqs(io.as_ref()).await);
3616        // Every batch extends the batch sealed before it, whichever was
3617        // durable first.
3618        for object_seq in 1..=6 {
3619            let header = decode_header(&io.get(object_seq).await.unwrap()).unwrap();
3620            assert_eq!(
3621                Some(object_seq - 1),
3622                header.prev.map(|link| link.object_seq)
3623            );
3624        }
3625        assert_eq!(
3626            expected_entries(
3627                region_id,
3628                &[
3629                    (id(1, 1), "a1"),
3630                    (id(2, 1), "a2"),
3631                    (id(3, 1), "a3"),
3632                    (id(4, 1), "a4"),
3633                    (id(5, 1), "a5"),
3634                    (id(6, 1), "a6")
3635                ]
3636            ),
3637            read_entries(&store, region_id, 1).await
3638        );
3639    }
3640
3641    #[tokio::test]
3642    async fn test_store_transient_failure_fails_every_batch_that_is_not_indexed() {
3643        let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await;
3644        let region_id = region(1);
3645        // Every slot is taken; the last batch waits for one.
3646        let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 1).await;
3647        let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await;
3648        round_trip_actor(&store).await;
3649        assert!(parked.try_recv().is_err());
3650
3651        // Object 1 fails while objects 2 to 4 are in flight and object 5 waits
3652        // for a slot: every batch fails at once.
3653        releases.remove(&1).unwrap().send(false).unwrap();
3654        for append in appends {
3655            let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
3656            assert!(
3657                matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
3658                "unexpected error: {error:?}"
3659            );
3660            assert_eq!(RetryHint::Retryable, error.retry_hint());
3661        }
3662        assert_eq!(0, latest(&store, region_id));
3663
3664        // The retried entries take the sequences after the last sealed one.
3665        // The creates of the failed batches still hold their slots, so only
3666        // one retry starts until they complete; what they store is off the
3667        // chain.
3668        let retries = spawn_appends(&store, region_id, 2).await;
3669        let (object_seq, first) = next_create(&mut parked).await;
3670        assert_eq!(6, object_seq);
3671        round_trip_actor(&store).await;
3672        assert!(parked.try_recv().is_err());
3673        for release in releases.into_values() {
3674            release.send(true).unwrap();
3675        }
3676        let (object_seq, second) = next_create(&mut parked).await;
3677        assert_eq!(7, object_seq);
3678        first.send(true).unwrap();
3679        second.send(true).unwrap();
3680        for (retry, object_seq) in retries.into_iter().zip(6..) {
3681            let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap();
3682            assert_eq!(
3683                HashMap::from([(region_id, id(object_seq, 1))]),
3684                response.last_entry_ids
3685            );
3686        }
3687        assert_eq!(vec![0, 2, 3, 4, 6, 7], object_seqs(io.as_ref()).await);
3688        assert_eq!(
3689            expected_entries(region_id, &[(id(6, 1), "a1"), (id(7, 1), "a2")]),
3690            read_entries(&store, region_id, 1).await
3691        );
3692    }
3693
3694    #[tokio::test]
3695    async fn test_store_transient_failure_before_a_created_object_keeps_it_off_the_chain() {
3696        let object_store = memory_store();
3697        let (store, io, mut parked) = open_parking_creates(object_store.clone(), &eager()).await;
3698        let region_id = region(1);
3699        let appends = spawn_appends(&store, region_id, 2).await;
3700        let mut releases = parked_creates(&mut parked, 2).await;
3701
3702        // Object 2 is created while object 1 failed: both batches fail, and
3703        // the store keeps serving.
3704        releases.remove(&2).unwrap().send(true).unwrap();
3705        timeout(WAIT, async {
3706            while object_seqs(io.as_ref()).await != [0, 2] {
3707                tokio::time::sleep(Duration::from_millis(1)).await;
3708            }
3709        })
3710        .await
3711        .unwrap();
3712        releases.remove(&1).unwrap().send(false).unwrap();
3713        for append in appends {
3714            let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
3715            assert!(
3716                matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
3717                "unexpected error: {error:?}"
3718            );
3719        }
3720        assert_eq!(0, latest(&store, region_id));
3721
3722        // The retry extends the start object, not the failed batches.
3723        let retry = spawn_appends(&store, region_id, 1).await;
3724        let (object_seq, release) = next_create(&mut parked).await;
3725        assert_eq!(3, object_seq);
3726        release.send(true).unwrap();
3727        let response = timeout(WAIT, retry.into_iter().next().unwrap())
3728            .await
3729            .unwrap()
3730            .unwrap()
3731            .unwrap();
3732        assert_eq!(
3733            HashMap::from([(region_id, id(3, 1))]),
3734            response.last_entry_ids
3735        );
3736        let bytes = io.get(3).await.unwrap();
3737        let header = decode_header(&bytes).unwrap();
3738        assert_eq!(
3739            Some(ChainLink {
3740                object_seq: 0,
3741                epoch: 1,
3742            }),
3743            header.prev
3744        );
3745        store.stop().await.unwrap();
3746
3747        // After a restart only the acknowledged entry replays: object 2 is
3748        // off the chain although it exists.
3749        let store = open(object_store, &eager()).await;
3750        assert_eq!(id(3, 1), latest(&store, region_id));
3751        assert_eq!(
3752            expected_entries(region_id, &[(id(3, 1), "a1")]),
3753            read_entries(&store, region_id, 1).await
3754        );
3755    }
3756
3757    #[tokio::test]
3758    async fn test_store_permanent_failure_before_a_durable_object_poisons() {
3759        let object_store = memory_store();
3760        let (store, _, mut parked) = open_parking_creates(object_store.clone(), &eager()).await;
3761        let region_id = region(1);
3762        let foreign = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
3763        // Another writer of the same epoch took sequence 1.
3764        put_foreign(&object_store, 1, 1).await;
3765        let appends = spawn_appends(&store, region_id, 2).await;
3766        let mut releases = parked_creates(&mut parked, 2).await;
3767
3768        // Object 2 is created; object 1 conflicts with the foreign object.
3769        releases.remove(&2).unwrap().send(true).unwrap();
3770        timeout(WAIT, async {
3771            while object_seqs(&foreign).await != [0, 1, 2] {
3772                tokio::time::sleep(Duration::from_millis(1)).await;
3773            }
3774        })
3775        .await
3776        .unwrap();
3777        round_trip_actor(&store).await;
3778        assert!(appends.iter().all(|append| !append.is_finished()));
3779        releases.remove(&1).unwrap().send(true).unwrap();
3780        for append in appends {
3781            let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
3782            assert!(
3783                matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }),
3784                "unexpected error: {error:?}"
3785            );
3786        }
3787        let error = store.latest_entry_id(&provider(region_id)).unwrap_err();
3788        assert!(
3789            matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }),
3790            "unexpected error: {error:?}"
3791        );
3792        store.stop().await.unwrap();
3793    }
3794
3795    #[tokio::test]
3796    async fn test_store_conflict_with_an_earlier_epoch_fails_the_batch_without_poisoning() {
3797        let object_store = memory_store();
3798        let (store, io, mut parked) = open_parking_creates(object_store.clone(), &eager()).await;
3799        let region_id = region(1);
3800        // A late object of an earlier epoch lands under sequence 1.
3801        put_foreign(&object_store, 1, 0).await;
3802        let append = spawn_appends(&store, region_id, 1).await;
3803        next_create(&mut parked).await.1.send(true).unwrap();
3804        let error = timeout(WAIT, append.into_iter().next().unwrap())
3805            .await
3806            .unwrap()
3807            .unwrap()
3808            .unwrap_err();
3809        assert!(
3810            matches!(
3811                unwrap_shared(&error),
3812                Error::StaleWalObject {
3813                    existing_epoch: 0,
3814                    epoch: 1,
3815                    ..
3816                }
3817            ),
3818            "unexpected error: {error:?}"
3819        );
3820        assert_eq!(RetryHint::Retryable, error.retry_hint());
3821
3822        // The store keeps serving: the retry takes the next sequence.
3823        let retry = spawn_appends(&store, region_id, 1).await;
3824        let (object_seq, release) = next_create(&mut parked).await;
3825        assert_eq!(2, object_seq);
3826        release.send(true).unwrap();
3827        let response = timeout(WAIT, retry.into_iter().next().unwrap())
3828            .await
3829            .unwrap()
3830            .unwrap()
3831            .unwrap();
3832        assert_eq!(
3833            HashMap::from([(region_id, id(2, 1))]),
3834            response.last_entry_ids
3835        );
3836        assert_eq!(vec![0, 1, 2], object_seqs(io.as_ref()).await);
3837        store.stop().await.unwrap();
3838
3839        // The stale object is off the chain after a restart.
3840        let store = open(object_store, &eager()).await;
3841        assert_eq!(
3842            expected_entries(region_id, &[(id(2, 1), "a1")]),
3843            read_entries(&store, region_id, 1).await
3844        );
3845    }
3846
3847    #[tokio::test]
3848    async fn test_store_obsolete_raises_the_sequence_floor() {
3849        let store = open(memory_store(), &manual()).await;
3850        let region_one = region(1);
3851        let region_two = region(2);
3852        let mut admitted = 0;
3853        let mut spawn_append = |region_id, data: &str| {
3854            admitted += 1;
3855            (
3856                spawn_append_batch(&store, vec![entry(&store, region_id, data)]),
3857                admitted,
3858            )
3859        };
3860        let (first, count) = spawn_append(region_one, "a1");
3861        store.wait_for_admitted_appends(count).await.unwrap();
3862        store.seal_open_batch().await.unwrap();
3863        let response = timeout(WAIT, first).await.unwrap().unwrap().unwrap();
3864        assert_eq!(
3865            HashMap::from([(region_one, id(1, 1))]),
3866            response.last_entry_ids
3867        );
3868        let (second, count) = spawn_append(region_one, "a2");
3869        store.wait_for_admitted_appends(count).await.unwrap();
3870
3871        // A durable watermark names an object below the next sequence and
3872        // changes nothing: the open batch goes on under sequence 2.
3873        store
3874            .obsolete(&provider(region_one), region_one, id(1, 1))
3875            .await
3876            .unwrap();
3877        let (third, count) = spawn_append(region_one, "a3");
3878        store.wait_for_admitted_appends(count).await.unwrap();
3879
3880        // A watermark that names a later object, as one inherited from
3881        // another prefix does, cannot move the sequence while the open batch
3882        // handed out ids under it: neither the sequence nor the watermark
3883        // moves.
3884        let error = store
3885            .obsolete(&provider(region_two), region_two, id(3, 7))
3886            .await
3887            .unwrap_err();
3888        assert!(
3889            matches!(
3890                error,
3891                Error::WalObjectSequenceUnsettled { object_seq: 2, .. }
3892            ),
3893            "unexpected error: {error:?}"
3894        );
3895        assert_eq!(RetryHint::Retryable, error.retry_hint());
3896        assert!(
3897            store
3898                .obsolete_entry_ids
3899                .lock()
3900                .unwrap()
3901                .get(&region_two)
3902                .is_none()
3903        );
3904
3905        // Once the batch is durable the sequence moves above that object, so
3906        // the region's next id is greater than the watermark.
3907        store.seal_open_batch().await.unwrap();
3908        let response = timeout(WAIT, second).await.unwrap().unwrap().unwrap();
3909        assert_eq!(
3910            HashMap::from([(region_one, id(2, 1))]),
3911            response.last_entry_ids
3912        );
3913        let response = timeout(WAIT, third).await.unwrap().unwrap().unwrap();
3914        assert_eq!(
3915            HashMap::from([(region_one, id(2, 2))]),
3916            response.last_entry_ids
3917        );
3918        assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await);
3919        store
3920            .obsolete(&provider(region_two), region_two, id(3, 7))
3921            .await
3922            .unwrap();
3923        assert_eq!(
3924            Some(&id(3, 7)),
3925            store.obsolete_entry_ids.lock().unwrap().get(&region_two)
3926        );
3927        let (fourth, count) = spawn_append(region_two, "b1");
3928        store.wait_for_admitted_appends(count).await.unwrap();
3929        store.seal_open_batch().await.unwrap();
3930        let response = timeout(WAIT, fourth).await.unwrap().unwrap().unwrap();
3931        assert_eq!(
3932            HashMap::from([(region_two, id(4, 1))]),
3933            response.last_entry_ids
3934        );
3935        assert_eq!(vec![0, 1, 2, 4], object_seqs(store.io.as_ref()).await);
3936        // The watermark hides nothing of this prefix; the region has no
3937        // entry at or below it.
3938        assert_eq!(
3939            expected_entries(region_two, &[(id(4, 1), "b1")]),
3940            read_entries(&store, region_two, 1).await
3941        );
3942        store.stop().await.unwrap();
3943    }
3944
3945    #[tokio::test]
3946    async fn test_store_does_not_reuse_the_sequence_of_a_failed_create() {
3947        let (io, _) = RecordingIo::over(memory_store());
3948        let store = open_over(io.clone(), &manual()).await;
3949        let region_one = region(1);
3950        let region_two = region(2);
3951
3952        // Two appends share object 1, which is written but reported as
3953        // failed: it exists and is not indexed.
3954        let first = spawn_append_batch(&store, vec![entry(&store, region_one, "a1")]);
3955        store.wait_for_admitted_appends(1).await.unwrap();
3956        let second = spawn_append_batch(&store, vec![entry(&store, region_two, "b1")]);
3957        store.wait_for_admitted_appends(2).await.unwrap();
3958        io.fail_after_next_put.store(true, Ordering::Relaxed);
3959        store.seal_open_batch().await.unwrap_err();
3960        for append in [first, second] {
3961            timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
3962        }
3963        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
3964        assert_eq!(0, latest(&store, region_one));
3965
3966        // Retried apart, the entries land in later objects rather than
3967        // conflicting with object 1.
3968        for (region_id, data, object_seq) in [(region_one, "a1", 2), (region_two, "b1", 3)] {
3969            let retry = spawn_append_batch(&store, vec![entry(&store, region_id, data)]);
3970            store
3971                .wait_for_admitted_appends(object_seq as usize + 1)
3972                .await
3973                .unwrap();
3974            store.seal_open_batch().await.unwrap();
3975            let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap();
3976            assert_eq!(
3977                HashMap::from([(region_id, id(object_seq, 1))]),
3978                response.last_entry_ids
3979            );
3980        }
3981        assert_eq!(vec![0, 1, 2, 3], object_seqs(io.as_ref()).await);
3982        store.stop().await.unwrap();
3983
3984        // Object 2 extends the start object, so recovery leaves object 1 off
3985        // the chain: only the retries replay.
3986        let store = open_over(io, &manual()).await;
3987        assert_eq!(
3988            expected_entries(region_one, &[(id(2, 1), "a1")]),
3989            read_entries(&store, region_one, 0).await
3990        );
3991        assert_eq!(
3992            expected_entries(region_two, &[(id(3, 1), "b1")]),
3993            read_entries(&store, region_two, 0).await
3994        );
3995        store.stop().await.unwrap();
3996    }
3997
3998    #[tokio::test]
3999    async fn test_store_floor_applies_while_a_create_is_in_flight() {
4000        let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await;
4001        let region_one = region(1);
4002        let region_two = region(2);
4003        let pending = spawn_append_batch(&store, vec![entry(&store, region_one, "a1")]);
4004        let (_, release) = next_create(&mut parked).await;
4005
4006        // Object 1 is in flight and the open batch is empty: the floor moves
4007        // the next sequence, which a failure of object 1 does not move back.
4008        store
4009            .obsolete(&provider(region_two), region_two, id(5, 1))
4010            .await
4011            .unwrap();
4012        release.send(false).unwrap();
4013        timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err();
4014
4015        let write = spawn_append_batch(&store, vec![entry(&store, region_two, "b1")]);
4016        next_create(&mut parked).await.1.send(true).unwrap();
4017        let response = timeout(WAIT, write).await.unwrap().unwrap().unwrap();
4018        assert_eq!(
4019            HashMap::from([(region_two, id(6, 1))]),
4020            response.last_entry_ids
4021        );
4022        assert_eq!(vec![0, 6], object_seqs(io.as_ref()).await);
4023    }
4024
4025    #[tokio::test]
4026    async fn test_store_obsolete_queued_behind_stop_records_the_watermark() {
4027        let (store, mut command_rx) = store_without_actor();
4028        let region_id = region(1);
4029        let provider_one = provider(region_id);
4030
4031        // The actor exits with the command still queued and never answers
4032        // it: nothing is assigned an id any more, so the watermark alone holds.
4033        let (result, ()) = tokio::join!(store.obsolete(&provider_one, region_id, 1), async {
4034            let command = timeout(WAIT, command_rx.recv()).await.unwrap();
4035            assert!(matches!(
4036                command,
4037                Some(Command::Obsolete { entry_id: 1, .. })
4038            ));
4039        });
4040        result.unwrap();
4041        assert_eq!(
4042            Some(&1),
4043            store.obsolete_entry_ids.lock().unwrap().get(&region_id)
4044        );
4045
4046        // The same holds for a call that finds the actor gone.
4047        drop(command_rx);
4048        store
4049            .obsolete(&provider(region(2)), region(2), id(1, 1))
4050            .await
4051            .unwrap();
4052        assert_eq!(
4053            Some(&id(1, 1)),
4054            store.obsolete_entry_ids.lock().unwrap().get(&region(2))
4055        );
4056    }
4057
4058    #[tokio::test]
4059    async fn test_store_hooks_hold_and_fail_creates() {
4060        let store = open(memory_store(), &eager()).await;
4061        let region_id = region(1);
4062        store.hold_creates();
4063        let held = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]);
4064        store.wait_for_admitted_appends(1).await.unwrap();
4065        round_trip_actor(&store).await;
4066        assert!(!held.is_finished());
4067        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
4068        store.release_creates();
4069        timeout(WAIT, held).await.unwrap().unwrap().unwrap();
4070        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
4071
4072        store.fail_creates();
4073        let error = append(&store, region_id, "a2").await.unwrap_err();
4074        assert!(
4075            matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
4076            "unexpected error: {error:?}"
4077        );
4078        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
4079
4080        store.begin_stop();
4081        assert_stopped(&append(&store, region_id, "a3").await.unwrap_err());
4082        store.stop().await.unwrap();
4083    }
4084
4085    #[tokio::test]
4086    async fn test_store_hook_fails_the_next_create_after_it_writes() {
4087        let store = open(memory_store(), &eager()).await;
4088        let region_id = region(1);
4089
4090        // Object 1 is stored, but its create reports a failure, so the
4091        // append fails and the region has no durable entry.
4092        store.fail_next_create_after_write();
4093        let error = append(&store, region_id, "a1").await.unwrap_err();
4094        assert!(
4095            matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
4096            "unexpected error: {error:?}"
4097        );
4098        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
4099        assert_eq!(0, latest(&store, region_id));
4100
4101        // Only that create failed.
4102        let response = append(&store, region_id, "a2").await.unwrap();
4103        assert_eq!(
4104            HashMap::from([(region_id, id(2, 1))]),
4105            response.last_entry_ids
4106        );
4107        assert_eq!(vec![0, 1, 2], object_seqs(store.io.as_ref()).await);
4108        store.stop().await.unwrap();
4109    }
4110
4111    #[tokio::test]
4112    async fn test_store_hooks_observe_parked_creates_waits_and_crash() {
4113        let store = open(memory_store(), &enqueued(eager())).await;
4114        let region_id = region(1);
4115
4116        // The append is acknowledged and its create parks; a wait for its
4117        // entry has to wait.
4118        store.hold_creates();
4119        append(&store, region_id, "a1").await.unwrap();
4120        timeout(WAIT, store.wait_for_parked_creates(1))
4121            .await
4122            .unwrap()
4123            .unwrap();
4124        let wait = {
4125            let store = store.clone();
4126            tokio::spawn(async move { store.wait_durable(&provider(region_id), id(1, 1)).await })
4127        };
4128        timeout(WAIT, store.wait_for_durability_waits(1))
4129            .await
4130            .unwrap()
4131            .unwrap();
4132        assert!(!wait.is_finished());
4133
4134        // The crash fails the waiter and the store, and the parked create
4135        // never writes its object.
4136        store.crash().await;
4137        assert_stopped(&timeout(WAIT, wait).await.unwrap().unwrap().unwrap_err());
4138        assert_stopped(&append(&store, region_id, "a2").await.unwrap_err());
4139        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
4140    }
4141
4142    #[tokio::test]
4143    async fn test_store_create_held_when_the_store_is_dropped_never_runs() {
4144        let (io, _) = RecordingIo::over(memory_store());
4145        let store = open_over(io.clone(), &eager()).await;
4146        store.hold_creates();
4147        let append = spawn_append_batch(&store, vec![entry(&store, region(1), "a1")]);
4148        store.wait_for_admitted_appends(1).await.unwrap();
4149        append.abort();
4150        let _ = append.await;
4151        // A command sender outlives the store, so the actor keeps running and
4152        // the parked create completes on the closed hold channel instead of
4153        // being dropped with the actor.
4154        let command_tx = store.command_tx.clone();
4155        let terminal_error = store.terminal_error.clone();
4156        drop(store);
4157        timeout(WAIT, async {
4158            while terminal(&terminal_error).is_none() && object_seqs(io.as_ref()).await == [0] {
4159                tokio::time::sleep(Duration::from_millis(1)).await;
4160            }
4161        })
4162        .await
4163        .unwrap();
4164        assert_eq!(vec![0], object_seqs(io.as_ref()).await);
4165        let error = terminal(&terminal_error).unwrap();
4166        assert!(
4167            matches!(error.as_ref(), Error::ObjectStoreWalStopped { .. }),
4168            "unexpected error: {error:?}"
4169        );
4170        drop(command_tx);
4171    }
4172
4173    #[tokio::test]
4174    async fn test_store_rollback_and_poison_drop_the_open_batch() {
4175        let object_store = memory_store();
4176        let (store, io, mut parked) = open_parking_creates(object_store.clone(), &manual()).await;
4177        let region_id = region(1);
4178        let spawn_seal = || {
4179            let store = store.clone();
4180            tokio::spawn(async move { store.seal_open_batch().await })
4181        };
4182
4183        // Object 1 is in flight and a2 waits in the open batch under
4184        // sequence 2 when the create fails: both fail.
4185        let sealed = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]);
4186        store.wait_for_admitted_appends(1).await.unwrap();
4187        let seal = spawn_seal();
4188        let (_, release) = next_create(&mut parked).await;
4189        let open = spawn_append_batch(&store, vec![entry(&store, region_id, "a2")]);
4190        store.wait_for_admitted_appends(2).await.unwrap();
4191        release.send(false).unwrap();
4192        for append in [sealed, open] {
4193            let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
4194            assert!(
4195                matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
4196                "unexpected error: {error:?}"
4197            );
4198        }
4199        timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err();
4200
4201        // The open batch had not taken its sequence, so the next entry takes
4202        // the first id of sequence 2 again, alone.
4203        let retry = spawn_append_batch(&store, vec![entry(&store, region_id, "b1")]);
4204        store.wait_for_admitted_appends(3).await.unwrap();
4205        let seal = spawn_seal();
4206        next_create(&mut parked).await.1.send(true).unwrap();
4207        timeout(WAIT, seal).await.unwrap().unwrap().unwrap();
4208        let response = timeout(WAIT, retry).await.unwrap().unwrap().unwrap();
4209        assert_eq!(
4210            HashMap::from([(region_id, id(2, 1))]),
4211            response.last_entry_ids
4212        );
4213        assert_eq!(
4214            expected_entries(region_id, &[(id(2, 1), "b1")]),
4215            read_entries(&store, region_id, 0).await
4216        );
4217
4218        // Object 3 conflicts while c2 waits in the open batch: both fail.
4219        let sealed = spawn_append_batch(&store, vec![entry(&store, region_id, "c1")]);
4220        store.wait_for_admitted_appends(4).await.unwrap();
4221        let seal = spawn_seal();
4222        let (_, release) = next_create(&mut parked).await;
4223        let open = spawn_append_batch(&store, vec![entry(&store, region_id, "c2")]);
4224        store.wait_for_admitted_appends(5).await.unwrap();
4225        put_foreign(&object_store, 3, 1).await;
4226        release.send(true).unwrap();
4227        for append in [sealed, open] {
4228            let error = timeout(WAIT, append).await.unwrap().unwrap().unwrap_err();
4229            assert!(
4230                matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }),
4231                "unexpected error: {error:?}"
4232            );
4233        }
4234        timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err();
4235        assert_eq!(vec![0, 2, 3], object_seqs(io.as_ref()).await);
4236        store.stop().await.unwrap();
4237    }
4238
4239    #[tokio::test]
4240    async fn test_store_poisons_once_the_last_sequence_is_taken() {
4241        let object_store = memory_store();
4242        let last = OBJECT_SEQ_LIMIT - 1;
4243        // The start object of the store takes the sequence before the last.
4244        put_object(&object_store, last - 2, region(1), &[id(last - 2, 1)]).await;
4245        let store = open(object_store, &eager()).await;
4246        let next = entry(&store, region(1), "next");
4247
4248        // The batch at the last sequence is created and acknowledged, and the
4249        // store is poisoned as soon as it takes that sequence.
4250        let response = append(&store, region(1), "last").await.unwrap();
4251        assert_eq!(
4252            HashMap::from([(region(1), id(last, 1))]),
4253            response.last_entry_ids
4254        );
4255        for error in [
4256            store.latest_entry_id(&provider(region(1))).unwrap_err(),
4257            store.append_batch(vec![next]).await.unwrap_err(),
4258        ] {
4259            assert!(
4260                matches!(
4261                    unwrap_shared(&error),
4262                    Error::WalObjectSequenceExhausted { last_object_seq, .. } if *last_object_seq == last
4263                ),
4264                "unexpected error: {error:?}"
4265            );
4266        }
4267        store.stop().await.unwrap();
4268    }
4269
4270    #[tokio::test]
4271    async fn test_store_holds_appends_back_while_the_sealed_batches_are_full() {
4272        let (store, _, mut parked) = open_parking_creates(memory_store(), &eager()).await;
4273        let region_id = region(1);
4274        let admitted = || *store.admitted_appends.borrow();
4275
4276        // Every create is parked: once no two more sealed batches fit,
4277        // further appends wait in the channel and are not admitted.
4278        let mut appends = spawn_appends(&store, region_id, MAX_SEALED_BATCHES - 1).await;
4279        for index in 0..3 {
4280            let entries = vec![entry(&store, region_id, &format!("w{index}"))];
4281            appends.push(spawn_append_batch(&store, entries));
4282        }
4283        let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await;
4284        round_trip_actor(&store).await;
4285        assert_eq!(MAX_SEALED_BATCHES - 1, admitted());
4286
4287        // A durable object frees a place, and one more append is admitted.
4288        releases.remove(&1).unwrap().send(true).unwrap();
4289        timeout(WAIT, store.wait_for_admitted_appends(MAX_SEALED_BATCHES))
4290            .await
4291            .unwrap()
4292            .unwrap();
4293        round_trip_actor(&store).await;
4294        assert_eq!(MAX_SEALED_BATCHES, admitted());
4295
4296        // Stop is handled while appends are held back: the batches whose create
4297        // has not started and the appends still waiting learn of the stop.
4298        let stop = begin_spawned_stop(&store).await;
4299        for release in releases.into_values() {
4300            release.send(true).unwrap();
4301        }
4302        next_create(&mut parked).await.1.send(true).unwrap();
4303        timeout(WAIT, stop).await.unwrap().unwrap().unwrap();
4304        let mut acknowledged = 0;
4305        for append in appends {
4306            match timeout(WAIT, append).await.unwrap().unwrap() {
4307                Ok(_) => acknowledged += 1,
4308                Err(error) => assert_stopped(&error),
4309            }
4310        }
4311        assert_eq!(MAX_IN_FLIGHT_CREATES + 1, acknowledged);
4312    }
4313
4314    /// Starts `stop` and waits until it has set the stopped flag.
4315    async fn begin_spawned_stop(
4316        store: &Arc<ObjectStoreLogStore>,
4317    ) -> tokio::task::JoinHandle<Result<()>> {
4318        let stop = {
4319            let store = store.clone();
4320            tokio::spawn(async move { store.stop().await })
4321        };
4322        while !store.stopped.load(Ordering::Acquire) {
4323            tokio::task::yield_now().await;
4324        }
4325        stop
4326    }
4327
4328    #[tokio::test]
4329    async fn test_store_stop_during_failed_flush_reports_stopped() {
4330        let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await;
4331        let pending = spawn_append_batch(&store, vec![entry(&store, region(1), "a1")]);
4332        let (_, release) = next_create(&mut parked).await;
4333        let stop = begin_spawned_stop(&store).await;
4334        assert!(!stop.is_finished());
4335
4336        // The create fails after stop began: its waiter learns of the stop.
4337        release.send(false).unwrap();
4338        timeout(WAIT, stop).await.unwrap().unwrap().unwrap();
4339        assert_stopped(&timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err());
4340        assert_eq!(vec![0], object_seqs(io.as_ref()).await);
4341        assert_stopped(&append(&store, region(1), "a2").await.unwrap_err());
4342    }
4343
4344    #[tokio::test]
4345    async fn test_store_stop_with_several_creates_in_flight() {
4346        let (store, io, mut parked) = open_parking_creates(memory_store(), &eager()).await;
4347        let region_id = region(1);
4348        // Four creates are in flight, the fifth batch waits for a slot.
4349        let appends = spawn_appends(&store, region_id, MAX_IN_FLIGHT_CREATES + 1).await;
4350        let mut releases = parked_creates(&mut parked, MAX_IN_FLIGHT_CREATES).await;
4351        assert!(parked.try_recv().is_err());
4352
4353        let stop = begin_spawned_stop(&store).await;
4354        assert!(!stop.is_finished());
4355        // Nothing is assigned an id after stop began, so a watermark needs no
4356        // floor even while creates are in flight.
4357        store
4358            .obsolete(&provider(region(2)), region(2), id(9, 1))
4359            .await
4360            .unwrap();
4361        assert_eq!(
4362            Some(&id(9, 1)),
4363            store.obsolete_entry_ids.lock().unwrap().get(&region(2))
4364        );
4365
4366        // The creates in flight run to completion in any order and are
4367        // acknowledged; the batch that never started learns of the stop.
4368        for object_seq in [3, 1, 4, 2] {
4369            releases.remove(&object_seq).unwrap().send(true).unwrap();
4370        }
4371        timeout(WAIT, stop).await.unwrap().unwrap().unwrap();
4372        let mut appends = appends.into_iter();
4373        for object_seq in 1..=MAX_IN_FLIGHT_CREATES as u64 {
4374            let response = timeout(WAIT, appends.next().unwrap())
4375                .await
4376                .unwrap()
4377                .unwrap()
4378                .unwrap();
4379            assert_eq!(
4380                HashMap::from([(region_id, id(object_seq, 1))]),
4381                response.last_entry_ids
4382            );
4383        }
4384        let error = timeout(WAIT, appends.next().unwrap())
4385            .await
4386            .unwrap()
4387            .unwrap()
4388            .unwrap_err();
4389        assert_stopped(&error);
4390        // No create was started after stop began.
4391        assert!(parked.try_recv().is_err());
4392        assert_eq!(vec![0, 1, 2, 3, 4], object_seqs(io.as_ref()).await);
4393        assert_eq!(id(4, 1), latest(&store, region_id));
4394        timeout(WAIT, store.command_tx.closed()).await.unwrap();
4395    }
4396
4397    #[tokio::test]
4398    async fn test_store_stop_during_conflicting_flush_reports_stopped_and_poisons() {
4399        let object_store = memory_store();
4400        let (store, _, mut parked) = open_parking_creates(object_store.clone(), &eager()).await;
4401        let region_id = region(1);
4402        let pending = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]);
4403        let (_, release) = next_create(&mut parked).await;
4404        let stop = begin_spawned_stop(&store).await;
4405        // Another writer of the same epoch takes the sequence before the
4406        // create runs.
4407        put_foreign(&object_store, 1, 1).await;
4408
4409        release.send(true).unwrap();
4410        timeout(WAIT, stop).await.unwrap().unwrap().unwrap();
4411        assert_stopped(&timeout(WAIT, pending).await.unwrap().unwrap().unwrap_err());
4412        for error in [
4413            store.latest_entry_id(&provider(region_id)).unwrap_err(),
4414            store
4415                .read(&provider(region_id), 1, None)
4416                .await
4417                .err()
4418                .unwrap(),
4419        ] {
4420            assert!(
4421                matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }),
4422                "unexpected error: {error:?}"
4423            );
4424        }
4425    }
4426
4427    fn injected_failure<T>(operation: &'static str, path: String) -> Result<T> {
4428        Err(
4429            object_store::Error::new(object_store::ErrorKind::Unexpected, "injected failure")
4430                .set_temporary(),
4431        )
4432        .context(WalObjectStoreSnafu { operation, path })
4433    }
4434
4435    fn chain_header(object_seq: u64, epoch: u64, prev: Option<(u64, u64)>) -> Header {
4436        Header {
4437            object_seq,
4438            epoch,
4439            prev: prev.map(|(object_seq, epoch)| ChainLink { object_seq, epoch }),
4440        }
4441    }
4442
4443    fn chain_of(headers: &[Header]) -> Vec<u64> {
4444        select_chain(
4445            &headers
4446                .iter()
4447                .map(|header| (header.object_seq, header))
4448                .collect(),
4449        )
4450    }
4451
4452    #[test]
4453    fn test_store_selects_the_chain_of_the_latest_complete_object() {
4454        assert!(chain_of(&[]).is_empty());
4455        // A linear chain is replayed whole.
4456        assert_eq!(
4457            vec![0, 1, 2],
4458            chain_of(&[
4459                chain_header(0, 1, None),
4460                chain_header(1, 1, Some((0, 1))),
4461                chain_header(2, 1, Some((1, 1))),
4462            ])
4463        );
4464        // Object 1 was reported as failed and object 2 links past it.
4465        assert_eq!(
4466            vec![0, 2],
4467            chain_of(&[
4468                chain_header(0, 1, None),
4469                chain_header(1, 1, Some((0, 1))),
4470                chain_header(2, 1, Some((0, 1))),
4471            ])
4472        );
4473        // Object 3 extends object 2, which never landed, so object 1 is the
4474        // tip although object 3 has a higher sequence.
4475        assert_eq!(
4476            vec![0, 1],
4477            chain_of(&[
4478                chain_header(0, 1, None),
4479                chain_header(1, 1, Some((0, 1))),
4480                chain_header(3, 1, Some((2, 1))),
4481            ])
4482        );
4483        // Object 1 links to an object of epoch 2 at sequence 0, where an
4484        // object of epoch 1 is.
4485        assert_eq!(
4486            vec![0],
4487            chain_of(&[chain_header(0, 1, None), chain_header(1, 2, Some((0, 2))),])
4488        );
4489        // Instance A left object 2 behind after instance B, a later epoch,
4490        // acknowledged object 1: the later epoch ends the chain.
4491        assert_eq!(
4492            vec![0, 1],
4493            chain_of(&[
4494                chain_header(0, 1, None),
4495                chain_header(1, 2, Some((0, 1))),
4496                chain_header(2, 1, Some((0, 1))),
4497            ])
4498        );
4499        // A missing predecessor breaks the link, whether or not it lies below
4500        // every present object, so an object that lands late below it cannot
4501        // change which links hold.
4502        assert!(
4503            chain_of(&[
4504                chain_header(3, 1, Some((2, 1))),
4505                chain_header(4, 1, Some((3, 1))),
4506            ])
4507            .is_empty()
4508        );
4509        assert!(
4510            chain_of(&[
4511                chain_header(2, 1, Some((1, 1))),
4512                chain_header(4, 2, Some((3, 1))),
4513                chain_header(5, 2, Some((4, 2))),
4514            ])
4515            .is_empty()
4516        );
4517        assert_eq!(
4518            vec![1],
4519            chain_of(&[chain_header(1, 1, None), chain_header(3, 1, Some((2, 1))),])
4520        );
4521        // A predecessor at or above the object itself never holds.
4522        assert!(chain_of(&[chain_header(5, 1, Some((5, 1)))]).is_empty());
4523    }
4524
4525    #[tokio::test]
4526    async fn test_store_recovery_indexes_only_the_chain() {
4527        let object_store = memory_store();
4528        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4529        let record = |entry_id: EntryId, data: &'static str| Record {
4530            region_id: region(1),
4531            entry_id,
4532            payload: Bytes::from_static(data.as_bytes()),
4533        };
4534        put_fixture(&io, 0, &[record(1, "a")]).await;
4535        let first = fixture_header(&io, 1).await;
4536        // Object 1 was reported as failed; object 2 reused its row sequences
4537        // and links past it.
4538        put_object_with_header(&io, first.clone(), &[record(id(1, 1), "old")]).await;
4539        let second = Header {
4540            object_seq: 2,
4541            ..first
4542        };
4543        put_object_with_header(&io, second, &[record(id(2, 1), "new")]).await;
4544
4545        for restart in 0..2 {
4546            let store = open(object_store.clone(), &eager()).await;
4547            assert_eq!(
4548                expected_entries(region(1), &[(1, "a"), (id(2, 1), "new")]),
4549                read_entries(&store, region(1), 0).await
4550            );
4551            assert_eq!(id(2, 1), latest(&store, region(1)));
4552            store.stop().await.unwrap();
4553            // Every open starts the epoch one above the sequence of its start
4554            // object, and the first restart's start object extends the chain
4555            // the second restart replays.
4556            let start = decode_header(&io.get(3 + restart).await.unwrap()).unwrap();
4557            assert_eq!(4 + restart, start.epoch);
4558            assert_eq!(Some(2 + restart), start.prev.map(|prev| prev.object_seq));
4559        }
4560        assert_eq!(vec![0, 1, 2, 3, 4], object_seqs(&io).await);
4561    }
4562
4563    #[tokio::test]
4564    async fn test_store_recovery_never_reuses_the_sequence_of_an_orphan() {
4565        let io = ObjectStoreIo::new(memory_store(), PREFIX).unwrap();
4566        put_fixture(&io, 0, &[]).await;
4567        // Object 2 extends object 1, which never landed.
4568        let orphan = Header {
4569            prev: Some(ChainLink {
4570                object_seq: 1,
4571                epoch: 1,
4572            }),
4573            ..fixture_header(&io, 2).await
4574        };
4575        put_object_with_header(&io, orphan, &[]).await;
4576        let recovered = recover(&io).await.unwrap();
4577        assert_eq!(Some(0), recovered.tip.map(|tip| tip.object_seq));
4578        assert_eq!(3, recovered.next_object_seq);
4579    }
4580
4581    #[tokio::test]
4582    async fn test_store_recovery_rejects_objects_without_a_complete_chain() {
4583        let object_store = memory_store();
4584        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4585        // Object 1 extends object 0, which is gone.
4586        put_object_with_header(&io, chain_header(1, 1, Some((0, 1))), &[]).await;
4587        let error = ObjectStoreLogStore::try_new(object_store, &eager(), 1, 2)
4588            .await
4589            .unwrap_err();
4590        assert!(
4591            matches!(&error, Error::CorruptedWalObject { reason, .. }
4592                if reason.contains("no object among 1 present objects completes a chain")),
4593            "unexpected error: {error:?}"
4594        );
4595        assert_eq!(vec![1], object_seqs(&io).await);
4596    }
4597
4598    #[tokio::test]
4599    async fn test_store_recovery_keeps_the_later_epoch_over_a_late_object() {
4600        let object_store = memory_store();
4601        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4602        let record = |entry_id: EntryId, data: &'static str| Record {
4603            region_id: region(1),
4604            entry_id,
4605            payload: Bytes::from_static(data.as_bytes()),
4606        };
4607        put_object_with_header(&io, chain_header(0, 1, None), &[record(1, "a")]).await;
4608        // Instance B acknowledged object 1 before the create instance A issued
4609        // for object 2 landed.
4610        put_object_with_header(
4611            &io,
4612            chain_header(1, 2, Some((0, 1))),
4613            &[record(id(1, 1), "b")],
4614        )
4615        .await;
4616        put_object_with_header(
4617            &io,
4618            chain_header(2, 1, Some((0, 1))),
4619            &[record(id(2, 1), "late")],
4620        )
4621        .await;
4622
4623        let store = open(object_store, &eager()).await;
4624        assert_eq!(
4625            expected_entries(region(1), &[(1, "a"), (id(1, 1), "b")]),
4626            read_entries(&store, region(1), 0).await
4627        );
4628        store.stop().await.unwrap();
4629        let start = decode_header(&io.get(3).await.unwrap()).unwrap();
4630        assert_eq!(4, start.epoch);
4631        assert_eq!(
4632            Some(ChainLink {
4633                object_seq: 1,
4634                epoch: 2,
4635            }),
4636            start.prev
4637        );
4638    }
4639
4640    #[tokio::test]
4641    async fn test_store_open_starts_the_first_epoch_on_an_empty_prefix() {
4642        let object_store = memory_store();
4643        let store = open(object_store.clone(), &eager()).await;
4644        store.stop().await.unwrap();
4645        let io = ObjectStoreIo::new(object_store, PREFIX).unwrap();
4646        let start = decode_header(&io.get(0).await.unwrap()).unwrap();
4647        assert_eq!(1, start.epoch);
4648        assert_eq!(None, start.prev);
4649        assert_eq!(vec![0], object_seqs(&io).await);
4650    }
4651
4652    #[tokio::test]
4653    async fn test_store_start_object_moves_past_an_earlier_epoch_only() {
4654        for (epoch, moves) in [(1, true), (2, false)] {
4655            let object_store = memory_store();
4656            let fixtures = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4657            put_fixture(&fixtures, 0, &[]).await;
4658            // An object lands at the sequence of the start object after
4659            // recovery listed the prefix.
4660            let late = encode_object(chain_header(1, epoch, None), &[])
4661                .unwrap()
4662                .bytes;
4663            let io = Arc::new(RacingIo {
4664                inner: ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap(),
4665                late: Mutex::new(Some((1, late))),
4666            });
4667            let result = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string()).await;
4668            if moves {
4669                result.unwrap().stop().await.unwrap();
4670                let start = decode_header(&fixtures.get(2).await.unwrap()).unwrap();
4671                assert_eq!(3, start.epoch);
4672                assert_eq!(Some(0), start.prev.map(|prev| prev.object_seq));
4673            } else {
4674                let error = result.unwrap_err();
4675                assert!(
4676                    matches!(&error, Error::WalObjectConflict { path, .. } if path == &fixtures.object_path(1)),
4677                    "unexpected error: {error:?}"
4678                );
4679                assert_eq!(vec![0, 1], object_seqs(&fixtures).await);
4680            }
4681        }
4682    }
4683
4684    #[tokio::test]
4685    async fn test_store_open_fails_when_its_start_object_is_already_present() {
4686        // On an empty prefix and after object 0 of epoch 1, another open that
4687        // recovered the same objects wrote the very start object this open
4688        // writes.
4689        for present in [false, true] {
4690            let object_store = memory_store();
4691            let fixtures = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4692            let (start_seq, epoch, tip) = if present {
4693                put_fixture(&fixtures, 0, &[]).await;
4694                (1, 2, Some((0, 1)))
4695            } else {
4696                (0, 1, None)
4697            };
4698            let same = encode_object(chain_header(start_seq, epoch, tip), &[])
4699                .unwrap()
4700                .bytes;
4701            let io = Arc::new(RacingIo {
4702                inner: ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap(),
4703                late: Mutex::new(Some((start_seq, same))),
4704            });
4705            let error = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string())
4706                .await
4707                .unwrap_err();
4708            assert!(
4709                matches!(&error, Error::UnconfirmedWalEpochStart { epoch: actual, path, .. }
4710                    if *actual == epoch && path == &fixtures.object_path(start_seq)),
4711                "unexpected error: {error:?}"
4712            );
4713            assert_eq!(RetryHint::Retryable, error.retry_hint());
4714
4715            // The next open counts that object and starts a later epoch.
4716            open(object_store, &eager()).await.stop().await.unwrap();
4717            let start = decode_header(&fixtures.get(start_seq + 1).await.unwrap()).unwrap();
4718            assert_eq!(epoch + 1, start.epoch);
4719            assert_eq!(
4720                Some(ChainLink {
4721                    object_seq: start_seq,
4722                    epoch,
4723                }),
4724                start.prev
4725            );
4726        }
4727    }
4728
4729    #[tokio::test]
4730    async fn test_store_opens_racing_across_a_late_object_take_distinct_epochs() {
4731        let object_store = memory_store();
4732        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4733        put_fixture(&io, 0, &[]).await;
4734        // Open A recovers object 0 and pauses before its start object.
4735        let recovered = recover(&io).await.unwrap();
4736        // A late object of epoch 1 lands at sequence 2, across the gap at
4737        // sequence 1, and open B, which lists it, completes.
4738        put_object_with_header(&io, chain_header(2, 1, Some((0, 1))), &[]).await;
4739        open(object_store, &eager()).await.stop().await.unwrap();
4740        let b = decode_header(&io.get(3).await.unwrap()).unwrap();
4741        assert_eq!(4, b.epoch);
4742
4743        // Open A resumes at the sequence it recovered: the epoch follows from
4744        // the sequence it claims, so it differs from the epoch of open B.
4745        let a = start_epoch(&io, recovered.next_object_seq, recovered.tip)
4746            .await
4747            .unwrap();
4748        assert_eq!(
4749            ChainLink {
4750                object_seq: 1,
4751                epoch: 2,
4752            },
4753            a
4754        );
4755    }
4756
4757    #[tokio::test]
4758    async fn test_store_open_rejects_an_epoch_above_the_next_sequence() {
4759        let object_store = memory_store();
4760        let io = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4761        // No instance that starts at or below sequence 0 writes epoch 5.
4762        put_object_with_header(&io, chain_header(0, 5, None), &[]).await;
4763        let error = ObjectStoreLogStore::try_new(object_store, &eager(), 1, 2)
4764            .await
4765            .unwrap_err();
4766        assert!(
4767            matches!(&error, Error::CorruptedWalObject { reason, .. }
4768                if reason.contains("carries epoch 5, above the next sequence 1")),
4769            "unexpected error: {error:?}"
4770        );
4771        assert_eq!(vec![0], object_seqs(&io).await);
4772    }
4773
4774    #[tokio::test]
4775    async fn test_store_open_fails_when_its_start_object_lands_with_an_unknown_outcome() {
4776        let object_store = memory_store();
4777        let fixtures = ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap();
4778        put_fixture(&fixtures, 0, &[]).await;
4779        // The start object is stored, but its create reports an error.
4780        let io = Arc::new(LostResponseIo {
4781            inner: ObjectStoreIo::new(object_store.clone(), PREFIX).unwrap(),
4782            lose_next: AtomicBool::new(true),
4783        });
4784        let error = ObjectStoreLogStore::open(io, &eager(), PREFIX.to_string())
4785            .await
4786            .unwrap_err();
4787        assert!(
4788            matches!(error, Error::WalObjectStore { .. }),
4789            "unexpected error: {error:?}"
4790        );
4791        // The open does not move to another sequence.
4792        assert_eq!(vec![0, 1], object_seqs(&fixtures).await);
4793        let stored = decode_header(&fixtures.get(1).await.unwrap()).unwrap();
4794        assert_eq!(2, stored.epoch);
4795
4796        // The next open counts the stored start object and claims a later
4797        // epoch.
4798        open(object_store, &eager()).await.stop().await.unwrap();
4799        let start = decode_header(&fixtures.get(2).await.unwrap()).unwrap();
4800        assert_eq!(3, start.epoch);
4801        assert_eq!(
4802            Some(ChainLink {
4803                object_seq: 1,
4804                epoch: 2,
4805            }),
4806            start.prev
4807        );
4808    }
4809
4810    #[tokio::test]
4811    async fn test_store_enqueued_append_returns_before_the_object_exists() {
4812        let store = open(memory_store(), &enqueued(manual())).await;
4813        let region_id = region(1);
4814
4815        let response = timeout(WAIT, append(&store, region_id, "a1"))
4816            .await
4817            .unwrap()
4818            .unwrap();
4819        assert_eq!(
4820            HashMap::from([(region_id, id(1, 1))]),
4821            response.last_entry_ids
4822        );
4823        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
4824        assert_eq!(0, latest(&store, region_id));
4825        assert_eq!(0, store.durable_entry_id(&provider(region_id)).unwrap());
4826        assert!(read_entries(&store, region_id, 1).await.is_empty());
4827        let response = append(&store, region_id, "a2").await.unwrap();
4828        assert_eq!(
4829            HashMap::from([(region_id, id(1, 2))]),
4830            response.last_entry_ids
4831        );
4832
4833        store.seal_open_batch().await.unwrap();
4834        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
4835        assert_eq!(id(1, 2), latest(&store, region_id));
4836        assert_eq!(
4837            id(1, 2),
4838            store.durable_entry_id(&provider(region_id)).unwrap()
4839        );
4840        assert_eq!(
4841            expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]),
4842            read_entries(&store, region_id, 1).await
4843        );
4844    }
4845
4846    #[tokio::test]
4847    async fn test_store_enqueued_durable_id_advances_after_create_and_indexing() {
4848        let (store, io, mut parked) =
4849            open_parking_creates(memory_store(), &enqueued(eager())).await;
4850        let region_one = region(1);
4851        let region_two = region(2);
4852
4853        let response = timeout(WAIT, append(&store, region_one, "a1"))
4854            .await
4855            .unwrap()
4856            .unwrap();
4857        assert_eq!(
4858            HashMap::from([(region_one, id(1, 1))]),
4859            response.last_entry_ids
4860        );
4861        let (object_seq, release) = next_create(&mut parked).await;
4862        assert_eq!(1, object_seq);
4863        assert_eq!(0, store.durable_entry_id(&provider(region_one)).unwrap());
4864        let wait = |entry_id| {
4865            let store = store.clone();
4866            tokio::spawn(async move { store.wait_durable(&provider(region_one), entry_id).await })
4867        };
4868        // An id the store never handed out waits for the ids of the region
4869        // that were handed out below it.
4870        let waits = [wait(id(1, 1)), wait(id(1, 7))];
4871        // The other region has nothing to wait for.
4872        timeout(WAIT, store.wait_durable(&provider(region_two), id(1, 1)))
4873            .await
4874            .unwrap()
4875            .unwrap();
4876        for _ in 0..16 {
4877            tokio::task::yield_now().await;
4878        }
4879        assert!(waits.iter().all(|wait| !wait.is_finished()));
4880
4881        release.send(true).unwrap();
4882        for wait in waits {
4883            timeout(WAIT, wait).await.unwrap().unwrap().unwrap();
4884        }
4885        assert_eq!(
4886            id(1, 1),
4887            store.durable_entry_id(&provider(region_one)).unwrap()
4888        );
4889        assert_eq!(0, store.durable_entry_id(&provider(region_two)).unwrap());
4890        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
4891    }
4892
4893    #[tokio::test]
4894    async fn test_store_enqueued_wait_ignores_later_pending_entries_of_the_region() {
4895        let (store, _io, mut parked) =
4896            open_parking_creates(memory_store(), &enqueued(eager())).await;
4897        let region_a = region(1);
4898        let region_b = region(2);
4899        // Object 1 holds region A's id(1, 1) and object 2 region B's id(2, 1).
4900        for (region_id, data) in [(region_a, "a1"), (region_b, "b1")] {
4901            let response = append(&store, region_id, data).await.unwrap();
4902            let (_, release) = next_create(&mut parked).await;
4903            release.send(true).unwrap();
4904            let entry_id = response.last_entry_ids[&region_id];
4905            timeout(WAIT, store.wait_durable(&provider(region_id), entry_id))
4906                .await
4907                .unwrap()
4908                .unwrap();
4909        }
4910        let response = append(&store, region_a, "a2").await.unwrap();
4911        assert_eq!(
4912            HashMap::from([(region_a, id(3, 1))]),
4913            response.last_entry_ids
4914        );
4915        let (object_seq, release) = next_create(&mut parked).await;
4916        assert_eq!(3, object_seq);
4917
4918        // Every entry of region A up to id(2, 1) is durable, so the wait does
4919        // not depend on the parked create of object 3.
4920        timeout(WAIT, store.wait_durable(&provider(region_a), id(2, 1)))
4921            .await
4922            .unwrap()
4923            .unwrap();
4924        // A wait for id(3, 1) itself is answered only once object 3 is
4925        // indexed; the round trip orders the check after the actor handled it.
4926        let (response_tx, mut response_rx) = oneshot::channel();
4927        store
4928            .command_tx
4929            .send(Command::WaitDurable {
4930                region_id: region_a,
4931                entry_id: id(3, 1),
4932                response: response_tx,
4933            })
4934            .await
4935            .unwrap();
4936        round_trip_actor(&store).await;
4937        assert!(matches!(
4938            response_rx.try_recv(),
4939            Err(oneshot::error::TryRecvError::Empty)
4940        ));
4941        release.send(true).unwrap();
4942        timeout(WAIT, response_rx).await.unwrap().unwrap().unwrap();
4943    }
4944
4945    #[tokio::test]
4946    async fn test_store_enqueued_obsolete_never_passes_the_durable_id() {
4947        let store = open(memory_store(), &enqueued(manual())).await;
4948        let region_id = region(1);
4949        append(&store, region_id, "a1").await.unwrap();
4950        append(&store, region_id, "a2").await.unwrap();
4951
4952        store
4953            .obsolete(&provider(region_id), region_id, id(1, 2))
4954            .await
4955            .unwrap();
4956        assert_eq!(
4957            Some(&0),
4958            store.obsolete_entry_ids.lock().unwrap().get(&region_id)
4959        );
4960        store.seal_open_batch().await.unwrap();
4961        assert_eq!(
4962            expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]),
4963            read_entries(&store, region_id, 1).await
4964        );
4965        store
4966            .obsolete(&provider(region_id), region_id, id(1, 2))
4967            .await
4968            .unwrap();
4969        assert_eq!(
4970            Some(&id(1, 2)),
4971            store.obsolete_entry_ids.lock().unwrap().get(&region_id)
4972        );
4973        assert!(read_entries(&store, region_id, 1).await.is_empty());
4974    }
4975
4976    /// Appends one entry in the enqueued mode with `config`, whose create
4977    /// parks, then appends a second one that the backlog threshold must
4978    /// stall. Returns the store, the parked creates, the stalled append and
4979    /// the release of the first create.
4980    async fn stall_second_append(
4981        config: ObjectStoreWalConfig,
4982    ) -> (
4983        Arc<ObjectStoreLogStore>,
4984        Arc<ParkedIo>,
4985        mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
4986        tokio::task::JoinHandle<Result<AppendBatchResponse>>,
4987        oneshot::Sender<bool>,
4988    ) {
4989        let (store, io, mut parked) = open_parking_creates(memory_store(), &config).await;
4990        let region_id = region(1);
4991        let response = timeout(WAIT, append(&store, region_id, "a1"))
4992            .await
4993            .unwrap()
4994            .unwrap();
4995        assert_eq!(
4996            HashMap::from([(region_id, id(1, 1))]),
4997            response.last_entry_ids
4998        );
4999        if config.max_unpersisted_age < Duration::from_secs(1) {
5000            tokio::time::sleep(config.max_unpersisted_age * 2).await;
5001        }
5002
5003        let stalled = spawn_append_batch(&store, vec![entry(&store, region_id, "a2")]);
5004        // The stall seals the open batch so that an upload is in flight.
5005        let (object_seq, release) = next_create(&mut parked).await;
5006        assert_eq!(1, object_seq);
5007        for _ in 0..16 {
5008            tokio::task::yield_now().await;
5009        }
5010        assert!(!stalled.is_finished());
5011        assert!(parked.try_recv().is_err());
5012        (store, io, parked, stalled, release)
5013    }
5014
5015    #[tokio::test]
5016    async fn test_store_enqueued_backlog_bytes_stall_admission_until_an_upload_completes() {
5017        let config = ObjectStoreWalConfig {
5018            max_unpersisted_bytes: ReadableSize(1),
5019            ..enqueued(manual())
5020        };
5021        let (store, io, mut parked, stalled, release) = stall_second_append(config).await;
5022        let region_id = region(1);
5023        // Two more appends queue behind the stalled one.
5024        let third = spawn_append_batch(&store, vec![entry(&store, region_id, "a3")]);
5025        let fourth = spawn_append_batch(&store, vec![entry(&store, region_id, "a4")]);
5026        for _ in 0..16 {
5027            tokio::task::yield_now().await;
5028        }
5029        assert!(!third.is_finished() && !fourth.is_finished());
5030        assert!(parked.try_recv().is_err());
5031
5032        // The upload releases the second append, whose entry reaches the
5033        // threshold again; the stall seals it so that the next upload can
5034        // release the third, and so on.
5035        release.send(true).unwrap();
5036        let response = timeout(WAIT, stalled).await.unwrap().unwrap().unwrap();
5037        assert_eq!(
5038            HashMap::from([(region_id, id(2, 1))]),
5039            response.last_entry_ids
5040        );
5041        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
5042        assert_eq!(id(1, 1), latest(&store, region_id));
5043        for (append, object_seq) in [(third, 2), (fourth, 3)] {
5044            let (parked_seq, release) = next_create(&mut parked).await;
5045            assert_eq!(object_seq, parked_seq);
5046            release.send(true).unwrap();
5047            let response = timeout(WAIT, append).await.unwrap().unwrap().unwrap();
5048            assert_eq!(
5049                HashMap::from([(region_id, id(object_seq + 1, 1))]),
5050                response.last_entry_ids
5051            );
5052        }
5053        assert_eq!(vec![0, 1, 2, 3], object_seqs(io.as_ref()).await);
5054        assert_eq!(
5055            expected_entries(
5056                region_id,
5057                &[(id(1, 1), "a1"), (id(2, 1), "a2"), (id(3, 1), "a3")]
5058            ),
5059            read_entries(&store, region_id, 1).await
5060        );
5061    }
5062
5063    #[tokio::test]
5064    async fn test_store_enqueued_backlog_age_stalls_admission_until_an_upload_completes() {
5065        let config = ObjectStoreWalConfig {
5066            max_unpersisted_age: Duration::from_millis(50),
5067            ..enqueued(manual())
5068        };
5069        let (store, io, _parked, stalled, release) = stall_second_append(config).await;
5070        let region_id = region(1);
5071
5072        release.send(true).unwrap();
5073        let response = timeout(WAIT, stalled).await.unwrap().unwrap().unwrap();
5074        assert_eq!(
5075            HashMap::from([(region_id, id(2, 1))]),
5076            response.last_entry_ids
5077        );
5078        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
5079        assert_eq!(
5080            expected_entries(region_id, &[(id(1, 1), "a1")]),
5081            read_entries(&store, region_id, 1).await
5082        );
5083    }
5084
5085    #[tokio::test]
5086    async fn test_store_enqueued_stop_fails_a_stalled_append_and_uploads_the_backlog() {
5087        let config = ObjectStoreWalConfig {
5088            max_unpersisted_bytes: ReadableSize(1),
5089            ..enqueued(manual())
5090        };
5091        let (store, io, mut parked, stalled, release) = stall_second_append(config).await;
5092
5093        let stop = {
5094            let store = store.clone();
5095            tokio::spawn(async move { store.stop().await })
5096        };
5097        assert_stopped(&timeout(WAIT, stalled).await.unwrap().unwrap().unwrap_err());
5098        release.send(true).unwrap();
5099        timeout(WAIT, stop).await.unwrap().unwrap().unwrap();
5100        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
5101        assert!(parked.try_recv().is_err());
5102    }
5103
5104    #[tokio::test]
5105    async fn test_store_enqueued_repeats_a_create_that_failed_transiently() {
5106        let (io, _) = RecordingIo::over(memory_store());
5107        let store = open_over(io.clone(), &enqueued(eager())).await;
5108        let region_id = region(1);
5109
5110        // The create stores the object but reports a failure; the repeat
5111        // writes the same bytes under the same sequence.
5112        io.fail_after_next_put.store(true, Ordering::Relaxed);
5113        let response = append(&store, region_id, "a1").await.unwrap();
5114        assert_eq!(
5115            HashMap::from([(region_id, id(1, 1))]),
5116            response.last_entry_ids
5117        );
5118        timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1)))
5119            .await
5120            .unwrap()
5121            .unwrap();
5122        assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await);
5123        assert_eq!(
5124            expected_entries(region_id, &[(id(1, 1), "a1")]),
5125            read_entries(&store, region_id, 1).await
5126        );
5127
5128        // An object store with a retry layer reports a temporary error that
5129        // outlasted its retries as persistent; the create is repeated too.
5130        io.fail_next_put_persistently();
5131        append(&store, region_id, "a2").await.unwrap();
5132        timeout(WAIT, store.wait_durable(&provider(region_id), id(2, 1)))
5133            .await
5134            .unwrap()
5135            .unwrap();
5136        assert_eq!(vec![0, 1, 2], object_seqs(io.as_ref()).await);
5137    }
5138
5139    #[tokio::test]
5140    async fn test_store_enqueued_permanent_failure_poisons() {
5141        // A conflict of any epoch poisons the store without reading the
5142        // existing object, since the acknowledged entries cannot move to
5143        // another sequence. So does a permanent storage error.
5144        for foreign_epoch in [Some(0), Some(2), None] {
5145            let object_store = memory_store();
5146            let (io, reads) = RecordingIo::over(object_store.clone());
5147            let store = open_over(io.clone(), &enqueued(eager())).await;
5148            let region_id = region(1);
5149            match foreign_epoch {
5150                Some(epoch) => put_foreign(&object_store, 1, epoch).await,
5151                None => io.fail_next_put_permanently(),
5152            }
5153            let assert_poisoned = |error: Error| {
5154                let poisoned = match foreign_epoch {
5155                    Some(_) => matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }),
5156                    None => matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
5157                };
5158                assert!(poisoned, "unexpected error: {error:?}");
5159            };
5160
5161            // The append was acknowledged; the failure surfaces afterwards.
5162            let second = entry(&store, region_id, "a2");
5163            append(&store, region_id, "a1").await.unwrap();
5164            let error = timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1)))
5165                .await
5166                .unwrap()
5167                .unwrap_err();
5168            assert_poisoned(error);
5169            assert_poisoned(store.append_batch(vec![second]).await.unwrap_err());
5170            assert!(store.latest_entry_id(&provider(region_id)).is_err());
5171            // The acknowledged entry was dropped, which stop reports.
5172            assert_poisoned(store.stop().await.unwrap_err());
5173            assert!(reads.lock().unwrap().is_empty());
5174        }
5175    }
5176
5177    #[tokio::test]
5178    async fn test_store_enqueued_stop_reports_a_backlog_lost_before_the_stop_command() {
5179        let (store, io, mut parked) =
5180            open_parking_creates(memory_store(), &enqueued(manual())).await;
5181        let region_id = region(1);
5182        append(&store, region_id, "a1").await.unwrap();
5183        let seal = {
5184            let store = store.clone();
5185            tokio::spawn(async move { store.seal_open_batch().await })
5186        };
5187        let (object_seq, release) = next_create(&mut parked).await;
5188        assert_eq!(1, object_seq);
5189
5190        // Stop began, but the actor has not received the stop command when
5191        // the create fails: the backlog is dropped and the failure is kept
5192        // for the stop that follows.
5193        store.begin_stop();
5194        release.send(false).unwrap();
5195        assert_stopped(&timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err());
5196        // The lost entry cannot be certified as durable before the stop
5197        // command is handled; a durable id still is.
5198        let error = timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1)))
5199            .await
5200            .unwrap()
5201            .unwrap_err();
5202        assert!(
5203            matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
5204            "unexpected error: {error:?}"
5205        );
5206        timeout(WAIT, store.wait_durable(&provider(region_id), 0))
5207            .await
5208            .unwrap()
5209            .unwrap();
5210        let error = store.stop().await.unwrap_err();
5211        assert!(
5212            matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
5213            "unexpected error: {error:?}"
5214        );
5215        assert_eq!(vec![0], object_seqs(io.as_ref()).await);
5216        assert!(parked.try_recv().is_err());
5217    }
5218
5219    #[tokio::test]
5220    async fn test_store_enqueued_stop_uploads_the_backlog() {
5221        let store = open(memory_store(), &enqueued(manual())).await;
5222        let region_id = region(1);
5223        append(&store, region_id, "a1").await.unwrap();
5224        append(&store, region_id, "a2").await.unwrap();
5225        assert_eq!(vec![0], object_seqs(store.io.as_ref()).await);
5226
5227        store.stop().await.unwrap();
5228        assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await);
5229        assert_eq!(id(1, 2), latest(&store, region_id));
5230        assert_eq!(
5231            expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]),
5232            read_entries(&store, region_id, 1).await
5233        );
5234
5235        // A backlog that cannot be uploaded is reported by stop. A transient
5236        // failure, temporary or persistent, drops it; a permanent one also
5237        // poisons the store.
5238        for failure in ["temporary", "persistent", "permanent"] {
5239            let (io, _) = RecordingIo::over(memory_store());
5240            let store = open_over(io.clone(), &enqueued(manual())).await;
5241            append(&store, region_id, "a1").await.unwrap();
5242            match failure {
5243                "temporary" => store.fail_creates(),
5244                "persistent" => io.fail_next_put_persistently(),
5245                _ => io.fail_next_put_permanently(),
5246            }
5247            let error = store.stop().await.unwrap_err();
5248            assert!(
5249                matches!(unwrap_shared(&error), Error::WalObjectStore { .. }),
5250                "unexpected error: {error:?}"
5251            );
5252            assert_eq!(vec![0], object_seqs(io.as_ref()).await);
5253            assert_eq!(
5254                failure == "permanent",
5255                store.latest_entry_id(&provider(region_id)).is_err()
5256            );
5257            store.stop().await.unwrap();
5258        }
5259    }
5260
5261    #[tokio::test]
5262    async fn test_store_enqueued_issued_id_needs_no_floor_before_it_is_durable() {
5263        let store = open(memory_store(), &enqueued(manual())).await;
5264        let region_id = region(1);
5265        append(&store, region_id, "a1").await.unwrap();
5266
5267        // The open batch holds id(1, 1) under the next sequence. The id was
5268        // handed out, so it is accepted without a floor, with the watermark
5269        // capped to what is durable; an id the store never handed out under
5270        // that sequence is refused.
5271        store
5272            .obsolete(&provider(region_id), region_id, id(1, 1))
5273            .await
5274            .unwrap();
5275        assert_eq!(
5276            Some(&0),
5277            store.obsolete_entry_ids.lock().unwrap().get(&region_id)
5278        );
5279        let error = store
5280            .obsolete(&provider(region_id), region_id, id(1, 2))
5281            .await
5282            .unwrap_err();
5283        assert!(
5284            matches!(
5285                error,
5286                Error::WalObjectSequenceUnsettled { object_seq: 1, .. }
5287            ),
5288            "unexpected error: {error:?}"
5289        );
5290
5291        store.seal_open_batch().await.unwrap();
5292        assert_eq!(
5293            expected_entries(region_id, &[(id(1, 1), "a1")]),
5294            read_entries(&store, region_id, 0).await
5295        );
5296        store
5297            .obsolete(&provider(region_id), region_id, id(1, 1))
5298            .await
5299            .unwrap();
5300        assert!(read_entries(&store, region_id, 0).await.is_empty());
5301        let response = append(&store, region_id, "a2").await.unwrap();
5302        assert_eq!(
5303            HashMap::from([(region_id, id(2, 1))]),
5304            response.last_entry_ids
5305        );
5306    }
5307
5308    #[tokio::test]
5309    async fn test_store_enqueued_inherited_watermark_raises_the_floor() {
5310        let object_store = memory_store();
5311        let store = open(object_store.clone(), &enqueued(manual())).await;
5312        let region_id = region(1);
5313
5314        // The region has nothing durable here, so the recorded watermark
5315        // stays at zero, but its ids must still start above the watermark.
5316        store
5317            .obsolete(&provider(region_id), region_id, id(5, 1))
5318            .await
5319            .unwrap();
5320        assert_eq!(
5321            Some(&0),
5322            store.obsolete_entry_ids.lock().unwrap().get(&region_id)
5323        );
5324        let response = append(&store, region_id, "r1").await.unwrap();
5325        assert_eq!(
5326            HashMap::from([(region_id, id(6, 1))]),
5327            response.last_entry_ids
5328        );
5329        store.seal_open_batch().await.unwrap();
5330        timeout(WAIT, store.wait_durable(&provider(region_id), id(6, 1)))
5331            .await
5332            .unwrap()
5333            .unwrap();
5334        assert_eq!(id(6, 1), latest(&store, region_id));
5335        store.stop().await.unwrap();
5336
5337        // Replay from the watermark sees the entry.
5338        let store = open(object_store, &enqueued(manual())).await;
5339        assert_eq!(
5340            expected_entries(region_id, &[(id(6, 1), "r1")]),
5341            read_entries(&store, region_id, id(5, 1) + 1).await
5342        );
5343    }
5344
5345    #[tokio::test]
5346    async fn test_store_durability_wait_pending_at_stop_fails_with_stopped() {
5347        let store = open(memory_store(), &manual()).await;
5348        let region_id = region(1);
5349        let append = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]);
5350        store.wait_for_admitted_appends(1).await.unwrap();
5351        let wait = {
5352            let store = store.clone();
5353            tokio::spawn(async move { store.wait_durable(&provider(region_id), id(1, 1)).await })
5354        };
5355        for _ in 0..16 {
5356            tokio::task::yield_now().await;
5357        }
5358        assert!(!wait.is_finished());
5359
5360        store.stop().await.unwrap();
5361        assert_stopped(&timeout(WAIT, append).await.unwrap().unwrap().unwrap_err());
5362        assert_stopped(&timeout(WAIT, wait).await.unwrap().unwrap().unwrap_err());
5363    }
5364
5365    /// Object access whose next create stores the object but reports a
5366    /// transient failure, as a create whose response is lost.
5367    struct LostResponseIo {
5368        inner: ObjectStoreIo,
5369        lose_next: AtomicBool,
5370    }
5371
5372    #[async_trait::async_trait]
5373    impl WalObjectIo for LostResponseIo {
5374        async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult> {
5375            let result = self.inner.put_if_absent(object_seq, content).await?;
5376            if !self.lose_next.swap(false, Ordering::SeqCst) {
5377                return Ok(result);
5378            }
5379            Err(
5380                object_store::Error::new(object_store::ErrorKind::Unexpected, "lost response")
5381                    .set_temporary(),
5382            )
5383            .context(WalObjectStoreSnafu {
5384                operation: "write",
5385                path: self.object_path(object_seq),
5386            })
5387        }
5388
5389        async fn get(&self, object_seq: u64) -> Result<Bytes> {
5390            self.inner.get(object_seq).await
5391        }
5392
5393        async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes> {
5394            self.inner.get_range(object_seq, offset, len).await
5395        }
5396
5397        async fn list(&self) -> Result<Vec<ListedObject>> {
5398            self.inner.list().await
5399        }
5400
5401        fn object_path(&self, object_seq: u64) -> String {
5402            self.inner.object_path(object_seq)
5403        }
5404    }
5405
5406    /// Object access whose first create lets another object land at a given
5407    /// sequence first, as a create issued before recovery would.
5408    struct RacingIo {
5409        inner: ObjectStoreIo,
5410        late: Mutex<Option<(u64, Bytes)>>,
5411    }
5412
5413    #[async_trait::async_trait]
5414    impl WalObjectIo for RacingIo {
5415        async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult> {
5416            let late = self.late.lock().unwrap().take();
5417            if let Some((late_seq, late)) = late {
5418                self.inner.put_if_absent(late_seq, late).await?;
5419            }
5420            self.inner.put_if_absent(object_seq, content).await
5421        }
5422
5423        async fn get(&self, object_seq: u64) -> Result<Bytes> {
5424            self.inner.get(object_seq).await
5425        }
5426
5427        async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes> {
5428            self.inner.get_range(object_seq, offset, len).await
5429        }
5430
5431        async fn list(&self) -> Result<Vec<ListedObject>> {
5432            self.inner.list().await
5433        }
5434
5435        fn object_path(&self, object_seq: u64) -> String {
5436            self.inner.object_path(object_seq)
5437        }
5438    }
5439
5440    /// Object access whose whole-object reads, or conditional creates, park
5441    /// until the test releases them, counting how many are in flight. A
5442    /// release of false fails the operation with a transient error before it
5443    /// reaches the object store.
5444    struct ParkedIo {
5445        inner: ObjectStoreIo,
5446        parked: mpsc::UnboundedSender<(u64, oneshot::Sender<bool>)>,
5447        park_creates: bool,
5448        /// Whether creates park yet; set once the store opened.
5449        creates_parked: AtomicBool,
5450        in_flight: AtomicUsize,
5451        max_in_flight: AtomicUsize,
5452    }
5453
5454    impl ParkedIo {
5455        fn over(
5456            object_store: ObjectStore,
5457        ) -> (
5458            Arc<Self>,
5459            mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
5460        ) {
5461            Self::new(object_store, false)
5462        }
5463
5464        fn parking_creates(
5465            object_store: ObjectStore,
5466        ) -> (
5467            Arc<Self>,
5468            mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
5469        ) {
5470            Self::new(object_store, true)
5471        }
5472
5473        fn new(
5474            object_store: ObjectStore,
5475            park_creates: bool,
5476        ) -> (
5477            Arc<Self>,
5478            mpsc::UnboundedReceiver<(u64, oneshot::Sender<bool>)>,
5479        ) {
5480            let (parked, parked_rx) = mpsc::unbounded_channel();
5481            (
5482                Arc::new(Self {
5483                    inner: ObjectStoreIo::new(object_store, PREFIX).unwrap(),
5484                    parked,
5485                    park_creates,
5486                    creates_parked: AtomicBool::new(false),
5487                    in_flight: AtomicUsize::new(0),
5488                    max_in_flight: AtomicUsize::new(0),
5489                }),
5490                parked_rx,
5491            )
5492        }
5493
5494        /// Parks until the test releases `object_seq` and returns whether the
5495        /// operation proceeds.
5496        async fn park(&self, object_seq: u64) -> bool {
5497            let in_flight = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1;
5498            self.max_in_flight.fetch_max(in_flight, Ordering::SeqCst);
5499            let (release, released) = oneshot::channel();
5500            self.parked.send((object_seq, release)).unwrap();
5501            let proceed = released.await.unwrap();
5502            self.in_flight.fetch_sub(1, Ordering::SeqCst);
5503            proceed
5504        }
5505    }
5506
5507    #[async_trait::async_trait]
5508    impl WalObjectIo for ParkedIo {
5509        async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult> {
5510            if self.park_creates
5511                && self.creates_parked.load(Ordering::SeqCst)
5512                && !self.park(object_seq).await
5513            {
5514                return injected_failure("write", self.object_path(object_seq));
5515            }
5516            self.inner.put_if_absent(object_seq, content).await
5517        }
5518
5519        async fn get(&self, object_seq: u64) -> Result<Bytes> {
5520            if !self.park_creates && !self.park(object_seq).await {
5521                return injected_failure("read", self.object_path(object_seq));
5522            }
5523            self.inner.get(object_seq).await
5524        }
5525
5526        async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes> {
5527            self.inner.get_range(object_seq, offset, len).await
5528        }
5529
5530        async fn list(&self) -> Result<Vec<ListedObject>> {
5531            self.inner.list().await
5532        }
5533
5534        fn object_path(&self, object_seq: u64) -> String {
5535            self.inner.object_path(object_seq)
5536        }
5537    }
5538
5539    type RangeReads = Arc<Mutex<Vec<(u64, u64, u64)>>>;
5540
5541    /// Object access that records every range read as (sequence, offset,
5542    /// length) and, on request, reports the next conditional create as failed
5543    /// after it wrote the object, or fails it with a given error before it
5544    /// writes.
5545    struct RecordingIo {
5546        inner: ObjectStoreIo,
5547        reads: RangeReads,
5548        fail_after_next_put: AtomicBool,
5549        fail_next_put: Mutex<Option<object_store::Error>>,
5550    }
5551
5552    impl RecordingIo {
5553        fn over(object_store: ObjectStore) -> (Arc<Self>, RangeReads) {
5554            let reads = RangeReads::default();
5555            let io = Self {
5556                inner: ObjectStoreIo::new(object_store, PREFIX).unwrap(),
5557                reads: reads.clone(),
5558                fail_after_next_put: AtomicBool::new(false),
5559                fail_next_put: Mutex::new(None),
5560            };
5561            (Arc::new(io), reads)
5562        }
5563
5564        fn fail_next_put_permanently(&self) {
5565            let error = object_store::Error::new(
5566                object_store::ErrorKind::PermissionDenied,
5567                "injected failure",
5568            );
5569            *self.fail_next_put.lock().unwrap() = Some(error);
5570        }
5571
5572        /// Fails the next create as a retry layer reports a temporary error
5573        /// that outlasted its retries.
5574        fn fail_next_put_persistently(&self) {
5575            let error = object_store::Error::new(object_store::ErrorKind::Unexpected, "injected")
5576                .set_temporary()
5577                .set_persistent();
5578            *self.fail_next_put.lock().unwrap() = Some(error);
5579        }
5580    }
5581
5582    #[async_trait::async_trait]
5583    impl WalObjectIo for RecordingIo {
5584        async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result<PutResult> {
5585            if let Some(error) = self.fail_next_put.lock().unwrap().take() {
5586                return Err(error).context(WalObjectStoreSnafu {
5587                    operation: "write",
5588                    path: self.object_path(object_seq),
5589                });
5590            }
5591            let result = self.inner.put_if_absent(object_seq, content).await?;
5592            if self.fail_after_next_put.swap(false, Ordering::Relaxed) {
5593                return injected_failure("write", self.object_path(object_seq));
5594            }
5595            Ok(result)
5596        }
5597
5598        async fn get(&self, object_seq: u64) -> Result<Bytes> {
5599            self.inner.get(object_seq).await
5600        }
5601
5602        async fn get_range(&self, object_seq: u64, offset: u64, len: u64) -> Result<Bytes> {
5603            self.reads.lock().unwrap().push((object_seq, offset, len));
5604            self.inner.get_range(object_seq, offset, len).await
5605        }
5606
5607        async fn list(&self) -> Result<Vec<ListedObject>> {
5608            self.inner.list().await
5609        }
5610
5611        fn object_path(&self, object_seq: u64) -> String {
5612            self.inner.object_path(object_seq)
5613        }
5614    }
5615}