use std::collections::HashMap;
use futures::stream::{self, StreamExt};
use lunaris_core::keyspace::{chunk_key, episode_key, fact_key};
use lunaris_core::{BiTemporal, Chunk, Episode, Hlc, HlcClock, LunarisError, Scope, StoragePort};
use ulid::Ulid;
use crate::types::{Hit, RawHit};
const HYDRATE_CONCURRENCY: usize = 32;
fn chunk_lookup_key(scope: &Scope, id_bytes: &[u8]) -> Option<Vec<u8>> {
let ulid = Ulid::from_bytes(id_bytes.try_into().ok()?);
Some(chunk_key(scope, ulid))
}
fn episode_lookup_key(scope: &Scope, id: Ulid) -> Vec<u8> {
episode_key(scope, id)
}
fn sys_closed(bt: &BiTemporal) -> bool {
bt.sys.1.is_some()
}
type EpisodeGateEntry = Option<(Ulid, String, bool)>;
pub async fn hydrate(
storage: &dyn StoragePort,
scope: &Scope,
hits: Vec<RawHit>,
as_of: Option<Hlc>,
initial_degraded: bool,
) -> Result<Vec<Hit>, LunarisError> {
let live_clock = HlcClock::new(0);
let snapshot = as_of.unwrap_or_else(|| live_clock.tick());
let chunks: Vec<(RawHit, Chunk)> = {
let results: Vec<Result<Option<(RawHit, Chunk)>, LunarisError>> = stream::iter(hits)
.map(|raw| {
let scope = scope.clone();
async move {
let key = match chunk_lookup_key(&scope, &raw.id) {
Some(k) => k,
None => return Ok(None), };
match storage.read_as_of(&scope, &key, snapshot).await? {
Some(row) => {
let chunk: Chunk = match serde_json::from_slice(&row.value) {
Ok(c) => c,
Err(_) => return Ok(None),
};
if sys_closed(&chunk.bt) {
return Ok(None); }
Ok(Some((raw, chunk)))
}
None => Ok(None), }
}
})
.buffered(HYDRATE_CONCURRENCY)
.collect()
.await;
results.into_iter().collect::<Result<Vec<_>, _>>()?.into_iter().flatten().collect()
};
let unique_ep: Vec<Ulid> = {
let mut s = std::collections::HashSet::new();
for (_, c) in &chunks {
s.insert(c.episode_id);
}
s.into_iter().collect()
};
let mut episode_sources: HashMap<Ulid, (String, bool)> = HashMap::new();
let ep_results: Vec<Result<EpisodeGateEntry, LunarisError>> = stream::iter(unique_ep)
.map(|ep_id| {
let scope = scope.clone();
async move {
let key = episode_lookup_key(&scope, ep_id);
if let Some(row) = storage.read_as_of(&scope, &key, snapshot).await?
&& let Ok(ep) = serde_json::from_slice::<Episode>(&row.value)
{
return Ok(Some((ep_id, ep.source, sys_closed(&ep.bt))));
}
Ok(None)
}
})
.buffer_unordered(HYDRATE_CONCURRENCY)
.collect()
.await;
for entry in ep_results.into_iter().collect::<Result<Vec<_>, _>>()?.into_iter().flatten() {
episode_sources.insert(entry.0, (entry.1, entry.2));
}
Ok(chunks
.into_iter()
.filter(|(_, chunk)| !matches!(episode_sources.get(&chunk.episode_id), Some((_, true))))
.map(|(raw, chunk)| Hit {
id: raw.id,
episode_id: chunk.episode_id.to_bytes().to_vec(),
score: raw.score,
text: chunk.text,
source: episode_sources
.get(&chunk.episode_id)
.map(|(s, _)| s.clone())
.unwrap_or_default(),
heading_path: chunk.heading_path,
valid_from: chunk.bt.valid.0,
valid_to: chunk.bt.valid.1,
degraded: raw.degraded || initial_degraded,
rerank_applied: raw.rerank_applied,
source_op: raw.source_op,
})
.collect())
}
enum MixedResolved {
Chunk(RawHit, Chunk),
Fact(RawHit, lunaris_extract::Fact, BiTemporal),
}
pub async fn hydrate_mixed(
storage: &dyn StoragePort,
scope: &Scope,
hits: Vec<RawHit>,
as_of: Option<Hlc>,
initial_degraded: bool,
) -> Result<Vec<Hit>, LunarisError> {
let live_clock = HlcClock::new(0);
let snapshot = as_of.unwrap_or_else(|| live_clock.tick());
let resolved: Vec<MixedResolved> = {
let results: Vec<Result<Option<MixedResolved>, LunarisError>> = stream::iter(hits)
.map(|raw| {
let scope = scope.clone();
async move {
let ulid = match <[u8; 16]>::try_from(raw.id.as_slice()) {
Ok(bytes) => Ulid::from_bytes(bytes),
Err(_) => return Ok(None), };
if let Some(row) =
storage.read_as_of(&scope, &chunk_key(&scope, ulid), snapshot).await?
&& let Ok(chunk) = serde_json::from_slice::<Chunk>(&row.value)
{
if sys_closed(&chunk.bt) {
return Ok(None); }
return Ok(Some(MixedResolved::Chunk(raw, chunk)));
}
if let Some(row) =
storage.read_as_of(&scope, &fact_key(&scope, ulid), snapshot).await?
&& let Ok(fact) =
serde_json::from_slice::<lunaris_extract::Fact>(&row.value)
{
if sys_closed(&row.bt) {
return Ok(None); }
return Ok(Some(MixedResolved::Fact(raw, fact, row.bt)));
}
Ok(None) }
})
.buffered(HYDRATE_CONCURRENCY)
.collect()
.await;
results.into_iter().collect::<Result<Vec<_>, _>>()?.into_iter().flatten().collect()
};
let unique_ep: Vec<Ulid> = {
let mut s = std::collections::HashSet::new();
for r in &resolved {
if let MixedResolved::Chunk(_, c) = r {
s.insert(c.episode_id);
}
}
s.into_iter().collect()
};
let mut episode_sources: HashMap<Ulid, (String, bool)> = HashMap::new();
let ep_results: Vec<Result<EpisodeGateEntry, LunarisError>> = stream::iter(unique_ep)
.map(|ep_id| {
let scope = scope.clone();
async move {
let key = episode_lookup_key(&scope, ep_id);
if let Some(row) = storage.read_as_of(&scope, &key, snapshot).await?
&& let Ok(ep) = serde_json::from_slice::<Episode>(&row.value)
{
return Ok(Some((ep_id, ep.source, sys_closed(&ep.bt))));
}
Ok(None)
}
})
.buffer_unordered(HYDRATE_CONCURRENCY)
.collect()
.await;
for entry in ep_results.into_iter().collect::<Result<Vec<_>, _>>()?.into_iter().flatten() {
episode_sources.insert(entry.0, (entry.1, entry.2));
}
Ok(resolved
.into_iter()
.filter(|r| match r {
MixedResolved::Chunk(_, chunk) => {
!matches!(episode_sources.get(&chunk.episode_id), Some((_, true)))
}
MixedResolved::Fact(..) => true, })
.map(|r| match r {
MixedResolved::Chunk(raw, chunk) => Hit {
id: raw.id,
episode_id: chunk.episode_id.to_bytes().to_vec(),
score: raw.score,
text: chunk.text,
source: episode_sources
.get(&chunk.episode_id)
.map(|(s, _)| s.clone())
.unwrap_or_default(),
heading_path: chunk.heading_path,
valid_from: chunk.bt.valid.0,
valid_to: chunk.bt.valid.1,
degraded: raw.degraded || initial_degraded,
rerank_applied: raw.rerank_applied,
source_op: raw.source_op,
},
MixedResolved::Fact(raw, fact, bt) => Hit {
id: raw.id,
episode_id: Vec::new(),
score: raw.score,
text: fact.fact_text,
source: format!("fact:{}", fact.predicate),
heading_path: Vec::new(),
valid_from: bt.valid.0,
valid_to: bt.valid.1,
degraded: raw.degraded || initial_degraded,
rerank_applied: raw.rerank_applied,
source_op: raw.source_op,
},
})
.collect())
}
pub async fn partial_hydrate_text(
storage: &dyn StoragePort,
scope: &Scope,
hits: &[RawHit],
as_of: Option<Hlc>,
) -> Result<HashMap<Vec<u8>, String>, LunarisError> {
let live_clock = HlcClock::new(0);
let snapshot = as_of.unwrap_or_else(|| live_clock.tick());
let owned: Vec<(Vec<u8>, Ulid)> = hits
.iter()
.filter_map(|raw| {
let bytes = <[u8; 16]>::try_from(raw.id.as_slice()).ok()?;
Some((raw.id.clone(), Ulid::from_bytes(bytes)))
})
.collect();
let pairs: Vec<(Vec<u8>, String)> = {
#[allow(clippy::type_complexity)]
let results: Vec<Result<Option<(Vec<u8>, String)>, LunarisError>> = stream::iter(owned)
.map(|(id, ulid)| {
let scope = scope.clone();
async move {
if let Some(row) =
storage.read_as_of(&scope, &chunk_key(&scope, ulid), snapshot).await?
&& let Ok(chunk) = serde_json::from_slice::<Chunk>(&row.value)
{
return Ok(Some((id, chunk.text)));
}
if let Some(row) =
storage.read_as_of(&scope, &fact_key(&scope, ulid), snapshot).await?
&& let Ok(fact) =
serde_json::from_slice::<lunaris_extract::Fact>(&row.value)
{
return Ok(Some((id, fact.fact_text)));
}
Ok(None) }
})
.buffered(HYDRATE_CONCURRENCY)
.collect()
.await;
results.into_iter().collect::<Result<Vec<_>, _>>()?.into_iter().flatten().collect()
};
Ok(pairs.into_iter().collect())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn chunk_lookup_key_round_trips_ulid() {
let id = Ulid::new();
let bytes = id.to_bytes().to_vec();
let scope = Scope::dev();
let key = chunk_lookup_key(&scope, &bytes).unwrap();
let s = String::from_utf8(key).unwrap();
assert!(s.contains(":chunk:"));
assert!(s.contains(&id.to_string()));
}
#[test]
fn chunk_lookup_key_rejects_wrong_size() {
let scope = Scope::dev();
assert!(chunk_lookup_key(&scope, b"too-short").is_none());
}
}