use std::cmp::Ordering;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use lunaris_core::storage::types::Filter;
use lunaris_core::{Embedder, Hlc, KeywordPort, LunarisError, Scope, StoragePort};
use ulid::Ulid;
use crate::operators::recency::{RecencyConfig, rescore_recency};
pub use crate::operators::graph::{DEFAULT_GRAPH_HOPS, DEFAULT_GRAPH_K, Graph, MAX_GRAPH_HOPS};
use crate::hydrate::hydrate_mixed;
use crate::operators::modifiers::FilterParseError;
use crate::operators::vector::Vector;
use crate::operators::{QueryContext, Retriever};
use crate::types::RawHit;
use crate::types::{Hit, Query};
#[must_use = "RetrievalBuilder must be terminated with .vector(..) / .keyword(..) / .graph(..) and awaited to actually retrieve hits"]
pub struct RetrievalBuilder {
pub(crate) root: Box<dyn Retriever>,
pub(crate) embedder: Arc<dyn Embedder>,
pub(crate) storage: Arc<dyn StoragePort>,
pub(crate) keyword: Arc<dyn KeywordPort>,
pub(crate) moon_storage: Option<Arc<lunaris_storage_moon::MoonStorage>>,
pub(crate) base_filter: Option<Filter>,
pub(crate) base_as_of: Option<Hlc>,
pub(crate) scope: Scope,
pub(crate) initial_degraded: bool,
pub(crate) recency_config: Option<RecencyConfig>,
#[allow(clippy::type_complexity)]
pub(crate) boost_cache: Option<Arc<parking_lot::RwLock<lru::LruCache<(Scope, Ulid), f32>>>>,
pub(crate) boost_provider: Option<Arc<dyn crate::boost_provider::BoostProvider>>,
}
impl RetrievalBuilder {
pub fn new(
storage: Arc<dyn StoragePort>,
keyword: Arc<dyn KeywordPort>,
embedder: Arc<dyn Embedder>,
) -> Self {
Self {
root: Box::new(Vector::new("chunks", 30)),
storage,
keyword,
embedder,
moon_storage: None,
base_filter: None,
base_as_of: None,
scope: Scope::dev(),
initial_degraded: false,
recency_config: None,
boost_cache: None,
boost_provider: None,
}
}
pub fn from_handle(
storage: Arc<dyn StoragePort>,
keyword: Arc<dyn KeywordPort>,
embedder: Arc<dyn Embedder>,
) -> Self {
Self::new(storage, keyword, embedder)
}
pub fn with_scope(mut self, scope: Scope) -> Self {
self.scope = scope;
self
}
pub fn with_moon_storage(mut self, moon: Arc<lunaris_storage_moon::MoonStorage>) -> Self {
self.moon_storage = Some(moon);
self
}
pub fn root_plan(&self) -> String {
crate::composition::plan_repr(self.root.as_ref())
}
pub fn with_root<R: Retriever + 'static>(mut self, root: R) -> Self {
self.root = Box::new(root);
self
}
pub fn with_root_boxed(mut self, root: Box<dyn Retriever>) -> Self {
self.root = root;
self
}
pub fn filter(mut self, f: Filter) -> Self {
self.base_filter = Some(f);
self
}
pub fn filter_str(mut self, s: &str) -> Result<Self, FilterParseError> {
let f = crate::operators::modifiers::filter_str(s)?;
self.base_filter = Some(f);
Ok(self)
}
pub fn as_of(mut self, ts: Hlc) -> Self {
self.base_as_of = Some(ts);
self
}
pub fn with_initial_degraded(mut self, deg: bool) -> Self {
self.initial_degraded = deg;
self
}
pub fn top(self, n: usize) -> Self {
let new_root = crate::operators::modifiers::TopRetriever::new(self.root, n);
Self { root: Box::new(new_root), ..self }
}
pub fn rerank(self, reranker: Arc<dyn lunaris_rerank::Reranker>) -> Self {
let new_root = crate::operators::rerank::RerankRetriever::new(self.root, reranker);
Self { root: Box::new(new_root), ..self }
}
pub fn rerank_with_threshold(
self,
reranker: Arc<dyn lunaris_rerank::Reranker>,
min_score: f32,
) -> Self {
let new_root = crate::operators::rerank::RerankRetriever::new(self.root, reranker)
.with_min_score(min_score);
Self { root: Box::new(new_root), ..self }
}
pub fn degraded_fallback<R: Retriever + 'static>(self, fallback: R) -> Self {
let new_root = crate::operators::degraded::DegradedFallbackRetriever::new(
self.root,
Box::new(fallback),
);
Self { root: Box::new(new_root), ..self }
}
pub fn tree(mut self, index: impl Into<String>, k: usize, depth: usize) -> Self {
let root = crate::operators::tree::Tree::new(index, k).with_depth(depth);
self.root = Box::new(root);
self
}
pub fn recency(mut self, config: RecencyConfig) -> Self {
self.recency_config = Some(config);
self
}
pub fn with_boost_cache(
mut self,
cache: Arc<parking_lot::RwLock<lru::LruCache<(Scope, Ulid), f32>>>,
) -> Self {
self.boost_cache = Some(cache);
self
}
pub fn with_boost_provider(
mut self,
provider: Arc<dyn crate::boost_provider::BoostProvider>,
) -> Self {
self.boost_provider = Some(provider);
self
}
pub async fn execute(mut self, mut query: Query) -> Result<Vec<Hit>, LunarisError> {
if self.scope == Scope::dev() {
tracing::warn!(
"RetrievalBuilder using Scope::dev() default — migrate to engine.scoped(scope).dsl()"
);
}
if query.filter.is_none() {
query.filter = self.base_filter.clone();
}
if query.as_of.is_none() {
query.as_of = self.base_as_of;
}
if query.as_of.is_none() {
query.as_of = Some(lunaris_core::HlcClock::new(0).tick());
}
let as_of = query.as_of;
let initial_degraded = self.initial_degraded;
let scope = self.scope.clone();
let recency_config = self.recency_config;
let boost_cache = self.boost_cache.take();
let boost_provider = self.boost_provider.take();
let ctx = match self.moon_storage.clone() {
Some(moon) => QueryContext::with_moon(
query,
scope,
self.embedder,
self.storage.clone(),
self.keyword,
moon,
),
None => {
QueryContext::new(query, scope, self.embedder, self.storage.clone(), self.keyword)
}
};
let raw = self.root.retrieve(&ctx).await?;
let mut hits =
hydrate_mixed(self.storage.as_ref(), &ctx.scope, raw, as_of, initial_degraded).await?;
if let Some(cfg) = recency_config {
let now = as_of.unwrap_or_else(|| {
let wall_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
Hlc::from_parts(wall_ms, 0, 0)
});
rescore_recency(&mut hits, now, &cfg);
}
if let Some(provider) = boost_provider {
fn ledger_id(hit: &Hit) -> Option<Ulid> {
let bytes = if hit.episode_id.is_empty() { &hit.id } else { &hit.episode_id };
<[u8; 16]>::try_from(bytes.as_slice()).ok().map(Ulid::from_bytes)
}
let ids: Vec<Ulid> = hits.iter().filter_map(ledger_id).collect();
if !ids.is_empty() {
let priors = provider.priors(&ctx.scope, &ids).await;
if !priors.is_empty() {
for hit in &mut hits {
if let Some(id) = ledger_id(hit)
&& let Some(&prior) = priors.get(&id)
{
hit.score += prior;
}
}
hits.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(Ordering::Equal));
}
}
}
if let Some(cache) = boost_cache {
let snapshot: HashMap<Ulid, f32> = {
let guard = cache.read();
hits.iter()
.filter_map(|h| {
let id = Ulid::from_bytes(h.id.as_slice().try_into().ok()?);
guard.peek(&(ctx.scope.clone(), id)).map(|&v| (id, v))
})
.collect()
};
if !snapshot.is_empty() {
for hit in &mut hits {
if let Ok(arr) = <[u8; 16]>::try_from(hit.id.as_slice()) {
let id = Ulid::from_bytes(arr);
if let Some(&delta) = snapshot.get(&id) {
hit.score += delta;
}
}
}
hits.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(Ordering::Equal));
}
}
Ok(hits)
}
pub async fn execute_raw(self, mut query: Query) -> Result<Vec<RawHit>, LunarisError> {
if query.filter.is_none() {
query.filter = self.base_filter.clone();
}
if query.as_of.is_none() {
query.as_of = self.base_as_of;
}
if query.as_of.is_none() {
query.as_of = Some(lunaris_core::HlcClock::new(0).tick());
}
let scope = self.scope.clone();
let ctx = match self.moon_storage.clone() {
Some(moon) => QueryContext::with_moon(
query,
scope,
self.embedder,
self.storage.clone(),
self.keyword,
moon,
),
None => {
QueryContext::new(query, scope, self.embedder, self.storage.clone(), self.keyword)
}
};
self.root.retrieve(&ctx).await
}
}