mod belief;
mod cache;
mod intent;
mod action;
mod evaluator;
mod policy;
mod suggest;
mod agenda;
mod temporal;
mod hawkes;
mod receptivity;
mod tick;
mod surfacing;
mod observer;
mod flywheel;
mod world_model;
mod experimenter;
mod skills;
mod extractor;
mod calibration;
mod introspection;
mod causal;
mod planner;
mod cognition;
mod coherence;
mod metacognition;
mod personality_bias;
mod query_dsl;
mod conflict;
mod analogy_engine;
mod schema_induction_engine;
mod narrative_engine;
mod counterfactual_engine;
mod belief_network_engine;
mod replay_engine;
mod perspective_engine;
mod feedback;
pub mod graph_state;
mod graph_ops;
mod indices;
mod learning;
mod lifecycle;
mod recall;
mod record;
mod procedural;
mod session;
mod stats;
mod storage;
mod temporal_helpers;
pub mod tenant;
#[cfg(test)]
mod tests;
use std::collections::HashMap;
use std::sync::{Mutex, RwLock, MutexGuard};
use base64::Engine;
use rand::Rng;
use rusqlite::{params, Connection};
use crate::encryption::{self, EncryptionProvider};
use crate::error::{YantrikDbError, Result};
use crate::graph_index::GraphIndex;
use crate::hlc::{HLCTimestamp, HLC};
use crate::hnsw::HnswIndex;
use crate::schema::{
MIGRATE_V1_TO_V2, MIGRATE_V2_TO_V3, MIGRATE_V3_TO_V4, MIGRATE_V4_TO_V5,
MIGRATE_V5_TO_V6, MIGRATE_V6_TO_V7, MIGRATE_V7_TO_V8, MIGRATE_V8_TO_V9,
MIGRATE_V9_TO_V10, MIGRATE_V10_TO_V11, MIGRATE_V11_TO_V12, MIGRATE_V12_TO_V13,
MIGRATE_V13_TO_V14,
SCHEMA_SQL, SCHEMA_VERSION,
};
use crate::types::*;
pub struct YantrikDB {
pub(crate) conn: Mutex<Connection>,
pub(crate) embedding_dim: usize,
pub(crate) hlc: Mutex<HLC>,
pub(crate) actor_id: String,
pub(crate) scoring_cache: RwLock<HashMap<String, ScoringRow>>,
pub(crate) vec_index: RwLock<HnswIndex>,
pub(crate) graph_index: RwLock<GraphIndex>,
pub(crate) enc: Option<EncryptionProvider>,
embedder: Option<Box<dyn crate::types::Embedder + Send + Sync>>,
pub(crate) active_sessions: RwLock<HashMap<String, String>>,
}
const _: () = {
fn _assert_send<T: Send>() {}
fn _assert_sync<T: Sync>() {}
fn _check() {
_assert_send::<YantrikDB>();
_assert_sync::<YantrikDB>();
}
};
pub(crate) fn now() -> f64 {
crate::time::now_secs()
}
pub(crate) fn embedding_hash(embedding: &[f32]) -> Vec<u8> {
let blob = crate::serde_helpers::serialize_f32(embedding);
blake3::hash(&blob).as_bytes().to_vec()
}
pub(crate) struct TextMetadataRow {
pub rid: String,
pub text: String,
pub metadata: String,
}
impl YantrikDB {
pub fn new(db_path: &str, embedding_dim: usize) -> Result<Self> {
Self::open(db_path, embedding_dim, None, None)
}
pub fn new_with_actor(db_path: &str, embedding_dim: usize, actor_id: &str) -> Result<Self> {
Self::open(db_path, embedding_dim, Some(actor_id.to_string()), None)
}
pub fn new_encrypted(db_path: &str, embedding_dim: usize, master_key: &[u8; 32]) -> Result<Self> {
Self::open(db_path, embedding_dim, None, Some(master_key))
}
fn open(
db_path: &str,
embedding_dim: usize,
actor_id: Option<String>,
master_key: Option<&[u8; 32]>,
) -> Result<Self> {
let conn = Connection::open(db_path)?;
conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
let existing_version = Self::get_schema_version(&conn);
let migrations: &[(i32, &str)] = &[
(1, MIGRATE_V1_TO_V2),
(2, MIGRATE_V2_TO_V3),
(3, MIGRATE_V3_TO_V4),
(4, MIGRATE_V4_TO_V5),
(5, MIGRATE_V5_TO_V6),
(6, MIGRATE_V6_TO_V7),
(7, MIGRATE_V7_TO_V8),
(8, MIGRATE_V8_TO_V9),
(9, MIGRATE_V9_TO_V10),
(10, MIGRATE_V10_TO_V11),
(11, MIGRATE_V11_TO_V12),
(12, MIGRATE_V12_TO_V13),
(13, MIGRATE_V13_TO_V14),
];
if let Some(v) = existing_version {
for &(from_v, sql) in migrations {
if v <= from_v {
conn.execute_batch(sql)?;
}
}
}
conn.execute_batch(SCHEMA_SQL)?;
crate::distributed::seed_categories::populate_seed_categories(&conn)?;
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', ?1)",
params![SCHEMA_VERSION.to_string()],
)?;
let actor_id = if let Some(id) = actor_id {
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('actor_id', ?1)",
params![id],
)?;
id
} else {
match Self::get_meta(&conn, "actor_id")? {
Some(id) => id,
None => {
let id = crate::id::new_id();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('actor_id', ?1)",
params![id],
)?;
id
}
}
};
let node_id: u32 = match Self::get_meta(&conn, "node_id")? {
Some(s) => s.parse().unwrap_or_else(|_| {
let id: u32 = rand::thread_rng().gen();
id
}),
None => {
let id: u32 = rand::thread_rng().gen();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('node_id', ?1)",
params![id.to_string()],
)?;
id
}
};
let enc = if let Some(mk) = master_key {
let provider = match Self::get_meta(&conn, "encrypted_dek")? {
Some(wrapped_b64) => {
let wrapped = base64::engine::general_purpose::STANDARD
.decode(&wrapped_b64)
.map_err(|e| YantrikDbError::Encryption(format!("DEK base64: {e}")))?;
let dek = encryption::unwrap_dek(mk, &wrapped)?;
EncryptionProvider::from_dek(&dek)
}
None => {
let dek = encryption::generate_key();
let wrapped = encryption::wrap_dek(mk, &dek)?;
let wrapped_b64 = base64::engine::general_purpose::STANDARD.encode(&wrapped);
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('encrypted_dek', ?1)",
params![wrapped_b64],
)?;
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('encryption_enabled', '1')",
[],
)?;
EncryptionProvider::from_dek(&dek)
}
};
Some(provider)
} else {
if Self::get_meta(&conn, "encryption_enabled")?.as_deref() == Some("1") {
return Err(YantrikDbError::Encryption(
"database is encrypted but no master_key provided".into(),
));
}
None
};
let scoring_cache = Self::load_scoring_cache(&conn)?;
let vec_index = Self::build_vec_index_with_enc(&conn, embedding_dim, enc.as_ref())?;
let graph_index = GraphIndex::build_from_db(&conn)?;
let active_sessions = Self::load_active_sessions(&conn)?;
Ok(Self {
conn: Mutex::new(conn),
embedding_dim,
hlc: Mutex::new(HLC::new(node_id)),
actor_id,
scoring_cache: RwLock::new(scoring_cache),
vec_index: RwLock::new(vec_index),
graph_index: RwLock::new(graph_index),
enc,
embedder: None,
active_sessions: RwLock::new(active_sessions),
})
}
fn get_schema_version(conn: &Connection) -> Option<i32> {
conn.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|row| {
let v: String = row.get(0)?;
Ok(v.parse::<i32>().unwrap_or(0))
},
)
.ok()
}
fn load_active_sessions(conn: &Connection) -> Result<HashMap<String, String>> {
let mut map = HashMap::new();
let mut stmt = match conn.prepare(
"SELECT namespace, session_id FROM sessions WHERE status = 'active'",
) {
Ok(s) => s,
Err(_) => return Ok(map),
};
let rows = stmt.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})?;
for row in rows {
let (ns, sid) = row?;
map.insert(ns, sid);
}
Ok(map)
}
fn get_meta(conn: &Connection, key: &str) -> Result<Option<String>> {
match conn.query_row(
"SELECT value FROM meta WHERE key = ?1",
params![key],
|row| row.get(0),
) {
Ok(v) => Ok(Some(v)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e.into()),
}
}
pub fn tick_hlc(&self) -> HLCTimestamp {
self.hlc.lock().unwrap().now()
}
pub fn merge_hlc(&self, remote: HLCTimestamp) -> HLCTimestamp {
self.hlc.lock().unwrap().recv(remote)
}
pub fn actor_id(&self) -> &str {
&self.actor_id
}
pub fn embedding_dim(&self) -> usize {
self.embedding_dim
}
pub fn conn(&self) -> MutexGuard<'_, Connection> {
self.conn.lock().unwrap()
}
pub fn is_encrypted(&self) -> bool {
self.enc.is_some()
}
pub fn encryption(&self) -> Option<&EncryptionProvider> {
self.enc.as_ref()
}
pub(crate) fn encrypt_text(&self, plaintext: &str) -> Result<String> {
match &self.enc {
Some(e) => e.encrypt_string(plaintext),
None => Ok(plaintext.to_string()),
}
}
pub(crate) fn decrypt_text(&self, stored: &str) -> Result<String> {
match &self.enc {
Some(e) => e.decrypt_string(stored),
None => Ok(stored.to_string()),
}
}
pub(crate) fn encrypt_embedding(&self, emb_blob: &[u8]) -> Result<Vec<u8>> {
match &self.enc {
Some(e) => e.encrypt_bytes(emb_blob),
None => Ok(emb_blob.to_vec()),
}
}
pub(crate) fn decrypt_embedding(&self, stored: &[u8]) -> Result<Vec<u8>> {
match &self.enc {
Some(e) => e.decrypt_bytes(stored),
None => Ok(stored.to_vec()),
}
}
pub fn close(self) -> Result<()> {
self.conn
.into_inner()
.unwrap()
.close()
.map_err(|(_, e)| YantrikDbError::Database(e))
}
pub fn set_embedder(&mut self, embedder: Box<dyn crate::types::Embedder + Send + Sync>) {
self.embedder = Some(embedder);
}
pub fn has_embedder(&self) -> bool {
self.embedder.is_some()
}
pub fn embed(&self, text: &str) -> Result<Vec<f32>> {
self.embedder
.as_ref()
.ok_or(YantrikDbError::NoEmbedder)?
.embed(text)
.map_err(|e| YantrikDbError::Inference(e.to_string()))
}
pub fn record_text(
&self,
text: &str,
memory_type: &str,
importance: f64,
valence: f64,
half_life: f64,
metadata: &serde_json::Value,
namespace: &str,
certainty: f64,
domain: &str,
source: &str,
emotional_state: Option<&str>,
) -> Result<String> {
let embedding = self.embed(text)?;
self.record(
text,
memory_type,
importance,
valence,
half_life,
metadata,
&embedding,
namespace,
certainty,
domain,
source,
emotional_state,
)
}
pub fn recall_text(
&self,
query: &str,
top_k: usize,
) -> Result<Vec<RecallResult>> {
let embedding = self.embed(query)?;
self.recall(
&embedding,
top_k,
None, None, false, true, Some(query),
false, None, None, None, )
}
pub fn recall_text_filtered(
&self,
query: &str,
top_k: usize,
domain: Option<&str>,
source: Option<&str>,
) -> Result<Vec<RecallResult>> {
let embedding = self.embed(query)?;
self.recall(
&embedding,
top_k,
None, None, false, true, Some(query),
false, None, domain,
source,
)
}
}