use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use car_engine::ToolExecutor;
use car_eventlog::Event;
use car_memgine::{
MemgineEngine, ProactiveMaintenanceRequest, ProactiveMemoryDecision, ProactiveMemoryRequest,
ProactiveMemorySave,
};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
#[derive(Clone, Serialize, Deserialize)]
struct Note {
subject: String,
body: String,
#[serde(default)]
kind: NoteKind,
}
#[derive(Clone, Copy, Default, PartialEq, Eq, Debug, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum NoteKind {
#[default]
Fact,
Preference,
Procedure,
}
impl NoteKind {
fn parse(raw: Option<&str>) -> Result<Self, String> {
match raw.map(str::trim) {
None | Some("") | Some("fact") => Ok(Self::Fact),
Some("preference") => Ok(Self::Preference),
Some("procedure") => Ok(Self::Procedure),
Some(other) => Err(format!(
"remember: unknown kind '{other}' (expected fact, preference, or procedure)"
)),
}
}
pub fn as_str(self) -> &'static str {
match self {
Self::Fact => "fact",
Self::Preference => "preference",
Self::Procedure => "procedure",
}
}
pub fn from_wire(raw: Option<&str>) -> Self {
Self::parse(raw).unwrap_or_default()
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct SyncedFact {
pub subject: String,
pub body: String,
pub kind: NoteKind,
}
#[async_trait]
pub trait MemorySync: Send + Sync {
async fn append_knowledge(
&self,
subject: &str,
body: &str,
kind: NoteKind,
) -> Result<(), String>;
async fn pull_knowledge(&self) -> Result<Vec<SyncedFact>, String>;
}
pub struct MemoryTools {
inner: Mutex<Inner>,
path: PathBuf,
sync: Option<Arc<dyn MemorySync>>,
}
struct Inner {
engine: MemgineEngine,
notes: Vec<Note>,
}
impl MemoryTools {
pub fn open(path: PathBuf) -> Self {
let notes: Vec<Note> = std::fs::read_to_string(&path)
.ok()
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default();
let mut engine = MemgineEngine::new(None);
for (i, n) in notes.iter().enumerate() {
ingest(&mut engine, i, n);
}
Self {
inner: Mutex::new(Inner { engine, notes }),
path,
sync: None,
}
}
pub fn with_sync(mut self, sync: Option<Arc<dyn MemorySync>>) -> Self {
self.sync = sync;
self
}
pub fn tool_defs() -> Vec<Value> {
vec![
json!({
"name": "remember",
"description": "Save a durable fact about the user or task so you can recall it in \
future sessions. Use for things worth persisting, not transient \
details. `kind` decides how the fact is used later, so pick it \
deliberately.",
"parameters": {
"type": "object",
"properties": {
"subject": { "type": "string", "description": "Short label for the fact (e.g. 'project name')." },
"body": { "type": "string", "description": "The fact to remember." },
"kind": {
"type": "string",
"enum": ["fact", "preference", "procedure"],
"description": "'fact' (default): something true about the user, \
project, or environment; retrieved when relevant to \
the query. 'preference': a standing instruction about \
how the user wants you to work — surfaced in EVERY \
future session, so reserve it for rules you should \
never violate. 'procedure': how a task was done and \
whether it worked, for reuse next time."
}
},
"required": ["subject", "body"]
},
"mutating": true,
"tier": "full_access"
}),
json!({
"name": "recall",
"description": "Retrieve previously remembered facts relevant to a query. Use at the \
start of a task to recall what you know about the user or project.",
"parameters": {
"type": "object",
"properties": {
"query": { "type": "string", "description": "What to recall." }
},
"required": ["query"]
}
}),
]
}
fn remember(&self, params: &Value) -> Result<(Value, NoteKind), String> {
let subject = normalize_subject(
params
.get("subject")
.and_then(Value::as_str)
.ok_or("remember requires a 'subject' string")?,
);
let body = params
.get("body")
.and_then(Value::as_str)
.ok_or("remember requires a 'body' string")?;
let kind = NoteKind::parse(params.get("kind").and_then(Value::as_str))?;
let mut g = self.inner.lock().map_err(|_| "memory lock poisoned")?;
let note = Note {
subject: subject.to_string(),
body: body.to_string(),
kind,
};
match g
.notes
.iter_mut()
.find(|n| n.subject.eq_ignore_ascii_case(subject))
{
Some(existing) => {
existing.body = body.to_string();
existing.kind = kind;
}
None => g.notes.push(note),
}
let mut engine = MemgineEngine::new(None);
for (i, n) in g.notes.iter().enumerate() {
ingest(&mut engine, i, n);
}
g.engine = engine;
if let Some(parent) = self.path.parent() {
let _ = std::fs::create_dir_all(parent);
}
let serialized =
serde_json::to_string_pretty(&g.notes).map_err(|e| format!("serialize: {e}"))?;
std::fs::write(&self.path, serialized).map_err(|e| format!("persist memory: {e}"))?;
Ok((
json!({ "remembered": subject, "total_facts": g.notes.len() }),
kind,
))
}
async fn mirror_to_sync(&self, params: &Value, kind: NoteKind) {
let Some(sync) = &self.sync else {
return;
};
let (Some(subject), Some(body)) = (
params.get("subject").and_then(Value::as_str),
params.get("body").and_then(Value::as_str),
) else {
return;
};
if let Err(e) = sync
.append_knowledge(normalize_subject(subject), body, kind)
.await
{
tracing::debug!(
target: "assistant.memory",
"knowledge sync append failed (non-fatal): {e}"
);
}
}
async fn pull_peer_knowledge(&self) -> Vec<SyncedFact> {
let Some(sync) = &self.sync else {
return Vec::new();
};
match sync.pull_knowledge().await {
Ok(facts) => facts,
Err(e) => {
tracing::debug!(
target: "assistant.memory",
"peer knowledge pull failed (non-fatal): {e}"
);
Vec::new()
}
}
}
async fn recall(&self, params: &Value) -> Result<Value, String> {
let query = params
.get("query")
.and_then(Value::as_str)
.ok_or("recall requires a 'query' string")?;
let peers = self.pull_peer_knowledge().await;
let mut g = self.inner.lock().map_err(|_| "memory lock poisoned")?;
if g.notes.is_empty() && peers.is_empty() {
return Ok(json!({ "query": query, "context": "", "note": "no facts remembered yet" }));
}
let local_keys: std::collections::HashSet<String> =
g.notes.iter().map(|n| subject_key(&n.subject)).collect();
let gaps: Vec<&SyncedFact> = peers
.iter()
.filter(|fact| !local_keys.contains(&subject_key(&fact.subject)))
.collect();
if gaps.is_empty() {
let context = g.engine.build_context(query);
return Ok(json!({ "query": query, "context": context }));
}
let mut engine = MemgineEngine::new(None);
let mut idx = 0;
for n in &g.notes {
ingest(&mut engine, idx, n);
idx += 1;
}
for fact in gaps {
ingest(
&mut engine,
idx,
&Note {
subject: fact.subject.clone(),
body: fact.body.clone(),
kind: fact.kind,
},
);
idx += 1;
}
let context = engine.build_context(query);
Ok(json!({ "query": query, "context": context }))
}
pub async fn proactive_intervention(
&self,
query: &str,
recent: Vec<String>,
events: &[Event],
) -> Result<
(
car_memgine::ProactiveMaintenanceReport,
ProactiveMemoryDecision,
),
String,
> {
let peers = self.pull_peer_knowledge().await;
let mut g = self.inner.lock().map_err(|_| "memory lock poisoned")?;
let local_keys: std::collections::HashSet<String> =
g.notes.iter().map(|n| subject_key(&n.subject)).collect();
let gaps: Vec<&SyncedFact> = peers
.iter()
.filter(|fact| !local_keys.contains(&subject_key(&fact.subject)))
.collect();
let maintenance_request = ProactiveMaintenanceRequest {
max_recent: 32,
tenant_id: None,
};
let mut request = ProactiveMemoryRequest {
query: query.to_string(),
recent,
..Default::default()
};
if gaps.is_empty() {
let maintenance = g
.engine
.maintain_proactive_memory_from_events(events, &maintenance_request);
request.trigger.merge(maintenance.trigger.clone());
let decision = g.engine.proactive_intervention(&request);
return Ok((maintenance, decision));
}
let mut engine = MemgineEngine::new(None);
let mut idx = 0;
for n in &g.notes {
ingest(&mut engine, idx, n);
idx += 1;
}
for fact in gaps {
ingest(
&mut engine,
idx,
&Note {
subject: fact.subject.clone(),
body: fact.body.clone(),
kind: fact.kind,
},
);
idx += 1;
}
let maintenance =
engine.maintain_proactive_memory_from_events(events, &maintenance_request);
request.trigger.merge(maintenance.trigger.clone());
let decision = engine.proactive_intervention(&request);
Ok((maintenance, decision))
}
}
fn normalize_subject(subject: &str) -> &str {
subject.trim()
}
fn subject_key(subject: &str) -> String {
normalize_subject(subject).to_ascii_lowercase()
}
fn ingest(engine: &mut MemgineEngine, idx: usize, n: &Note) {
let save = ProactiveMemorySave {
id: Some(format!("assistant-note-{idx}")),
subject: n.subject.clone(),
body: n.body.clone(),
tags: vec!["assistant_memory".to_string()],
confidence: Some("high".to_string()),
tenant_id: None,
is_constraint: n.kind == NoteKind::Preference,
};
match n.kind {
NoteKind::Procedure => engine.save_proactive_procedural(save),
NoteKind::Fact | NoteKind::Preference => engine.save_proactive_knowledge(save),
};
}
#[async_trait]
impl ToolExecutor for MemoryTools {
async fn execute(&self, tool: &str, params: &Value) -> Result<Value, String> {
match tool {
"remember" => {
let (out, kind) = self.remember(params)?;
self.mirror_to_sync(params, kind).await;
Ok(out)
}
"recall" => self.recall(params).await,
other => Err(format!("unknown tool: '{other}'")),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn remember_declares_mutating_and_full_access() {
let defs = MemoryTools::tool_defs();
let remember = defs
.iter()
.find(|def| def["name"] == "remember")
.expect("remember tool def");
assert_eq!(remember["mutating"], true);
assert_eq!(remember["tier"], "full_access");
}
#[test]
fn remember_advertises_the_memory_taxonomy() {
let defs = MemoryTools::tool_defs();
let remember = defs.iter().find(|def| def["name"] == "remember").unwrap();
let kinds = remember["parameters"]["properties"]["kind"]["enum"]
.as_array()
.expect("kind is an enum — the schema is what teaches the model");
for expected in ["fact", "preference", "procedure"] {
assert!(kinds.iter().any(|k| k == expected), "missing {expected}");
}
let required = remember["parameters"]["required"].as_array().unwrap();
assert!(!required.iter().any(|r| r == "kind"));
}
#[tokio::test]
async fn preference_becomes_an_always_included_constraint() {
let dir = tempfile::tempdir().unwrap();
let mem = MemoryTools::open(dir.path().join("assistant.json"));
mem.execute(
"remember",
&json!({
"subject": "code review",
"body": "Never merge without a green CI run.",
"kind": "preference"
}),
)
.await
.unwrap();
mem.execute(
"remember",
&json!({ "subject": "project name", "body": "The project is called Zephyr." }),
)
.await
.unwrap();
let out = mem
.execute("recall", &json!({ "query": "what time is the standup" }))
.await
.unwrap();
let context = out["context"].as_str().unwrap();
assert!(
context.contains("Active Constraints") && context.contains("green CI run"),
"preference must land in the always-included constraints layer: {context}"
);
}
#[tokio::test]
async fn unknown_kind_is_rejected_rather_than_silently_stored() {
let dir = tempfile::tempdir().unwrap();
let mem = MemoryTools::open(dir.path().join("assistant.json"));
let err = mem
.execute(
"remember",
&json!({ "subject": "s", "body": "b", "kind": "identity" }),
)
.await
.unwrap_err();
assert!(err.contains("unknown kind 'identity'"), "{err}");
}
#[tokio::test]
async fn remembers_across_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("assistant.json");
{
let mem = MemoryTools::open(path.clone());
mem.execute(
"remember",
&json!({ "subject": "project name", "body": "The project is called Zephyr." }),
)
.await
.unwrap();
}
let mem = MemoryTools::open(path.clone());
let out = mem
.execute("recall", &json!({ "query": "what is my project called?" }))
.await
.unwrap();
let ctx = out["context"].as_str().unwrap();
assert!(
ctx.contains("Zephyr"),
"recall should surface the fact: {ctx}"
);
}
#[tokio::test]
async fn remember_supersedes_same_subject() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("m.json");
let mem = MemoryTools::open(path.clone());
mem.execute("remember", &json!({ "subject": "editor", "body": "vim" }))
.await
.unwrap();
let out = mem
.execute("remember", &json!({ "subject": "Editor", "body": "emacs" }))
.await
.unwrap();
assert_eq!(out["total_facts"], 1);
let notes: Vec<Note> =
serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(notes.len(), 1);
assert_eq!(notes[0].body, "emacs");
}
#[tokio::test]
async fn recall_on_empty_is_graceful() {
let dir = tempfile::tempdir().unwrap();
let mem = MemoryTools::open(dir.path().join("m.json"));
let out = mem
.execute("recall", &json!({ "query": "anything" }))
.await
.unwrap();
assert_eq!(out["context"], "");
}
#[tokio::test]
async fn recall_ranks_by_relevance_not_creation_order() {
let dir = tempfile::tempdir().unwrap();
let mem = MemoryTools::open(dir.path().join("m.json"));
for (subject, body) in [
("deployment target", "Unclear which region we deploy to."),
("db", "Postgres 16 in prod."),
("owner", "Matt owns the release process."),
] {
mem.execute("remember", &json!({ "subject": subject, "body": body }))
.await
.unwrap();
}
let out = mem
.execute("recall", &json!({ "query": "deployment target" }))
.await
.unwrap();
let ctx = out["context"].as_str().unwrap_or_default();
let facts: Vec<&str> = ctx
.lines()
.skip_while(|l| !l.starts_with("## Current Facts"))
.filter(|l| l.starts_with("- "))
.collect();
assert!(facts.len() >= 3, "expected all facts back, got: {ctx}");
assert!(
facts.last().unwrap().contains("deployment target"),
"the query-matching fact must rank last (relevance-ascending); \
got creation order, which means Fast mode is running.\nfacts: {facts:#?}"
);
}
struct RecordingSync(Mutex<Vec<SyncedFact>>);
#[async_trait]
impl MemorySync for RecordingSync {
async fn append_knowledge(
&self,
subject: &str,
body: &str,
kind: NoteKind,
) -> Result<(), String> {
self.0.lock().unwrap().push(SyncedFact {
subject: subject.into(),
body: body.into(),
kind,
});
Ok(())
}
async fn pull_knowledge(&self) -> Result<Vec<SyncedFact>, String> {
Ok(Vec::new())
}
}
struct PeerSync(Vec<SyncedFact>);
#[async_trait]
impl MemorySync for PeerSync {
async fn append_knowledge(&self, _: &str, _: &str, _: NoteKind) -> Result<(), String> {
Ok(())
}
async fn pull_knowledge(&self) -> Result<Vec<SyncedFact>, String> {
Ok(self.0.clone())
}
}
fn peer_fact(subject: &str, body: &str, kind: NoteKind) -> SyncedFact {
SyncedFact {
subject: subject.into(),
body: body.into(),
kind,
}
}
#[tokio::test]
async fn remember_mirrors_normalized_subject_and_body_to_sync() {
let dir = tempfile::tempdir().unwrap();
let rec = Arc::new(RecordingSync(Mutex::new(Vec::new())));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(rec.clone()));
mem.execute(
"remember",
&json!({ "subject": " Editor ", "body": "vim" }),
)
.await
.unwrap();
let got = rec.0.lock().unwrap().clone();
assert_eq!(got, vec![peer_fact("Editor", "vim", NoteKind::Fact)]);
}
#[tokio::test]
async fn remember_without_sync_stays_local_only() {
let dir = tempfile::tempdir().unwrap();
let mem = MemoryTools::open(dir.path().join("m.json"));
let out = mem
.execute("remember", &json!({ "subject": "x", "body": "y" }))
.await
.unwrap();
assert_eq!(out["remembered"], "x");
}
struct FailingSync;
#[async_trait]
impl MemorySync for FailingSync {
async fn append_knowledge(&self, _: &str, _: &str, _: NoteKind) -> Result<(), String> {
Err("oplog unreachable".into())
}
async fn pull_knowledge(&self) -> Result<Vec<SyncedFact>, String> {
Err("oplog unreachable".into())
}
}
#[tokio::test]
async fn recall_surfaces_a_peer_fact_not_held_locally() {
let dir = tempfile::tempdir().unwrap();
let peer = Arc::new(PeerSync(vec![peer_fact(
"project",
"the project is Zephyr",
NoteKind::Fact,
)]));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(peer));
let out = mem
.execute("recall", &json!({ "query": "what is my project?" }))
.await
.unwrap();
assert!(
out["context"].as_str().unwrap().contains("Zephyr"),
"{}",
out["context"]
);
}
#[tokio::test]
async fn a_peer_preference_stays_a_standing_rule_and_a_peer_fact_does_not() {
let rule = "Never merge without a green CI run.";
let unrelated = "what time is the standup";
let recall_with = |kind| async move {
let dir = tempfile::tempdir().unwrap();
let peer = Arc::new(PeerSync(vec![peer_fact("code review", rule, kind)]));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(peer));
let out = mem
.execute("recall", &json!({ "query": unrelated }))
.await
.unwrap();
out["context"].as_str().unwrap().to_string()
};
let as_preference = recall_with(NoteKind::Preference).await;
assert!(
as_preference.contains("Active Constraints") && as_preference.contains(rule),
"a synced preference must land in the always-included constraints \
layer, exactly as it does on the device that set it: {as_preference}"
);
let as_fact = recall_with(NoteKind::Fact).await;
assert!(
!as_fact.contains("Active Constraints"),
"control: a plain fact is not a standing rule — if it reaches the \
constraints layer the test cannot tell the two kinds apart: {as_fact}"
);
}
#[tokio::test]
async fn remember_mirrors_the_kind_so_a_preference_crosses_the_wire() {
let dir = tempfile::tempdir().unwrap();
let rec = Arc::new(RecordingSync(Mutex::new(Vec::new())));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(rec.clone()));
mem.execute(
"remember",
&json!({
"subject": "code review",
"body": "Never merge without a green CI run.",
"kind": "preference"
}),
)
.await
.unwrap();
assert_eq!(
rec.0.lock().unwrap().clone(),
vec![peer_fact(
"code review",
"Never merge without a green CI run.",
NoteKind::Preference
)],
"the mirror must send the kind `remember` validated, not a default"
);
}
#[tokio::test]
async fn recall_keeps_local_over_a_peer_edit_of_the_same_subject() {
let dir = tempfile::tempdir().unwrap();
let peer = Arc::new(PeerSync(vec![peer_fact("editor", "emacs", NoteKind::Fact)]));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(peer));
mem.execute("remember", &json!({ "subject": "editor", "body": "vim" }))
.await
.unwrap();
let out = mem
.execute("recall", &json!({ "query": "which editor?" }))
.await
.unwrap();
let ctx = out["context"].as_str().unwrap().to_string();
assert!(ctx.contains("vim"), "{ctx}");
assert!(
!ctx.contains("emacs"),
"peer edit must not override local under Option B: {ctx}"
);
}
#[tokio::test]
async fn recall_peer_gate_is_case_insensitive() {
let dir = tempfile::tempdir().unwrap();
let peer = Arc::new(PeerSync(vec![peer_fact(
"pet",
"a peer dog",
NoteKind::Fact,
)]));
let mem = MemoryTools::open(dir.path().join("m.json")).with_sync(Some(peer));
mem.execute(
"remember",
&json!({ "subject": "Pet", "body": "my cat Mittens" }),
)
.await
.unwrap();
let out = mem
.execute("recall", &json!({ "query": "what pet?" }))
.await
.unwrap();
let ctx = out["context"].as_str().unwrap().to_string();
assert!(ctx.contains("Mittens"), "{ctx}");
assert!(
!ctx.contains("peer dog"),
"a case-variant peer subject must be gated: {ctx}"
);
}
#[tokio::test]
async fn recall_pull_failure_degrades_to_local() {
let dir = tempfile::tempdir().unwrap();
let mem =
MemoryTools::open(dir.path().join("m.json")).with_sync(Some(Arc::new(FailingSync)));
mem.execute("remember", &json!({ "subject": "city", "body": "Denver" }))
.await
.unwrap();
let out = mem
.execute("recall", &json!({ "query": "which city?" }))
.await
.unwrap();
assert!(out["context"].as_str().unwrap().contains("Denver"));
}
#[tokio::test]
async fn remember_survives_sync_failure() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("m.json");
let mem = MemoryTools::open(path.clone()).with_sync(Some(Arc::new(FailingSync)));
let out = mem
.execute("remember", &json!({ "subject": "a", "body": "b" }))
.await
.unwrap();
assert_eq!(out["remembered"], "a");
let notes: Vec<Note> =
serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(notes.len(), 1);
assert_eq!(notes[0].body, "b");
}
#[test]
fn synced_knowledge_survives_edits_and_converges_to_newest_body() {
use car_sync::fold::fold;
use car_sync::oplog::{DeviceLog, Scope, Surface};
let mut log = DeviceLog::new("device-a");
let vim = log.append(
Scope::Personal,
Surface::Knowledge,
json!({ "subject": "editor", "body": "vim" }),
);
let emacs = log.append(
Scope::Personal,
Surface::Knowledge,
json!({ "subject": "editor", "body": "emacs" }),
);
let state = fold(&[vim, emacs]);
let entries = state.log_entries("knowledge");
assert_eq!(entries.len(), 2);
let newest = entries
.iter()
.filter(|r| r.payload["subject"] == "editor")
.max_by(|a, b| a.hlc.cmp(&b.hlc))
.unwrap();
assert_eq!(newest.payload["body"], "emacs");
}
#[test]
fn identical_reremember_dedups_in_the_synced_fold() {
use car_sync::fold::fold;
use car_sync::oplog::{DeviceLog, Scope, Surface};
let mut log = DeviceLog::new("device-a");
let a = log.append(
Scope::Personal,
Surface::Knowledge,
json!({ "subject": "editor", "body": "vim" }),
);
let b = log.append(
Scope::Personal,
Surface::Knowledge,
json!({ "subject": "editor", "body": "vim" }),
);
let state = fold(&[a, b]);
assert_eq!(state.log_entries("knowledge").len(), 1);
}
}