use super::list::{BucketSource, IndexBucket, Scan};
use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
use crate::{NamespaceStore, Partition, RepoId, StoreError};
use mkit_core::hash::Hash;
pub type PublishedRows = Vec<(String, Hash)>;
pub type PublishedBucket = Result<Option<PublishedRows>, StoreError>;
pub trait PublishedSource: MaybeSend + MaybeSync {
fn inspection_configured(&self) -> bool;
fn uses_published_values(&self) -> bool {
false
}
fn read_ref_enabled(&self) -> bool {
false
}
fn bucket<'a>(
&'a self,
repo: &'a RepoId,
partition: &'a Partition,
now_ms: u64,
) -> BoxFuture<'a, PublishedBucket>;
}
pub(super) struct ReaderBucket<'a, N> {
pub store: &'a N,
pub partition: &'a Partition,
pub source: Option<&'a dyn PublishedSource>,
pub now_ms: u64,
}
impl<N: NamespaceStore> BucketSource for ReaderBucket<'_, N> {
async fn scan(
&self,
repo: &RepoId,
prefix: &str,
last: Option<&str>,
limit: u32,
) -> Result<Scan, StoreError> {
if let Some(source) = self.source {
if source.inspection_configured() && !source.uses_published_values() {
return Err(StoreError::unavailable("published view unavailable"));
}
if let Some(rows) = source.bucket(repo, self.partition, self.now_ms).await? {
if self.store.capabilities().atomic_multi_key && !source.uses_published_values() {
return Err(StoreError::unavailable("published view unavailable"));
}
let mut selected = rows
.into_iter()
.filter(|(name, _)| {
name.starts_with(prefix) && last.is_none_or(|last| name.as_str() > last)
})
.take(limit as usize + 1)
.collect::<Vec<_>>();
let more = selected.len() > limit as usize;
selected.truncate(limit as usize);
return Ok(Scan {
rows: selected,
more,
});
}
}
IndexBucket {
store: self.store,
partition: self.partition,
}
.scan(repo, prefix, last, limit)
.await
}
}