use crate::persistence::Persistence;
use crate::retrieval::hybrid::RetrievalAbstention;
use crate::semantic_bridge::{CanonicalConcept, ConceptGraph};
use chrono::{DateTime, Utc};
use csm_core_lib::error::{MemoryError, Result};
use libsql::params;
impl Persistence {
pub async fn save_canonical_concept(&self, ns: &str, concept: &CanonicalConcept) -> Result<()> {
let _permit = self.acquire_remote_slot().await?;
let conn = self.connect().await?;
let labels_json = serde_json::to_string(&concept.labels)?;
let related_json = serde_json::to_string(&concept.related)?;
conn.execute(
"INSERT INTO csm_canonical (namespace, id, version, labels_json, related_json)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT(namespace, id) DO UPDATE SET
version = excluded.version,
labels_json = excluded.labels_json,
related_json = excluded.related_json",
params![
ns.to_string(),
concept.id.clone(),
concept.version as i64,
labels_json,
related_json
],
)
.await
.map_err(|e| MemoryError::database(format!("Failed to save canonical concept: {e}")))?;
Ok(())
}
pub async fn delete_canonical_concept(&self, ns: &str, id: &str) -> Result<()> {
let _permit = self.acquire_remote_slot().await?;
let conn = self.connect().await?;
conn.execute(
"DELETE FROM csm_canonical WHERE namespace = ?1 AND id = ?2",
params![ns.to_string(), id],
)
.await
.map_err(|e| MemoryError::database(format!("Failed to delete canonical concept: {e}")))?;
Ok(())
}
pub async fn load_canonical_concept(
&self,
ns: &str,
id: &str,
) -> Result<Option<CanonicalConcept>> {
let _permit = self.acquire_remote_slot().await?;
let conn = self.connect().await?;
let mut rows = conn
.query(
"SELECT id, version, labels_json, related_json FROM csm_canonical WHERE namespace = ?1 AND id = ?2",
params![ns.to_string(), id],
)
.await
.map_err(|e| MemoryError::database(format!("Failed to load canonical concept: {e}")))?;
if let Some(row) = rows.next().await.map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept row: {e}"))
})? {
let id: String = row.get(0).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept id: {e}"))
})?;
let version: i64 = row.get(1).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept version: {e}"))
})?;
let labels_json: String = row.get(2).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept labels: {e}"))
})?;
let related_json: String = row.get(3).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept related: {e}"))
})?;
let labels: Vec<String> = serde_json::from_str(&labels_json)?;
let related: Vec<String> = serde_json::from_str(&related_json)?;
Ok(Some(CanonicalConcept {
id,
version: u32::try_from(version).unwrap_or(0),
labels,
related,
}))
} else {
Ok(None)
}
}
pub async fn load_all_canonical_concepts(&self, ns: &str) -> Result<Vec<CanonicalConcept>> {
let _permit = self.acquire_remote_slot().await?;
let conn = self.connect().await?;
let mut rows = conn
.query(
"SELECT id, version, labels_json, related_json FROM csm_canonical WHERE namespace = ?1",
params![ns.to_string()],
)
.await
.map_err(|e| {
MemoryError::database(format!("Failed to load canonical concepts: {e}"))
})?;
let mut concepts = Vec::new();
while let Some(row) = rows.next().await.map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept row: {e}"))
})? {
let id: String = row.get(0).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept id: {e}"))
})?;
let version: i64 = row.get(1).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept version: {e}"))
})?;
let labels_json: String = row.get(2).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept labels: {e}"))
})?;
let related_json: String = row.get(3).map_err(|e| {
MemoryError::database(format!("Failed to read canonical concept related: {e}"))
})?;
let labels: Vec<String> = serde_json::from_str(&labels_json)?;
let related: Vec<String> = serde_json::from_str(&related_json)?;
concepts.push(CanonicalConcept {
id,
version: u32::try_from(version).unwrap_or(0),
labels,
related,
});
}
Ok(concepts)
}
pub async fn save_concept_graph(&self, ns: &str, graph: &ConceptGraph) -> Result<()> {
let _permit = self.acquire_remote_slot().await?;
let conn = self.connect().await?;
conn.execute("BEGIN", ())
.await
.map_err(|e| MemoryError::database(format!("Failed to begin transaction: {e}")))?;
if let Err(e) = conn
.execute(
"DELETE FROM csm_canonical WHERE namespace = ?1",
params![ns.to_string()],
)
.await
{
let _ = conn.execute("ROLLBACK", ()).await;
return Err(MemoryError::database(format!(
"Failed to clear canonical concepts: {e}"
)));
}
let mut first_error: Option<MemoryError> = None;
for concept in graph.all_concepts() {
let labels_json = match serde_json::to_string(&concept.labels) {
Ok(j) => j,
Err(e) => {
first_error = Some(MemoryError::Serialization(e));
break;
}
};
let related_json = match serde_json::to_string(&concept.related) {
Ok(j) => j,
Err(e) => {
first_error = Some(MemoryError::Serialization(e));
break;
}
};
if let Err(e) = conn
.execute(
"INSERT INTO csm_canonical (namespace, id, version, labels_json, related_json)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![
ns.to_string(),
concept.id.clone(),
concept.version as i64,
labels_json,
related_json
],
)
.await
{
first_error = Some(MemoryError::database(format!(
"Failed to save canonical concept: {e}"
)));
break;
}
}
if let Some(error) = first_error {
let _ = conn.execute("ROLLBACK", ()).await;
return Err(error);
}
conn.execute("COMMIT", ())
.await
.map_err(|e| MemoryError::database(format!("Failed to commit transaction: {e}")))?;
Ok(())
}
pub async fn load_concept_graph(&self, ns: &str) -> Result<ConceptGraph> {
let concepts = self.load_all_canonical_concepts(ns).await?;
let mut graph = ConceptGraph::new();
for concept in concepts {
graph.add_concept(concept);
}
Ok(graph)
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AbsenceEntry {
pub id: String,
pub query: String,
pub normalized_query: String,
pub attempt_count: u32,
pub last_threshold: f32,
pub best_score_ever: Option<f32>,
pub first_seen: DateTime<Utc>,
pub last_seen: DateTime<Utc>,
}
impl AbsenceEntry {
pub fn normalize(query: &str) -> String {
query.trim().to_lowercase()
}
pub fn id_for(query: &str) -> String {
let normalized = Self::normalize(query);
let hash = Self::fnv1a_hash(normalized.as_bytes());
format!("absence:{hash:016x}")
}
pub(crate) fn fnv1a_hash(bytes: &[u8]) -> u64 {
const OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = OFFSET_BASIS;
for &byte in bytes {
hash ^= byte as u64;
hash = hash.wrapping_mul(PRIME);
}
hash
}
pub fn from_abstention(abstention: &RetrievalAbstention) -> Self {
let normalized = Self::normalize(&abstention.query);
AbsenceEntry {
id: Self::id_for(&abstention.query),
query: abstention.query.clone(),
normalized_query: normalized,
attempt_count: 1,
last_threshold: abstention.min_score_threshold,
best_score_ever: abstention.best_score_seen,
first_seen: abstention.timestamp,
last_seen: abstention.timestamp,
}
}
pub fn merge_with(&mut self, abstention: &RetrievalAbstention) {
self.attempt_count += 1;
self.last_seen = abstention.timestamp;
self.last_threshold = abstention.min_score_threshold;
match (abstention.best_score_seen, self.best_score_ever) {
(Some(new), Some(existing)) => {
if new > existing {
self.best_score_ever = Some(new);
}
}
(Some(new), None) => {
self.best_score_ever = Some(new);
}
_ => {}
}
}
}
#[async_trait::async_trait]
pub trait AbsenceStore: Send + Sync {
async fn get_absence(&self, id: &str) -> Result<Option<AbsenceEntry>>;
async fn upsert_absence(&self, entry: &AbsenceEntry) -> Result<()>;
async fn list_absences(&self, min_attempts: u32) -> Result<Vec<AbsenceEntry>>;
}
pub async fn persist_absence(
abstention: &RetrievalAbstention,
store: &dyn AbsenceStore,
) -> Result<AbsenceEntry> {
let id = AbsenceEntry::id_for(&abstention.query);
match store.get_absence(&id).await? {
Some(mut existing) => {
existing.merge_with(abstention);
store.upsert_absence(&existing).await?;
Ok(existing)
}
None => {
let entry = AbsenceEntry::from_abstention(abstention);
store.upsert_absence(&entry).await?;
Ok(entry)
}
}
}