mkit_server/pipeline/
published.rs1use super::list::{BucketSource, IndexBucket, Scan};
3use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
4use crate::{NamespaceStore, Partition, RepoId, StoreError};
5use mkit_core::hash::Hash;
6
7pub type PublishedRows = Vec<(String, Hash)>;
9pub type PublishedBucket = Result<Option<PublishedRows>, StoreError>;
11
12pub trait PublishedSource: MaybeSend + MaybeSync {
14 fn inspection_configured(&self) -> bool;
16 fn uses_published_values(&self) -> bool {
18 false
19 }
20 fn read_ref_enabled(&self) -> bool {
22 false
23 }
24 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}