1use 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#[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 pub fn previous(self) -> Option<SeqNo> {
89 Some(SeqNo(self.0.checked_sub(1)?))
90 }
91
92 pub fn next(self) -> SeqNo {
94 SeqNo(self.0 + 1)
95 }
96
97 pub fn minimum() -> Self {
99 SeqNo(0)
100 }
101
102 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#[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 pub fn new(inner: anyhow::Error) -> Self {
149 Determinate { inner }
150 }
151
152 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#[derive(Debug)]
164pub struct Indeterminate {
165 pub(crate) inner: anyhow::Error,
166}
167
168impl Indeterminate {
169 pub fn new(inner: anyhow::Error) -> Self {
173 Indeterminate { inner }
174 }
175
176 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#[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#[derive(Debug)]
209pub enum ExternalError {
210 Determinate(Determinate),
212 Indeterminate(Indeterminate),
214}
215
216impl ExternalError {
217 #[track_caller]
222 pub fn new_timeout(deadline: Instant) -> Self {
223 ExternalError::Indeterminate(Indeterminate {
224 inner: anyhow!("timeout at {:?}", deadline),
225 })
226 }
227
228 pub fn is_timeout(&self) -> bool {
233 self.to_string().contains("timeout")
235 }
236
237 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#[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 &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 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 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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
392pub struct VersionedData {
393 pub seqno: SeqNo,
395 pub data: Bytes,
397}
398
399#[allow(clippy::as_conversions)]
403pub const SCAN_ALL: usize = u64_to_usize(i64::MAX as u64);
404
405pub const CONSENSUS_HEAD_LIVENESS_KEY: &str = "LIVENESS";
407
408#[derive(Debug, PartialEq, Serialize, Deserialize)]
410pub enum CaSResult {
411 Committed,
413 ExpectationMismatch,
415}
416
417#[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
429pub type ResultStream<'a, T> = Pin<Box<dyn Stream<Item = Result<T, ExternalError>> + Send + 'a>>;
432
433#[async_trait]
442pub trait Consensus: std::fmt::Debug + Send + Sync {
443 fn list_keys(&self) -> ResultStream<'_, String>;
445
446 async fn head(&self, key: &str) -> Result<Option<VersionedData>, ExternalError>;
449
450 async fn compare_and_set(
455 &self,
456 key: &str,
457 new: VersionedData,
458 ) -> Result<CaSResult, ExternalError>;
459
460 async fn scan(
466 &self,
467 key: &str,
468 from: SeqNo,
469 limit: usize,
470 ) -> Result<Vec<VersionedData>, ExternalError>;
471
472 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 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#[derive(Debug)]
544pub struct BlobMetadata<'a> {
545 pub key: &'a str,
547 pub size_in_bytes: u64,
549}
550
551pub const BLOB_GET_LIVENESS_KEY: &str = "LIVENESS";
553
554#[async_trait]
566pub trait Blob: std::fmt::Debug + Send + Sync {
567 async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError>;
569
570 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 async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError>;
585
586 async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError>;
591
592 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 async fn list_keys_and_metadata(
621 &self,
622 key_prefix: &str,
623 f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
624 ) -> Result<(), ExternalError> {
625 self.0.list_keys_and_metadata(key_prefix, f).await
628 }
629
630 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 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#[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 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 let _ = new_fn("path1").await?;
718
719 let blob1 = new_fn("path0").await?;
721
722 let k0 = "foo/bar/k0";
723
724 assert_eq!(blob0.get(k0).await?, None);
726 assert_eq!(blob1.get(k0).await?, None);
727
728 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 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 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 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 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 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 assert_eq!(blob0.delete(k0).await, Ok(Some(2)));
787 assert_eq!(blob0.get(k0).await?, None);
789 assert_eq!(blob1.get(k0).await?, None);
790 assert_eq!(blob0.delete(k0).await, Ok(None));
792 assert_eq!(blob0.delete("nope").await, Ok(None));
794 blob0.set("empty", Bytes::new()).await?;
797 assert_eq!(blob0.delete("empty").await, Ok(Some(0)));
798
799 blob0.set("undelete", Bytes::from("data")).await?;
802 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 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 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 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 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 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 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 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 let key = Uuid::new_v4().to_string();
889
890 assert_eq!(consensus.head(&key).await, Ok(None));
892
893 assert_eq!(consensus.scan(&key, SeqNo(0), SCAN_ALL).await, Ok(vec![]));
895
896 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 assert_eq!(
906 consensus.compare_and_set(&key, state_at(1)).await,
907 Ok(CaSResult::ExpectationMismatch),
908 );
909
910 assert_eq!(
912 consensus.compare_and_set(&key, state_at(0)).await,
913 Ok(CaSResult::Committed),
914 );
915
916 let keys: Vec<_> = consensus.list_keys().try_collect().await?;
918 assert_eq!(keys, vec![key.to_owned()]);
919
920 assert_eq!(consensus.head(&key).await, Ok(Some(state_at(0))));
922
923 assert_eq!(
925 consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
926 Ok(vec![state_at(0)])
927 );
928
929 assert_eq!(
931 consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
932 Ok(vec![state_at(0)])
933 );
934
935 assert_eq!(consensus.scan(&key, SeqNo(1), SCAN_ALL).await, Ok(vec![]));
938
939 assert_ok!(consensus.truncate(&key, SeqNo(0)).await);
942
943 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 assert_eq!(
953 consensus.compare_and_set(&key, new_state_at(3)).await,
954 Ok(CaSResult::ExpectationMismatch),
955 );
956
957 assert_eq!(
959 consensus.compare_and_set(&key, new_state_at(0)).await,
960 Ok(CaSResult::ExpectationMismatch),
961 );
962
963 assert_eq!(
965 consensus.compare_and_set(&key, new_state_at(1)).await,
966 Ok(CaSResult::Committed),
967 );
968
969 assert_eq!(consensus.head(&key).await, Ok(Some(new_state_at(1))));
971
972 assert_eq!(
975 consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
976 Ok(vec![state_at(0), new_state_at(1)])
977 );
978
979 assert_eq!(
982 consensus.scan(&key, SeqNo(1), SCAN_ALL).await,
983 Ok(vec![new_state_at(1)])
984 );
985
986 assert_eq!(consensus.scan(&key, SeqNo(2), SCAN_ALL).await, Ok(vec![]));
988
989 assert_eq!(
991 consensus.scan(&key, SeqNo::minimum(), 1).await,
992 Ok(vec![state_at(0)])
993 );
994
995 assert_eq!(
997 consensus.scan(&key, SeqNo::minimum(), 2).await,
998 Ok(vec![state_at(0), new_state_at(1)])
999 );
1000
1001 assert_eq!(
1003 consensus.scan(&key, SeqNo(0), 100).await,
1004 Ok(vec![state_at(0), new_state_at(1)])
1005 );
1006
1007 assert_ok!(consensus.truncate(&key, SeqNo(1)).await);
1009
1010 assert_eq!(
1012 consensus.scan(&key, SeqNo(0), SCAN_ALL).await,
1013 Ok(vec![new_state_at(1)])
1014 );
1015
1016 assert_ok!(consensus.truncate(&key, SeqNo(1)).await);
1019
1020 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 assert_eq!(consensus.head(&key).await, Ok(Some(new_state_at(1))));
1039
1040 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 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 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 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}