mod embedding;
mod hydrate;
mod lease;
mod project;
mod recall;
mod schema;
mod util;
mod write;
use std::collections::HashSet;
use std::fs;
use std::path::{Path, PathBuf};
use anyhow::Context as _;
use rusqlite::{Connection, OptionalExtension as _, TransactionBehavior, params};
use crate::engine::Client;
use crate::engine::display;
use crate::engine::model::{EMBEDDING_MODEL_ID, Embedder, Embedding};
use super::entity::{EntityCandidate, extract_entities};
use super::extract::{
Extractor, KnownClaim, KnownEvidence, LocalExtractor, PendingSources, SourceCandidate,
extraction_text,
};
use super::types::{
ClaimRelationship, EntityReference, ForgetReport, MatchKind, MemoryError, MemoryFilter,
MemoryListItem, MemoryRecord, MemoryStats, MemoryType, RecallHit, RelationshipDirection,
RememberInput, RememberReport, SourceReference,
};
use super::{CONSOLIDATION_EVIDENCE_CHARS, clipped_chars, external_id};
use embedding::{EmbeddingProvider, backfill_project_embeddings, embed_query};
use hydrate::hydrate_claims;
use lease::{
ClaimedEvidence, claim_evidence, complete_evidence, leased_evidence_ids, release_evidence_lease,
};
use project::{coalesce_projects, normalized_path, prune_orphans};
use recall::{
MemoryEmbedding, RecallFill, SemanticCandidate, entity_match_detail, load_entity_related_rows,
load_memory_rows, load_query_entity_rows, map_memory_row, recall_match_reasons,
relation_candidates, relation_match_detail, relation_score_factor, store_embeddings,
};
#[cfg(unix)]
use schema::secure_directory;
use schema::{
SCHEMA_VERSION, database_path, initialize, is_lock_error, open_database_connection,
prepare_database_path, register_sqlite_vec,
};
use util::{
CONSOLIDATION_CANDIDATE_LIMIT, DISPLAY_ID_CHARS, MAX_RECALL_LIMIT, SEMANTIC_MIN_SIMILARITY,
SQLITE_VEC_MAX_K, claim_id, consolidation_fts_query, count, count_two, count_where, fts_query,
memory_token_estimate, memory_type_and_text_conn, memory_type_counts, now_millis,
parse_id_prefix, session_tombstone_key, sha256_hex, split_keywords,
};
use write::{ensure_project, insert_claims};
pub struct Memory {
conn: Connection,
path: PathBuf,
}
impl Memory {
pub fn open() -> anyhow::Result<Self> {
let path = database_path()?;
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
#[cfg(unix)]
secure_directory(parent)?;
}
Self::open_path(&path)
}
pub fn open_path(path: &Path) -> anyhow::Result<Self> {
register_sqlite_vec()?;
prepare_database_path(path)?;
let mut conn = open_database_connection(path)?;
initialize(&mut conn)?;
if let Err(error) = conn.pragma_update(None, "journal_mode", "WAL")
&& !is_lock_error(&error)
{
return Err(error.into());
}
coalesce_projects(&mut conn)?;
Ok(Self {
conn,
path: path.to_path_buf(),
})
}
pub fn remember(&mut self, input: &RememberInput<'_>) -> anyhow::Result<RememberReport> {
let pending = self.pending_evidence(input)?;
if pending.evidence.is_empty() || pending.context_forgotten {
return Ok(pending.empty_report());
}
let mut claimed = claim_evidence(&mut self.conn, input, pending)?;
if claimed.pending.evidence.is_empty() || claimed.pending.context_forgotten {
return Ok(claimed.empty_report());
}
let result = (|| {
let embedder = match Embedder::load() {
Ok(embedder) => Some(embedder),
Err(error) => {
eprintln!("goosedump: warning: could not load semantic memory model: {error}");
None
}
};
let known = self.consolidation_candidates(&claimed.pending, embedder.as_ref())?;
let mut extractor = LocalExtractor::load()?;
let report = self.remember_claimed(input, &mut claimed, &known, &mut extractor)?;
if let Some(embedder) = embedder.as_ref() {
self.maintain_embeddings(&report, &claimed.pending.project, embedder);
}
Ok(report)
})();
if result.is_err()
&& let Err(error) = release_evidence_lease(&self.conn, &claimed.owner)
{
eprintln!("goosedump: warning: could not release failed learning lease: {error}");
}
result
}
pub fn recall(
&self,
query: &str,
filter: &MemoryFilter,
limit: usize,
max_tokens: usize,
history: bool,
) -> anyhow::Result<Vec<RecallHit>> {
if query.trim().is_empty() || limit == 0 || max_tokens == 0 {
return Ok(Vec::new());
}
let fts = fts_query(query);
let query_entities = extract_entities(query, &[], &[]);
let project = filter.project.as_deref().map(normalized_path).transpose()?;
let memory_type = filter.memory_type.map(MemoryType::as_str);
let hit_limit = limit.min(MAX_RECALL_LIMIT);
let mut hits = Vec::new();
let mut used_tokens: usize = 0;
if !fts.is_empty() {
let mut stmt = self.conn.prepare(
"SELECT claims.id, claims.memory_type, claims.statement,
claims.keywords, projects.path, claims.valid_from,
claims.valid_until, claims.created_at, claims.status,
claims.superseded_by,
bm25(claims_fts, 1.0, 0.5) AS rank
FROM claims_fts
JOIN claims ON claims.rowid = claims_fts.rowid
JOIN projects ON projects.id = claims.project_id
WHERE claims_fts MATCH ?1
AND (?2 IS NULL OR projects.path = ?2)
AND (?3 IS NULL OR claims.memory_type = ?3)
AND (?4 OR claims.status = 'active')
ORDER BY rank, claims.created_at DESC
LIMIT ?5",
)?;
let rows = stmt
.query_map(
params![
fts,
project,
memory_type,
history,
i64::try_from(hit_limit)?
],
map_memory_row,
)?
.collect::<rusqlite::Result<Vec<_>>>()?;
let claim_ids = rows.iter().map(|row| row.id.clone()).collect::<Vec<_>>();
let mut hydrated = hydrate_claims(&self.conn, &claim_ids)?;
for row in rows {
let hydration = hydrated
.remove(&row.id)
.context("lexical row not hydrated")?;
let evidence = hydration.evidence;
let relationships = hydration.relationships;
let match_reasons = recall_match_reasons(
MatchKind::Lexical,
"full-text query",
row.status,
&relationships,
);
let estimated = memory_token_estimate(
&row.statement,
&relationships,
&match_reasons,
&evidence,
);
if !hits.is_empty() && used_tokens.saturating_add(estimated) > max_tokens {
break;
}
used_tokens = used_tokens.saturating_add(estimated);
let display_id = hydration.display_id;
hits.push(RecallHit {
id: external_id(&row.id),
display_id,
memory_type: row.memory_type,
text: row.statement,
keywords: split_keywords(&row.keywords),
status: row.status,
score: 1.0 / (1.0 + row.rank.abs()),
project: PathBuf::from(row.project),
valid_from: row.valid_from,
valid_until: row.valid_until,
superseded_by: row.superseded_by.map(|id| external_id(&id)),
related_by: Vec::new(),
relationships,
match_reasons,
evidence,
});
if hits.len() >= hit_limit {
break;
}
}
}
let mut fill = RecallFill {
hits: &mut hits,
used_tokens: &mut used_tokens,
limit: hit_limit,
max_tokens,
};
self.expand_structured_hits(&query_entities, filter, history, &mut fill)?;
if *fill.used_tokens < fill.max_tokens
&& let Err(error) = self.fill_semantic_hits(query, filter, history, &mut fill)
{
eprintln!("goosedump: warning: semantic memory recall unavailable: {error}");
}
self.expand_relation_hits(filter, history, &mut fill)?;
Ok(hits)
}
pub fn list(
&self,
filter: &MemoryFilter,
limit: usize,
history: bool,
) -> anyhow::Result<Vec<MemoryListItem>> {
if limit == 0 {
return Ok(Vec::new());
}
let project = filter.project.as_deref().map(normalized_path).transpose()?;
let memory_type = filter.memory_type.map(MemoryType::as_str);
let limit = i64::try_from(limit.min(MAX_RECALL_LIMIT))?;
let mut stmt = self.conn.prepare(
"SELECT claims.id, claims.memory_type, claims.statement,
claims.keywords, projects.path, claims.valid_from,
claims.valid_until, claims.created_at, claims.status,
(SELECT count(*) FROM claim_evidence
WHERE claim_evidence.claim_id = claims.id)
FROM claims
JOIN projects ON projects.id = claims.project_id
WHERE (?1 IS NULL OR projects.path = ?1)
AND (?2 IS NULL OR claims.memory_type = ?2)
AND (?3 OR claims.status = 'active')
ORDER BY coalesce(claims.valid_from, claims.created_at) DESC,
claims.created_at DESC
LIMIT ?4",
)?;
let rows = stmt.query_map(params![project, memory_type, history, limit], |row| {
let raw_type = row.get::<_, String>(1)?;
let memory_type = raw_type.parse().map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
1,
rusqlite::types::Type::Text,
Box::new(error),
)
})?;
let raw_status = row.get::<_, String>(8)?;
let status = raw_status.parse().map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
8,
rusqlite::types::Type::Text,
Box::new(error),
)
})?;
Ok((
row.get::<_, String>(0)?,
memory_type,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, String>(4)?,
row.get::<_, Option<i64>>(5)?,
row.get::<_, Option<i64>>(6)?,
row.get::<_, i64>(7)?,
status,
row.get::<_, i64>(9)?,
))
})?;
let mut items = Vec::new();
for row in rows {
let (
id,
memory_type,
statement,
keywords,
project,
valid_from,
valid_until,
created_at,
status,
evidence_count,
) = row?;
items.push(MemoryListItem {
id: external_id(&id),
display_id: self.display_id(&id)?,
memory_type,
text: statement,
keywords: split_keywords(&keywords),
status,
project: PathBuf::from(project),
valid_from,
valid_until,
created_at,
evidence_count: usize::try_from(evidence_count)?,
});
}
Ok(items)
}
pub fn show(&self, target: &str) -> anyhow::Result<MemoryRecord> {
let id = self.resolve_claim_id(target)?;
let row = self
.conn
.query_row(
"SELECT claims.id, claims.memory_type, claims.statement,
claims.keywords, projects.path, claims.valid_from,
claims.valid_until, claims.created_at, claims.status,
claims.superseded_by, 0.0
FROM claims
JOIN projects ON projects.id = claims.project_id
WHERE claims.id = ?1",
params![id],
map_memory_row,
)
.optional()?
.with_context(|| format!("memory '{target}' not found"))?;
let evidence = self.evidence_for(&row.id)?;
let supersedes = self.supersedes_for(&row.id)?;
Ok(MemoryRecord {
id: external_id(&row.id),
display_id: self.display_id(&row.id)?,
memory_type: row.memory_type,
text: row.statement,
keywords: split_keywords(&row.keywords),
status: row.status,
project: PathBuf::from(row.project),
valid_from: row.valid_from,
valid_until: row.valid_until,
created_at: row.created_at,
superseded_by: row.superseded_by.map(|id| external_id(&id)),
supersedes,
entities: self.entities_for(&row.id)?,
relationships: self.relationships_for(&row.id)?,
evidence,
})
}
fn entities_for(&self, claim_id: &str) -> anyhow::Result<Vec<EntityReference>> {
let mut stmt = self.conn.prepare(
"SELECT entities.kind, entities.value,
group_concat(claim_entities.origin, char(10))
FROM claim_entities
JOIN entities ON entities.id = claim_entities.entity_id
WHERE claim_entities.claim_id = ?1
GROUP BY entities.id
ORDER BY entities.kind, entities.value",
)?;
let rows = stmt
.query_map(params![claim_id], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
))
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(rows
.into_iter()
.map(|(kind, value, origins)| EntityReference {
kind,
value,
origins: origins
.lines()
.filter(|origin| !origin.is_empty())
.map(str::to_string)
.collect(),
})
.collect())
}
fn relationships_for(&self, claim_id: &str) -> anyhow::Result<Vec<ClaimRelationship>> {
let mut stmt = self.conn.prepare(
"SELECT 'outgoing', relations.predicate, claims.id, claims.memory_type,
claims.statement, claims.status, relations.rationale
FROM relations
JOIN claims ON claims.id = relations.object_claim_id
WHERE relations.subject_claim_id = ?1
UNION ALL
SELECT 'incoming', relations.predicate, claims.id, claims.memory_type,
claims.statement, claims.status, relations.rationale
FROM relations
JOIN claims ON claims.id = relations.subject_claim_id
WHERE relations.object_claim_id = ?1
ORDER BY 1, 2, 3",
)?;
let rows = stmt
.query_map(params![claim_id], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, String>(4)?,
row.get::<_, String>(5)?,
row.get::<_, String>(6)?,
))
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
rows.into_iter()
.map(
|(direction, kind, id, memory_type, text, status, rationale)| {
let direction = match direction.as_str() {
"outgoing" => RelationshipDirection::Outgoing,
"incoming" => RelationshipDirection::Incoming,
_ => return Err(MemoryError::InvalidRelationshipDirection.into()),
};
Ok(ClaimRelationship {
direction,
kind: kind.parse()?,
claim_id: external_id(&id),
display_id: self.display_id(&id)?,
memory_type: memory_type.parse()?,
text,
status: status.parse()?,
rationale,
})
},
)
.collect()
}
fn consolidation_candidates(
&self,
pending: &PendingSources,
embedder: Option<&Embedder>,
) -> anyhow::Result<Vec<KnownClaim>> {
let query = pending
.evidence
.iter()
.map(|source| source.extraction_text.as_str())
.filter(|text| !text.trim().is_empty())
.collect::<Vec<_>>()
.join("\n");
if query.trim().is_empty() {
return Ok(Vec::new());
}
let project = PathBuf::from(&pending.project);
let filter = MemoryFilter {
project: Some(project),
memory_type: None,
};
let entities = extract_entities(&query, &[], &[]);
let mut ids = load_query_entity_rows(
&self.conn,
&entities,
Some(pending.project.clone()),
None,
false,
CONSOLIDATION_CANDIDATE_LIMIT,
)?
.into_iter()
.map(|row| row.id)
.collect::<Vec<_>>();
if let Some(embedder) = embedder {
match embed_query(embedder, &query) {
Ok(query_embedding) => {
let excluded = ids.iter().cloned().collect();
for candidate in
self.semantic_candidates(&filter, false, &query_embedding, &excluded)?
{
if ids.len() == CONSOLIDATION_CANDIDATE_LIMIT {
break;
}
if !ids.contains(&candidate.row.id) {
ids.push(candidate.row.id);
}
}
}
Err(error) => {
eprintln!(
"goosedump: warning: semantic memory consolidation unavailable: {error}"
);
}
}
}
for id in
self.lexical_consolidation_ids(&query, &pending.project, CONSOLIDATION_CANDIDATE_LIMIT)?
{
if ids.len() == CONSOLIDATION_CANDIDATE_LIMIT {
break;
}
if !ids.contains(&id) {
ids.push(id);
}
}
ids.truncate(CONSOLIDATION_CANDIDATE_LIMIT);
ids.into_iter()
.enumerate()
.map(|(index, id)| self.known_claim(format!("m{index}"), &id))
.collect()
}
fn lexical_consolidation_ids(
&self,
query: &str,
project: &str,
limit: usize,
) -> anyhow::Result<Vec<String>> {
let fts = consolidation_fts_query(query);
if fts.is_empty() || limit == 0 {
return Ok(Vec::new());
}
let mut stmt = self.conn.prepare(
"SELECT claims.id
FROM claims_fts
JOIN claims ON claims.rowid = claims_fts.rowid
JOIN projects ON projects.id = claims.project_id
WHERE claims_fts MATCH ?1
AND projects.path = ?2
AND claims.status = 'active'
ORDER BY bm25(claims_fts, 1.0, 0.5), claims.created_at DESC
LIMIT ?3",
)?;
Ok(stmt
.query_map(params![fts, project, i64::try_from(limit)?], |row| {
row.get(0)
})?
.collect::<rusqlite::Result<Vec<_>>>()?)
}
fn known_claim(&self, prompt_id: String, id: &str) -> anyhow::Result<KnownClaim> {
let (memory_type, statement) = memory_type_and_text_conn(&self.conn, id)?;
let evidence = self
.evidence_for(id)?
.into_iter()
.map(|source| KnownEvidence {
observed_at: source.observed_at,
text: clipped_chars(&source.snippet, CONSOLIDATION_EVIDENCE_CHARS),
})
.collect();
Ok(KnownClaim {
prompt_id,
id: id.to_string(),
memory_type,
statement,
evidence,
})
}
fn expand_structured_hits(
&self,
query_entities: &[EntityCandidate],
filter: &MemoryFilter,
history: bool,
fill: &mut RecallFill<'_>,
) -> anyhow::Result<()> {
self.seed_entity_hits(query_entities, filter, history, fill)?;
self.expand_entity_hits(filter, history, fill)?;
self.expand_relation_hits(filter, history, fill)
}
fn seed_entity_hits(
&self,
query_entities: &[EntityCandidate],
filter: &MemoryFilter,
history: bool,
fill: &mut RecallFill<'_>,
) -> anyhow::Result<()> {
if query_entities.is_empty() || fill.hits.len() >= fill.limit {
return Ok(());
}
let project = filter.project.as_deref().map(normalized_path).transpose()?;
let memory_type = filter.memory_type.map(MemoryType::as_str);
let rows = load_query_entity_rows(
&self.conn,
query_entities,
project,
memory_type,
history,
fill.limit,
)?;
let claim_ids = rows.iter().map(|row| row.id.clone()).collect::<Vec<_>>();
let mut hydrated = hydrate_claims(&self.conn, &claim_ids)?;
let known: HashSet<String> = fill.hits.iter().map(|hit| hit.id.clone()).collect();
for row in rows {
if fill.hits.len() >= fill.limit {
break;
}
let external = external_id(&row.id);
if known.contains(&external) {
continue;
}
let hydration = hydrated
.remove(&row.id)
.context("entity recall row was not hydrated")?;
let evidence = hydration.evidence;
let relationships = hydration.relationships;
let match_reasons = recall_match_reasons(
MatchKind::Entity,
entity_match_detail(row.shared, &row.related_by),
row.status,
&relationships,
);
let estimated =
memory_token_estimate(&row.statement, &relationships, &match_reasons, &evidence);
if !fill.hits.is_empty() && fill.used_tokens.saturating_add(estimated) > fill.max_tokens
{
break;
}
*fill.used_tokens = fill.used_tokens.saturating_add(estimated);
let shared_factor = f64::from(u32::try_from(row.shared.clamp(1, 8)).unwrap_or(1)) / 8.0;
let score = (0.35 * shared_factor).min(0.45);
fill.hits.push(RecallHit {
id: external,
display_id: hydration.display_id,
memory_type: row.memory_type,
text: row.statement,
keywords: split_keywords(&row.keywords),
status: row.status,
score,
project: PathBuf::from(row.project),
valid_from: row.valid_from,
valid_until: row.valid_until,
superseded_by: row.superseded_by.map(|value| external_id(&value)),
related_by: row.related_by,
relationships,
match_reasons,
evidence,
});
}
Ok(())
}
fn expand_entity_hits(
&self,
filter: &MemoryFilter,
history: bool,
fill: &mut RecallFill<'_>,
) -> anyhow::Result<()> {
if fill.hits.is_empty() || fill.hits.len() >= fill.limit {
return Ok(());
}
let project = filter.project.as_deref().map(normalized_path).transpose()?;
let memory_type = filter.memory_type.map(MemoryType::as_str);
let seed_ids: Vec<String> = fill
.hits
.iter()
.filter_map(|hit| hit.id.strip_prefix("mem_").map(str::to_string))
.collect();
if seed_ids.is_empty() {
return Ok(());
}
let rows = load_entity_related_rows(
&self.conn,
&seed_ids,
project,
memory_type,
history,
fill.limit,
)?;
let claim_ids = rows.iter().map(|row| row.id.clone()).collect::<Vec<_>>();
let mut hydrated = hydrate_claims(&self.conn, &claim_ids)?;
let known: HashSet<String> = fill.hits.iter().map(|hit| hit.id.clone()).collect();
let min_lexical = fill
.hits
.iter()
.map(|hit| hit.score)
.fold(f64::INFINITY, f64::min);
let baseline = if min_lexical.is_finite() {
min_lexical
} else {
0.5
};
for row in rows {
if fill.hits.len() >= fill.limit {
break;
}
let external = external_id(&row.id);
if known.contains(&external) {
continue;
}
let hydration = hydrated
.remove(&row.id)
.context("related entity recall row was not hydrated")?;
let evidence = hydration.evidence;
let relationships = hydration.relationships;
let match_reasons = recall_match_reasons(
MatchKind::Entity,
entity_match_detail(row.shared, &row.related_by),
row.status,
&relationships,
);
let estimated =
memory_token_estimate(&row.statement, &relationships, &match_reasons, &evidence);
if !fill.hits.is_empty() && fill.used_tokens.saturating_add(estimated) > fill.max_tokens
{
break;
}
*fill.used_tokens = fill.used_tokens.saturating_add(estimated);
let shared_factor = f64::from(u32::try_from(row.shared.clamp(1, 8)).unwrap_or(1)) / 8.0;
let score = (baseline * 0.45 * shared_factor).min(baseline * 0.9);
fill.hits.push(RecallHit {
id: external,
display_id: hydration.display_id,
memory_type: row.memory_type,
text: row.statement,
keywords: split_keywords(&row.keywords),
status: row.status,
score,
project: PathBuf::from(row.project),
valid_from: row.valid_from,
valid_until: row.valid_until,
superseded_by: row.superseded_by.map(|value| external_id(&value)),
related_by: row.related_by,
relationships,
match_reasons,
evidence,
});
}
Ok(())
}
fn expand_relation_hits(
&self,
filter: &MemoryFilter,
history: bool,
fill: &mut RecallFill<'_>,
) -> anyhow::Result<()> {
if fill.hits.is_empty() || fill.hits.len() >= fill.limit {
return Ok(());
}
let project = filter.project.as_deref().map(normalized_path).transpose()?;
let memory_type = filter.memory_type.map(MemoryType::as_str);
let candidates = relation_candidates(fill.hits, history, filter.memory_type);
let mut known: HashSet<String> = fill.hits.iter().map(|hit| hit.id.clone()).collect();
let candidate_ids = candidates
.iter()
.map(|candidate| candidate.claim_id.clone())
.collect::<Vec<_>>();
let mut memory_rows = load_memory_rows(
&self.conn,
&candidate_ids,
project.as_deref(),
memory_type,
history,
)?;
let claim_ids = memory_rows.keys().cloned().collect::<Vec<_>>();
let mut hydrated = hydrate_claims(&self.conn, &claim_ids)?;
for candidate in candidates {
if fill.hits.len() >= fill.limit {
break;
}
let external = external_id(&candidate.claim_id);
if known.contains(&external) {
continue;
}
let Some(row) = memory_rows.remove(&candidate.claim_id) else {
continue;
};
let hydration = hydrated
.remove(&row.id)
.context("related claim recall row was not hydrated")?;
let evidence = hydration.evidence;
let relationships = hydration.relationships;
let match_reasons = recall_match_reasons(
MatchKind::Relation,
relation_match_detail(
candidate.kind,
&candidate.seed_display_id,
&candidate.rationale,
),
row.status,
&relationships,
);
let estimated =
memory_token_estimate(&row.statement, &relationships, &match_reasons, &evidence);
if fill.used_tokens.saturating_add(estimated) > fill.max_tokens {
continue;
}
*fill.used_tokens = fill.used_tokens.saturating_add(estimated);
known.insert(external.clone());
fill.hits.push(RecallHit {
id: external,
display_id: hydration.display_id,
memory_type: row.memory_type,
text: row.statement,
keywords: split_keywords(&row.keywords),
status: row.status,
score: candidate.seed_score * relation_score_factor(candidate.kind),
project: PathBuf::from(row.project),
valid_from: row.valid_from,
valid_until: row.valid_until,
superseded_by: row.superseded_by.map(|value| external_id(&value)),
related_by: Vec::new(),
relationships,
match_reasons,
evidence,
});
}
Ok(())
}
fn maintain_embeddings(
&self,
report: &RememberReport,
project: &str,
embedder: &impl EmbeddingProvider,
) {
if !report.added.is_empty()
&& let Err(error) = self.embed_added_claims(report, embedder)
{
eprintln!("goosedump: warning: could not index new claims: {error}");
}
if let Err(error) = backfill_project_embeddings(&self.conn, project, embedder) {
eprintln!("goosedump: warning: could not backfill claim embeddings: {error}");
}
}
fn embed_added_claims(
&self,
report: &RememberReport,
embedder: &impl EmbeddingProvider,
) -> anyhow::Result<()> {
let ids = report
.added
.iter()
.map(|memory| {
memory
.id
.strip_prefix("mem_")
.context("remember report has an invalid memory ID")
.map(str::to_string)
})
.collect::<anyhow::Result<Vec<_>>>()?;
let texts = report
.added
.iter()
.map(|memory| memory.text.as_str())
.collect::<Vec<_>>();
let vectors = embedder.embed_batch(&texts)?;
if vectors.len() != ids.len() {
return Err(MemoryError::EmbeddingBatchMismatch.into());
}
let embeddings = ids
.into_iter()
.zip(vectors)
.map(|(claim_id, embedding)| MemoryEmbedding {
claim_id,
embedding,
})
.collect::<Vec<_>>();
store_embeddings(&self.conn, &embeddings)
}
fn fill_semantic_hits(
&self,
query: &str,
filter: &MemoryFilter,
history: bool,
fill: &mut RecallFill<'_>,
) -> anyhow::Result<()> {
if fill.hits.len() >= fill.limit {
return Ok(());
}
let embedder = Embedder::load()?;
let query_embedding = embed_query(&embedder, query)?;
let mut excluded: HashSet<String> = fill
.hits
.iter()
.filter_map(|hit| hit.id.strip_prefix("mem_"))
.map(str::to_string)
.collect();
let baseline = fill
.hits
.iter()
.map(|hit| hit.score)
.fold(f64::INFINITY, f64::min);
let baseline = if baseline.is_finite() { baseline } else { 0.5 };
'pages: loop {
let candidates =
self.semantic_candidates(filter, history, &query_embedding, &excluded)?;
if candidates.is_empty() {
break;
}
let has_more = candidates.len() == MAX_RECALL_LIMIT;
let claim_ids = candidates
.iter()
.map(|candidate| candidate.row.id.clone())
.collect::<Vec<_>>();
let mut hydrated = hydrate_claims(&self.conn, &claim_ids)?;
for candidate in candidates {
if fill.hits.len() >= fill.limit {
break 'pages;
}
excluded.insert(candidate.row.id.clone());
let external = external_id(&candidate.row.id);
let hydration = hydrated
.remove(&candidate.row.id)
.context("semantic recall row was not hydrated")?;
let evidence = hydration.evidence;
let relationships = hydration.relationships;
let match_reasons = recall_match_reasons(
MatchKind::Vector,
format!("semantic similarity {:.3}", candidate.similarity),
candidate.row.status,
&relationships,
);
let estimated = memory_token_estimate(
&candidate.row.statement,
&relationships,
&match_reasons,
&evidence,
);
if !fill.hits.is_empty()
&& fill.used_tokens.saturating_add(estimated) > fill.max_tokens
{
continue;
}
*fill.used_tokens = fill.used_tokens.saturating_add(estimated);
fill.hits.push(RecallHit {
id: external,
display_id: hydration.display_id,
memory_type: candidate.row.memory_type,
text: candidate.row.statement,
keywords: split_keywords(&candidate.row.keywords),
status: candidate.row.status,
score: baseline * 0.4 * candidate.similarity.clamp(0.0, 1.0),
project: PathBuf::from(candidate.row.project),
valid_from: candidate.row.valid_from,
valid_until: candidate.row.valid_until,
superseded_by: candidate.row.superseded_by.map(|id| external_id(&id)),
related_by: Vec::new(),
relationships,
match_reasons,
evidence,
});
if *fill.used_tokens >= fill.max_tokens {
break 'pages;
}
}
if !has_more {
break;
}
}
Ok(())
}
fn semantic_candidates(
&self,
filter: &MemoryFilter,
history: bool,
query_embedding: &Embedding,
known: &HashSet<String>,
) -> anyhow::Result<Vec<SemanticCandidate>> {
if known.len() >= SQLITE_VEC_MAX_K {
return Ok(Vec::new());
}
let project_id = if let Some(path) = filter.project.as_deref() {
let path = normalized_path(path)?;
let id = self
.conn
.query_row(
"SELECT id FROM projects WHERE path = ?1",
params![path],
|row| row.get::<_, i64>(0),
)
.optional()?;
let Some(id) = id else {
return Ok(Vec::new());
};
Some(id)
} else {
None
};
let mut conditions = String::new();
let k = MAX_RECALL_LIMIT
.saturating_add(known.len())
.min(SQLITE_VEC_MAX_K);
let mut values = vec![
rusqlite::types::Value::Blob(Vec::<u8>::from(query_embedding)),
rusqlite::types::Value::Integer(i64::try_from(k)?),
rusqlite::types::Value::Text(EMBEDDING_MODEL_ID.to_string()),
];
if let Some(project_id) = project_id {
conditions.push_str(" AND project_id = ?");
values.push(rusqlite::types::Value::Integer(project_id));
}
if let Some(memory_type) = filter.memory_type {
conditions.push_str(" AND memory_type = ?");
values.push(rusqlite::types::Value::Text(
memory_type.as_str().to_string(),
));
}
if !history {
conditions.push_str(" AND memory_status = 'active'");
}
let sql = format!(
"WITH nearest AS (
SELECT claim_id, distance
FROM claim_embeddings
WHERE embedding MATCH ? AND k = ? AND embedding_model = ?{conditions}
)
SELECT claims.id, claims.memory_type, claims.statement,
claims.keywords, projects.path, claims.valid_from,
claims.valid_until, claims.created_at, claims.status,
claims.superseded_by, 1.0 - nearest.distance
FROM nearest
JOIN claims ON claims.id = nearest.claim_id
JOIN projects ON projects.id = claims.project_id
ORDER BY nearest.distance, claims.created_at DESC"
);
let mut stmt = self.conn.prepare(&sql)?;
let mut rows = stmt.query(rusqlite::params_from_iter(values))?;
let mut candidates = Vec::new();
while let Some(row) = rows.next()? {
let memory = map_memory_row(row)?;
if known.contains(&memory.id) {
continue;
}
let similarity = memory.rank;
if similarity < SEMANTIC_MIN_SIMILARITY {
break;
}
candidates.push(SemanticCandidate {
row: memory,
similarity,
});
if candidates.len() == MAX_RECALL_LIMIT {
break;
}
}
Ok(candidates)
}
fn supersedes_for(&self, claim_id: &str) -> anyhow::Result<Vec<String>> {
let mut stmt = self.conn.prepare(
"SELECT id FROM claims
WHERE superseded_by = ?1
ORDER BY coalesce(valid_until, created_at) DESC, created_at DESC",
)?;
let ids = stmt
.query_map(params![claim_id], |row| row.get::<_, String>(0))?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(ids.into_iter().map(|id| external_id(&id)).collect())
}
pub fn forget_memory(&mut self, target: &str, apply: bool) -> anyhow::Result<ForgetReport> {
let id = self.resolve_claim_id(target)?;
let (project, memory_type, statement): (String, String, String) = self.conn.query_row(
"SELECT projects.path, claims.memory_type, claims.statement
FROM claims
JOIN projects ON projects.id = claims.project_id
WHERE claims.id = ?1",
params![id],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
)?;
let canonical_id = claim_id(&project, memory_type.parse::<MemoryType>()?, &statement);
let claim_evidence = count_where(
&self.conn,
"SELECT count(*) FROM claim_evidence WHERE claim_id = ?1",
&id,
)?;
let evidence = count_where(
&self.conn,
"SELECT count(*) FROM evidence
WHERE EXISTS(
SELECT 1 FROM claim_evidence
WHERE claim_evidence.evidence_id = evidence.id
AND claim_evidence.claim_id = ?1
) AND NOT EXISTS(
SELECT 1 FROM claim_evidence
WHERE claim_evidence.evidence_id = evidence.id
AND claim_evidence.claim_id != ?1
)",
&id,
)?;
let existing_tombstones = u64::try_from(self.conn.query_row(
"SELECT count(*) FROM tombstones
WHERE kind = 'memory' AND key IN (?1, ?2)",
params![id, canonical_id],
|row| row.get::<_, i64>(0),
)?)?;
let tombstone_count: u64 = if id == canonical_id { 1 } else { 2 };
let mut report = ForgetReport {
target: external_id(&id),
claims: 1,
evidence,
claim_evidence,
tombstones: tombstone_count.saturating_sub(existing_tombstones),
applied: apply,
};
if !apply {
return Ok(report);
}
let exclusive_evidence_ids: Vec<i64> = {
let mut stmt = self.conn.prepare(
"SELECT evidence.id FROM evidence
WHERE EXISTS(
SELECT 1 FROM claim_evidence
WHERE claim_evidence.evidence_id = evidence.id
AND claim_evidence.claim_id = ?1
) AND NOT EXISTS(
SELECT 1 FROM claim_evidence
WHERE claim_evidence.evidence_id = evidence.id
AND claim_evidence.claim_id != ?1
)",
)?;
stmt.query_map(params![id], |row| row.get(0))?
.collect::<rusqlite::Result<Vec<_>>>()?
};
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
tx.execute(
"INSERT OR IGNORE INTO tombstones(kind, key, created_at)
VALUES ('memory', ?1, ?3), ('memory', ?2, ?3)",
params![id, canonical_id, now_millis()],
)?;
report.claims =
u64::try_from(tx.execute("DELETE FROM claims WHERE id = ?1", params![id])?)?;
let mut deleted_evidence = 0u64;
for evidence_id in exclusive_evidence_ids {
deleted_evidence += u64::try_from(
tx.execute("DELETE FROM evidence WHERE id = ?1", params![evidence_id])?,
)?;
}
report.evidence = deleted_evidence;
prune_orphans(&tx)?;
tx.commit()?;
Ok(report)
}
pub fn forget_session(
&mut self,
provider: Client,
session_id: &str,
apply: bool,
) -> anyhow::Result<ForgetReport> {
let provider_name = provider.as_str();
let key = session_tombstone_key(provider_name, session_id);
let evidence = count_two(
&self.conn,
"SELECT count(*) FROM evidence WHERE provider = ?1 AND session_id = ?2",
provider_name,
session_id,
)?;
let claim_evidence = count_two(
&self.conn,
"SELECT count(*)
FROM claim_evidence
JOIN evidence ON evidence.id = claim_evidence.evidence_id
WHERE evidence.provider = ?1 AND evidence.session_id = ?2",
provider_name,
session_id,
)?;
let claims = count_two(
&self.conn,
"SELECT count(*) FROM claims
WHERE EXISTS(
SELECT 1 FROM claim_evidence
JOIN evidence ON evidence.id = claim_evidence.evidence_id
WHERE claim_evidence.claim_id = claims.id
AND evidence.provider = ?1 AND evidence.session_id = ?2
) AND NOT EXISTS(
SELECT 1 FROM claim_evidence
JOIN evidence ON evidence.id = claim_evidence.evidence_id
WHERE claim_evidence.claim_id = claims.id
AND NOT (evidence.provider = ?1 AND evidence.session_id = ?2)
)",
provider_name,
session_id,
)?;
let tombstone_exists: bool = self.conn.query_row(
"SELECT EXISTS(SELECT 1 FROM tombstones WHERE kind = 'session' AND key = ?1)",
params![key],
|row| row.get(0),
)?;
let mut report = ForgetReport {
target: format!("{provider_name}:{session_id}"),
claims,
evidence,
claim_evidence,
tombstones: u64::from(!tombstone_exists),
applied: apply,
};
if !apply {
return Ok(report);
}
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
tx.execute(
"INSERT OR IGNORE INTO tombstones(kind, key, created_at)
VALUES ('session', ?1, ?2)",
params![key, now_millis()],
)?;
report.evidence = u64::try_from(tx.execute(
"DELETE FROM evidence WHERE provider = ?1 AND session_id = ?2",
params![provider_name, session_id],
)?)?;
report.claims = u64::try_from(tx.execute(
"DELETE FROM claims
WHERE NOT EXISTS(
SELECT 1 FROM claim_evidence
WHERE claim_evidence.claim_id = claims.id
)",
[],
)?)?;
prune_orphans(&tx)?;
tx.commit()?;
Ok(report)
}
pub fn stats(&self) -> anyhow::Result<MemoryStats> {
let claims = count(&self.conn, "claims")?;
let embeddings = count_where(
&self.conn,
"SELECT count(*) FROM claim_embeddings WHERE embedding_model = ?1",
EMBEDDING_MODEL_ID,
)?;
let pending_embeddings = count_where(
&self.conn,
"SELECT count(*)
FROM claims
LEFT JOIN claim_embeddings
ON claim_embeddings.claim_id = claims.id
AND claim_embeddings.embedding_model = ?1
WHERE claims.status = 'active'
AND claim_embeddings.claim_id IS NULL",
EMBEDDING_MODEL_ID,
)?;
Ok(MemoryStats {
schema_version: SCHEMA_VERSION,
database: self.path.clone(),
projects: count(&self.conn, "projects")?,
evidence: count(&self.conn, "evidence")?,
claims,
claim_evidence: count(&self.conn, "claim_evidence")?,
entities: count(&self.conn, "entities")?,
entity_links: count(&self.conn, "claim_entities")?,
relations: count(&self.conn, "relations")?,
embedding_model: EMBEDDING_MODEL_ID.to_string(),
embeddings,
pending_embeddings,
tombstones: count(&self.conn, "tombstones")?,
types: memory_type_counts(&self.conn)?,
last_learned_at: self.conn.query_row(
"SELECT max(created_at) FROM evidence",
[],
|row| row.get(0),
)?,
})
}
fn pending_evidence(&self, input: &RememberInput<'_>) -> anyhow::Result<PendingSources> {
let project = normalized_path(input.project)?;
let key = session_tombstone_key(input.provider.as_str(), input.session_id);
let context_forgotten: bool = self.conn.query_row(
"SELECT EXISTS(SELECT 1 FROM tombstones WHERE kind = 'session' AND key = ?1)",
params![key],
|row| row.get(0),
)?;
if context_forgotten {
return Ok(PendingSources {
project,
evidence_seen: input.context.messages.len(),
evidence: Vec::new(),
context_forgotten: true,
});
}
let mut evidence = Vec::new();
for (ordinal, message) in input.context.messages.iter().enumerate() {
let content_json = serde_json::to_string(message).context("encode memory evidence")?;
let content_hash = sha256_hex(content_json.as_bytes());
let entry_id = if message.entry_id.is_empty() {
format!("content:{content_hash}:{ordinal}")
} else {
message.entry_id.clone()
};
let existing = self
.conn
.query_row(
"SELECT projects.path, evidence.extraction_completed_at
FROM evidence
JOIN projects ON projects.id = evidence.project_id
WHERE evidence.provider = ?1 AND evidence.session_id = ?2
AND evidence.entry_id = ?3 AND evidence.content_hash = ?4",
params![
input.provider.as_str(),
input.session_id,
entry_id,
content_hash
],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?)),
)
.optional()?;
if let Some((existing_project, completed_at)) = existing {
if normalized_path(Path::new(&existing_project))? != project {
return Err(MemoryError::ProjectConflict {
session_id: input.session_id.to_string(),
entry_id: entry_id.clone(),
existing_project,
}
.into());
}
if completed_at.is_some() {
continue;
}
}
let observed_at = message
.timestamp
.map_or_else(now_millis, |timestamp| timestamp.timestamp_millis());
evidence.push(SourceCandidate {
prompt_id: format!("s{}", evidence.len()),
entry_id,
role: message.role_label(),
observed_at,
source_path: input.source_path.to_path_buf(),
content_hash,
content_json,
text: display::searchable_text(message),
extraction_text: extraction_text(message),
});
}
Ok(PendingSources {
project,
evidence_seen: input.context.messages.len(),
evidence,
context_forgotten: false,
})
}
fn remember_claimed<E: Extractor>(
&mut self,
input: &RememberInput<'_>,
claimed: &mut ClaimedEvidence,
known: &[KnownClaim],
extractor: &mut E,
) -> anyhow::Result<RememberReport> {
let extracted = match extractor.extract(&claimed.pending.evidence, known) {
Ok(extracted) => extracted,
Err(error) => {
if let Err(release_error) = release_evidence_lease(&self.conn, &claimed.owner) {
eprintln!(
"goosedump: warning: could not release failed extraction lease: {release_error}"
);
}
return Err(error);
}
};
let result = self.store_claimed(input, claimed, known, extracted);
if result.is_err()
&& let Err(error) = release_evidence_lease(&self.conn, &claimed.owner)
{
eprintln!("goosedump: warning: could not release extraction lease: {error}");
}
result
}
fn store_claimed(
&mut self,
input: &RememberInput<'_>,
claimed: &mut ClaimedEvidence,
known: &[KnownClaim],
extracted: Vec<super::extract::ExtractedMemory>,
) -> anyhow::Result<RememberReport> {
let now = now_millis();
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
let key = session_tombstone_key(input.provider.as_str(), input.session_id);
let forgotten: bool = tx.query_row(
"SELECT EXISTS(SELECT 1 FROM tombstones WHERE kind = 'session' AND key = ?1)",
params![key],
|row| row.get(0),
)?;
if forgotten {
release_evidence_lease(&tx, &claimed.owner)?;
tx.commit()?;
return Ok(RememberReport {
evidence_seen: claimed.pending.evidence_seen,
skipped_tombstones: claimed.pending.evidence_seen,
..RememberReport::default()
});
}
let owned_ids = leased_evidence_ids(&tx, &claimed.owner)?;
if owned_ids.is_empty() {
tx.commit()?;
return Ok(claimed.empty_report());
}
claimed.pending.evidence.retain(|source| {
claimed
.evidence_ids
.get(&source.prompt_id)
.is_some_and(|id| owned_ids.contains(id))
});
claimed
.evidence_ids
.retain(|_, evidence_id| owned_ids.contains(evidence_id));
let project_id = ensure_project(&tx, &claimed.pending.project, now)?;
let inserted = insert_claims(
&tx,
&claimed.pending,
&claimed.evidence_ids,
extracted,
known,
project_id,
now,
)?;
let completed = complete_evidence(&tx, &claimed.owner, now)?;
anyhow::ensure!(
completed == owned_ids.len(),
"evidence extraction lease changed during completion"
);
tx.commit()?;
Ok(RememberReport {
evidence_seen: claimed.pending.evidence_seen,
evidence_added: claimed.evidence_added,
claims_added: inserted.claims_added,
claims_superseded: inserted.claims_superseded,
claim_evidence_added: inserted.evidence_added,
relations_added: inserted.relations_added,
duplicates_merged: inserted.duplicates_merged,
contradictions_found: inserted.contradictions_found,
skipped_tombstones: 0,
added: inserted.added,
superseded: inserted.superseded,
})
}
fn evidence_for(&self, claim_id: &str) -> anyhow::Result<Vec<SourceReference>> {
let mut stmt = self.conn.prepare(
"SELECT evidence.provider, evidence.session_id, evidence.entry_id,
evidence.role, evidence.observed_at, projects.path,
evidence.source_path, evidence.content_hash,
substr(evidence.text, 1, 500), evidence.id
FROM claim_evidence
JOIN evidence ON evidence.id = claim_evidence.evidence_id
JOIN projects ON projects.id = evidence.project_id
WHERE claim_evidence.claim_id = ?1
ORDER BY evidence.observed_at, evidence.id",
)?;
Ok(stmt
.query_map(params![claim_id], |row| {
let raw_provider = row.get::<_, String>(0)?;
let provider = raw_provider.parse().map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
0,
rusqlite::types::Type::Text,
Box::new(error),
)
})?;
let evidence_id = row.get::<_, i64>(9)?;
Ok(SourceReference {
citation_id: format!("ev_{evidence_id}"),
provider,
session_id: row.get(1)?,
entry_id: row.get(2)?,
role: row.get(3)?,
observed_at: row.get(4)?,
project: PathBuf::from(row.get::<_, String>(5)?),
source_path: PathBuf::from(row.get::<_, String>(6)?),
content_hash: row.get(7)?,
snippet: row.get(8)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?)
}
fn resolve_claim_id(&self, target: &str) -> anyhow::Result<String> {
let prefix = parse_id_prefix(target)?;
let pattern = format!("{prefix}%");
let mut stmt = self
.conn
.prepare("SELECT id FROM claims WHERE id LIKE ?1 ORDER BY id LIMIT 2")?;
let ids = stmt
.query_map(params![pattern], |row| row.get::<_, String>(0))?
.collect::<rusqlite::Result<Vec<_>>>()?;
match ids.as_slice() {
[] => Err(MemoryError::NotFound(target.to_string()).into()),
[id] => Ok(id.clone()),
_ => Err(MemoryError::AmbiguousTarget(target.to_string()).into()),
}
}
fn display_id(&self, id: &str) -> anyhow::Result<String> {
let start = DISPLAY_ID_CHARS.min(id.len());
for chars in start..=id.len() {
let prefix = &id[..chars];
let matches: i64 = self.conn.query_row(
"SELECT count(*) FROM claims WHERE id LIKE ?1",
params![format!("{prefix}%")],
|row| row.get(0),
)?;
if matches <= 1 {
return Ok(format!("mem_{prefix}"));
}
}
Ok(external_id(id))
}
}