Skip to main content

mz_persist/
location.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Abstractions over files, cloud storage, etc used in persistence.
11
12use std::fmt;
13use std::pin::Pin;
14use std::sync::Arc;
15use std::time::Instant;
16
17use anyhow::anyhow;
18use async_trait::async_trait;
19use azure_core::http::StatusCode;
20use bytes::Bytes;
21use futures_util::Stream;
22use mz_ore::bytes::SegmentedBytes;
23use mz_ore::cast::u64_to_usize;
24use mz_postgres_client::error::PostgresError;
25use mz_proto::RustType;
26use proptest_derive::Arbitrary;
27use serde::{Deserialize, Serialize};
28use tracing::{Instrument, Span};
29
30use crate::error::Error;
31
32/// The "sequence number" of a persist state change.
33///
34/// Persist is a state machine, with all mutating requests modeled as input
35/// state changes sequenced into a log. This reflects that ordering.
36///
37/// This ordering also includes requests that were sequenced and applied to the
38/// persist state machine, but that application was deterministically made into
39/// a no-op because it was contextually invalid (a write or seal at a sealed
40/// timestamp, an allow_compactions at an unsealed timestamp, etc).
41///
42/// Read-only requests are assigned the SeqNo of a write, indicating that all
43/// mutating requests up to and including that one are reflected in the read
44/// state.
45#[derive(
46    Arbitrary,
47    Clone,
48    Copy,
49    Debug,
50    PartialOrd,
51    Ord,
52    PartialEq,
53    Eq,
54    Hash,
55    Serialize,
56    Deserialize
57)]
58pub struct SeqNo(pub u64);
59
60impl std::fmt::Display for SeqNo {
61    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
62        write!(f, "v{}", self.0)
63    }
64}
65
66impl timely::PartialOrder for SeqNo {
67    fn less_equal(&self, other: &Self) -> bool {
68        self <= other
69    }
70}
71
72impl std::str::FromStr for SeqNo {
73    type Err = String;
74
75    fn from_str(encoded: &str) -> Result<Self, Self::Err> {
76        let encoded = match encoded.strip_prefix('v') {
77            Some(x) => x,
78            None => return Err(format!("invalid SeqNo {}: incorrect prefix", encoded)),
79        };
80        let seqno =
81            u64::from_str(encoded).map_err(|err| format!("invalid SeqNo {}: {}", encoded, err))?;
82        Ok(SeqNo(seqno))
83    }
84}
85
86impl SeqNo {
87    /// Returns the previous SeqNo in the sequence, if there is one.
88    pub fn previous(self) -> Option<SeqNo> {
89        Some(SeqNo(self.0.checked_sub(1)?))
90    }
91
92    /// Returns the next SeqNo in the sequence.
93    pub fn next(self) -> SeqNo {
94        SeqNo(self.0 + 1)
95    }
96
97    /// A minimum value suitable as a default.
98    pub fn minimum() -> Self {
99        SeqNo(0)
100    }
101
102    /// A maximum value.
103    pub fn maximum() -> Self {
104        SeqNo(u64::MAX)
105    }
106}
107
108impl RustType<u64> for SeqNo {
109    fn into_proto(&self) -> u64 {
110        self.0
111    }
112
113    fn from_proto(proto: u64) -> Result<Self, mz_proto::TryFromProtoError> {
114        Ok(SeqNo(proto))
115    }
116}
117
118/// An error coming from an underlying durability system (e.g. s3) indicating
119/// that the operation _definitely did NOT succeed_ (e.g. permission denied).
120#[derive(Debug)]
121pub struct Determinate {
122    inner: anyhow::Error,
123}
124
125impl std::fmt::Display for Determinate {
126    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
127        write!(f, "determinate: ")?;
128        self.inner.fmt(f)
129    }
130}
131
132impl std::error::Error for Determinate {
133    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
134        self.inner.source()
135    }
136}
137
138impl From<anyhow::Error> for Determinate {
139    fn from(inner: anyhow::Error) -> Self {
140        Self::new(inner)
141    }
142}
143
144impl Determinate {
145    /// Return a new Determinate wrapping the given error.
146    ///
147    /// Exposed for testing via [crate::unreliable].
148    pub fn new(inner: anyhow::Error) -> Self {
149        Determinate { inner }
150    }
151
152    /// Adds context to the wrapped error. Mirrors [`anyhow::Error::context`].
153    pub fn context<C>(self, context: C) -> Self
154    where
155        C: fmt::Display + Send + Sync + 'static,
156    {
157        Determinate::new(self.inner.context(context))
158    }
159}
160
161/// An error coming from an underlying durability system (e.g. s3) indicating
162/// that the operation _might have succeeded_ (e.g. timeout).
163#[derive(Debug)]
164pub struct Indeterminate {
165    pub(crate) inner: anyhow::Error,
166}
167
168impl Indeterminate {
169    /// Return a new Indeterminate wrapping the given error.
170    ///
171    /// Exposed for testing.
172    pub fn new(inner: anyhow::Error) -> Self {
173        Indeterminate { inner }
174    }
175
176    /// Adds context to the wrapped error. Mirrors [`anyhow::Error::context`].
177    pub fn context<C>(self, context: C) -> Self
178    where
179        C: fmt::Display + Send + Sync + 'static,
180    {
181        Indeterminate::new(self.inner.context(context))
182    }
183}
184
185impl std::fmt::Display for Indeterminate {
186    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
187        write!(f, "indeterminate: ")?;
188        self.inner.fmt(f)
189    }
190}
191
192impl std::error::Error for Indeterminate {
193    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
194        self.inner.source()
195    }
196}
197
198/// An impl of PartialEq purely for convenience in tests and debug assertions.
199#[cfg(any(test, debug_assertions))]
200impl PartialEq for Indeterminate {
201    fn eq(&self, other: &Self) -> bool {
202        self.to_string() == other.to_string()
203    }
204}
205
206/// An error coming from an underlying durability system (e.g. s3) or from
207/// invalid data received from one.
208#[derive(Debug)]
209pub enum ExternalError {
210    /// A determinate error from an external system.
211    Determinate(Determinate),
212    /// An indeterminate error from an external system.
213    Indeterminate(Indeterminate),
214}
215
216impl ExternalError {
217    /// Returns a new error representing a timeout.
218    ///
219    /// TODO: When we overhaul errors, this presumably should instead be a type
220    /// that can be matched on.
221    #[track_caller]
222    pub fn new_timeout(deadline: Instant) -> Self {
223        ExternalError::Indeterminate(Indeterminate {
224            inner: anyhow!("timeout at {:?}", deadline),
225        })
226    }
227
228    /// Returns whether this error represents a timeout.
229    ///
230    /// TODO: When we overhaul errors, this presumably should instead be a type
231    /// that can be matched on.
232    pub fn is_timeout(&self) -> bool {
233        // Gross...
234        self.to_string().contains("timeout")
235    }
236
237    /// Adds context to the underlying error, preserving the determinate vs
238    /// indeterminate classification. Mirrors [`anyhow::Error::context`].
239    ///
240    /// Callers use this to record which resource an operation was acting on
241    /// (e.g. the blob key being fetched) so the error names it everywhere it is
242    /// displayed.
243    pub fn context<C>(self, context: C) -> Self
244    where
245        C: fmt::Display + Send + Sync + 'static,
246    {
247        match self {
248            ExternalError::Determinate(e) => ExternalError::Determinate(e.context(context)),
249            ExternalError::Indeterminate(e) => ExternalError::Indeterminate(e.context(context)),
250        }
251    }
252}
253
254impl std::fmt::Display for ExternalError {
255    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
256        match self {
257            ExternalError::Determinate(x) => std::fmt::Display::fmt(x, f),
258            ExternalError::Indeterminate(x) => std::fmt::Display::fmt(x, f),
259        }
260    }
261}
262
263impl std::error::Error for ExternalError {
264    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
265        match self {
266            ExternalError::Determinate(e) => e.source(),
267            ExternalError::Indeterminate(e) => e.source(),
268        }
269    }
270}
271
272/// An impl of PartialEq purely for convenience in tests and debug assertions.
273#[cfg(any(test, debug_assertions))]
274impl PartialEq for ExternalError {
275    fn eq(&self, other: &Self) -> bool {
276        self.to_string() == other.to_string()
277    }
278}
279
280impl From<PostgresError> for ExternalError {
281    fn from(x: PostgresError) -> Self {
282        match x {
283            PostgresError::Determinate(e) => ExternalError::Determinate(Determinate::new(e)),
284            PostgresError::Indeterminate(e) => ExternalError::Indeterminate(Indeterminate::new(e)),
285        }
286    }
287}
288
289impl From<Indeterminate> for ExternalError {
290    fn from(x: Indeterminate) -> Self {
291        ExternalError::Indeterminate(x)
292    }
293}
294
295impl From<Determinate> for ExternalError {
296    fn from(x: Determinate) -> Self {
297        ExternalError::Determinate(x)
298    }
299}
300
301impl From<anyhow::Error> for ExternalError {
302    fn from(inner: anyhow::Error) -> Self {
303        ExternalError::Indeterminate(Indeterminate { inner })
304    }
305}
306
307impl From<Error> for ExternalError {
308    fn from(x: Error) -> Self {
309        ExternalError::Indeterminate(Indeterminate {
310            inner: anyhow::Error::new(x),
311        })
312    }
313}
314
315impl From<std::io::Error> for ExternalError {
316    fn from(x: std::io::Error) -> Self {
317        ExternalError::Indeterminate(Indeterminate {
318            inner: anyhow::Error::new(x),
319        })
320    }
321}
322
323impl From<deadpool_postgres::tokio_postgres::Error> for ExternalError {
324    fn from(e: deadpool_postgres::tokio_postgres::Error) -> Self {
325        let code = match e.as_db_error().map(|x| x.code()) {
326            Some(x) => x,
327            None => {
328                return ExternalError::Indeterminate(Indeterminate {
329                    inner: anyhow::Error::new(e),
330                });
331            }
332        };
333        match code {
334            // Feel free to add more things to this allowlist as we encounter
335            // them as long as you're certain they're determinate.
336            &deadpool_postgres::tokio_postgres::error::SqlState::T_R_SERIALIZATION_FAILURE => {
337                ExternalError::Determinate(Determinate {
338                    inner: anyhow::Error::new(e),
339                })
340            }
341            _ => ExternalError::Indeterminate(Indeterminate {
342                inner: anyhow::Error::new(e),
343            }),
344        }
345    }
346}
347
348impl From<azure_core::Error> for ExternalError {
349    fn from(value: azure_core::Error) -> Self {
350        let definitely_determinate = match value.http_status() {
351            // There are many other status codes that _ought_ to be determinate, according to
352            // the HTTP spec, but this includes only codes that we've observed in practice for now.
353            Some(StatusCode::TooManyRequests) => true,
354            _ => false,
355        };
356        if definitely_determinate {
357            ExternalError::Determinate(Determinate {
358                inner: anyhow!(value),
359            })
360        } else {
361            ExternalError::Indeterminate(Indeterminate {
362                inner: anyhow!(value),
363            })
364        }
365    }
366}
367
368impl From<deadpool_postgres::PoolError> for ExternalError {
369    fn from(x: deadpool_postgres::PoolError) -> Self {
370        match x {
371            // We have logic for turning a postgres Error into an ExternalError,
372            // so use it.
373            deadpool_postgres::PoolError::Backend(x) => ExternalError::from(x),
374            x => ExternalError::Indeterminate(Indeterminate {
375                inner: anyhow::Error::new(x),
376            }),
377        }
378    }
379}
380
381impl From<tokio::task::JoinError> for ExternalError {
382    fn from(x: tokio::task::JoinError) -> Self {
383        ExternalError::Indeterminate(Indeterminate {
384            inner: anyhow::Error::new(x),
385        })
386    }
387}
388
389/// An abstraction for a single arbitrarily-sized binary blob and an associated
390/// version number (sequence number).
391#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
392pub struct VersionedData {
393    /// The sequence number of the data.
394    pub seqno: SeqNo,
395    /// The data itself.
396    pub data: Bytes,
397}
398
399/// Helper constant to scan all states in [Consensus::scan].
400/// The maximum possible SeqNo is i64::MAX.
401// TODO(benesch): find a way to express this without `as`.
402#[allow(clippy::as_conversions)]
403pub const SCAN_ALL: usize = u64_to_usize(i64::MAX as u64);
404
405/// A key usable for liveness checks via [Consensus::head].
406pub const CONSENSUS_HEAD_LIVENESS_KEY: &str = "LIVENESS";
407
408/// Return type to indicate whether [Consensus::compare_and_set] succeeded or failed.
409#[derive(Debug, PartialEq, Serialize, Deserialize)]
410pub enum CaSResult {
411    /// The compare-and-set succeeded and committed new state.
412    Committed,
413    /// The compare-and-set failed due to expectation mismatch.
414    ExpectationMismatch,
415}
416
417/// Wraps all calls to a backing store in a new tokio task. This adds extra overhead,
418/// but insulates the system from callers who fail to drive futures promptly to completion,
419/// which can cause timeouts or resource exhaustion in a store.
420#[derive(Debug)]
421pub struct Tasked<A>(pub Arc<A>);
422
423impl<A> Tasked<A> {
424    fn clone_backing(&self) -> Arc<A> {
425        Arc::clone(&self.0)
426    }
427}
428
429/// A boxed stream, similar to what `async_trait` desugars async functions to, but hardcoded
430/// to our standard result type.
431pub type ResultStream<'a, T> = Pin<Box<dyn Stream<Item = Result<T, ExternalError>> + Send + 'a>>;
432
433/// An abstraction for [VersionedData] held in a location in persistent storage
434/// where the data are conditionally updated by version.
435///
436/// Users are expected to use this API with consistently increasing sequence numbers
437/// to allow multiple processes across multiple machines to agree to a total order
438/// of the evolution of the data. To make roundtripping through various forms of durable
439/// storage easier, sequence numbers used with [Consensus] need to be restricted to the
440/// range [0, i64::MAX].
441#[async_trait]
442pub trait Consensus: std::fmt::Debug + Send + Sync {
443    /// Returns all the keys ever created in the consensus store.
444    fn list_keys(&self) -> ResultStream<'_, String>;
445
446    /// Returns a recent version of `data`, and the corresponding sequence number, if
447    /// one exists at this location.
448    async fn head(&self, key: &str) -> Result<Option<VersionedData>, ExternalError>;
449
450    /// Add the [VersionedData] to the log for the given key. If the sequence number is 0, the log
451    /// must be empty; otherwise, it must be one greater than the previous sequence number.
452    /// It is invalid to call
453    /// this function with a sequence number outside of the range `[0, i64::MAX]`.
454    async fn compare_and_set(
455        &self,
456        key: &str,
457        new: VersionedData,
458    ) -> Result<CaSResult, ExternalError>;
459
460    /// Return `limit` versions of data stored for this `key` at sequence numbers >= `from`,
461    /// in ascending order of sequence number.
462    ///
463    /// Returns an empty vec if `from` is greater than the current sequence
464    /// number or if there is no data at this key.
465    async fn scan(
466        &self,
467        key: &str,
468        from: SeqNo,
469        limit: usize,
470    ) -> Result<Vec<VersionedData>, ExternalError>;
471
472    /// Deletes all historical versions of the data stored at `key` that are <
473    /// `seqno`, iff `seqno` <= the current sequence number.
474    ///
475    /// Returns the number of versions deleted or `None` on success. Returns an error if
476    /// `seqno` is greater than the current sequence number, or if there is no
477    /// data at this key.
478    async fn truncate(&self, key: &str, seqno: SeqNo) -> Result<Option<usize>, ExternalError>;
479}
480
481#[async_trait]
482impl<A: Consensus + 'static> Consensus for Tasked<A> {
483    fn list_keys(&self) -> ResultStream<'_, String> {
484        // Similarly to Blob::list_keys_and_metadata, this is difficult to make into a task.
485        // (If we use an unbounded channel between the task and the caller, we can buffer forever;
486        // if we use a bounded channel, we lose the isolation benefits of Tasked.)
487        // However, this should only be called in administrative contexts
488        // and not in the main state-machine impl.
489        self.0.list_keys()
490    }
491
492    async fn head(&self, key: &str) -> Result<Option<VersionedData>, ExternalError> {
493        let backing = self.clone_backing();
494        let key = key.to_owned();
495        mz_ore::task::spawn(
496            || "persist::task::head",
497            async move { backing.head(&key).await }.instrument(Span::current()),
498        )
499        .await
500    }
501
502    async fn compare_and_set(
503        &self,
504        key: &str,
505        new: VersionedData,
506    ) -> Result<CaSResult, ExternalError> {
507        let backing = self.clone_backing();
508        let key = key.to_owned();
509        mz_ore::task::spawn(
510            || "persist::task::cas",
511            async move { backing.compare_and_set(&key, new).await }.instrument(Span::current()),
512        )
513        .await
514    }
515
516    async fn scan(
517        &self,
518        key: &str,
519        from: SeqNo,
520        limit: usize,
521    ) -> Result<Vec<VersionedData>, ExternalError> {
522        let backing = self.clone_backing();
523        let key = key.to_owned();
524        mz_ore::task::spawn(
525            || "persist::task::scan",
526            async move { backing.scan(&key, from, limit).await }.instrument(Span::current()),
527        )
528        .await
529    }
530
531    async fn truncate(&self, key: &str, seqno: SeqNo) -> Result<Option<usize>, ExternalError> {
532        let backing = self.clone_backing();
533        let key = key.to_owned();
534        mz_ore::task::spawn(
535            || "persist::task::truncate",
536            async move { backing.truncate(&key, seqno).await }.instrument(Span::current()),
537        )
538        .await
539    }
540}
541
542/// Metadata about a particular blob stored by persist
543#[derive(Debug)]
544pub struct BlobMetadata<'a> {
545    /// The key for the blob
546    pub key: &'a str,
547    /// Size of the blob
548    pub size_in_bytes: u64,
549}
550
551/// A key usable for liveness checks via [Blob::get].
552pub const BLOB_GET_LIVENESS_KEY: &str = "LIVENESS";
553
554/// An abstraction over read-write access to a `bytes key`->`bytes value` store.
555///
556/// Implementations are required to be _linearizable_.
557///
558/// TODO: Consider whether this can be relaxed. Since our usage is write-once
559/// modify-never, it certainly seems like we could by adding retries around
560/// `get` to wait for a non-linearizable `set` to show up. However, the tricky
561/// bit comes once we stop handing out seqno capabilities to readers and have to
562/// start reasoning about "this set hasn't show up yet" vs "the blob has already
563/// been deleted". Another tricky problem is the same but for a deletion when
564/// the first attempt timed out.
565#[async_trait]
566pub trait Blob: std::fmt::Debug + Send + Sync {
567    /// Returns a reference to the value corresponding to the key.
568    async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError>;
569
570    /// List all of the keys in the map with metadata about the entry.
571    ///
572    /// Can be optionally restricted to only list keys starting with a
573    /// given prefix.
574    async fn list_keys_and_metadata(
575        &self,
576        key_prefix: &str,
577        f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
578    ) -> Result<(), ExternalError>;
579
580    /// Inserts a key-value pair into the map.
581    ///
582    /// Writes must be atomic and either succeed or leave the previous value
583    /// intact.
584    async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError>;
585
586    /// Remove a key from the map.
587    ///
588    /// Returns Some and the size of the deleted blob if if exists. Succeeds and
589    /// returns None if it does not exist.
590    async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError>;
591
592    /// Restores a previously-deleted key to the map, if possible.
593    ///
594    /// Returns successfully if the key exists after this call: perhaps because it already existed
595    /// or was restored. (In particular, this makes restore idempotent.)
596    /// Fails if we were unable to restore any value for that key:
597    /// perhaps the key was never written, or was permanently deleted.
598    ///
599    /// It is acceptable for [Blob::restore] to be unable
600    /// to restore keys, in which case this method should succeed iff the key exists.
601    async fn restore(&self, key: &str) -> Result<(), ExternalError>;
602}
603
604#[async_trait]
605impl<A: Blob + 'static> Blob for Tasked<A> {
606    async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
607        let backing = self.clone_backing();
608        let key = key.to_owned();
609        mz_ore::task::spawn(
610            || "persist::task::get",
611            async move { backing.get(&key).await }.instrument(Span::current()),
612        )
613        .await
614    }
615
616    /// List all of the keys in the map with metadata about the entry.
617    ///
618    /// Can be optionally restricted to only list keys starting with a
619    /// given prefix.
620    async fn list_keys_and_metadata(
621        &self,
622        key_prefix: &str,
623        f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
624    ) -> Result<(), ExternalError> {
625        // TODO: No good way that I can see to make this one a task because of
626        // the closure and Blob needing to be object-safe.
627        self.0.list_keys_and_metadata(key_prefix, f).await
628    }
629
630    /// Inserts a key-value pair into the map.
631    async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError> {
632        let backing = self.clone_backing();
633        let key = key.to_owned();
634        mz_ore::task::spawn(
635            || "persist::task::set",
636            async move { backing.set(&key, value).await }.instrument(Span::current()),
637        )
638        .await
639    }
640
641    /// Remove a key from the map.
642    ///
643    /// Returns Some and the size of the deleted blob if if exists. Succeeds and
644    /// returns None if it does not exist.
645    async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError> {
646        let backing = self.clone_backing();
647        let key = key.to_owned();
648        mz_ore::task::spawn(
649            || "persist::task::delete",
650            async move { backing.delete(&key).await }.instrument(Span::current()),
651        )
652        .await
653    }
654
655    async fn restore(&self, key: &str) -> Result<(), ExternalError> {
656        let backing = self.clone_backing();
657        let key = key.to_owned();
658        mz_ore::task::spawn(
659            || "persist::task::restore",
660            async move { backing.restore(&key).await }.instrument(Span::current()),
661        )
662        .await
663    }
664}
665
666/// Test helpers for the crate.
667#[cfg(test)]
668pub mod tests {
669    use std::future::Future;
670
671    use anyhow::anyhow;
672    use futures_util::TryStreamExt;
673    use mz_ore::{assert_err, assert_ok};
674    use uuid::Uuid;
675
676    use crate::location::Blob;
677
678    use super::*;
679
680    fn keys(baseline: &[String], new: &[&str]) -> Vec<String> {
681        let mut ret = baseline.to_vec();
682        ret.extend(new.iter().map(|x| x.to_string()));
683        ret.sort();
684        ret
685    }
686
687    async fn get_keys(b: &impl Blob) -> Result<Vec<String>, ExternalError> {
688        let mut keys = vec![];
689        b.list_keys_and_metadata("", &mut |entry| keys.push(entry.key.to_string()))
690            .await?;
691        Ok(keys)
692    }
693
694    async fn get_keys_with_prefix(
695        b: &impl Blob,
696        prefix: &str,
697    ) -> Result<Vec<String>, ExternalError> {
698        let mut keys = vec![];
699        b.list_keys_and_metadata(prefix, &mut |entry| keys.push(entry.key.to_string()))
700            .await?;
701        Ok(keys)
702    }
703
704    /// Common test impl for different blob implementations.
705    pub async fn blob_impl_test<
706        B: Blob,
707        F: Future<Output = Result<B, ExternalError>>,
708        NewFn: Fn(&'static str) -> F,
709    >(
710        new_fn: NewFn,
711    ) -> Result<(), ExternalError> {
712        let values = ["v0".as_bytes().to_vec(), "v1".as_bytes().to_vec()];
713
714        let blob0 = new_fn("path0").await?;
715
716        // We can create a second blob writing to a different place.
717        let _ = new_fn("path1").await?;
718
719        // We can open two blobs to the same place, even.
720        let blob1 = new_fn("path0").await?;
721
722        let k0 = "foo/bar/k0";
723
724        // Empty key is empty.
725        assert_eq!(blob0.get(k0).await?, None);
726        assert_eq!(blob1.get(k0).await?, None);
727
728        // Empty list keys is empty.
729        let empty_keys = get_keys(&blob0).await?;
730        assert_eq!(empty_keys, Vec::<String>::new());
731        let empty_keys = get_keys(&blob1).await?;
732        assert_eq!(empty_keys, Vec::<String>::new());
733
734        // Set a key and get it back.
735        blob0.set(k0, values[0].clone().into()).await?;
736        assert_eq!(
737            blob0.get(k0).await?.map(|s| s.into_contiguous()),
738            Some(values[0].clone())
739        );
740        assert_eq!(
741            blob1.get(k0).await?.map(|s| s.into_contiguous()),
742            Some(values[0].clone())
743        );
744
745        // Set another key and get it back.
746        blob0.set("k0a", values[0].clone().into()).await?;
747        assert_eq!(
748            blob0.get("k0a").await?.map(|s| s.into_contiguous()),
749            Some(values[0].clone())
750        );
751        assert_eq!(
752            blob1.get("k0a").await?.map(|s| s.into_contiguous()),
753            Some(values[0].clone())
754        );
755
756        // Blob contains the key we just inserted.
757        let mut blob_keys = get_keys(&blob0).await?;
758        blob_keys.sort();
759        assert_eq!(blob_keys, keys(&empty_keys, &[k0, "k0a"]));
760        let mut blob_keys = get_keys(&blob1).await?;
761        blob_keys.sort();
762        assert_eq!(blob_keys, keys(&empty_keys, &[k0, "k0a"]));
763
764        // Can overwrite a key.
765        blob0.set(k0, values[1].clone().into()).await?;
766        assert_eq!(
767            blob0.get(k0).await?.map(|s| s.into_contiguous()),
768            Some(values[1].clone())
769        );
770        assert_eq!(
771            blob1.get(k0).await?.map(|s| s.into_contiguous()),
772            Some(values[1].clone())
773        );
774        // Can overwrite another key.
775        blob0.set("k0a", values[1].clone().into()).await?;
776        assert_eq!(
777            blob0.get("k0a").await?.map(|s| s.into_contiguous()),
778            Some(values[1].clone())
779        );
780        assert_eq!(
781            blob1.get("k0a").await?.map(|s| s.into_contiguous()),
782            Some(values[1].clone())
783        );
784
785        // Can delete a key.
786        assert_eq!(blob0.delete(k0).await, Ok(Some(2)));
787        // Can no longer get a deleted key.
788        assert_eq!(blob0.get(k0).await?, None);
789        assert_eq!(blob1.get(k0).await?, None);
790        // Double deleting a key succeeds but indicates that it did no work.
791        assert_eq!(blob0.delete(k0).await, Ok(None));
792        // Deleting a key that does not exist succeeds.
793        assert_eq!(blob0.delete("nope").await, Ok(None));
794        // Deleting a key with an empty value indicates it did work but deleted
795        // no bytes.
796        blob0.set("empty", Bytes::new()).await?;
797        assert_eq!(blob0.delete("empty").await, Ok(Some(0)));
798
799        // Attempt to restore a key. Not all backends will be able to restore, but
800        // we can confirm that our data is visible iff restore reported success.
801        blob0.set("undelete", Bytes::from("data")).await?;
802        // Restoring should always succeed when the key exists.
803        blob0.restore("undelete").await?;
804        assert_eq!(blob0.delete("undelete").await?, Some("data".len()));
805        let expected = match blob0.restore("undelete").await {
806            Ok(()) => Some(Bytes::from("data").into()),
807            Err(ExternalError::Determinate(_)) => None,
808            Err(other) => return Err(other),
809        };
810        assert_eq!(blob0.get("undelete").await?, expected);
811        blob0.delete("undelete").await?;
812
813        // Empty blob contains no keys.
814        blob0.delete("k0a").await?;
815        let mut blob_keys = get_keys(&blob0).await?;
816        blob_keys.sort();
817        assert_eq!(blob_keys, empty_keys);
818        let mut blob_keys = get_keys(&blob1).await?;
819        blob_keys.sort();
820        assert_eq!(blob_keys, empty_keys);
821        // Can reset a deleted key to some other value.
822        blob0.set(k0, values[1].clone().into()).await?;
823        assert_eq!(
824            blob1.get(k0).await?.map(|s| s.into_contiguous()),
825            Some(values[1].clone())
826        );
827        assert_eq!(
828            blob0.get(k0).await?.map(|s| s.into_contiguous()),
829            Some(values[1].clone())
830        );
831
832        // Insert multiple keys back to back and validate that we can list
833        // them all out.
834        let mut expected_keys = empty_keys;
835        for i in 1..=5 {
836            let key = format!("k{}", i);
837            blob0.set(&key, values[0].clone().into()).await?;
838            expected_keys.push(key);
839        }
840
841        // Blob contains the key we just inserted.
842        let mut blob_keys = get_keys(&blob0).await?;
843        blob_keys.sort();
844        assert_eq!(blob_keys, keys(&expected_keys, &[k0]));
845        let mut blob_keys = get_keys(&blob1).await?;
846        blob_keys.sort();
847        assert_eq!(blob_keys, keys(&expected_keys, &[k0]));
848
849        // Insert multiple keys with a different prefix and validate that we can
850        // list out keys by their prefix
851        let mut expected_prefix_keys = vec![];
852        for i in 1..=3 {
853            let key = format!("k-prefix-{}", i);
854            blob0.set(&key, values[0].clone().into()).await?;
855            expected_prefix_keys.push(key);
856        }
857        let mut blob_keys = get_keys_with_prefix(&blob0, "k-prefix").await?;
858        blob_keys.sort();
859        assert_eq!(blob_keys, expected_prefix_keys);
860        let mut blob_keys = get_keys_with_prefix(&blob0, "k").await?;
861        blob_keys.sort();
862        expected_keys.extend(expected_prefix_keys);
863        expected_keys.sort();
864        assert_eq!(blob_keys, expected_keys);
865
866        // We can open a new blob to the same path and use it.
867        let blob3 = new_fn("path0").await?;
868        assert_eq!(
869            blob3.get(k0).await?.map(|s| s.into_contiguous()),
870            Some(values[1].clone())
871        );
872
873        Ok(())
874    }
875
876    /// Common test impl for different consensus implementations.
877    pub async fn consensus_impl_test<
878        C: Consensus,
879        F: Future<Output = Result<C, ExternalError>>,
880        NewFn: FnMut() -> F,
881    >(
882        mut new_fn: NewFn,
883    ) -> Result<(), ExternalError> {
884        let consensus = new_fn().await?;
885
886        // Use a random key so independent runs of this test don't interfere
887        // with each other.
888        let key = Uuid::new_v4().to_string();
889
890        // Starting value of consensus data is None.
891        assert_eq!(consensus.head(&key).await, Ok(None));
892
893        // Can scan a key that has no data.
894        assert_eq!(consensus.scan(&key, SeqNo(0), SCAN_ALL).await, Ok(vec![]));
895
896        // Cannot truncate data from a key that doesn't have any data
897        assert_err!(consensus.truncate(&key, SeqNo(0)).await);
898
899        let state_at = |v| VersionedData {
900            seqno: SeqNo(v),
901            data: Bytes::from("abc"),
902        };
903
904        // Incorrectly setting the data with a non-initial seqno should fail.
905        assert_eq!(
906            consensus.compare_and_set(&key, state_at(1)).await,
907            Ok(CaSResult::ExpectationMismatch),
908        );
909
910        // Correctly updating the state with the correct expected value should succeed.
911        assert_eq!(
912            consensus.compare_and_set(&key, state_at(0)).await,
913            Ok(CaSResult::Committed),
914        );
915
916        // The new key is visible in state.
917        let keys: Vec<_> = consensus.list_keys().try_collect().await?;
918        assert_eq!(keys, vec![key.to_owned()]);
919
920        // We can observe the a recent value on successful update.
921        assert_eq!(consensus.head(&key).await, Ok(Some(state_at(0))));
922
923        // Can scan a key that has data with a lower bound sequence number < head.
924        assert_eq!(
925            consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
926            Ok(vec![state_at(0)])
927        );
928
929        // Can scan a key that has data with a lower bound sequence number == head.
930        assert_eq!(
931            consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
932            Ok(vec![state_at(0)])
933        );
934
935        // Can scan a key that has data with a lower bound sequence number >
936        // head.
937        assert_eq!(consensus.scan(&key, SeqNo(1), SCAN_ALL).await, Ok(vec![]));
938
939        // Can truncate data with an upper bound <= head, even if there is no data in the
940        // range [0, upper).
941        assert_ok!(consensus.truncate(&key, SeqNo(0)).await);
942
943        // Cannot truncate data with an upper bound > head.
944        assert_err!(consensus.truncate(&key, SeqNo(1)).await);
945
946        let new_state_at = |v| VersionedData {
947            seqno: SeqNo(v),
948            data: Bytes::from("def"),
949        };
950
951        // Trying to update without the correct expected seqno fails, (even if expected > current)
952        assert_eq!(
953            consensus.compare_and_set(&key, new_state_at(3)).await,
954            Ok(CaSResult::ExpectationMismatch),
955        );
956
957        // Trying to update without the correct expected seqno fails, (even if expected < current)
958        assert_eq!(
959            consensus.compare_and_set(&key, new_state_at(0)).await,
960            Ok(CaSResult::ExpectationMismatch),
961        );
962
963        // Can correctly update to a new state if we provide the right expected seqno
964        assert_eq!(
965            consensus.compare_and_set(&key, new_state_at(1)).await,
966            Ok(CaSResult::Committed),
967        );
968
969        // We can observe the a recent value on successful update.
970        assert_eq!(consensus.head(&key).await, Ok(Some(new_state_at(1))));
971
972        // We can observe both states in the correct order with scan if pass
973        // in a suitable lower bound.
974        assert_eq!(
975            consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
976            Ok(vec![state_at(0), new_state_at(1)])
977        );
978
979        // We can observe only the most recent state if the lower bound is higher
980        // than the previous insertion's sequence number.
981        assert_eq!(
982            consensus.scan(&key, SeqNo(1), SCAN_ALL).await,
983            Ok(vec![new_state_at(1)])
984        );
985
986        // We can scan if the provided lower bound > head's sequence number.
987        assert_eq!(consensus.scan(&key, SeqNo(2), SCAN_ALL).await, Ok(vec![]));
988
989        // We can scan with limits that don't cover all states
990        assert_eq!(
991            consensus.scan(&key, SeqNo::minimum(), 1).await,
992            Ok(vec![state_at(0)])
993        );
994
995        // We can scan with limits to cover exactly the number of states
996        assert_eq!(
997            consensus.scan(&key, SeqNo::minimum(), 2).await,
998            Ok(vec![state_at(0), new_state_at(1)])
999        );
1000
1001        // We can scan with a limit larger than the number of states
1002        assert_eq!(
1003            consensus.scan(&key, SeqNo(0), 100).await,
1004            Ok(vec![state_at(0), new_state_at(1)])
1005        );
1006
1007        // Can remove the previous write with the appropriate truncation.
1008        assert_ok!(consensus.truncate(&key, SeqNo(1)).await);
1009
1010        // Verify that the old write is indeed deleted.
1011        assert_eq!(
1012            consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
1013            Ok(vec![new_state_at(1)])
1014        );
1015
1016        // Truncate is idempotent and can be repeated. The return value
1017        // indicates we didn't do any work though.
1018        assert_ok!(consensus.truncate(&key, SeqNo(1)).await);
1019
1020        // Make sure entries under different keys don't clash.
1021        let other_key = Uuid::new_v4().to_string();
1022
1023        assert_eq!(consensus.head(&other_key).await, Ok(None));
1024
1025        let state = VersionedData {
1026            seqno: SeqNo(0),
1027            data: Bytes::from("einszweidrei"),
1028        };
1029
1030        assert_eq!(
1031            consensus.compare_and_set(&other_key, state.clone()).await,
1032            Ok(CaSResult::Committed),
1033        );
1034
1035        assert_eq!(consensus.head(&other_key).await, Ok(Some(state.clone())));
1036
1037        // State for the first key is still as expected.
1038        assert_eq!(consensus.head(&key).await, Ok(Some(new_state_at(1))));
1039
1040        // Trying to update from a stale version of current doesn't work.
1041        let invalid_jump_forward = VersionedData {
1042            seqno: SeqNo(11),
1043            data: Bytes::from("invalid"),
1044        };
1045        assert_eq!(
1046            consensus.compare_and_set(&key, invalid_jump_forward).await,
1047            Ok(CaSResult::ExpectationMismatch),
1048        );
1049
1050        // Writing a large (~10 KiB) amount of data works fine.
1051        let large_state = VersionedData {
1052            seqno: SeqNo(2),
1053            data: std::iter::repeat(b'a').take(10240).collect(),
1054        };
1055        assert_eq!(
1056            consensus.compare_and_set(&key, large_state).await,
1057            Ok(CaSResult::Committed),
1058        );
1059
1060        // Truncate can delete more than one version at a time.
1061        let v3 = VersionedData {
1062            seqno: SeqNo(3),
1063            data: Bytes::new(),
1064        };
1065        assert_eq!(
1066            consensus.compare_and_set(&key, v3).await,
1067            Ok(CaSResult::Committed),
1068        );
1069        assert_ok!(consensus.truncate(&key, SeqNo(3)).await);
1070
1071        // The shard now has head seqno 3 with seqnos 0..=2 truncated away. A re-initialization
1072        // (compare_and_set with an initial seqno, i.e. expected=None) must be rejected rather than
1073        // resurrecting seqno 0 in the gap below the live head.
1074        assert_eq!(
1075            consensus.compare_and_set(&key, state_at(0)).await,
1076            Ok(CaSResult::ExpectationMismatch),
1077        );
1078
1079        Ok(())
1080    }
1081
1082    #[mz_ore::test]
1083    fn timeout_error() {
1084        assert!(ExternalError::new_timeout(Instant::now()).is_timeout());
1085        assert!(!ExternalError::from(anyhow!("foo")).is_timeout());
1086    }
1087}