Skip to main content

mkit_server/pipeline/
watermark.rs

1//! Internal safety reads for native GC and takedown consumers. Workers use
2//! `store::watermark::namespace_relay_watermark_step` to resume bounded scans.
3
4use crate::error::ServerError;
5use crate::repo::NamespaceKey;
6use crate::store::watermark::{self, ActiveShardsPage, WatermarkError};
7use crate::store::{BlobStore, Cursor, KeyClasses, NamespaceStore, Partition};
8
9use super::{HookSet, Pipeline, Sharding, meta_error, ms};
10
11fn map_watermark(error: WatermarkError) -> ServerError {
12    match error {
13        WatermarkError::Recovering => {
14            ServerError::unavailable("coordinator lease table is recovering")
15        }
16        WatermarkError::Store(error) => meta_error(error),
17    }
18}
19
20impl<B: BlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
21    /// Namespace lower bound on commit time of undelivered relay rows for
22    /// native GC and takedown. Consumers compare it against
23    /// `T + MAX_APPLY_WINDOW + margin`. Workers use the resumable store step.
24    /// A new shard's stale-low first report or recovery can lower the value;
25    /// it is unavailable until the recovered lease table is reconciled.
26    ///
27    /// # Errors
28    /// Unavailable during lease-table recovery or on an unreadable store.
29    pub async fn namespace_relay_watermark(&self, ns: &NamespaceKey) -> Result<u64, ServerError> {
30        let now = ms(self.clock.now_ms());
31        if self.cfg.sharding == Sharding::Single {
32            if self.meta.capabilities().key_classes == KeyClasses::RefsOnly {
33                return Ok(now);
34            }
35            let partition = Partition::Namespace(ns.clone());
36            if self.meta.capabilities().key_classes == KeyClasses::All {
37                watermark::check_recovery(&self.meta, &partition, None)
38                    .await
39                    .map_err(map_watermark)?;
40            }
41            return crate::relay::relay_watermark(&self.meta, &partition, now)
42                .await
43                .map_err(meta_error);
44        }
45        watermark::namespace_relay_watermark(&self.meta, &self.shards.coordinator(ns), now)
46            .await
47            .map_err(map_watermark)
48    }
49
50    /// Enumerate coordinator ref-shard rows for a safety re-check. Expired
51    /// rows retained for relay backlog are included.
52    ///
53    /// # Errors
54    /// Unavailable during lease-table recovery or on an unreadable store.
55    pub async fn active_shards(
56        &self,
57        ns: &NamespaceKey,
58        cursor: Option<&Cursor>,
59        limit: u32,
60    ) -> Result<ActiveShardsPage, ServerError> {
61        if self.cfg.sharding == Sharding::Single {
62            if self.meta.capabilities().key_classes == KeyClasses::All {
63                watermark::check_recovery(&self.meta, &Partition::Namespace(ns.clone()), None)
64                    .await
65                    .map_err(map_watermark)?;
66            }
67            return Ok(ActiveShardsPage {
68                shards: if cursor.is_none() && limit > 0 {
69                    vec![Partition::Namespace(ns.clone())]
70                } else {
71                    Vec::new()
72                },
73                next: None,
74            });
75        }
76        watermark::active_shards(&self.meta, &self.shards.coordinator(ns), cursor, limit)
77            .await
78            .map_err(map_watermark)
79    }
80}