Skip to main content

mkit_server/store/
watermark.rs

1//! Coordinator relay watermark and shard enumeration. The value is a lower
2//! bound on the commit time of every undelivered relay row. Consumers compare
3//! it against `T + MAX_APPLY_WINDOW + margin`. A new shard's first report may
4//! be stale-low, so the namespace result can decrease without recovery;
5//! recovery can reset maxima too.
6
7use crate::store::{
8    Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, Precondition, ScanPage,
9    StoreError, codec, keys,
10};
11
12const PAGE_SIZE: u32 = 128;
13
14/// A recovery fence or a failed coordinator read.
15#[derive(Debug, thiserror::Error)]
16pub enum WatermarkError {
17    /// The lease table was restored and awaits R-116 reconciliation.
18    #[error("coordinator lease table is recovering")]
19    Recovering,
20    /// A storage or codec failure; consumers fail closed.
21    #[error(transparent)]
22    Store(#[from] StoreError),
23}
24
25/// A resumable scan. The cursor is an opaque store cursor for the `ls` range;
26/// the ceiling and minimum preserve the original scan-time bound.
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct WatermarkCheckpoint {
29    coordinator: Partition,
30    cursor: Cursor,
31    ceiling_ms: u64,
32    minimum_ms: u64,
33    recovery_generation: RecoveryGeneration,
34}
35
36/// The exact recovery and reconciliation marker values observed at scan start.
37pub type RecoveryGeneration = (Option<crate::store::Value>, Option<crate::store::Value>);
38
39fn encode_marker(bytes: &mut Vec<u8>, marker: Option<&crate::store::Value>) {
40    match marker {
41        Some(value) => {
42            bytes.extend_from_slice(
43                &u32::try_from(value.as_bytes().len())
44                    .expect("store value length fits u32")
45                    .to_be_bytes(),
46            );
47            bytes.extend_from_slice(value.as_bytes());
48        }
49        None => bytes.extend_from_slice(&u32::MAX.to_be_bytes()),
50    }
51}
52
53fn decode_marker(bytes: &[u8], at: &mut usize) -> Result<Option<crate::store::Value>, StoreError> {
54    let len = u32::from_be_bytes(
55        bytes
56            .get(*at..*at + 4)
57            .and_then(|slice| slice.try_into().ok())
58            .ok_or_else(|| StoreError::Invalid("truncated watermark generation".into()))?,
59    );
60    *at += 4;
61    if len == u32::MAX {
62        return Ok(None);
63    }
64    let end = at
65        .checked_add(
66            usize::try_from(len)
67                .map_err(|_| StoreError::Invalid("invalid watermark generation length".into()))?,
68        )
69        .ok_or_else(|| StoreError::Invalid("invalid watermark generation length".into()))?;
70    let value = bytes
71        .get(*at..end)
72        .ok_or_else(|| StoreError::Invalid("truncated watermark generation".into()))?;
73    *at = end;
74    Ok(Some(crate::store::Value::new(value.to_vec())))
75}
76
77impl WatermarkCheckpoint {
78    /// Stable binary form: version, partition length and bytes, scan ceiling,
79    /// partial minimum, exact recovery markers, then opaque cursor bytes.
80    ///
81    /// # Panics
82    /// A checkpoint created for an invalid or oversized partition cannot encode.
83    #[must_use]
84    pub fn encode(&self) -> Vec<u8> {
85        let partition = self
86            .coordinator
87            .encode()
88            .expect("coordinator identity encodes");
89        let mut bytes = Vec::with_capacity(27 + partition.len() + self.cursor.as_bytes().len());
90        bytes.push(2);
91        bytes.extend_from_slice(
92            &u16::try_from(partition.len())
93                .expect("partition length fits")
94                .to_be_bytes(),
95        );
96        bytes.extend_from_slice(&partition);
97        bytes.extend_from_slice(&self.ceiling_ms.to_be_bytes());
98        bytes.extend_from_slice(&self.minimum_ms.to_be_bytes());
99        encode_marker(&mut bytes, self.recovery_generation.0.as_ref());
100        encode_marker(&mut bytes, self.recovery_generation.1.as_ref());
101        bytes.extend_from_slice(self.cursor.as_bytes());
102        bytes
103    }
104
105    /// Decode a checkpoint supplied by the previous scan step.
106    pub fn decode(bytes: &[u8]) -> Result<Self, StoreError> {
107        if bytes.len() < 28 || bytes[0] != 2 {
108            return Err(StoreError::Invalid("invalid watermark checkpoint".into()));
109        }
110        let part_len = usize::from(u16::from_be_bytes([bytes[1], bytes[2]]));
111        let end = 3usize
112            .checked_add(part_len)
113            .and_then(|n| n.checked_add(16))
114            .ok_or_else(|| StoreError::Invalid("invalid watermark checkpoint length".into()))?;
115        if bytes.len() < end + 9 {
116            return Err(StoreError::Invalid("truncated watermark checkpoint".into()));
117        }
118        let partition = bytes
119            .get(3..3 + part_len)
120            .ok_or_else(|| StoreError::Invalid("truncated coordinator identity".into()))?;
121        let coordinator = Partition::decode(partition)?;
122        let ceiling_ms = u64::from_be_bytes(
123            bytes
124                .get(end - 16..end - 8)
125                .and_then(|slice| slice.try_into().ok())
126                .ok_or_else(|| StoreError::Invalid("truncated watermark ceiling".into()))?,
127        );
128        let minimum_ms = u64::from_be_bytes(
129            bytes
130                .get(end - 8..end)
131                .and_then(|slice| slice.try_into().ok())
132                .ok_or_else(|| StoreError::Invalid("truncated watermark minimum".into()))?,
133        );
134        if minimum_ms > ceiling_ms {
135            return Err(StoreError::Invalid("invalid watermark minimum".into()));
136        }
137        let mut at = end;
138        let recovery = decode_marker(bytes, &mut at)?;
139        let reconcile = decode_marker(bytes, &mut at)?;
140        if bytes.len() <= at {
141            return Err(StoreError::Invalid("truncated watermark cursor".into()));
142        }
143        recovery
144            .as_ref()
145            .map(codec::decode_lease_recovery)
146            .transpose()?;
147        reconcile.as_ref().map(codec::decode_u64).transpose()?;
148        Ok(Self {
149            coordinator,
150            cursor: Cursor::new(bytes[at..].to_vec()),
151            ceiling_ms,
152            minimum_ms,
153            recovery_generation: (recovery, reconcile),
154        })
155    }
156}
157
158/// One bounded scan step.
159#[derive(Debug, Clone, PartialEq, Eq)]
160pub enum WatermarkStep {
161    /// More pages remain; resume with this checkpoint.
162    Pending(WatermarkCheckpoint),
163    /// The namespace lower bound, capped at scan start.
164    Complete(u64),
165}
166
167/// Check the durable recovery fence before reading a watermark or shard set.
168pub async fn check_recovery<S: NamespaceStore>(
169    store: &S,
170    coordinator: &Partition,
171    expected: Option<&RecoveryGeneration>,
172) -> Result<RecoveryGeneration, WatermarkError> {
173    if !matches!(
174        coordinator,
175        Partition::Coordinator(_) | Partition::Namespace(_)
176    ) {
177        return Err(WatermarkError::Store(StoreError::Invalid(
178            "watermark requires namespace coordinator".into(),
179        )));
180    }
181    let rows = store
182        .get_many(
183            coordinator,
184            &[keys::lease_recovery(), keys::lease_reconcile()],
185        )
186        .await?;
187    let [recovery, reconcile] = rows.as_slice() else {
188        return Err(WatermarkError::Store(StoreError::Corrupt(
189            "recovery get_many length".into(),
190        )));
191    };
192    let generation = (recovery.clone(), reconcile.clone());
193    if expected.is_some_and(|prior| *prior != generation) {
194        return Err(WatermarkError::Store(StoreError::Invalid(
195            "watermark checkpoint recovery generation changed".into(),
196        )));
197    }
198    let recovered = recovery
199        .as_ref()
200        .map(codec::decode_lease_recovery)
201        .transpose()?;
202    let reconciled = reconcile.as_ref().map(codec::decode_u64).transpose()?;
203    if recovered
204        .and_then(codec::LeaseRecovery::recovery_time)
205        .is_some_and(|resumed| reconciled.is_none_or(|at| at <= resumed))
206    {
207        return Err(WatermarkError::Recovering);
208    }
209    Ok(generation)
210}
211
212/// Mark the current recovery generation reconciled, only after an external
213/// R-116 driver has rebuilt the lease and index tables. A concurrent new
214/// recovery defeats the guarded write.
215pub async fn mark_lease_table_reconciled<S: NamespaceStore>(
216    store: &S,
217    coordinator: &Partition,
218    at_ms: u64,
219) -> Result<(), StoreError> {
220    if !matches!(
221        coordinator,
222        Partition::Coordinator(_) | Partition::Namespace(_)
223    ) {
224        return Err(StoreError::Invalid(
225            "reconciliation requires namespace coordinator".into(),
226        ));
227    }
228    let key = keys::lease_recovery();
229    let value = store
230        .get(coordinator, &key)
231        .await?
232        .ok_or_else(|| StoreError::Invalid("no lease recovery to reconcile".into()))?;
233    let recovery = codec::decode_lease_recovery(&value)?;
234    let later = recovery
235        .recovery_time()
236        .ok_or_else(|| StoreError::Invalid("no actual lease recovery to reconcile".into()))?
237        .checked_add(1)
238        .ok_or_else(|| StoreError::Invalid("lease recovery time has no successor".into()))?;
239    let timestamp = at_ms.max(later);
240    let batch = Batch::new()
241        .require(Precondition::Equals(key, value))
242        .put(keys::lease_reconcile(), codec::encode_u64(timestamp));
243    match store.apply(coordinator, batch).await? {
244        BatchOutcome::Committed => Ok(()),
245        _ => Err(StoreError::Corrupt(
246            "lease reconciliation raced recovery".into(),
247        )),
248    }
249}
250
251fn decode_shard(
252    key: &Key,
253    value: &crate::store::Value,
254    coordinator: &Partition,
255) -> Result<(Partition, u64), StoreError> {
256    let Some(keys::ParsedKey::LeasedShard { repo, shard_ref }) = keys::parse(key) else {
257        return Err(StoreError::Corrupt("invalid leased shard key".into()));
258    };
259    let Partition::Coordinator(ns) = coordinator else {
260        return Err(StoreError::Invalid(
261            "watermark requires coordinator partition".into(),
262        ));
263    };
264    let row = codec::decode_leased_shard(value)?;
265    Ok((
266        Partition::Ref {
267            ns: ns.clone(),
268            repo,
269            shard_ref,
270        },
271        row.relay_watermark_ms,
272    ))
273}
274
275/// Scan one page of the coordinator table. Workers consumers call this
276/// resumable step. A new shard can enter behind the cursor only after a new
277/// lease grant, so its first relay commit is after the checkpoint ceiling.
278/// Existing rows retain their running maximum.
279pub async fn namespace_relay_watermark_step<S: NamespaceStore>(
280    store: &S,
281    coordinator: &Partition,
282    now_ms: u64,
283    checkpoint: Option<WatermarkCheckpoint>,
284    limit: u32,
285) -> Result<WatermarkStep, WatermarkError> {
286    if !matches!(coordinator, Partition::Coordinator(_)) {
287        return Err(WatermarkError::Store(StoreError::Invalid(
288            "watermark scan requires coordinator partition".into(),
289        )));
290    }
291    let generation = check_recovery(
292        store,
293        coordinator,
294        checkpoint.as_ref().map(|c| &c.recovery_generation),
295    )
296    .await?;
297    let (start, end) = keys::class_range(keys::TAG_LEASED_SHARD);
298    let (cursor, ceiling_ms, mut minimum_ms) = match checkpoint {
299        Some(c) if c.coordinator == *coordinator => (Some(c.cursor), c.ceiling_ms, c.minimum_ms),
300        Some(_) => {
301            return Err(WatermarkError::Store(StoreError::Invalid(
302                "watermark checkpoint belongs to another coordinator".into(),
303            )));
304        }
305        None => (None, now_ms, now_ms),
306    };
307    let page = store
308        .scan(coordinator, &start, &end, cursor.as_ref(), limit.max(1))
309        .await?;
310    for (key, value) in &page.entries {
311        minimum_ms = minimum_ms.min(decode_shard(key, value, coordinator)?.1);
312    }
313    Ok(if let Some(cursor) = page.next {
314        WatermarkStep::Pending(WatermarkCheckpoint {
315            coordinator: coordinator.clone(),
316            cursor,
317            ceiling_ms,
318            minimum_ms,
319            recovery_generation: generation,
320        })
321    } else {
322        check_recovery(store, coordinator, Some(&generation)).await?;
323        WatermarkStep::Complete(minimum_ms)
324    })
325}
326
327/// Minimum of scan-start `now_ms` and all shard maxima. Native consumers may
328/// use this full scan; Workers use [`namespace_relay_watermark_step`]. The
329/// caller compares it against `T + MAX_APPLY_WINDOW + margin`. It can decrease
330/// when a new shard enters with a stale-low report or on restore.
331pub async fn namespace_relay_watermark<S: NamespaceStore>(
332    store: &S,
333    coordinator: &Partition,
334    now_ms: u64,
335) -> Result<u64, WatermarkError> {
336    let mut checkpoint = None;
337    loop {
338        match namespace_relay_watermark_step(store, coordinator, now_ms, checkpoint, PAGE_SIZE)
339            .await?
340        {
341            WatermarkStep::Pending(next) => checkpoint = Some(next),
342            WatermarkStep::Complete(value) => return Ok(value),
343        }
344    }
345}
346
347/// A page of all coordinator ref-shard rows, including expired rows kept for
348/// an undelivered outbox. This is the shard set GC must include.
349#[derive(Debug, Clone, PartialEq, Eq)]
350pub struct ActiveShardsPage {
351    /// Ref shard identities.
352    pub shards: Vec<Partition>,
353    /// Opaque continuation cursor.
354    pub next: Option<Cursor>,
355}
356
357/// Scan coordinator shard rows, failing on recovery or any undecodable row.
358pub async fn active_shards<S: NamespaceStore>(
359    store: &S,
360    coordinator: &Partition,
361    cursor: Option<&Cursor>,
362    limit: u32,
363) -> Result<ActiveShardsPage, WatermarkError> {
364    if !matches!(coordinator, Partition::Coordinator(_)) {
365        return Err(WatermarkError::Store(StoreError::Invalid(
366            "shard scan requires coordinator partition".into(),
367        )));
368    }
369    check_recovery(store, coordinator, None).await?;
370    let (start, end) = keys::class_range(keys::TAG_LEASED_SHARD);
371    let ScanPage { entries, next } = store
372        .scan(coordinator, &start, &end, cursor, limit.max(1))
373        .await?;
374    let shards = entries
375        .iter()
376        .map(|(key, value)| decode_shard(key, value, coordinator).map(|(p, _)| p))
377        .collect::<Result<_, _>>()?;
378    Ok(ActiveShardsPage { shards, next })
379}
380
381#[cfg(all(test, feature = "memory"))]
382mod tests {
383    use super::*;
384    use crate::memory::MemoryKv;
385    use crate::repo::{NamespaceKey, RepoName};
386    use crate::store::{Batch, BatchOutcome, PartitionStats, StoreCapabilities, Value};
387    use std::sync::{
388        Arc,
389        atomic::{AtomicBool, Ordering},
390    };
391
392    struct RecoverDuringScan {
393        inner: Arc<MemoryKv>,
394        once: AtomicBool,
395    }
396
397    impl NamespaceStore for RecoverDuringScan {
398        fn capabilities(&self) -> StoreCapabilities {
399            self.inner.capabilities()
400        }
401        async fn get(&self, p: &Partition, k: &Key) -> Result<Option<Value>, StoreError> {
402            self.inner.get(p, k).await
403        }
404        async fn get_many(
405            &self,
406            p: &Partition,
407            keys: &[Key],
408        ) -> Result<Vec<Option<Value>>, StoreError> {
409            self.inner.get_many(p, keys).await
410        }
411        async fn scan(
412            &self,
413            p: &Partition,
414            start: &Key,
415            end: &Key,
416            after: Option<&Cursor>,
417            limit: u32,
418        ) -> Result<ScanPage, StoreError> {
419            let page = self.inner.scan(p, start, end, after, limit).await?;
420            if self.once.swap(false, Ordering::SeqCst) {
421                self.inner
422                    .apply(
423                        p,
424                        Batch::new().put(
425                            keys::lease_recovery(),
426                            codec::encode_lease_recovery(&codec::LeaseRecovery {
427                                authority_fence: None,
428                                authority_ready: None,
429                                activation_only: None,
430                                resumed_at_ms: 110,
431                            }),
432                        ),
433                    )
434                    .await?;
435            }
436            Ok(page)
437        }
438        async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
439            self.inner.apply(p, batch).await
440        }
441        async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
442            self.inner.stats(p).await
443        }
444        async fn probe(&self) -> Result<(), StoreError> {
445            self.inner.probe().await
446        }
447    }
448
449    fn coordinator() -> Partition {
450        Partition::Coordinator(NamespaceKey::deployment_default())
451    }
452    fn repo(n: u8) -> RepoName {
453        RepoName::new(format!("r{n}")).expect("test repository name")
454    }
455    fn source(n: u8) -> Partition {
456        Partition::Ref {
457            ns: NamespaceKey::deployment_default(),
458            repo: repo(n),
459            shard_ref: "refs/heads/main".into(),
460        }
461    }
462    fn lease(watermark: u64, expiry: u64) -> codec::LeasedShard {
463        codec::LeasedShard {
464            authority_generation: None,
465            acked_authority_generation: None,
466            epoch: 1,
467            expires_at_ms: expiry,
468            acked_epoch: 1,
469            relay_watermark_ms: watermark,
470            sweep_due_ms: expiry,
471        }
472    }
473    async fn put_lease(store: &MemoryKv, n: u8, watermark: u64, expiry: u64) {
474        store
475            .apply(
476                &coordinator(),
477                Batch::new().put(
478                    keys::leased_shard(&repo(n), "refs/heads/main"),
479                    codec::encode_leased_shard(&lease(watermark, expiry)),
480                ),
481            )
482            .await
483            .expect("insert test lease");
484    }
485    #[tokio::test]
486    async fn pages_and_inserts_preserve_scan_start_ceiling() {
487        let store = MemoryKv::default();
488        put_lease(&store, 0, 40, 200).await;
489        put_lease(&store, 2, 90, 200).await;
490        let WatermarkStep::Pending(checkpoint) =
491            namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
492                .await
493                .unwrap()
494        else {
495            panic!("first page");
496        };
497        let encoded = checkpoint.encode();
498        let checkpoint = WatermarkCheckpoint::decode(&encoded).unwrap();
499        let other = Partition::Coordinator(NamespaceKey::from_stored("another".into()));
500        assert!(matches!(
501            namespace_relay_watermark_step(&store, &other, 150, Some(checkpoint.clone()), 1).await,
502            Err(WatermarkError::Store(StoreError::Invalid(_)))
503        ));
504        put_lease(&store, 1, 70, 200).await;
505        let WatermarkStep::Pending(checkpoint) =
506            namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1)
507                .await
508                .unwrap()
509        else {
510            panic!("second page");
511        };
512        let WatermarkStep::Complete(value) =
513            namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1)
514                .await
515                .unwrap()
516        else {
517            panic!("third page");
518        };
519        assert_eq!(value, 40);
520        let page = active_shards(&store, &coordinator(), None, 1)
521            .await
522            .unwrap();
523        assert_eq!(page.shards, vec![source(0)]);
524        assert!(page.next.is_some());
525    }
526
527    #[tokio::test]
528    async fn resume_rejects_a_new_recovery_generation() {
529        let store = MemoryKv::default();
530        put_lease(&store, 0, 40, 200).await;
531        put_lease(&store, 1, 90, 200).await;
532        store
533            .apply(
534                &coordinator(),
535                Batch::new()
536                    .put(
537                        keys::lease_recovery(),
538                        codec::encode_lease_recovery(&codec::LeaseRecovery {
539                            authority_fence: None,
540                            authority_ready: None,
541                            activation_only: None,
542                            resumed_at_ms: 90,
543                        }),
544                    )
545                    .put(keys::lease_reconcile(), codec::encode_u64(91)),
546            )
547            .await
548            .unwrap();
549        let WatermarkStep::Pending(checkpoint) =
550            namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
551                .await
552                .unwrap()
553        else {
554            panic!("first page");
555        };
556        let checkpoint = WatermarkCheckpoint::decode(&checkpoint.encode()).unwrap();
557        store
558            .apply(
559                &coordinator(),
560                Batch::new()
561                    .put(
562                        keys::lease_recovery(),
563                        codec::encode_lease_recovery(&codec::LeaseRecovery {
564                            authority_fence: None,
565                            authority_ready: None,
566                            activation_only: None,
567                            resumed_at_ms: 110,
568                        }),
569                    )
570                    .put(keys::lease_reconcile(), codec::encode_u64(111)),
571            )
572            .await
573            .unwrap();
574        assert!(matches!(
575            namespace_relay_watermark_step(
576                &store,
577                &coordinator(),
578                120,
579                Some(checkpoint.clone()),
580                1
581            )
582            .await,
583            Err(WatermarkError::Store(StoreError::Invalid(_)))
584        ));
585        store
586            .apply(
587                &coordinator(),
588                Batch::new()
589                    .put(
590                        keys::lease_recovery(),
591                        codec::encode_lease_recovery(&codec::LeaseRecovery {
592                            authority_fence: None,
593                            authority_ready: None,
594                            activation_only: None,
595                            resumed_at_ms: 130,
596                        }),
597                    )
598                    .delete(keys::lease_reconcile()),
599            )
600            .await
601            .unwrap();
602        assert!(matches!(
603            namespace_relay_watermark_step(
604                &store,
605                &coordinator(),
606                140,
607                Some(checkpoint.clone()),
608                1
609            )
610            .await,
611            Err(WatermarkError::Store(StoreError::Invalid(_)))
612        ));
613        store
614            .apply(
615                &coordinator(),
616                Batch::new().put(keys::lease_recovery(), Value::new(&b"corrupt marker"[..])),
617            )
618            .await
619            .unwrap();
620        assert!(matches!(
621            namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1).await,
622            Err(WatermarkError::Store(StoreError::Invalid(_)))
623        ));
624    }
625
626    #[tokio::test]
627    async fn final_page_rechecks_recovery_generation() {
628        let inner = Arc::new(MemoryKv::default());
629        put_lease(inner.as_ref(), 0, 40, 200).await;
630        let store = RecoverDuringScan {
631            inner,
632            once: AtomicBool::new(true),
633        };
634        assert!(matches!(
635            namespace_relay_watermark_step(&store, &coordinator(), 100, None, 10).await,
636            Err(WatermarkError::Store(StoreError::Invalid(_)))
637        ));
638    }
639
640    #[tokio::test]
641    async fn a_new_row_behind_the_cursor_is_excluded_by_the_scan_ceiling() {
642        let store = MemoryKv::default();
643        put_lease(&store, 1, 80, 200).await;
644        put_lease(&store, 2, 90, 200).await;
645        let WatermarkStep::Pending(checkpoint) =
646            namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
647                .await
648                .unwrap()
649        else {
650            panic!("first page");
651        };
652        put_lease(&store, 0, 0, 200).await;
653        let WatermarkStep::Complete(value) =
654            namespace_relay_watermark_step(&store, &coordinator(), 120, Some(checkpoint), 1)
655                .await
656                .unwrap()
657        else {
658            panic!("last page");
659        };
660        assert_eq!(value, 80);
661        assert!(value <= 100, "the scan-start ceiling remains binding");
662    }
663
664    #[tokio::test]
665    async fn corruption_and_recovery_fail_closed() {
666        let store = MemoryKv::default();
667        store
668            .apply(
669                &coordinator(),
670                Batch::new().put(
671                    keys::leased_shard(&repo(0), "refs/heads/main"),
672                    Value::new(&b"bad"[..]),
673                ),
674            )
675            .await
676            .unwrap();
677        assert!(matches!(
678            namespace_relay_watermark(&store, &coordinator(), 100).await,
679            Err(WatermarkError::Store(StoreError::Corrupt(_)))
680        ));
681        assert!(
682            active_shards(&store, &coordinator(), None, 10)
683                .await
684                .is_err()
685        );
686        store
687            .apply(
688                &coordinator(),
689                Batch::new().put(
690                    keys::lease_recovery(),
691                    codec::encode_lease_recovery(&codec::LeaseRecovery {
692                        authority_fence: None,
693                        authority_ready: None,
694                        activation_only: None,
695                        resumed_at_ms: 100,
696                    }),
697                ),
698            )
699            .await
700            .unwrap();
701        assert!(matches!(
702            namespace_relay_watermark(&store, &coordinator(), 200).await,
703            Err(WatermarkError::Recovering)
704        ));
705        mark_lease_table_reconciled(&store, &coordinator(), 101)
706            .await
707            .unwrap();
708        assert!(matches!(
709            namespace_relay_watermark(&store, &coordinator(), 200).await,
710            Err(WatermarkError::Store(StoreError::Corrupt(_)))
711        ));
712    }
713}