use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::json;
use ulid::Ulid;
use std::collections::HashMap;
use lunaris_core::keyspace::{chunk_key, episode_key, fact_key as scoped_fact_key, fact_spo_key};
use lunaris_core::{
Chunk, Embedder, Hlc, HlcClock, Lsn, LunarisError, Scope, StorageError, StoragePort, WriteOp,
sanitize_graph_ident,
};
use lunaris_extract::types::{EntityId, Fact, FactId};
use lunaris_extract::validator::{NeedsReviewItem, NeedsReviewReason};
use lunaris_ingest::chunk_markdown;
use crate::episode_builder::EpisodeBuilder;
use crate::reconcile::{FactDecision, FactTriple, SpoEntry, classify_fact};
const CHUNK_VECTOR_INDEX: &str = "chunks";
const ENTITIES_INDEX: &str = "entities";
const FACTS_INDEX: &str = "facts";
const GRAPH_NAME: &str = "lunaris_graph";
const DEFAULT_TARGET_TOKENS: usize = 256;
const DEFAULT_OVERLAP_TOKENS: usize = 32;
fn default_confidence() -> f32 {
1.0
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct EntityInput {
pub name: String,
pub entity_type: String,
#[serde(default)]
pub aliases: Vec<String>,
#[serde(default = "default_confidence")]
pub confidence: f32,
pub valid_from: DateTime<Utc>,
#[serde(default)]
pub valid_to: Option<DateTime<Utc>>,
#[serde(default)]
pub embedding: Option<Vec<f32>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct RelationInput {
pub subject_name: String,
pub subject_type: String,
pub predicate: String,
pub object_name: String,
pub object_type: String,
#[serde(default = "default_confidence")]
pub confidence: f32,
pub valid_from: DateTime<Utc>,
#[serde(default)]
pub valid_to: Option<DateTime<Utc>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct FactInput {
pub fact_text: String,
pub subject_name: String,
pub subject_type: String,
pub predicate: String,
pub object_name: String,
pub object_type: String,
#[serde(default = "default_confidence")]
pub confidence: f32,
pub valid_from: DateTime<Utc>,
#[serde(default)]
pub valid_to: Option<DateTime<Utc>>,
}
pub struct StructuredIngest {
pub episode: EpisodeBuilder,
pub entities: Vec<EntityInput>,
pub relations: Vec<RelationInput>,
pub facts: Vec<FactInput>,
}
impl StructuredIngest {
#[must_use]
pub fn new(episode: EpisodeBuilder) -> Self {
Self { episode, entities: Vec::new(), relations: Vec::new(), facts: Vec::new() }
}
#[must_use]
pub fn with_entities(mut self, entities: Vec<EntityInput>) -> Self {
self.entities = entities;
self
}
#[must_use]
pub fn with_relations(mut self, relations: Vec<RelationInput>) -> Self {
self.relations = relations;
self
}
#[must_use]
pub fn with_facts(mut self, facts: Vec<FactInput>) -> Self {
self.facts = facts;
self
}
}
#[doc(hidden)]
pub async fn ingest_structured_inner(
storage: &dyn StoragePort,
embedder: &dyn Embedder,
clock: &HlcClock,
payload: StructuredIngest,
scope: Scope,
) -> Result<Lsn, LunarisError> {
let episode = payload.episode.into_episode(scope, clock);
let embedder_dim = embedder.dim();
let drafts = chunk_markdown(&episode.content, DEFAULT_TARGET_TOKENS, DEFAULT_OVERLAP_TOKENS);
let chunk_embeddings: Vec<Vec<f32>> = if drafts.is_empty() {
Vec::new()
} else {
let texts: Vec<&str> = drafts.iter().map(|d| d.text.as_str()).collect();
let rows = embedder.embed_batch(&texts).await?;
if rows.len() != texts.len() {
return Err(LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: chunk embed returned {} rows for {} chunks",
rows.len(),
texts.len()
))));
}
rows
};
let mut entity_embeds: Vec<Vec<f32>> = vec![Vec::new(); payload.entities.len()];
let mut to_embed_idx: Vec<usize> = Vec::new();
let mut to_embed_text: Vec<String> = Vec::new();
for (i, e) in payload.entities.iter().enumerate() {
if let Some(emb) = &e.embedding {
if emb.len() != embedder_dim {
return Err(LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: EntityInput {:?} supplied embedding has dim {} but \
handle expects {}",
e.name,
emb.len(),
embedder_dim
))));
}
entity_embeds[i] = emb.clone();
} else {
to_embed_idx.push(i);
to_embed_text.push(e.name.clone());
}
}
if !to_embed_text.is_empty() {
let texts: Vec<&str> = to_embed_text.iter().map(String::as_str).collect();
let rows = embedder.embed_batch(&texts).await?;
if rows.len() != to_embed_idx.len() {
return Err(LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: entity embed returned {} rows for {} entities",
rows.len(),
to_embed_idx.len()
))));
}
for (idx, emb) in to_embed_idx.into_iter().zip(rows.into_iter()) {
entity_embeds[idx] = emb;
}
}
let fact_embeds: Vec<Vec<f32>> = if payload.facts.is_empty() {
Vec::new()
} else {
let texts: Vec<&str> = payload.facts.iter().map(|f| f.fact_text.as_str()).collect();
let rows = embedder.embed_batch(&texts).await?;
if rows.len() != texts.len() {
return Err(LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: fact embed returned {} rows for {} facts",
rows.len(),
texts.len()
))));
}
rows
};
let mut ops: Vec<WriteOp> = Vec::with_capacity(
1 + 2 * drafts.len()
+ 2 * payload.entities.len()
+ payload.relations.len()
+ 2 * payload.facts.len(),
);
let episode_value = serde_json::to_vec(&episode).map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: episode serialize: {e}"
)))
})?;
ops.push(WriteOp::KvPut { key: episode_key(&episode.scope, episode.id), value: episode_value });
let mut chunks: Vec<Chunk> = Vec::with_capacity(drafts.len());
for (draft, emb) in drafts.into_iter().zip(chunk_embeddings.into_iter()) {
let mut c = draft.into_chunk_valid_from(
episode.scope.clone(),
episode.id,
clock,
episode.bt.valid.0,
);
c.embedding = Some(emb.clone());
let chunk_value = serde_json::to_vec(&c).map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: chunk serialize: {e}"
)))
})?;
ops.push(WriteOp::KvPut { key: chunk_key(&episode.scope, c.id), value: chunk_value });
ops.push(WriteOp::VectorUpsert {
index: CHUNK_VECTOR_INDEX.into(),
id: c.id.to_bytes().to_vec(),
embedding: emb,
metadata: json!({
"episode_id": c.episode_id.to_string(),
"heading_path": c.heading_path,
"offset": c.offset,
"text": c.text,
"source": &episode.source,
}),
});
chunks.push(c);
}
let episode_id_str = episode.id.to_string();
for (e, emb) in payload.entities.iter().zip(entity_embeds.iter()) {
let eid = EntityId::from_name_and_type(&e.name, &e.entity_type);
let id_bytes = eid.0.to_vec();
ops.push(WriteOp::GraphNode {
graph: GRAPH_NAME.into(),
id: id_bytes.clone(),
label: sanitize_graph_ident(&e.entity_type, "Entity"),
props: json!({
"id_hex": format!("{eid}"),
"name": e.name,
"type": e.entity_type,
"aliases": e.aliases,
"confidence": e.confidence,
"valid_from_iso": e.valid_from.to_rfc3339(),
"valid_to_iso": e.valid_to.map(|t| t.to_rfc3339()),
"source_episode_id": episode_id_str,
}),
index_kind: "entities".into(),
});
ops.push(WriteOp::VectorUpsert {
index: ENTITIES_INDEX.into(),
id: id_bytes,
embedding: emb.clone(),
metadata: json!({"entity_type": e.entity_type, "name": e.name}),
});
}
for r in &payload.relations {
let sid = EntityId::from_name_and_type(&r.subject_name, &r.subject_type);
let oid = EntityId::from_name_and_type(&r.object_name, &r.object_type);
ops.push(WriteOp::GraphEdge {
graph: GRAPH_NAME.into(),
src: sid.0.to_vec(),
dst: oid.0.to_vec(),
rel: sanitize_graph_ident(&r.predicate, "RELATED_TO"),
props: json!({
"confidence": r.confidence,
"valid_from_iso": r.valid_from.to_rfc3339(),
"valid_to_iso": r.valid_to.map(|t| t.to_rfc3339()),
"source_episode_id": episode_id_str,
}),
});
}
let now_hlc = clock.tick();
let mut spo_index: HashMap<Vec<u8>, Vec<SpoEntry>> = HashMap::new();
let mut needs_review: Vec<NeedsReviewItem> = Vec::new();
for (f, emb) in payload.facts.iter().zip(fact_embeds.iter()) {
let sid = EntityId::from_name_and_type(&f.subject_name, &f.subject_type);
let oid = EntityId::from_name_and_type(&f.object_name, &f.object_type);
let fact_id = Ulid::from_bytes(FactId::from_triple(sid, &f.predicate, oid).0);
let spo_key = fact_spo_key(&episode.scope, &sid.0, &f.predicate);
if !spo_index.contains_key(&spo_key) {
let prior = read_spo_index(storage, &episode.scope, &spo_key, now_hlc).await?;
spo_index.insert(spo_key.clone(), prior);
}
let new_triple = FactTriple {
subject_id: sid,
predicate: f.predicate.clone(),
object_id: oid,
valid_from: f.valid_from,
valid_to: f.valid_to,
};
let prior = &spo_index[&spo_key];
match classify_fact(&new_triple, prior) {
FactDecision::Noop => {
if let Some(entry) = spo_index
.get_mut(&spo_key)
.and_then(|v| v.iter_mut().find(|e| e.object_id == oid))
{
entry.valid_from = f.valid_from;
entry.valid_to = f.valid_to;
}
}
FactDecision::Append => {
spo_index.get_mut(&spo_key).expect("seeded above").push(SpoEntry {
object_id: oid,
fact_id,
valid_from: f.valid_from,
valid_to: f.valid_to,
});
}
FactDecision::Supersede { loser_fact_id } => {
let existing_object =
prior.iter().find(|p| p.fact_id == loser_fact_id).map_or(oid, |p| p.object_id);
needs_review.push(NeedsReviewItem::Fact {
reason: NeedsReviewReason::CrossEpisodeContradiction {
subject: sid,
predicate: f.predicate.clone(),
existing_fact_id: loser_fact_id,
existing_object,
new_fact_id: fact_id,
new_object: oid,
},
raw: Fact {
id: fact_id,
subject_id: sid,
predicate: f.predicate.clone(),
object_id: oid,
fact_text: f.fact_text.clone(),
confidence: f.confidence,
valid_from_iso: f.valid_from.to_rfc3339(),
valid_to_iso: f.valid_to.map(|t| t.to_rfc3339()),
},
});
spo_index.get_mut(&spo_key).expect("seeded above").push(SpoEntry {
object_id: oid,
fact_id,
valid_from: f.valid_from,
valid_to: f.valid_to,
});
}
}
let fact_value = serde_json::to_vec(&serde_json::json!({
"id": fact_id.to_string(),
"subject_id": sid.0,
"predicate": f.predicate,
"object_id": oid.0,
"fact_text": f.fact_text,
"confidence": f.confidence,
"valid_from_iso": f.valid_from.to_rfc3339(),
"valid_to_iso": f.valid_to.map(|t| t.to_rfc3339()),
"source_episode_id": episode_id_str,
}))
.map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: fact serialize: {e}"
)))
})?;
ops.push(WriteOp::KvPut {
key: scoped_fact_key(&episode.scope, fact_id),
value: fact_value,
});
ops.push(WriteOp::VectorUpsert {
index: FACTS_INDEX.into(),
id: fact_id.to_bytes().to_vec(),
embedding: emb.clone(),
metadata: json!({"predicate": f.predicate, "fact_text": f.fact_text}),
});
let fact_id_bytes = fact_id.to_bytes().to_vec();
ops.push(WriteOp::GraphNode {
graph: GRAPH_NAME.into(),
id: fact_id_bytes.clone(),
label: "Fact".into(),
props: json!({
"id_hex": fact_id_bytes.iter().map(|b| format!("{b:02x}")).collect::<String>(),
"predicate": f.predicate,
"confidence": f.confidence,
"valid_from_iso": f.valid_from.to_rfc3339(),
"valid_to_iso": f.valid_to.map(|t| t.to_rfc3339()),
}),
index_kind: "facts".into(),
});
ops.push(WriteOp::GraphEdge {
graph: GRAPH_NAME.into(),
src: sid.0.to_vec(),
dst: fact_id_bytes.clone(),
rel: "HAS_FACT".into(),
props: json!({}),
});
ops.push(WriteOp::GraphEdge {
graph: GRAPH_NAME.into(),
src: fact_id_bytes,
dst: oid.0.to_vec(),
rel: "FACT_ABOUT".into(),
props: json!({}),
});
}
for (key, entries) in &spo_index {
let value = serde_json::to_vec(&spo_entries_to_json(entries)).map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!(
"structured_ingest: spo-index serialize: {e}"
)))
})?;
ops.push(WriteOp::KvPut { key: key.clone(), value });
}
let lsn = storage.atomic_write(&episode.scope, &ops).await?;
if !needs_review.is_empty() {
crate::ingest::publish_needs_review(storage, &episode.scope, &needs_review).await;
}
Ok(lsn)
}
pub(crate) async fn read_spo_index(
storage: &dyn StoragePort,
scope: &Scope,
key: &[u8],
as_of: Hlc,
) -> Result<Vec<SpoEntry>, LunarisError> {
let Some(row) = storage.read_as_of(scope, key, as_of).await.map_err(LunarisError::Storage)?
else {
return Ok(Vec::new());
};
let arr: Vec<serde_json::Value> = serde_json::from_slice(&row.value).unwrap_or_default();
let mut out = Vec::with_capacity(arr.len());
for v in arr {
let (Some(obj_hex), Some(fid_str), Some(vf_str)) = (
v.get("object_id").and_then(|x| x.as_str()),
v.get("fact_id").and_then(|x| x.as_str()),
v.get("valid_from").and_then(|x| x.as_str()),
) else {
continue;
};
let (Some(object_id), Some(fact_id), Some(valid_from)) = (
EntityId::from_hex(obj_hex),
Ulid::from_string(fid_str).ok(),
DateTime::parse_from_rfc3339(vf_str).ok().map(|d| d.with_timezone(&Utc)),
) else {
continue;
};
let valid_to = v
.get("valid_to")
.and_then(|x| x.as_str())
.and_then(|s| DateTime::parse_from_rfc3339(s).ok())
.map(|d| d.with_timezone(&Utc));
out.push(SpoEntry { object_id, fact_id, valid_from, valid_to });
}
Ok(out)
}
pub(crate) fn spo_entries_to_json(entries: &[SpoEntry]) -> Vec<serde_json::Value> {
entries
.iter()
.map(|e| {
json!({
"object_id": format!("{}", e.object_id),
"fact_id": e.fact_id.to_string(),
"valid_from": e.valid_from.to_rfc3339(),
"valid_to": e.valid_to.map(|t| t.to_rfc3339()),
})
})
.collect()
}