mz_timestamp_oracle/
lib.rs1use async_trait::async_trait;
18use mz_ore::now::{EpochMillis, NowFn};
19
20pub mod batching_oracle;
21pub mod config;
22pub mod metrics;
23pub mod postgres_oracle;
24pub mod retry;
25
26pub use config::TimestampOracleConfig;
27
28#[derive(Debug)]
30pub struct WriteTimestamp<T = mz_repr::Timestamp> {
31 pub timestamp: T,
33 pub advance_to: T,
35}
36
37#[async_trait]
44pub trait TimestampOracle<T>: std::fmt::Debug {
45 async fn write_ts(&self) -> WriteTimestamp<T>;
50
51 async fn peek_write_ts(&self) -> T;
53
54 async fn read_ts(&self) -> T;
60
61 async fn apply_write(&self, lower_bound: T);
66}
67
68pub trait GenericNowFn<T>: Clone + Send + Sync {
74 fn now(&self) -> T;
75}
76
77impl GenericNowFn<mz_repr::Timestamp> for NowFn<EpochMillis> {
78 fn now(&self) -> mz_repr::Timestamp {
79 (self)().into()
80 }
81}
82
83impl<T: Clone + Send + Sync> GenericNowFn<T> for NowFn<T> {
84 fn now(&self) -> T {
85 (self)()
86 }
87}
88
89pub mod tests {
92 use std::sync::Arc;
93
94 use futures::Future;
95 use mz_repr::Timestamp;
96
97 use super::*;
98
99 pub async fn timestamp_oracle_impl_test<F, NewFn>(
103 mut new_fn: NewFn,
104 ) -> Result<(), anyhow::Error>
105 where
106 F: Future<Output = Arc<dyn TimestampOracle<Timestamp> + Send + Sync>>,
107 NewFn: FnMut(String, NowFn, Timestamp) -> F,
108 {
109 let timeline = uuid::Uuid::new_v4().to_string();
115 let oracle = new_fn(timeline, NowFn::from(|| 0u64), Timestamp::MIN).await;
116 assert_eq!(oracle.read_ts().await, Timestamp::MIN);
117 assert_eq!(oracle.peek_write_ts().await, Timestamp::MIN);
118
119 let timeline = uuid::Uuid::new_v4().to_string();
121 let oracle = new_fn(timeline, NowFn::from(|| 0u64), Timestamp::MAX).await;
122 assert_eq!(oracle.read_ts().await, Timestamp::MAX);
123 assert_eq!(oracle.peek_write_ts().await, Timestamp::MAX);
124
125 let timeline = uuid::Uuid::new_v4().to_string();
128 let oracle = new_fn(
129 timeline,
130 NowFn::from(|| Timestamp::MAX.step_back().expect("known to work").into()),
131 Timestamp::MIN,
132 )
133 .await;
134 assert_eq!(oracle.read_ts().await, Timestamp::MIN);
136 assert_eq!(oracle.peek_write_ts().await, Timestamp::MIN);
137 assert_eq!(
138 oracle.write_ts().await.timestamp,
139 Timestamp::MAX.step_back().expect("known to work")
140 );
141 assert_eq!(oracle.read_ts().await, Timestamp::MIN);
143 assert_eq!(
144 oracle.peek_write_ts().await,
145 Timestamp::MAX.step_back().expect("known to work")
146 );
147
148 let timeline = uuid::Uuid::new_v4().to_string();
150 let oracle = new_fn(timeline, NowFn::from(|| 0u64), Timestamp::MIN).await;
151 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(1u64));
152 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(2u64));
153
154 let timeline = uuid::Uuid::new_v4().to_string();
156 let oracle = new_fn(timeline, NowFn::from(|| 0u64), Timestamp::MIN).await;
157 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(0u64));
158 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(0u64));
159
160 let timeline = uuid::Uuid::new_v4().to_string();
165 let oracle = new_fn(timeline, NowFn::from(|| 0u64), 10u64.into()).await;
166 oracle.apply_write(5u64.into()).await;
167 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(10u64));
168 assert_eq!(oracle.read_ts().await, Timestamp::from(10u64));
169
170 let timeline = uuid::Uuid::new_v4().to_string();
173 let oracle = new_fn(timeline, NowFn::from(|| 0u64), 0u64.into()).await;
174 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(1u64));
176 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(2u64));
177 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(3u64));
178 assert_eq!(oracle.write_ts().await.timestamp, Timestamp::from(4u64));
179 oracle.apply_write(2u64.into()).await;
180 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(4u64));
181 assert_eq!(oracle.read_ts().await, Timestamp::from(2u64));
182
183 let timeline = uuid::Uuid::new_v4().to_string();
186 let oracle = new_fn(timeline, NowFn::from(|| 0u64), 0u64.into()).await;
187 oracle.apply_write(2u64.into()).await;
188 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(2u64));
189 assert_eq!(oracle.read_ts().await, Timestamp::from(2u64));
190 oracle.apply_write(4u64.into()).await;
191 assert_eq!(oracle.peek_write_ts().await, Timestamp::from(4u64));
192 assert_eq!(oracle.read_ts().await, Timestamp::from(4u64));
193
194 Ok(())
195 }
196}