1pub const CREATE_MEMORIES_TABLE: &str = "
2CREATE TABLE IF NOT EXISTS memories (
3 id VARCHAR PRIMARY KEY,
4 agent_id VARCHAR NOT NULL,
5 content TEXT NOT NULL,
6 memory_type VARCHAR NOT NULL,
7 scope VARCHAR NOT NULL DEFAULT 'private',
8 importance FLOAT NOT NULL DEFAULT 0.5,
9 tags JSON,
10 metadata JSON,
11 embedding BLOB,
12 content_hash BLOB NOT NULL,
13 prev_hash BLOB,
14 source_type VARCHAR NOT NULL DEFAULT 'agent',
15 source_id VARCHAR,
16 consolidation_state VARCHAR NOT NULL DEFAULT 'raw',
17 access_count BIGINT NOT NULL DEFAULT 0,
18 org_id VARCHAR,
19 thread_id VARCHAR,
20 created_at VARCHAR NOT NULL,
21 updated_at VARCHAR NOT NULL,
22 last_accessed_at VARCHAR,
23 expires_at VARCHAR,
24 deleted_at VARCHAR,
25 decay_rate FLOAT,
26 created_by VARCHAR,
27 version INTEGER NOT NULL DEFAULT 1,
28 prev_version_id VARCHAR,
29 quarantined BOOLEAN NOT NULL DEFAULT false,
30 quarantine_reason VARCHAR,
31 decay_function VARCHAR
32);
33CREATE INDEX IF NOT EXISTS idx_memories_agent_id ON memories(agent_id);
34CREATE INDEX IF NOT EXISTS idx_memories_scope ON memories(scope);
35CREATE INDEX IF NOT EXISTS idx_memories_memory_type ON memories(memory_type);
36CREATE INDEX IF NOT EXISTS idx_memories_org_id ON memories(org_id);
37CREATE INDEX IF NOT EXISTS idx_memories_created_at ON memories(created_at);
38CREATE INDEX IF NOT EXISTS idx_memories_deleted_at ON memories(deleted_at);
39CREATE INDEX IF NOT EXISTS idx_memories_thread_id ON memories(thread_id);
40CREATE INDEX IF NOT EXISTS idx_memories_expires_at ON memories(expires_at);
41CREATE INDEX IF NOT EXISTS idx_memories_consolidation_state ON memories(consolidation_state);
42";
43
44pub const CREATE_ACLS_TABLE: &str = "
45CREATE TABLE IF NOT EXISTS acls (
46 id VARCHAR PRIMARY KEY,
47 memory_id VARCHAR NOT NULL,
48 principal_type VARCHAR NOT NULL,
49 principal_id VARCHAR NOT NULL,
50 permission VARCHAR NOT NULL,
51 granted_by VARCHAR NOT NULL,
52 created_at VARCHAR NOT NULL,
53 expires_at VARCHAR
54);
55CREATE INDEX IF NOT EXISTS idx_acls_memory_id ON acls(memory_id);
56CREATE INDEX IF NOT EXISTS idx_acls_principal ON acls(principal_type, principal_id);
57";
58
59pub const CREATE_RELATIONS_TABLE: &str = "
60CREATE TABLE IF NOT EXISTS relations (
61 id VARCHAR PRIMARY KEY,
62 source_id VARCHAR NOT NULL,
63 target_id VARCHAR NOT NULL,
64 relation_type VARCHAR NOT NULL,
65 weight FLOAT NOT NULL DEFAULT 1.0,
66 metadata JSON,
67 created_at VARCHAR NOT NULL
68);
69CREATE INDEX IF NOT EXISTS idx_relations_source ON relations(source_id);
70CREATE INDEX IF NOT EXISTS idx_relations_target ON relations(target_id);
71";
72
73pub const CREATE_AGENT_EVENTS_TABLE: &str = "
79CREATE TABLE IF NOT EXISTS agent_events (
80 id VARCHAR PRIMARY KEY,
81 agent_id VARCHAR NOT NULL,
82 thread_id VARCHAR,
83 run_id VARCHAR,
84 parent_event_id VARCHAR,
85 event_type VARCHAR NOT NULL,
86 payload JSON,
87 trace_id VARCHAR,
88 span_id VARCHAR,
89 model VARCHAR,
90 tokens_input BIGINT,
91 tokens_output BIGINT,
92 latency_ms BIGINT,
93 cost_usd DOUBLE,
94 timestamp VARCHAR NOT NULL,
95 logical_clock BIGINT NOT NULL DEFAULT 0,
96 content_hash BLOB NOT NULL,
97 prev_hash BLOB,
98 embedding BLOB
99);
100CREATE INDEX IF NOT EXISTS idx_events_agent_id ON agent_events(agent_id);
101CREATE INDEX IF NOT EXISTS idx_events_thread_id ON agent_events(thread_id);
102CREATE INDEX IF NOT EXISTS idx_events_event_type ON agent_events(event_type);
103CREATE INDEX IF NOT EXISTS idx_events_timestamp ON agent_events(timestamp);
104CREATE INDEX IF NOT EXISTS idx_events_trace_id ON agent_events(trace_id);
105CREATE INDEX IF NOT EXISTS idx_events_parent ON agent_events(parent_event_id);
106";
107
108pub const CREATE_CHECKPOINTS_TABLE: &str = "
109CREATE TABLE IF NOT EXISTS checkpoints (
110 id VARCHAR PRIMARY KEY,
111 thread_id VARCHAR NOT NULL,
112 agent_id VARCHAR NOT NULL,
113 parent_id VARCHAR,
114 branch_name VARCHAR NOT NULL DEFAULT 'main',
115 state_snapshot JSON,
116 state_diff JSON,
117 memory_refs JSON,
118 event_cursor VARCHAR,
119 label VARCHAR,
120 created_at VARCHAR NOT NULL,
121 metadata JSON
122);
123CREATE INDEX IF NOT EXISTS idx_checkpoints_thread_id ON checkpoints(thread_id);
124CREATE INDEX IF NOT EXISTS idx_checkpoints_branch ON checkpoints(thread_id, branch_name);
125CREATE INDEX IF NOT EXISTS idx_checkpoints_agent ON checkpoints(agent_id);
126CREATE INDEX IF NOT EXISTS idx_checkpoints_created_at ON checkpoints(created_at);
127";
128
129pub const SPRINT3_COLUMN_ALTERS: &[&str] = &[
134 "ALTER TABLE memories ADD COLUMN decay_rate FLOAT",
135 "ALTER TABLE memories ADD COLUMN created_by VARCHAR",
136 "ALTER TABLE memories ADD COLUMN version INTEGER DEFAULT 1",
137 "ALTER TABLE memories ADD COLUMN prev_version_id VARCHAR",
138 "ALTER TABLE memories ADD COLUMN quarantined BOOLEAN DEFAULT false",
139 "ALTER TABLE memories ADD COLUMN quarantine_reason VARCHAR",
140];
141
142pub const SPRINT4_COLUMN_ALTERS: &[&str] = &[
144 "ALTER TABLE agent_events ADD COLUMN embedding BLOB",
145 "ALTER TABLE memories ADD COLUMN decay_function VARCHAR",
146];
147
148pub const CREATE_DELEGATIONS_TABLE: &str = "
149CREATE TABLE IF NOT EXISTS delegations (
150 id VARCHAR PRIMARY KEY,
151 delegator_id VARCHAR NOT NULL,
152 delegate_id VARCHAR NOT NULL,
153 permission VARCHAR NOT NULL,
154 scope_type VARCHAR NOT NULL DEFAULT 'all_memories',
155 scope_value JSON,
156 max_depth INTEGER NOT NULL DEFAULT 0,
157 current_depth INTEGER NOT NULL DEFAULT 0,
158 parent_delegation_id VARCHAR,
159 created_at VARCHAR NOT NULL,
160 expires_at VARCHAR,
161 revoked_at VARCHAR
162);
163CREATE INDEX IF NOT EXISTS idx_delegations_delegator ON delegations(delegator_id);
164CREATE INDEX IF NOT EXISTS idx_delegations_delegate ON delegations(delegate_id);
165";
166
167pub const CREATE_AGENT_PROFILES_TABLE: &str = "
168CREATE TABLE IF NOT EXISTS agent_profiles (
169 agent_id VARCHAR PRIMARY KEY,
170 avg_importance DOUBLE NOT NULL DEFAULT 0.5,
171 avg_content_length DOUBLE NOT NULL DEFAULT 100,
172 total_memories BIGINT NOT NULL DEFAULT 0,
173 last_updated VARCHAR NOT NULL
174);
175";
176
177pub const CREATE_SYNC_METADATA_TABLE: &str = "
178CREATE TABLE IF NOT EXISTS sync_metadata (
179 key VARCHAR PRIMARY KEY,
180 value VARCHAR NOT NULL,
181 updated_at VARCHAR NOT NULL DEFAULT CURRENT_TIMESTAMP
182);
183";
184
185pub const CREATE_MNEMO_META_TABLE: &str = "
189CREATE TABLE IF NOT EXISTS mnemo_meta (
190 key VARCHAR PRIMARY KEY,
191 value VARCHAR NOT NULL,
192 updated_at VARCHAR NOT NULL DEFAULT CURRENT_TIMESTAMP
193);
194";
195
196pub const CREATE_EMBEDDING_BASELINE_TABLE: &str = "
201CREATE TABLE IF NOT EXISTS embedding_baseline (
202 agent_id VARCHAR PRIMARY KEY,
203 mu JSON NOT NULL,
204 cov_diag JSON NOT NULL,
205 n BIGINT NOT NULL,
206 updated_at VARCHAR NOT NULL
207);
208";
209
210pub const CURRENT_PERSISTENCE_VERSION: u32 = 4;
213
214pub fn run_migrations(conn: &duckdb::Connection) -> duckdb::Result<()> {
215 conn.execute_batch(CREATE_MEMORIES_TABLE)?;
216 conn.execute_batch(CREATE_ACLS_TABLE)?;
217 conn.execute_batch(CREATE_RELATIONS_TABLE)?;
218 conn.execute_batch(CREATE_AGENT_EVENTS_TABLE)?;
219 conn.execute_batch(CREATE_CHECKPOINTS_TABLE)?;
220 apply_alters_idempotent(conn, SPRINT3_COLUMN_ALTERS)?;
227 conn.execute_batch(CREATE_DELEGATIONS_TABLE)?;
228 conn.execute_batch(CREATE_AGENT_PROFILES_TABLE)?;
229 apply_alters_idempotent(conn, SPRINT4_COLUMN_ALTERS)?;
231 conn.execute(
234 "CREATE INDEX IF NOT EXISTS idx_events_parent ON agent_events(parent_event_id)",
235 [],
236 )?;
237 conn.execute_batch(CREATE_SYNC_METADATA_TABLE)?;
239 conn.execute_batch(CREATE_MNEMO_META_TABLE)?;
241 conn.execute_batch(CREATE_EMBEDDING_BASELINE_TABLE)?;
243 stamp_persistence_version(conn)?;
244 Ok(())
245}
246
247fn parse_alter_table_add_column(sql: &str) -> Option<(&str, &str)> {
251 let trimmed = sql.trim();
252 let lower = trimmed.to_ascii_lowercase();
253 let prefix = "alter table ";
254 let rest = lower.strip_prefix(prefix)?;
255 let after_table_keyword_idx = prefix.len();
256 let table_end = rest.find(' ')?;
257 let table = &trimmed[after_table_keyword_idx..after_table_keyword_idx + table_end];
258 let after_table = &lower[after_table_keyword_idx + table_end..];
259 let add_column = " add column ";
260 let after_add = after_table.strip_prefix(add_column)?;
261 let col_end = after_add.find(' ')?;
262 let col_start_in_full = after_table_keyword_idx + table_end + add_column.len();
263 let col = &trimmed[col_start_in_full..col_start_in_full + col_end];
264 Some((table, col))
265}
266
267fn column_exists(conn: &duckdb::Connection, table: &str, column: &str) -> duckdb::Result<bool> {
269 let mut stmt = conn.prepare(
270 "SELECT 1 FROM information_schema.columns \
271 WHERE lower(table_name) = lower(?) AND lower(column_name) = lower(?) LIMIT 1",
272 )?;
273 let mut rows = stmt.query(duckdb::params![table, column])?;
274 Ok(rows.next()?.is_some())
275}
276
277fn apply_alters_idempotent(conn: &duckdb::Connection, alters: &[&str]) -> duckdb::Result<()> {
282 for sql in alters {
283 let Some((table, column)) = parse_alter_table_add_column(sql) else {
284 return Err(duckdb::Error::ToSqlConversionFailure(
285 format!("migration ALTER did not match expected shape: {sql:?}").into(),
286 ));
287 };
288 if column_exists(conn, table, column)? {
289 continue;
290 }
291 conn.execute(sql, [])?;
292 }
293 Ok(())
294}
295
296pub fn read_persistence_version(conn: &duckdb::Connection) -> duckdb::Result<Option<u32>> {
299 let mut stmt =
300 conn.prepare("SELECT value FROM mnemo_meta WHERE key = 'persistence_version'")?;
301 let mut rows = stmt.query([])?;
302 if let Some(row) = rows.next()? {
303 let raw: String = row.get(0)?;
304 Ok(raw.parse::<u32>().ok())
305 } else {
306 Ok(None)
307 }
308}
309
310fn stamp_persistence_version(conn: &duckdb::Connection) -> duckdb::Result<()> {
320 let existing = read_persistence_version(conn)?;
321 let current = CURRENT_PERSISTENCE_VERSION;
322 if let Some(v) = existing
323 && v == current
324 {
325 return Ok(());
326 }
327 let now = chrono::Utc::now().to_rfc3339();
328 conn.execute(
331 "DELETE FROM mnemo_meta WHERE key = 'persistence_version'",
332 [],
333 )?;
334 conn.execute(
335 "INSERT INTO mnemo_meta(key, value, updated_at) VALUES ('persistence_version', ?, ?)",
336 duckdb::params![current.to_string(), now],
337 )?;
338 Ok(())
339}
340
341#[cfg(test)]
342mod tests {
343 use super::*;
344
345 #[test]
346 fn test_fresh_db_stamps_current_persistence_version() {
347 let conn = duckdb::Connection::open_in_memory().unwrap();
348 run_migrations(&conn).unwrap();
349 let v = read_persistence_version(&conn).unwrap();
350 assert_eq!(v, Some(CURRENT_PERSISTENCE_VERSION));
351 }
352
353 #[test]
358 fn test_legacy_db_gets_stamped_on_open() {
359 let conn = duckdb::Connection::open_in_memory().unwrap();
360 conn.execute_batch(CREATE_MEMORIES_TABLE).unwrap();
362 conn.execute_batch(CREATE_ACLS_TABLE).unwrap();
363 conn.execute_batch(CREATE_RELATIONS_TABLE).unwrap();
364 conn.execute_batch(CREATE_AGENT_EVENTS_TABLE).unwrap();
365 conn.execute_batch(CREATE_CHECKPOINTS_TABLE).unwrap();
366 conn.execute_batch(CREATE_DELEGATIONS_TABLE).unwrap();
367 conn.execute_batch(CREATE_AGENT_PROFILES_TABLE).unwrap();
368 conn.execute_batch(CREATE_SYNC_METADATA_TABLE).unwrap();
369
370 assert!(
371 read_persistence_version(&conn).is_err()
372 || read_persistence_version(&conn).unwrap().is_none(),
373 "pre-migration legacy file should have no stamp"
374 );
375
376 run_migrations(&conn).unwrap();
377 assert_eq!(
378 read_persistence_version(&conn).unwrap(),
379 Some(CURRENT_PERSISTENCE_VERSION)
380 );
381
382 run_migrations(&conn).unwrap();
384 assert_eq!(
385 read_persistence_version(&conn).unwrap(),
386 Some(CURRENT_PERSISTENCE_VERSION)
387 );
388 }
389
390 #[test]
391 fn sprint_alters_match_expected_shape() {
392 for sql in SPRINT3_COLUMN_ALTERS
397 .iter()
398 .chain(SPRINT4_COLUMN_ALTERS.iter())
399 {
400 let parsed = parse_alter_table_add_column(sql);
401 assert!(
402 parsed.is_some(),
403 "ALTER does not match `ALTER TABLE <t> ADD COLUMN <c> ...`: {sql:?}"
404 );
405 }
406 }
407
408 #[test]
409 fn alters_are_idempotent_under_duckdb_152() {
410 let conn = duckdb::Connection::open_in_memory().unwrap();
415 run_migrations(&conn).unwrap();
416 run_migrations(&conn).unwrap();
417 let mut stmt = conn.prepare("SELECT COUNT(*) FROM memories").unwrap();
419 let n: i64 = stmt.query_row([], |row| row.get(0)).unwrap();
420 assert_eq!(n, 0);
421 }
422
423 #[test]
424 fn test_migrations_run_on_in_memory_db() {
425 let conn = duckdb::Connection::open_in_memory().unwrap();
426 run_migrations(&conn).unwrap();
427 let mut stmt = conn.prepare("SELECT COUNT(*) FROM memories").unwrap();
429 let count: i64 = stmt.query_row([], |row| row.get(0)).unwrap();
430 assert_eq!(count, 0);
431
432 let mut stmt = conn.prepare("SELECT COUNT(*) FROM agent_events").unwrap();
433 let count: i64 = stmt.query_row([], |row| row.get(0)).unwrap();
434 assert_eq!(count, 0);
435
436 let mut stmt = conn.prepare("SELECT COUNT(*) FROM checkpoints").unwrap();
437 let count: i64 = stmt.query_row([], |row| row.get(0)).unwrap();
438 assert_eq!(count, 0);
439 }
440}