Skip to main content

mkit_server/pipeline/
published.rs

1//! Explicit published-value source. Inspection must never fall back to live rows.
2use super::list::{BucketSource, IndexBucket, Scan};
3use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
4use crate::{NamespaceStore, Partition, RepoId, StoreError};
5use mkit_core::hash::Hash;
6
7/// Sorted full ref names and raw ids.
8pub type PublishedRows = Vec<(String, Hash)>;
9/// A validated snapshot, a live fallback, or a fail-closed storage error.
10pub type PublishedBucket = Result<Option<PublishedRows>, StoreError>;
11
12/// Ref data only: the pipeline authorizes each request before invoking this seam.
13pub trait PublishedSource: MaybeSend + MaybeSync {
14    /// Whether the deployment configures inspection.
15    fn inspection_configured(&self) -> bool;
16    /// Whether inputs contain exclusively published values, using the current snapshot format.
17    fn uses_published_values(&self) -> bool {
18        false
19    }
20    /// Optional unsigned `ReadRef` acceleration; false by default.
21    fn read_ref_enabled(&self) -> bool {
22        false
23    }
24    /// A validated bucket, or `None` to use the live published-equivalent index.
25    fn bucket<'a>(
26        &'a self,
27        repo: &'a RepoId,
28        partition: &'a Partition,
29        now_ms: u64,
30    ) -> BoxFuture<'a, PublishedBucket>;
31}
32
33pub(super) struct ReaderBucket<'a, N> {
34    pub store: &'a N,
35    pub partition: &'a Partition,
36    pub source: Option<&'a dyn PublishedSource>,
37    pub now_ms: u64,
38}
39impl<N: NamespaceStore> BucketSource for ReaderBucket<'_, N> {
40    async fn scan(
41        &self,
42        repo: &RepoId,
43        prefix: &str,
44        last: Option<&str>,
45        limit: u32,
46    ) -> Result<Scan, StoreError> {
47        if let Some(source) = self.source {
48            if source.inspection_configured() && !source.uses_published_values() {
49                return Err(StoreError::unavailable("published view unavailable"));
50            }
51            if let Some(rows) = source.bucket(repo, self.partition, self.now_ms).await? {
52                if self.store.capabilities().atomic_multi_key && !source.uses_published_values() {
53                    return Err(StoreError::unavailable("published view unavailable"));
54                }
55                let mut selected = rows
56                    .into_iter()
57                    .filter(|(name, _)| {
58                        name.starts_with(prefix) && last.is_none_or(|last| name.as_str() > last)
59                    })
60                    .take(limit as usize + 1)
61                    .collect::<Vec<_>>();
62                let more = selected.len() > limit as usize;
63                selected.truncate(limit as usize);
64                return Ok(Scan {
65                    rows: selected,
66                    more,
67                });
68            }
69        }
70        IndexBucket {
71            store: self.store,
72            partition: self.partition,
73        }
74        .scan(repo, prefix, last, limit)
75        .await
76    }
77}