1use 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;
59const APPEND_BUFFER: usize = 16;
61const MIN_FLUSH_INTERVAL: Duration = Duration::from_millis(10);
62const MAX_IN_FLIGHT_CREATES: usize = 4;
63const CREATE_RETRY_DELAY: Duration = Duration::from_millis(100);
66const MAX_SEALED_BATCHES: usize = 2 * MAX_IN_FLIGHT_CREATES;
72const RECOVERY_CONCURRENCY: usize = 8;
74const RECOVERY_TAIL_WINDOW: usize = 64 * 1024;
78
79pub 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 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 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 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 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 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#[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
354fn 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 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 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 pub fn hold_creates(&self) {
410 self.creates_held.send_replace(true);
411 }
412
413 pub fn release_creates(&self) {
415 self.creates_held.send_replace(false);
416 }
417
418 pub fn fail_creates(&self) {
421 self.creates_fail.store(true, Ordering::Release);
422 }
423
424 pub fn fail_next_create_after_write(&self) {
428 self.next_create_fails_after_write
429 .store(true, Ordering::Release);
430 }
431
432 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 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 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 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 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 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 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(®ion_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 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 None => {
670 record_obsolete(&self.obsolete_entry_ids, region_id, watermark);
671 Ok(())
672 }
673 }
674 }
675
676 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 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 WaitDurable {
734 region_id: RegionId,
735 entry_id: EntryId,
736 response: oneshot::Sender<Result<()>>,
737 },
738 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
758struct PendingAppend {
760 last_entry_ids: HashMap<RegionId, EntryId>,
761 response: AppendResponse,
762}
763
764struct DurableWaiter {
766 region_id: RegionId,
767 entry_id: EntryId,
768 response: oneshot::Sender<Result<()>>,
769}
770
771enum CreateState {
773 Pending,
776 InFlight,
777 Created,
778}
779
780struct SealedBatch {
783 object_seq: u64,
784 bytes: Bytes,
785 footer: Vec<FooterEntry>,
786 first_admitted_at: Instant,
787 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
813struct 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 issued_entry_ids: HashMap<RegionId, EntryId>,
860 pending: Vec<PendingAppend>,
862 sealed: VecDeque<SealedBatch>,
864 creates: FuturesUnordered<BoxFuture<'static, CreateOutcome>>,
866 stalled: Option<QueuedAppend>,
869 durable_waiters: Vec<DurableWaiter>,
870 stop: Vec<oneshot::Sender<Result<()>>>,
872 stop_error: Option<Arc<Error>>,
875 next_object_seq: Option<u64>,
879 epoch: u64,
881 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 interval.tick().await;
905
906 loop {
907 tokio::select! {
908 _ = interval.tick() => {
909 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 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 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 fn admit(&mut self, entries: Vec<Entry>, response: AppendResponse) {
975 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 if self.ack_mode == AckMode::Enqueued {
1375 self.stop_error.get_or_insert(error);
1376 }
1377 }
1378
1379 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 fn reset_open_batch(&mut self) {
1410 self.open_batch.reset();
1411 }
1412
1413 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 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 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 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 fn raise_sequence_floor(&mut self, region_id: RegionId, entry_id: EntryId) -> Result<()> {
1488 if self.is_stopped() {
1490 return Ok(());
1491 }
1492 if let Some(error) = terminal(&self.terminal_error) {
1493 return Err(shared(&error));
1494 }
1495 if self.ack_mode == AckMode::Enqueued
1501 && entry_id <= self.issued_entry_ids.get(®ion_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 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 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 #[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
1627fn 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
1637fn is_transient(error: &Error) -> bool {
1641 matches!(error, Error::WalObjectStore { error, .. } if !error.is_permanent())
1642}
1643
1644fn 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#[derive(Debug)]
1663struct Recovered {
1664 catalog: ObjectCatalog,
1665 next_object_seq: u64,
1667 durable_entry_ids: HashMap<RegionId, EntryId>,
1668 tip: Option<ChainLink>,
1670 max_epoch: u64,
1672}
1673
1674struct FetchedObject {
1676 object: ListedObject,
1677 header: Header,
1678 footer: Vec<FooterEntry>,
1679}
1680
1681async 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
1693fn 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 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
1756fn select_chain(headers: &BTreeMap<u64, &Header>) -> Vec<u64> {
1769 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
1801async 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
1864async 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
1886async 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
1911async 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
1983fn 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#[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 fn enqueued(config: ObjectStoreWalConfig) -> ObjectStoreWalConfig {
2090 ObjectStoreWalConfig {
2091 ack_mode: AckMode::Enqueued,
2092 ..config
2093 }
2094 }
2095
2096 fn eager() -> ObjectStoreWalConfig {
2098 config(Duration::from_secs(3600), 1)
2099 }
2100
2101 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 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 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 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 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 io.creates_parked.store(true, Ordering::SeqCst);
2241 (store, io, parked)
2242 }
2243
2244 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 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 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 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 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 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[®ion_id]);
2518 assert_eq!(durable[®ion_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 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 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 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 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 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 assert_stopped(&store.append_batch(Vec::new()).await.unwrap_err());
2866 }
2867
2868 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 for _ in 0..16 {
2983 tokio::task::yield_now().await;
2984 }
2985 assert!(parked.try_recv().is_err());
2986
2987 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 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[®ion(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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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(®ion_two)
3902 .is_none()
3903 );
3904
3905 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(®ion_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 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 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 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 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 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 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(®ion_id)
4044 );
4045
4046 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(®ion(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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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(®ion(2))
4364 );
4365
4366 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 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 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 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 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 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 assert_eq!(
4486 vec![0],
4487 chain_of(&[chain_header(0, 1, None), chain_header(1, 2, Some((0, 2))),])
4488 );
4489 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 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 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 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 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 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 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 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 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 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 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 let recovered = recover(&io).await.unwrap();
4736 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 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 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 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 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 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 let waits = [wait(id(1, 1)), wait(id(1, 7))];
4871 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 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[®ion_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 timeout(WAIT, store.wait_durable(&provider(region_a), id(2, 1)))
4921 .await
4922 .unwrap()
4923 .unwrap();
4924 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(®ion_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(®ion_id)
4972 );
4973 assert!(read_entries(&store, region_id, 1).await.is_empty());
4974 }
4975
4976 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 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 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 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 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 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 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 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 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 store.begin_stop();
5194 release.send(false).unwrap();
5195 assert_stopped(&timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err());
5196 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 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 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(®ion_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 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(®ion_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 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 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 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 struct ParkedIo {
5445 inner: ObjectStoreIo,
5446 parked: mpsc::UnboundedSender<(u64, oneshot::Sender<bool>)>,
5447 park_creates: bool,
5448 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 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 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 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}