use std::env;
use std::fs::{self, OpenOptions};
use std::path::{Path, PathBuf};
use std::sync::OnceLock;
use std::time::Duration;
#[cfg(unix)]
use std::os::unix::fs::{OpenOptionsExt as _, PermissionsExt as _};
use anyhow::Context as _;
use rusqlite::{Connection, TransactionBehavior, ffi::sqlite3_auto_extension};
use sqlite_vec::sqlite3_vec_init;
use crate::engine::memory::types::MemoryError;
pub(super) const SCHEMA_VERSION: i64 = 19;
static SQLITE_VEC_REGISTRATION: OnceLock<i32> = OnceLock::new();
pub(super) const SCHEMA_SQL: &str = r"CREATE TABLE projects(
id INTEGER PRIMARY KEY,
path TEXT NOT NULL UNIQUE,
created_at INTEGER NOT NULL
);
CREATE TABLE evidence(
id INTEGER PRIMARY KEY,
project_id INTEGER NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
provider TEXT NOT NULL,
session_id TEXT NOT NULL,
entry_id TEXT NOT NULL,
role TEXT NOT NULL,
observed_at INTEGER NOT NULL,
source_path TEXT NOT NULL,
content_hash TEXT NOT NULL,
content_json TEXT NOT NULL,
text TEXT NOT NULL,
created_at INTEGER NOT NULL,
extraction_completed_at INTEGER,
extraction_lease_owner TEXT,
extraction_lease_until INTEGER,
UNIQUE(provider, session_id, entry_id, content_hash)
);
CREATE INDEX evidence_project ON evidence(project_id, observed_at);
CREATE INDEX evidence_session ON evidence(provider, session_id);
CREATE INDEX evidence_entry ON evidence(provider, session_id, entry_id, created_at DESC);
CREATE INDEX evidence_extraction_lease ON evidence(extraction_lease_owner)
WHERE extraction_completed_at IS NULL;
CREATE TABLE claims(
id TEXT PRIMARY KEY,
project_id INTEGER NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
memory_type TEXT NOT NULL CHECK(memory_type IN (
'decision', 'fact', 'preference', 'procedure', 'lesson'
)),
statement TEXT NOT NULL,
keywords TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active' CHECK(status IN ('active', 'superseded')),
superseded_by TEXT REFERENCES claims(id) ON DELETE SET NULL,
valid_from INTEGER,
valid_until INTEGER,
created_at INTEGER NOT NULL,
CHECK(valid_from IS NULL OR valid_until IS NULL OR valid_until >= valid_from)
);
CREATE INDEX claims_project_type ON claims(project_id, memory_type, created_at);
CREATE INDEX claims_active ON claims(project_id, status, memory_type);
CREATE TABLE claim_evidence(
claim_id TEXT NOT NULL REFERENCES claims(id) ON DELETE CASCADE,
evidence_id INTEGER NOT NULL REFERENCES evidence(id) ON DELETE CASCADE,
PRIMARY KEY(claim_id, evidence_id)
) WITHOUT ROWID;
CREATE INDEX claim_evidence_source ON claim_evidence(evidence_id);
CREATE TABLE entities(
id INTEGER PRIMARY KEY,
project_id INTEGER NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
kind TEXT NOT NULL CHECK(kind IN ('path', 'crate', 'symbol', 'command', 'concept')),
value TEXT NOT NULL,
normalized TEXT NOT NULL,
UNIQUE(project_id, kind, normalized)
);
CREATE INDEX entities_lookup ON entities(project_id, normalized);
CREATE TABLE claim_entities(
claim_id TEXT NOT NULL REFERENCES claims(id) ON DELETE CASCADE,
entity_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
origin TEXT NOT NULL CHECK(origin IN ('statement', 'keyword', 'evidence')),
PRIMARY KEY(claim_id, entity_id, origin)
) WITHOUT ROWID;
CREATE INDEX claim_entities_entity ON claim_entities(entity_id, claim_id);
CREATE TABLE relations(
subject_claim_id TEXT NOT NULL REFERENCES claims(id) ON DELETE CASCADE,
predicate TEXT NOT NULL CHECK(predicate IN (
'duplicates', 'supports', 'revises', 'contradicts'
)),
object_claim_id TEXT NOT NULL REFERENCES claims(id) ON DELETE CASCADE,
rationale TEXT NOT NULL,
created_at INTEGER NOT NULL,
CHECK(subject_claim_id != object_claim_id),
PRIMARY KEY(subject_claim_id, predicate, object_claim_id)
) WITHOUT ROWID;
CREATE INDEX relations_object ON relations(object_claim_id, predicate);
CREATE VIRTUAL TABLE claim_embeddings USING vec0(
claim_id TEXT PRIMARY KEY,
project_id INTEGER PARTITION KEY,
embedding_model TEXT,
memory_type TEXT,
memory_status TEXT,
embedding FLOAT[384] distance_metric=cosine
);
CREATE TRIGGER claims_ad_embedding AFTER DELETE ON claims BEGIN
DELETE FROM claim_embeddings WHERE claim_id = old.id;
END;
CREATE TRIGGER claims_au_embedding
AFTER UPDATE OF memory_type, status ON claims BEGIN
UPDATE claim_embeddings
SET memory_type = new.memory_type,
memory_status = new.status
WHERE claim_id = old.id;
END;
CREATE TABLE tombstones(
kind TEXT NOT NULL CHECK(kind IN ('memory', 'session')),
key TEXT NOT NULL,
created_at INTEGER NOT NULL,
PRIMARY KEY(kind, key)
) WITHOUT ROWID;
CREATE VIRTUAL TABLE claims_fts USING fts5(
statement,
keywords,
content='claims',
content_rowid='rowid'
);
CREATE TRIGGER claims_clear_superseded_by BEFORE DELETE ON claims BEGIN
UPDATE claims SET superseded_by = NULL WHERE superseded_by = old.id;
END;
CREATE TRIGGER claims_ai AFTER INSERT ON claims BEGIN
INSERT INTO claims_fts(rowid, statement, keywords)
VALUES (new.rowid, new.statement, new.keywords);
END;
CREATE TRIGGER claims_ad AFTER DELETE ON claims BEGIN
INSERT INTO claims_fts(claims_fts, rowid, statement, keywords)
VALUES ('delete', old.rowid, old.statement, old.keywords);
END;
CREATE TRIGGER claims_au AFTER UPDATE ON claims BEGIN
INSERT INTO claims_fts(claims_fts, rowid, statement, keywords)
VALUES ('delete', old.rowid, old.statement, old.keywords);
INSERT INTO claims_fts(rowid, statement, keywords)
VALUES (new.rowid, new.statement, new.keywords);
END;
PRAGMA user_version = 19;";
pub(super) fn register_sqlite_vec() -> anyhow::Result<()> {
let result = *SQLITE_VEC_REGISTRATION.get_or_init(|| {
unsafe {
sqlite3_auto_extension(Some(std::mem::transmute::<
*const (),
unsafe extern "C" fn(
*mut rusqlite::ffi::sqlite3,
*mut *mut std::ffi::c_char,
*const rusqlite::ffi::sqlite3_api_routines,
) -> std::ffi::c_int,
>(sqlite3_vec_init as *const ())))
}
});
if result != rusqlite::ffi::SQLITE_OK {
return Err(MemoryError::SqliteVecRegistration(result).into());
}
Ok(())
}
pub(super) fn initialize(conn: &mut Connection) -> anyhow::Result<()> {
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let version: i64 = tx.pragma_query_value(None, "user_version", |row| row.get(0))?;
if version == 0 {
let objects: i64 = tx.query_row(
"SELECT count(*) FROM sqlite_schema WHERE name NOT LIKE 'sqlite_%'",
[],
|row| row.get(0),
)?;
if objects != 0 {
return Err(MemoryError::SchemaVersionMismatch {
found: version,
expected: SCHEMA_VERSION,
}
.into());
}
tx.execute_batch(SCHEMA_SQL)?;
} else if version == 18 {
tx.execute_batch(
"ALTER TABLE evidence ADD COLUMN extraction_completed_at INTEGER;
ALTER TABLE evidence ADD COLUMN extraction_lease_owner TEXT;
ALTER TABLE evidence ADD COLUMN extraction_lease_until INTEGER;
CREATE INDEX evidence_extraction_lease ON evidence(extraction_lease_owner)
WHERE extraction_completed_at IS NULL;
UPDATE evidence SET extraction_completed_at = created_at;
PRAGMA user_version = 19;",
)?;
} else if version != SCHEMA_VERSION {
return Err(MemoryError::SchemaVersionMismatch {
found: version,
expected: SCHEMA_VERSION,
}
.into());
}
tx.commit()?;
Ok(())
}
pub(super) fn database_path() -> anyhow::Result<PathBuf> {
if let Some(path) = env::var_os("GOOSEDUMP_STATE_DIR").filter(|path| !path.is_empty()) {
return Ok(PathBuf::from(path).join("memory.sqlite3"));
}
let dir = dirs::state_dir()
.or_else(dirs::data_local_dir)
.context("state or local data directory not found; set GOOSEDUMP_STATE_DIR")?;
Ok(dir.join("goosedump").join("memory.sqlite3"))
}
pub(super) fn open_database_connection(path: &Path) -> anyhow::Result<Connection> {
let conn = Connection::open(path).with_context(|| format!("open {}", path.display()))?;
conn.busy_timeout(Duration::from_secs(5))?;
conn.pragma_update(None, "foreign_keys", true)?;
Ok(conn)
}
pub(super) fn prepare_database_path(path: &Path) -> anyhow::Result<()> {
if let Some(parent) = path.parent()
&& !parent.exists()
{
fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
#[cfg(unix)]
secure_directory(parent)?;
}
prepare_database_file(path)
}
#[cfg(unix)]
pub(super) fn secure_directory(path: &Path) -> anyhow::Result<()> {
fs::set_permissions(path, fs::Permissions::from_mode(0o700))
.with_context(|| format!("secure {}", path.display()))
}
#[cfg(unix)]
pub(super) fn prepare_database_file(path: &Path) -> anyhow::Result<()> {
let _file = OpenOptions::new()
.create(true)
.append(true)
.mode(0o600)
.open(path)
.with_context(|| format!("prepare {}", path.display()))?;
fs::set_permissions(path, fs::Permissions::from_mode(0o600))
.with_context(|| format!("secure {}", path.display()))
}
#[cfg(not(unix))]
pub(super) fn prepare_database_file(path: &Path) -> anyhow::Result<()> {
let _file = OpenOptions::new()
.create(true)
.append(true)
.open(path)
.with_context(|| format!("prepare {}", path.display()))?;
Ok(())
}
pub(super) fn is_lock_error(error: &rusqlite::Error) -> bool {
matches!(
error,
rusqlite::Error::SqliteFailure(code, _)
if matches!(code.code, rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked)
)
}