1use mnemo_core::error::{Error, Result};
2
3pub async fn run_migrations(pool: &sqlx::PgPool, dimensions: usize) -> Result<()> {
8 sqlx::query("CREATE EXTENSION IF NOT EXISTS vector")
10 .execute(pool)
11 .await
12 .map_err(|e| Error::Storage(format!("failed to create vector extension: {e}")))?;
13
14 let create_memories = format!(
16 r#"
17CREATE TABLE IF NOT EXISTS memories (
18 id UUID PRIMARY KEY,
19 agent_id VARCHAR NOT NULL,
20 content TEXT NOT NULL,
21 memory_type VARCHAR NOT NULL,
22 scope VARCHAR NOT NULL DEFAULT 'private',
23 importance REAL NOT NULL DEFAULT 0.5,
24 tags TEXT[],
25 metadata JSONB,
26 embedding vector({dimensions}),
27 content_hash BYTEA NOT NULL,
28 prev_hash BYTEA,
29 source_type VARCHAR NOT NULL DEFAULT 'agent',
30 source_id VARCHAR,
31 consolidation_state VARCHAR NOT NULL DEFAULT 'raw',
32 access_count BIGINT NOT NULL DEFAULT 0,
33 org_id VARCHAR,
34 thread_id VARCHAR,
35 created_at VARCHAR NOT NULL,
36 updated_at VARCHAR NOT NULL,
37 last_accessed_at VARCHAR,
38 expires_at VARCHAR,
39 deleted_at VARCHAR,
40 decay_rate REAL,
41 created_by VARCHAR,
42 version INTEGER NOT NULL DEFAULT 1,
43 prev_version_id UUID,
44 quarantined BOOLEAN NOT NULL DEFAULT FALSE,
45 quarantine_reason VARCHAR,
46 decay_function VARCHAR
47)
48"#
49 );
50 sqlx::query(sqlx::AssertSqlSafe(create_memories.as_str()))
51 .execute(pool)
52 .await
53 .map_err(|e| Error::Storage(format!("create memories: {e}")))?;
54
55 sqlx::query(
57 r#"
58CREATE TABLE IF NOT EXISTS acls (
59 id UUID PRIMARY KEY,
60 memory_id UUID NOT NULL,
61 principal_type VARCHAR NOT NULL,
62 principal_id VARCHAR NOT NULL,
63 permission VARCHAR NOT NULL,
64 granted_by VARCHAR NOT NULL,
65 created_at VARCHAR NOT NULL,
66 expires_at VARCHAR
67)
68"#,
69 )
70 .execute(pool)
71 .await
72 .map_err(|e| Error::Storage(format!("create acls: {e}")))?;
73
74 sqlx::query(
76 r#"
77CREATE TABLE IF NOT EXISTS relations (
78 id UUID PRIMARY KEY,
79 source_id UUID NOT NULL,
80 target_id UUID NOT NULL,
81 relation_type VARCHAR NOT NULL,
82 weight REAL NOT NULL DEFAULT 1.0,
83 metadata JSONB,
84 created_at VARCHAR NOT NULL
85)
86"#,
87 )
88 .execute(pool)
89 .await
90 .map_err(|e| Error::Storage(format!("create relations: {e}")))?;
91
92 sqlx::query(
94 r#"
95CREATE TABLE IF NOT EXISTS agent_events (
96 id UUID PRIMARY KEY,
97 agent_id VARCHAR NOT NULL,
98 thread_id VARCHAR,
99 run_id VARCHAR,
100 parent_event_id UUID,
101 event_type VARCHAR NOT NULL,
102 payload JSONB,
103 trace_id VARCHAR,
104 span_id VARCHAR,
105 model VARCHAR,
106 tokens_input BIGINT,
107 tokens_output BIGINT,
108 latency_ms BIGINT,
109 cost_usd DOUBLE PRECISION,
110 "timestamp" VARCHAR NOT NULL,
111 logical_clock BIGINT NOT NULL DEFAULT 0,
112 content_hash BYTEA NOT NULL,
113 prev_hash BYTEA,
114 embedding BYTEA
115)
116"#,
117 )
118 .execute(pool)
119 .await
120 .map_err(|e| Error::Storage(format!("create agent_events: {e}")))?;
121
122 sqlx::query(
124 r#"
125CREATE TABLE IF NOT EXISTS checkpoints (
126 id UUID PRIMARY KEY,
127 thread_id VARCHAR NOT NULL,
128 agent_id VARCHAR NOT NULL,
129 parent_id UUID,
130 branch_name VARCHAR NOT NULL DEFAULT 'main',
131 state_snapshot JSONB,
132 state_diff JSONB,
133 memory_refs TEXT[],
134 event_cursor UUID,
135 label VARCHAR,
136 created_at VARCHAR NOT NULL,
137 metadata JSONB
138)
139"#,
140 )
141 .execute(pool)
142 .await
143 .map_err(|e| Error::Storage(format!("create checkpoints: {e}")))?;
144
145 sqlx::query(
147 r#"
148CREATE TABLE IF NOT EXISTS delegations (
149 id UUID PRIMARY KEY,
150 delegator_id VARCHAR NOT NULL,
151 delegate_id VARCHAR NOT NULL,
152 permission VARCHAR NOT NULL,
153 scope_type VARCHAR NOT NULL DEFAULT 'all_memories',
154 scope_value JSONB,
155 max_depth INTEGER NOT NULL DEFAULT 0,
156 current_depth INTEGER NOT NULL DEFAULT 0,
157 parent_delegation_id UUID,
158 created_at VARCHAR NOT NULL,
159 expires_at VARCHAR,
160 revoked_at VARCHAR
161)
162"#,
163 )
164 .execute(pool)
165 .await
166 .map_err(|e| Error::Storage(format!("create delegations: {e}")))?;
167
168 sqlx::query(
170 r#"
171CREATE TABLE IF NOT EXISTS agent_profiles (
172 agent_id VARCHAR PRIMARY KEY,
173 avg_importance DOUBLE PRECISION NOT NULL DEFAULT 0.5,
174 avg_content_length DOUBLE PRECISION NOT NULL DEFAULT 100,
175 total_memories BIGINT NOT NULL DEFAULT 0,
176 last_updated VARCHAR NOT NULL
177)
178"#,
179 )
180 .execute(pool)
181 .await
182 .map_err(|e| Error::Storage(format!("create agent_profiles: {e}")))?;
183
184 sqlx::query(
186 r#"
187CREATE TABLE IF NOT EXISTS sync_metadata (
188 key VARCHAR PRIMARY KEY,
189 value TEXT NOT NULL,
190 updated_at VARCHAR NOT NULL
191)
192"#,
193 )
194 .execute(pool)
195 .await
196 .map_err(|e| Error::Storage(format!("create sync_metadata: {e}")))?;
197
198 sqlx::query(
200 r#"
201CREATE TABLE IF NOT EXISTS embedding_baseline (
202 agent_id VARCHAR PRIMARY KEY,
203 mu JSONB NOT NULL,
204 cov_diag JSONB NOT NULL,
205 n BIGINT NOT NULL,
206 updated_at VARCHAR NOT NULL
207)
208"#,
209 )
210 .execute(pool)
211 .await
212 .map_err(|e| Error::Storage(format!("create embedding_baseline: {e}")))?;
213
214 sqlx::query(
217 r#"
218CREATE TABLE IF NOT EXISTS write_provenance (
219 id UUID PRIMARY KEY,
220 memory_id UUID NOT NULL,
221 principal VARCHAR NOT NULL,
222 capability_id UUID,
223 session_id VARCHAR,
224 op VARCHAR NOT NULL,
225 authored_at VARCHAR NOT NULL,
226 flags VARCHAR,
227 prev_hash BYTEA,
228 content_hash BYTEA NOT NULL
229)
230"#,
231 )
232 .execute(pool)
233 .await
234 .map_err(|e| Error::Storage(format!("create write_provenance: {e}")))?;
235
236 sqlx::query("ALTER TABLE write_provenance ADD COLUMN IF NOT EXISTS flags VARCHAR")
241 .execute(pool)
242 .await
243 .map_err(|e| Error::Storage(format!("add write_provenance.flags: {e}")))?;
244
245 let index_stmts: &[&str] = &[
249 "CREATE INDEX IF NOT EXISTS idx_memories_agent ON memories(agent_id)",
250 "CREATE INDEX IF NOT EXISTS idx_memories_thread ON memories(agent_id, thread_id)",
251 "CREATE INDEX IF NOT EXISTS idx_acls_memory ON acls(memory_id)",
252 "CREATE INDEX IF NOT EXISTS idx_acls_principal ON acls(principal_id)",
253 "CREATE INDEX IF NOT EXISTS idx_relations_source ON relations(source_id)",
254 "CREATE INDEX IF NOT EXISTS idx_relations_target ON relations(target_id)",
255 "CREATE INDEX IF NOT EXISTS idx_events_agent ON agent_events(agent_id)",
256 "CREATE INDEX IF NOT EXISTS idx_events_thread ON agent_events(thread_id)",
257 "CREATE INDEX IF NOT EXISTS idx_events_parent ON agent_events(parent_event_id)",
258 "CREATE INDEX IF NOT EXISTS idx_checkpoints_thread ON checkpoints(thread_id, branch_name)",
259 "CREATE INDEX IF NOT EXISTS idx_delegations_delegator ON delegations(delegator_id)",
260 "CREATE INDEX IF NOT EXISTS idx_delegations_delegate ON delegations(delegate_id)",
261 "CREATE INDEX IF NOT EXISTS idx_write_provenance_memory ON write_provenance(memory_id)",
262 "CREATE INDEX IF NOT EXISTS idx_write_provenance_principal ON write_provenance(principal)",
263 "CREATE INDEX IF NOT EXISTS idx_write_provenance_session ON write_provenance(session_id)",
264 ];
265
266 for stmt in index_stmts {
267 sqlx::query(sqlx::AssertSqlSafe(*stmt))
268 .execute(pool)
269 .await
270 .map_err(|e| Error::Storage(format!("create index: {e}")))?;
271 }
272
273 sqlx::query(
275 r#"
276CREATE OR REPLACE FUNCTION prevent_event_modification() RETURNS trigger AS $$
277BEGIN
278 RAISE EXCEPTION 'agent_events is append-only: % not allowed', TG_OP;
279 RETURN NULL;
280END;
281$$ LANGUAGE plpgsql
282"#,
283 )
284 .execute(pool)
285 .await
286 .map_err(|e| Error::Storage(format!("create append-only function: {e}")))?;
287
288 sqlx::query(
289 r#"
290DO $$ BEGIN
291 IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'enforce_append_only_events') THEN
292 CREATE TRIGGER enforce_append_only_events
293 BEFORE UPDATE OR DELETE ON agent_events
294 FOR EACH ROW EXECUTE FUNCTION prevent_event_modification();
295 END IF;
296END $$
297"#,
298 )
299 .execute(pool)
300 .await
301 .map_err(|e| Error::Storage(format!("create append-only trigger: {e}")))?;
302
303 let hnsw_sql = r#"
307DO $$
308BEGIN
309 IF NOT EXISTS (
310 SELECT 1 FROM pg_indexes WHERE indexname = 'idx_memories_embedding_hnsw'
311 ) THEN
312 CREATE INDEX idx_memories_embedding_hnsw ON memories USING hnsw (embedding vector_cosine_ops);
313 END IF;
314END
315$$
316"#;
317 sqlx::query(hnsw_sql)
318 .execute(pool)
319 .await
320 .map_err(|e| Error::Storage(format!("create hnsw index: {e}")))?;
321
322 Ok(())
323}