1use std::collections::{BTreeMap, HashMap, HashSet};
9use std::fmt;
10use std::path::Path;
11use std::sync::{Mutex, MutexGuard};
12
13use rusqlite::{params, Connection, OptionalExtension, Transaction};
14use serde::{Deserialize, Serialize};
15use time::format_description::well_known::Rfc3339;
16use time::OffsetDateTime;
17
18use crate::runtime_sqlite::{
19 initialize_runtime_sqlite, initialize_transient_runtime_sqlite, RuntimeSqliteSchema,
20 DEFAULT_BUSY_TIMEOUT,
21};
22
23use super::backend::{AtomRef, FlowSlice, GitExportReceipt, ShipReceipt, VcsBackend};
24use super::{Atom, AtomId, Intent, IntentId, Slice as DerivedSlice, SliceId, VcsBackendError};
25
26const SQLITE_ATOM_REF_PREFIX: &str = "sqlite://atoms";
27const SQLITE_SLICE_REF_PREFIX: &str = "sqlite://slices";
28const SQLITE_SCHEMA: RuntimeSqliteSchema =
29 RuntimeSqliteSchema::new("harn_flow", 1, FLOW_SCHEMA_SQL);
30
31#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
33pub struct StateVector {
34 clocks: BTreeMap<String, u64>,
35}
36
37impl StateVector {
38 pub fn new() -> Self {
39 Self::default()
40 }
41
42 pub fn insert(&mut self, site_id: impl Into<String>, clock: u64) {
43 self.clocks.insert(site_id.into(), clock);
44 }
45
46 pub fn clock(&self, site_id: &str) -> u64 {
47 self.clocks.get(site_id).copied().unwrap_or(0)
48 }
49
50 pub fn iter(&self) -> impl Iterator<Item = (&str, u64)> {
51 self.clocks
52 .iter()
53 .map(|(site_id, clock)| (site_id.as_str(), *clock))
54 }
55}
56
57#[derive(Clone, Debug, PartialEq, Eq)]
59pub struct AtomDelta {
60 pub atom: Atom,
61 pub site_id: String,
62 pub clock: u64,
63}
64
65#[derive(Clone, Debug, PartialEq, Eq)]
67pub struct StoredDerivedSlice {
68 pub slice: DerivedSlice,
69 pub created_at: String,
70}
71
72pub struct SqliteFlowStore {
74 site_id: String,
75 conn: Mutex<Connection>,
76}
77
78impl fmt::Debug for SqliteFlowStore {
79 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
80 f.debug_struct("SqliteFlowStore")
81 .field("site_id", &self.site_id)
82 .finish_non_exhaustive()
83 }
84}
85
86impl SqliteFlowStore {
87 pub fn open(
90 path: impl AsRef<Path>,
91 site_id: impl Into<String>,
92 ) -> Result<Self, VcsBackendError> {
93 let site_id = normalize_site_id(site_id.into())?;
94 let path = path.as_ref();
95 let conn = Connection::open(path)?;
96 initialize_file_schema(&conn)?;
97 Ok(Self {
98 site_id,
99 conn: Mutex::new(conn),
100 })
101 }
102
103 pub fn in_memory(site_id: impl Into<String>) -> Result<Self, VcsBackendError> {
105 let site_id = normalize_site_id(site_id.into())?;
106 let conn = Connection::open_in_memory()?;
107 initialize_transient_schema(&conn)?;
108 Ok(Self {
109 site_id,
110 conn: Mutex::new(conn),
111 })
112 }
113
114 pub fn site_id(&self) -> &str {
115 &self.site_id
116 }
117
118 pub fn emit_atoms(&self, atoms: &[Atom]) -> Result<Vec<AtomRef>, VcsBackendError> {
120 self.emit_atoms_inner(atoms, true)
121 }
122
123 pub fn emit_preverified_atoms(&self, atoms: &[Atom]) -> Result<Vec<AtomRef>, VcsBackendError> {
130 let mut conn = self.lock_conn()?;
131 let tx = conn.transaction()?;
132 let mut clocks: HashMap<(String, String), u64> = HashMap::new();
133 let mut refs = Vec::with_capacity(atoms.len());
134
135 {
136 let mut insert_atom = tx.prepare_cached(
137 "INSERT INTO atoms (
138 id, principal, persona, timestamp_ns, timestamp_rfc3339,
139 site_id, site_clock, inverse_of, body_binary
140 )
141 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
142 )?;
143 let mut insert_parent = tx.prepare_cached(
144 "INSERT INTO atom_parents (child_id, parent_id, ordinal)
145 VALUES (?1, ?2, ?3)",
146 )?;
147
148 for atom in atoms {
149 let key = (
150 atom.provenance.principal.clone(),
151 atom.provenance.persona.clone(),
152 );
153 if !clocks.contains_key(&key) {
154 let current = state_vector_clock_tx(
155 &tx,
156 &atom.provenance.principal,
157 &atom.provenance.persona,
158 &self.site_id,
159 )?;
160 clocks.insert(key.clone(), current);
161 }
162 let clock = clocks
163 .get_mut(&key)
164 .expect("clock was inserted before increment");
165 *clock = clock
166 .checked_add(1)
167 .ok_or_else(|| VcsBackendError::Invalid("site clock overflow".to_string()))?;
168
169 let body = atom.to_binary()?;
170 let timestamp_ns = atom_timestamp_ns(atom)?;
171 let timestamp_rfc3339 = atom_timestamp_rfc3339(atom)?;
172 let inverse_of = atom.inverse_of.map(|id| id.0.to_vec());
173 insert_atom.execute(params![
174 atom.id.0.as_slice(),
175 atom.provenance.principal,
176 atom.provenance.persona,
177 timestamp_ns,
178 timestamp_rfc3339,
179 self.site_id.as_str(),
180 i64_from_u64(*clock, "atom site clock")?,
181 inverse_of.as_deref(),
182 body.as_slice(),
183 ])?;
184
185 for (ordinal, parent) in atom.parents.iter().enumerate() {
186 insert_parent.execute(params![
187 atom.id.0.as_slice(),
188 parent.0.as_slice(),
189 i64_from_usize(ordinal, "atom parent ordinal")?
190 ])?;
191 }
192 refs.push(sqlite_atom_ref(atom.id, &self.site_id, *clock));
193 }
194 }
195
196 for ((principal, persona), clock) in clocks {
197 advance_state_vector_tx(&tx, &principal, &persona, &self.site_id, clock)?;
198 }
199 tx.commit()?;
200 Ok(refs)
201 }
202
203 fn emit_atoms_inner(
204 &self,
205 atoms: &[Atom],
206 verify: bool,
207 ) -> Result<Vec<AtomRef>, VcsBackendError> {
208 let mut conn = self.lock_conn()?;
209 let tx = conn.transaction()?;
210 let mut refs = Vec::with_capacity(atoms.len());
211 for atom in atoms {
212 if verify {
213 atom.verify()?;
214 }
215 refs.push(insert_atom_tx(&tx, atom, &self.site_id, None)?);
216 }
217 tx.commit()?;
218 Ok(refs)
219 }
220
221 pub fn insert_remote_atom(
223 &self,
224 atom: &Atom,
225 site_id: &str,
226 clock: u64,
227 ) -> Result<AtomRef, VcsBackendError> {
228 atom.verify()?;
229 if clock == 0 {
230 return Err(VcsBackendError::Invalid(
231 "remote atom clock must be greater than zero".to_string(),
232 ));
233 }
234 let site_id = normalize_site_id(site_id.to_string())?;
235 let mut conn = self.lock_conn()?;
236 let tx = conn.transaction()?;
237 let atom_ref = insert_atom_tx(&tx, atom, &site_id, Some(clock))?;
238 tx.commit()?;
239 Ok(atom_ref)
240 }
241
242 pub fn get_atom(&self, atom_id: AtomId) -> Result<Atom, VcsBackendError> {
244 let conn = self.lock_conn()?;
245 load_atom(&conn, atom_id)
246 }
247
248 pub fn atom_by_content_hash(
251 &self,
252 content_hash: AtomId,
253 ) -> Result<Option<Atom>, VcsBackendError> {
254 let conn = self.lock_conn()?;
255 conn.query_row(
256 "SELECT body_binary FROM atoms WHERE id = ?1",
257 params![content_hash.0.as_slice()],
258 |row| row.get::<_, Vec<u8>>(0),
259 )
260 .optional()?
261 .map(|body| Atom::from_binary_slice(&body).map_err(Into::into))
262 .transpose()
263 }
264
265 pub fn atoms_for_principal_persona(
267 &self,
268 principal: &str,
269 persona: &str,
270 ) -> Result<Vec<Atom>, VcsBackendError> {
271 let conn = self.lock_conn()?;
272 let mut stmt = conn.prepare(
273 "SELECT id FROM atoms
274 WHERE principal = ?1 AND persona = ?2
275 ORDER BY timestamp_ns, id",
276 )?;
277 let rows = stmt.query_map(params![principal, persona], |row| row.get::<_, Vec<u8>>(0))?;
278 let mut atoms = Vec::new();
279 for row in rows {
280 atoms.push(load_atom(&conn, atom_id_from_blob(row?)?)?);
281 }
282 Ok(atoms)
283 }
284
285 pub fn atom_count_for_principal_persona(
287 &self,
288 principal: &str,
289 persona: &str,
290 ) -> Result<u64, VcsBackendError> {
291 let conn = self.lock_conn()?;
292 let count = conn.query_row(
293 "SELECT COUNT(*) FROM atoms WHERE principal = ?1 AND persona = ?2",
294 params![principal, persona],
295 |row| row.get::<_, i64>(0),
296 )?;
297 u64_from_i64(count, "atom count")
298 }
299
300 pub fn atoms_with_parent(&self, parent: AtomId) -> Result<Vec<Atom>, VcsBackendError> {
302 let conn = self.lock_conn()?;
303 let mut stmt = conn.prepare(
304 "SELECT child_id FROM atom_parents
305 WHERE parent_id = ?1
306 ORDER BY child_id",
307 )?;
308 let rows = stmt.query_map(params![parent.0.as_slice()], |row| row.get::<_, Vec<u8>>(0))?;
309 let mut atoms = Vec::new();
310 for row in rows {
311 atoms.push(load_atom(&conn, atom_id_from_blob(row?)?)?);
312 }
313 Ok(atoms)
314 }
315
316 pub fn state_vector(
318 &self,
319 principal: &str,
320 persona: &str,
321 ) -> Result<StateVector, VcsBackendError> {
322 let conn = self.lock_conn()?;
323 let mut stmt = conn.prepare(
324 "SELECT site_id, clock FROM state_vectors
325 WHERE principal = ?1 AND persona = ?2
326 ORDER BY site_id",
327 )?;
328 let rows = stmt.query_map(params![principal, persona], |row| {
329 Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
330 })?;
331 let mut vector = StateVector::new();
332 for row in rows {
333 let (site_id, clock) = row?;
334 vector.insert(site_id, u64_from_i64(clock, "state vector clock")?);
335 }
336 Ok(vector)
337 }
338
339 pub fn causal_delta(
341 &self,
342 principal: &str,
343 persona: &str,
344 remote: &StateVector,
345 ) -> Result<Vec<AtomDelta>, VcsBackendError> {
346 let conn = self.lock_conn()?;
347 let mut stmt = conn.prepare(
348 "SELECT id, site_id, site_clock FROM atoms
349 WHERE principal = ?1 AND persona = ?2
350 ORDER BY site_id, site_clock, id",
351 )?;
352 let rows = stmt.query_map(params![principal, persona], |row| {
353 Ok((
354 row.get::<_, Vec<u8>>(0)?,
355 row.get::<_, String>(1)?,
356 row.get::<_, i64>(2)?,
357 ))
358 })?;
359 let mut delta = Vec::new();
360 for row in rows {
361 let (id_blob, site_id, clock_raw) = row?;
362 let clock = u64_from_i64(clock_raw, "atom site clock")?;
363 if clock > remote.clock(&site_id) {
364 delta.push(AtomDelta {
365 atom: load_atom(&conn, atom_id_from_blob(id_blob)?)?,
366 site_id,
367 clock,
368 });
369 }
370 }
371 Ok(delta)
372 }
373
374 pub fn put_intent(&self, intent: &Intent) -> Result<(), VcsBackendError> {
376 let body = serde_json::to_vec(intent)?;
377 let mut conn = self.lock_conn()?;
378 let tx = conn.transaction()?;
379 tx.execute(
380 "INSERT OR IGNORE INTO intents (id, body_json, goal_description, confidence)
381 VALUES (?1, ?2, ?3, ?4)",
382 params![
383 intent.id.0.as_slice(),
384 body.as_slice(),
385 intent.goal_description,
386 f64::from(intent.confidence)
387 ],
388 )?;
389 for (ordinal, atom_id) in intent.atoms.iter().enumerate() {
390 tx.execute(
391 "INSERT OR IGNORE INTO intent_atoms (intent_id, atom_id, ordinal)
392 VALUES (?1, ?2, ?3)",
393 params![
394 intent.id.0.as_slice(),
395 atom_id.0.as_slice(),
396 i64_from_usize(ordinal, "intent atom ordinal")?
397 ],
398 )?;
399 }
400 tx.commit()?;
401 Ok(())
402 }
403
404 pub fn get_intent(&self, intent_id: IntentId) -> Result<Intent, VcsBackendError> {
405 let conn = self.lock_conn()?;
406 let body = conn
407 .query_row(
408 "SELECT body_json FROM intents WHERE id = ?1",
409 params![intent_id.0.as_slice()],
410 |row| row.get::<_, Vec<u8>>(0),
411 )
412 .optional()?
413 .ok_or_else(|| VcsBackendError::NotFound(format!("intent {intent_id} not found")))?;
414 serde_json::from_slice(&body).map_err(Into::into)
415 }
416
417 pub fn put_derived_slice(&self, slice: &DerivedSlice) -> Result<(), VcsBackendError> {
419 let body = serde_json::to_vec(slice)?;
420 self.insert_slice_record(slice.id, &slice.atoms, "derived", body, false)
421 }
422
423 pub fn put_shipped_derived_slice(&self, slice: &DerivedSlice) -> Result<(), VcsBackendError> {
428 let body = serde_json::to_vec(slice)?;
429 self.insert_slice_record(slice.id, &slice.atoms, "derived", body, true)
430 }
431
432 pub fn get_derived_slice(&self, slice_id: SliceId) -> Result<DerivedSlice, VcsBackendError> {
433 let conn = self.lock_conn()?;
434 let body = conn
435 .query_row(
436 "SELECT body_json FROM slices WHERE id = ?1 AND slice_kind = 'derived'",
437 params![slice_id.0.as_slice()],
438 |row| row.get::<_, Vec<u8>>(0),
439 )
440 .optional()?
441 .ok_or_else(|| VcsBackendError::NotFound(format!("slice {slice_id} not found")))?;
442 serde_json::from_slice(&body).map_err(Into::into)
443 }
444
445 pub fn shipped_derived_slices_since(
448 &self,
449 since: Option<OffsetDateTime>,
450 ) -> Result<Vec<StoredDerivedSlice>, VcsBackendError> {
451 let since = since
452 .map(|value| value.format(&Rfc3339))
453 .transpose()
454 .map_err(|error| VcsBackendError::Invalid(format!("timestamp format: {error}")))?;
455 let conn = self.lock_conn()?;
456 let mut stmt = conn.prepare(
457 "SELECT body_json, created_at FROM slices
458 WHERE slice_kind = 'derived'
459 AND shipped = 1
460 AND (?1 IS NULL OR created_at >= datetime(?1))
461 ORDER BY created_at, id",
462 )?;
463 let rows = stmt.query_map(params![since.as_deref()], |row| {
464 Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, String>(1)?))
465 })?;
466 let mut slices = Vec::new();
467 for row in rows {
468 let (body, created_at) = row?;
469 slices.push(StoredDerivedSlice {
470 slice: serde_json::from_slice(&body)?,
471 created_at,
472 });
473 }
474 Ok(slices)
475 }
476
477 fn insert_flow_slice(&self, slice: &FlowSlice, shipped: bool) -> Result<(), VcsBackendError> {
478 let body = serde_json::to_vec(slice)?;
479 self.insert_slice_record(slice.id, &slice.atoms, "flow", body, shipped)
480 }
481
482 fn insert_slice_record(
483 &self,
484 slice_id: SliceId,
485 atoms: &[AtomId],
486 kind: &str,
487 body: Vec<u8>,
488 shipped: bool,
489 ) -> Result<(), VcsBackendError> {
490 let mut conn = self.lock_conn()?;
491 let tx = conn.transaction()?;
492 insert_slice_record_tx(&tx, slice_id, atoms, kind, &body, shipped)?;
493 tx.commit()?;
494 Ok(())
495 }
496
497 fn atom_closure(&self, roots: &[AtomId]) -> Result<Vec<AtomId>, VcsBackendError> {
498 let mut opened = HashSet::new();
499 let mut emitted = HashSet::new();
500 let mut out = Vec::new();
501 let mut stack: Vec<(AtomId, bool)> = roots
502 .iter()
503 .rev()
504 .copied()
505 .map(|atom_id| (atom_id, false))
506 .collect();
507
508 while let Some((atom_id, emit)) = stack.pop() {
509 if emit {
510 if emitted.insert(atom_id) {
511 out.push(atom_id);
512 }
513 continue;
514 }
515 if emitted.contains(&atom_id) || !opened.insert(atom_id) {
516 continue;
517 }
518
519 let atom = self.get_atom(atom_id)?;
520 stack.push((atom_id, true));
521 for parent in atom.parents.iter().rev() {
522 if !emitted.contains(parent) {
523 stack.push((*parent, false));
524 }
525 }
526 }
527
528 Ok(out)
529 }
530
531 fn lock_conn(&self) -> Result<MutexGuard<'_, Connection>, VcsBackendError> {
532 self.conn
533 .lock()
534 .map_err(|_| VcsBackendError::Io("sqlite flow store lock poisoned".to_string()))
535 }
536}
537
538impl VcsBackend for SqliteFlowStore {
539 fn emit_atom(&self, atom: &Atom) -> Result<AtomRef, VcsBackendError> {
540 self.emit_atoms(std::slice::from_ref(atom))
541 .map(|mut refs| refs.remove(0))
542 }
543
544 fn derive_slice(&self, atoms: &[AtomId]) -> Result<FlowSlice, VcsBackendError> {
545 FlowSlice::new(self.atom_closure(atoms)?)
546 }
547
548 fn ship_slice(&self, slice: &FlowSlice) -> Result<ShipReceipt, VcsBackendError> {
549 self.insert_flow_slice(slice, true)?;
550 Ok(ShipReceipt {
551 slice_id: slice.id,
552 commit: slice.id.to_string(),
553 ref_name: format!("{SQLITE_SLICE_REF_PREFIX}/{}", slice.id),
554 })
555 }
556
557 fn list_atoms(&self) -> Result<Vec<AtomRef>, VcsBackendError> {
558 let conn = self.lock_conn()?;
559 let mut stmt = conn.prepare(
560 "SELECT id, site_id, site_clock FROM atoms
561 ORDER BY principal, persona, timestamp_ns, id",
562 )?;
563 let rows = stmt.query_map([], |row| {
564 Ok((
565 row.get::<_, Vec<u8>>(0)?,
566 row.get::<_, String>(1)?,
567 row.get::<_, i64>(2)?,
568 ))
569 })?;
570 let mut atoms = Vec::new();
571 for row in rows {
572 let (id_blob, site_id, clock_raw) = row?;
573 atoms.push(sqlite_atom_ref(
574 atom_id_from_blob(id_blob)?,
575 &site_id,
576 u64_from_i64(clock_raw, "atom site clock")?,
577 ));
578 }
579 Ok(atoms)
580 }
581
582 fn replay_slice(&self, slice: &FlowSlice) -> Result<Vec<Atom>, VcsBackendError> {
583 slice
584 .atoms
585 .iter()
586 .map(|atom_id| self.get_atom(*atom_id))
587 .collect()
588 }
589
590 fn export_git(
591 &self,
592 _slice: &FlowSlice,
593 _ref_name: &str,
594 ) -> Result<GitExportReceipt, VcsBackendError> {
595 Err(VcsBackendError::Unsupported(
596 "SqliteFlowStore cannot export git refs; use ShadowGitBackend for git export"
597 .to_string(),
598 ))
599 }
600
601 fn import_git(&self, _ref_name: &str) -> Result<FlowSlice, VcsBackendError> {
602 Err(VcsBackendError::Unsupported(
603 "SqliteFlowStore cannot import git refs; use ShadowGitBackend for git import"
604 .to_string(),
605 ))
606 }
607}
608
609const FLOW_SCHEMA_SQL: &str = r"
610 CREATE TABLE IF NOT EXISTS atoms (
611 id BLOB PRIMARY KEY CHECK(length(id) = 32),
612 principal TEXT NOT NULL,
613 persona TEXT NOT NULL,
614 timestamp_ns INTEGER NOT NULL,
615 timestamp_rfc3339 TEXT NOT NULL,
616 site_id TEXT NOT NULL,
617 site_clock INTEGER NOT NULL CHECK(site_clock > 0),
618 inverse_of BLOB CHECK(inverse_of IS NULL OR length(inverse_of) = 32),
619 body_binary BLOB NOT NULL,
620 created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
621 UNIQUE(principal, persona, site_id, site_clock)
622 );
623
624 CREATE INDEX IF NOT EXISTS atoms_principal_persona_timestamp_idx
625 ON atoms(principal, persona, timestamp_ns, id);
626 CREATE INDEX IF NOT EXISTS atoms_principal_persona_site_clock_idx
627 ON atoms(principal, persona, site_id, site_clock);
628 CREATE INDEX IF NOT EXISTS atoms_inverse_of_idx ON atoms(inverse_of);
629
630 CREATE TABLE IF NOT EXISTS atom_parents (
631 child_id BLOB NOT NULL CHECK(length(child_id) = 32),
632 parent_id BLOB NOT NULL CHECK(length(parent_id) = 32),
633 ordinal INTEGER NOT NULL CHECK(ordinal >= 0),
634 PRIMARY KEY(child_id, ordinal),
635 UNIQUE(child_id, parent_id),
636 FOREIGN KEY(child_id) REFERENCES atoms(id)
637 );
638 CREATE INDEX IF NOT EXISTS atom_parents_parent_idx
639 ON atom_parents(parent_id, child_id);
640
641 CREATE TABLE IF NOT EXISTS intents (
642 id BLOB PRIMARY KEY CHECK(length(id) = 32),
643 body_json BLOB NOT NULL,
644 goal_description TEXT NOT NULL,
645 confidence REAL NOT NULL,
646 created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
647 );
648
649 CREATE TABLE IF NOT EXISTS intent_atoms (
650 intent_id BLOB NOT NULL CHECK(length(intent_id) = 32),
651 atom_id BLOB NOT NULL CHECK(length(atom_id) = 32),
652 ordinal INTEGER NOT NULL CHECK(ordinal >= 0),
653 PRIMARY KEY(intent_id, ordinal),
654 UNIQUE(intent_id, atom_id),
655 FOREIGN KEY(intent_id) REFERENCES intents(id)
656 );
657 CREATE INDEX IF NOT EXISTS intent_atoms_atom_idx
658 ON intent_atoms(atom_id, intent_id);
659
660 CREATE TABLE IF NOT EXISTS slices (
661 id BLOB PRIMARY KEY CHECK(length(id) = 32),
662 slice_kind TEXT NOT NULL,
663 body_json BLOB NOT NULL,
664 shipped INTEGER NOT NULL DEFAULT 0,
665 ref_name TEXT,
666 created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
667 );
668
669 CREATE TABLE IF NOT EXISTS slice_atoms (
670 slice_id BLOB NOT NULL CHECK(length(slice_id) = 32),
671 atom_id BLOB NOT NULL CHECK(length(atom_id) = 32),
672 ordinal INTEGER NOT NULL CHECK(ordinal >= 0),
673 PRIMARY KEY(slice_id, ordinal),
674 UNIQUE(slice_id, atom_id),
675 FOREIGN KEY(slice_id) REFERENCES slices(id)
676 );
677 CREATE INDEX IF NOT EXISTS slice_atoms_atom_idx
678 ON slice_atoms(atom_id, slice_id);
679
680 CREATE TABLE IF NOT EXISTS state_vectors (
681 principal TEXT NOT NULL,
682 persona TEXT NOT NULL,
683 site_id TEXT NOT NULL,
684 clock INTEGER NOT NULL CHECK(clock >= 0),
685 updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
686 PRIMARY KEY(principal, persona, site_id)
687 );
688
689 CREATE TRIGGER IF NOT EXISTS atoms_no_update
690 BEFORE UPDATE ON atoms
691 BEGIN
692 SELECT RAISE(ABORT, 'atoms are append-only');
693 END;
694
695 CREATE TRIGGER IF NOT EXISTS atoms_no_delete
696 BEFORE DELETE ON atoms
697 BEGIN
698 SELECT RAISE(ABORT, 'atoms are append-only');
699 END;
700
701 CREATE TRIGGER IF NOT EXISTS atom_parents_no_update
702 BEFORE UPDATE ON atom_parents
703 BEGIN
704 SELECT RAISE(ABORT, 'atom parent edges are append-only');
705 END;
706
707 CREATE TRIGGER IF NOT EXISTS atom_parents_no_delete
708 BEFORE DELETE ON atom_parents
709 BEGIN
710 SELECT RAISE(ABORT, 'atom parent edges are append-only');
711 END;
712
713 CREATE TRIGGER IF NOT EXISTS slices_no_update
714 BEFORE UPDATE ON slices
715 BEGIN
716 SELECT RAISE(ABORT, 'slices are append-only');
717 END;
718
719 CREATE TRIGGER IF NOT EXISTS slices_no_delete
720 BEFORE DELETE ON slices
721 BEGIN
722 SELECT RAISE(ABORT, 'slices are append-only');
723 END;
724
725 CREATE TRIGGER IF NOT EXISTS slice_atoms_no_update
726 BEFORE UPDATE ON slice_atoms
727 BEGIN
728 SELECT RAISE(ABORT, 'slice atom edges are append-only');
729 END;
730
731 CREATE TRIGGER IF NOT EXISTS slice_atoms_no_delete
732 BEFORE DELETE ON slice_atoms
733 BEGIN
734 SELECT RAISE(ABORT, 'slice atom edges are append-only');
735 END;
736 ";
737
738fn initialize_file_schema(conn: &Connection) -> Result<(), VcsBackendError> {
739 configure_flow_connection(conn)?;
740 initialize_runtime_sqlite(conn, DEFAULT_BUSY_TIMEOUT, &SQLITE_SCHEMA)
741 .map_err(|error| VcsBackendError::Sqlite(error.to_string()))
742}
743
744fn initialize_transient_schema(conn: &Connection) -> Result<(), VcsBackendError> {
745 configure_flow_connection(conn)?;
746 initialize_transient_runtime_sqlite(conn, DEFAULT_BUSY_TIMEOUT, &SQLITE_SCHEMA)
747 .map_err(|error| VcsBackendError::Sqlite(error.to_string()))
748}
749
750fn configure_flow_connection(conn: &Connection) -> Result<(), VcsBackendError> {
751 conn.pragma_update(None, "foreign_keys", true)?;
752 Ok(())
753}
754
755fn insert_atom_tx(
756 tx: &Transaction<'_>,
757 atom: &Atom,
758 site_id: &str,
759 explicit_clock: Option<u64>,
760) -> Result<AtomRef, VcsBackendError> {
761 if let Some((existing_site, existing_clock)) = atom_clock_tx(tx, atom.id)? {
762 return Ok(sqlite_atom_ref(atom.id, &existing_site, existing_clock));
763 }
764
765 let clock = match explicit_clock {
766 Some(clock) => {
767 reject_site_clock_conflict(
768 tx,
769 &atom.provenance.principal,
770 &atom.provenance.persona,
771 site_id,
772 clock,
773 atom.id,
774 )?;
775 advance_state_vector_tx(
776 tx,
777 &atom.provenance.principal,
778 &atom.provenance.persona,
779 site_id,
780 clock,
781 )?;
782 clock
783 }
784 None => reserve_next_clock_tx(
785 tx,
786 &atom.provenance.principal,
787 &atom.provenance.persona,
788 site_id,
789 )?,
790 };
791
792 let body = atom.to_binary()?;
793 let timestamp_ns = atom_timestamp_ns(atom)?;
794 let timestamp_rfc3339 = atom_timestamp_rfc3339(atom)?;
795 let inverse_of = atom.inverse_of.map(|id| id.0.to_vec());
796 tx.execute(
797 "INSERT INTO atoms (
798 id, principal, persona, timestamp_ns, timestamp_rfc3339,
799 site_id, site_clock, inverse_of, body_binary
800 )
801 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
802 params![
803 atom.id.0.as_slice(),
804 atom.provenance.principal,
805 atom.provenance.persona,
806 timestamp_ns,
807 timestamp_rfc3339,
808 site_id,
809 i64_from_u64(clock, "atom site clock")?,
810 inverse_of.as_deref(),
811 body.as_slice(),
812 ],
813 )?;
814
815 for (ordinal, parent) in atom.parents.iter().enumerate() {
816 tx.execute(
817 "INSERT INTO atom_parents (child_id, parent_id, ordinal)
818 VALUES (?1, ?2, ?3)",
819 params![
820 atom.id.0.as_slice(),
821 parent.0.as_slice(),
822 i64_from_usize(ordinal, "atom parent ordinal")?
823 ],
824 )?;
825 }
826
827 Ok(sqlite_atom_ref(atom.id, site_id, clock))
828}
829
830fn insert_slice_record_tx(
831 tx: &Transaction<'_>,
832 slice_id: SliceId,
833 atoms: &[AtomId],
834 kind: &str,
835 body: &[u8],
836 shipped: bool,
837) -> Result<(), VcsBackendError> {
838 tx.execute(
839 "INSERT OR IGNORE INTO slices (id, slice_kind, body_json, shipped, ref_name)
840 VALUES (?1, ?2, ?3, ?4, ?5)",
841 params![
842 slice_id.0.as_slice(),
843 kind,
844 body,
845 i32::from(shipped),
846 if shipped {
847 Some(format!("{SQLITE_SLICE_REF_PREFIX}/{slice_id}"))
848 } else {
849 None
850 }
851 ],
852 )?;
853 for (ordinal, atom_id) in atoms.iter().enumerate() {
854 tx.execute(
855 "INSERT OR IGNORE INTO slice_atoms (slice_id, atom_id, ordinal)
856 VALUES (?1, ?2, ?3)",
857 params![
858 slice_id.0.as_slice(),
859 atom_id.0.as_slice(),
860 i64_from_usize(ordinal, "slice atom ordinal")?
861 ],
862 )?;
863 }
864 Ok(())
865}
866
867fn load_atom(conn: &Connection, atom_id: AtomId) -> Result<Atom, VcsBackendError> {
868 let body = conn
869 .query_row(
870 "SELECT body_binary FROM atoms WHERE id = ?1",
871 params![atom_id.0.as_slice()],
872 |row| row.get::<_, Vec<u8>>(0),
873 )
874 .optional()?
875 .ok_or_else(|| VcsBackendError::NotFound(format!("atom {atom_id} not found")))?;
876 Atom::from_binary_slice(&body).map_err(Into::into)
877}
878
879fn atom_clock_tx(
880 tx: &Transaction<'_>,
881 atom_id: AtomId,
882) -> Result<Option<(String, u64)>, VcsBackendError> {
883 tx.query_row(
884 "SELECT site_id, site_clock FROM atoms WHERE id = ?1",
885 params![atom_id.0.as_slice()],
886 |row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)),
887 )
888 .optional()?
889 .map(|(site_id, clock)| Ok((site_id, u64_from_i64(clock, "atom site clock")?)))
890 .transpose()
891}
892
893fn reserve_next_clock_tx(
894 tx: &Transaction<'_>,
895 principal: &str,
896 persona: &str,
897 site_id: &str,
898) -> Result<u64, VcsBackendError> {
899 let current = state_vector_clock_tx(tx, principal, persona, site_id)?;
900 let next = current
901 .checked_add(1)
902 .ok_or_else(|| VcsBackendError::Invalid("state vector clock overflow".to_string()))?;
903 advance_state_vector_tx(tx, principal, persona, site_id, next)?;
904 Ok(next)
905}
906
907fn state_vector_clock_tx(
908 tx: &Transaction<'_>,
909 principal: &str,
910 persona: &str,
911 site_id: &str,
912) -> Result<u64, VcsBackendError> {
913 tx.query_row(
914 "SELECT clock FROM state_vectors
915 WHERE principal = ?1 AND persona = ?2 AND site_id = ?3",
916 params![principal, persona, site_id],
917 |row| row.get::<_, i64>(0),
918 )
919 .optional()?
920 .map(|clock| u64_from_i64(clock, "state vector clock"))
921 .transpose()
922 .map(|clock| clock.unwrap_or(0))
923}
924
925fn advance_state_vector_tx(
926 tx: &Transaction<'_>,
927 principal: &str,
928 persona: &str,
929 site_id: &str,
930 clock: u64,
931) -> Result<(), VcsBackendError> {
932 tx.execute(
933 "INSERT INTO state_vectors (principal, persona, site_id, clock, updated_at)
934 VALUES (?1, ?2, ?3, ?4, CURRENT_TIMESTAMP)
935 ON CONFLICT(principal, persona, site_id) DO UPDATE SET
936 clock = CASE
937 WHEN excluded.clock > state_vectors.clock THEN excluded.clock
938 ELSE state_vectors.clock
939 END,
940 updated_at = CURRENT_TIMESTAMP",
941 params![
942 principal,
943 persona,
944 site_id,
945 i64_from_u64(clock, "state vector clock")?
946 ],
947 )?;
948 Ok(())
949}
950
951fn reject_site_clock_conflict(
952 tx: &Transaction<'_>,
953 principal: &str,
954 persona: &str,
955 site_id: &str,
956 clock: u64,
957 atom_id: AtomId,
958) -> Result<(), VcsBackendError> {
959 let existing = tx
960 .query_row(
961 "SELECT id FROM atoms
962 WHERE principal = ?1 AND persona = ?2 AND site_id = ?3 AND site_clock = ?4",
963 params![
964 principal,
965 persona,
966 site_id,
967 i64_from_u64(clock, "atom site clock")?
968 ],
969 |row| row.get::<_, Vec<u8>>(0),
970 )
971 .optional()?;
972 if let Some(existing) = existing {
973 let existing = atom_id_from_blob(existing)?;
974 if existing != atom_id {
975 return Err(VcsBackendError::Invalid(format!(
976 "site clock conflict for {site_id}@{clock}: existing atom {existing}, new atom {atom_id}"
977 )));
978 }
979 }
980 Ok(())
981}
982
983fn sqlite_atom_ref(atom_id: AtomId, site_id: &str, clock: u64) -> AtomRef {
984 AtomRef {
985 atom_id,
986 commit: format!("{site_id}:{clock}"),
987 ref_name: format!("{SQLITE_ATOM_REF_PREFIX}/{atom_id}"),
988 }
989}
990
991fn atom_timestamp_ns(atom: &Atom) -> Result<i64, VcsBackendError> {
992 i64::try_from(atom.provenance.timestamp.unix_timestamp_nanos())
993 .map_err(|_| VcsBackendError::Invalid("atom timestamp is out of SQLite range".to_string()))
994}
995
996fn atom_timestamp_rfc3339(atom: &Atom) -> Result<String, VcsBackendError> {
997 atom.provenance
998 .timestamp
999 .format(&Rfc3339)
1000 .map_err(|error| VcsBackendError::Invalid(format!("atom timestamp format: {error}")))
1001}
1002
1003fn atom_id_from_blob(blob: Vec<u8>) -> Result<AtomId, VcsBackendError> {
1004 if blob.len() != 32 {
1005 return Err(VcsBackendError::Invalid(format!(
1006 "atom id blob must be 32 bytes, got {}",
1007 blob.len()
1008 )));
1009 }
1010 let mut out = [0u8; 32];
1011 out.copy_from_slice(&blob);
1012 Ok(AtomId(out))
1013}
1014
1015fn normalize_site_id(site_id: String) -> Result<String, VcsBackendError> {
1016 if site_id.trim().is_empty() {
1017 return Err(VcsBackendError::Invalid(
1018 "flow store site_id must not be empty".to_string(),
1019 ));
1020 }
1021 Ok(site_id)
1022}
1023
1024fn i64_from_u64(value: u64, field: &str) -> Result<i64, VcsBackendError> {
1025 i64::try_from(value)
1026 .map_err(|_| VcsBackendError::Invalid(format!("{field} exceeds SQLite i64 range")))
1027}
1028
1029fn i64_from_usize(value: usize, field: &str) -> Result<i64, VcsBackendError> {
1030 i64::try_from(value)
1031 .map_err(|_| VcsBackendError::Invalid(format!("{field} exceeds SQLite i64 range")))
1032}
1033
1034fn u64_from_i64(value: i64, field: &str) -> Result<u64, VcsBackendError> {
1035 u64::try_from(value).map_err(|_| VcsBackendError::Invalid(format!("{field} is negative")))
1036}
1037
1038#[cfg(test)]
1039mod tests {
1040 use super::*;
1041 use crate::flow::{Approval, CoverageMap, PredicateHash, Slice, SliceStatus, TestId, TextOp};
1042 use ed25519_dalek::SigningKey;
1043 use tempfile::TempDir;
1044 use time::OffsetDateTime;
1045
1046 fn key(seed: u8) -> SigningKey {
1047 SigningKey::from_bytes(&[seed; 32])
1048 }
1049
1050 fn atom(index: u64, parents: Vec<AtomId>) -> Atom {
1051 let principal = key(1);
1052 let persona = key(2);
1053 let timestamp = OffsetDateTime::from_unix_timestamp(1_775_000_000 + index as i64).unwrap();
1054 Atom::sign(
1055 vec![TextOp::Insert {
1056 offset: index,
1057 content: format!("atom-{index}"),
1058 }],
1059 parents,
1060 crate::flow::Provenance {
1061 principal: "user:alice".to_string(),
1062 persona: "ship-captain".to_string(),
1063 agent_run_id: format!("run-{index}"),
1064 tool_call_id: Some(format!("tool-{index}")),
1065 trace_id: "trace-1".to_string(),
1066 transcript_ref: "transcript-1".to_string(),
1067 timestamp,
1068 },
1069 None,
1070 &principal,
1071 &persona,
1072 )
1073 .unwrap()
1074 }
1075
1076 #[test]
1077 fn file_store_uses_versioned_runtime_sqlite_contract() {
1078 let temp = TempDir::new().unwrap();
1079 let store = SqliteFlowStore::open(temp.path().join("flow.sqlite"), "site-a").unwrap();
1080 let conn = store.lock_conn().unwrap();
1081
1082 let journal_mode = conn
1083 .query_row("PRAGMA journal_mode", [], |row| row.get::<_, String>(0))
1084 .unwrap();
1085 let foreign_keys = conn
1086 .query_row("PRAGMA foreign_keys", [], |row| row.get::<_, i64>(0))
1087 .unwrap();
1088 let schema_version = conn
1089 .query_row(
1090 "SELECT version FROM _harn_sqlite_schema_versions WHERE name = ?1",
1091 ["harn_flow"],
1092 |row| row.get::<_, i64>(0),
1093 )
1094 .unwrap();
1095
1096 assert_eq!(
1097 (journal_mode, foreign_keys, schema_version),
1098 ("wal".to_string(), 1, 1)
1099 );
1100 }
1101
1102 #[test]
1103 fn emits_replays_and_queries_atoms() {
1104 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1105 let first = atom(1, vec![]);
1106 let second = atom(2, vec![first.id]);
1107
1108 let refs = store.emit_atoms(&[first.clone(), second.clone()]).unwrap();
1109 assert_eq!(refs.len(), 2);
1110 assert_eq!(refs[0].commit, "site-a:1");
1111 assert_eq!(refs[1].commit, "site-a:2");
1112 assert_eq!(store.get_atom(first.id).unwrap(), first);
1113 assert_eq!(
1114 store.atom_by_content_hash(second.id).unwrap(),
1115 Some(second.clone())
1116 );
1117 assert_eq!(
1118 store.atoms_with_parent(first.id).unwrap(),
1119 vec![second.clone()]
1120 );
1121 assert_eq!(
1122 store
1123 .atoms_for_principal_persona("user:alice", "ship-captain")
1124 .unwrap(),
1125 vec![first, second]
1126 );
1127 }
1128
1129 #[test]
1130 fn derives_and_replays_parent_closed_slices() {
1131 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1132 let first = atom(1, vec![]);
1133 let second = atom(2, vec![first.id]);
1134 store.emit_atoms(&[first.clone(), second.clone()]).unwrap();
1135
1136 let slice = store.derive_slice(&[second.id]).unwrap();
1137 assert_eq!(slice.atoms, vec![first.id, second.id]);
1138 let receipt = store.ship_slice(&slice).unwrap();
1139 assert_eq!(receipt.slice_id, slice.id);
1140 assert_eq!(store.replay_slice(&slice).unwrap(), vec![first, second]);
1141 }
1142
1143 #[test]
1144 fn state_vector_delta_round_trips_between_replicas() {
1145 let source = SqliteFlowStore::in_memory("site-a").unwrap();
1146 let replica = SqliteFlowStore::in_memory("site-b").unwrap();
1147 let first = atom(1, vec![]);
1148 let second = atom(2, vec![first.id]);
1149 source.emit_atoms(&[first, second.clone()]).unwrap();
1150
1151 let empty = replica.state_vector("user:alice", "ship-captain").unwrap();
1152 let delta = source
1153 .causal_delta("user:alice", "ship-captain", &empty)
1154 .unwrap();
1155 assert_eq!(delta.len(), 2);
1156 for item in &delta {
1157 replica
1158 .insert_remote_atom(&item.atom, &item.site_id, item.clock)
1159 .unwrap();
1160 }
1161
1162 let vector = replica.state_vector("user:alice", "ship-captain").unwrap();
1163 assert_eq!(vector.clock("site-a"), 2);
1164 assert!(source
1165 .causal_delta("user:alice", "ship-captain", &vector)
1166 .unwrap()
1167 .is_empty());
1168 assert_eq!(replica.get_atom(second.id).unwrap(), second);
1169 }
1170
1171 #[test]
1172 fn persists_intents_and_derived_slices() {
1173 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1174 let first = atom(1, vec![]);
1175 store.emit_atom(&first).unwrap();
1176
1177 let intent = Intent::new(
1178 vec![first.id],
1179 "ship the smallest possible change",
1180 crate::flow::TranscriptSpan::new("transcript-1", 1, 1).unwrap(),
1181 0.9,
1182 )
1183 .unwrap();
1184 store.put_intent(&intent).unwrap();
1185 assert_eq!(store.get_intent(intent.id).unwrap(), intent);
1186
1187 let mut coverage = CoverageMap::new();
1188 coverage.insert(first.id, TestId::new("flow-store"));
1189 let slice = Slice {
1190 id: SliceId([3; 32]),
1191 atoms: vec![first.id],
1192 intents: vec![intent.id],
1193 invariants_applied: vec![(
1194 PredicateHash::new("pred"),
1195 crate::flow::InvariantResult::allow(),
1196 )],
1197 required_tests: vec![TestId::new("flow-store")],
1198 approval_chain: vec![Approval {
1199 reviewer: "alice".to_string(),
1200 approved_at: "2026-04-25T00:00:00Z".to_string(),
1201 reason: None,
1202 signature: None,
1203 }],
1204 base_ref: first.id,
1205 status: SliceStatus::Ready,
1206 };
1207 store.put_derived_slice(&slice).unwrap();
1208 assert_eq!(store.get_derived_slice(slice.id).unwrap(), slice);
1209 }
1210
1211 #[test]
1212 fn lists_only_shipped_derived_slices_for_replay_audit() {
1213 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1214 let first = atom(1, vec![]);
1215 store.emit_atom(&first).unwrap();
1216
1217 let shipped = Slice {
1218 id: SliceId([4; 32]),
1219 atoms: vec![first.id],
1220 intents: Vec::new(),
1221 invariants_applied: vec![(
1222 PredicateHash::new("sha256:retro"),
1223 crate::flow::InvariantResult::allow(),
1224 )],
1225 required_tests: vec![TestId::new("flow-store")],
1226 approval_chain: Vec::new(),
1227 base_ref: first.id,
1228 status: SliceStatus::Ready,
1229 };
1230 let unshipped = Slice {
1231 id: SliceId([5; 32]),
1232 atoms: vec![first.id],
1233 intents: Vec::new(),
1234 invariants_applied: Vec::new(),
1235 required_tests: Vec::new(),
1236 approval_chain: Vec::new(),
1237 base_ref: first.id,
1238 status: SliceStatus::Ready,
1239 };
1240
1241 store.put_shipped_derived_slice(&shipped).unwrap();
1242 store.put_derived_slice(&unshipped).unwrap();
1243
1244 let rows = store.shipped_derived_slices_since(None).unwrap();
1245 assert_eq!(rows.len(), 1);
1246 assert_eq!(rows[0].slice, shipped);
1247 assert!(!rows[0].created_at.is_empty());
1248 }
1249
1250 #[test]
1251 fn atoms_are_append_only_at_sql_boundary() {
1252 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1253 let first = atom(1, vec![]);
1254 store.emit_atom(&first).unwrap();
1255
1256 let conn = store.lock_conn().unwrap();
1257 let error = conn
1258 .execute(
1259 "DELETE FROM atoms WHERE id = ?1",
1260 params![first.id.0.as_slice()],
1261 )
1262 .unwrap_err();
1263 assert!(error.to_string().contains("atoms are append-only"));
1264 }
1265
1266 #[test]
1267 fn slices_are_append_only_at_sql_boundary() {
1268 let store = SqliteFlowStore::in_memory("site-a").unwrap();
1269 let first = atom(1, vec![]);
1270 store.emit_atom(&first).unwrap();
1271 let slice = Slice {
1272 id: SliceId([6; 32]),
1273 atoms: vec![first.id],
1274 intents: Vec::new(),
1275 invariants_applied: Vec::new(),
1276 required_tests: Vec::new(),
1277 approval_chain: Vec::new(),
1278 base_ref: first.id,
1279 status: SliceStatus::Ready,
1280 };
1281 store.put_shipped_derived_slice(&slice).unwrap();
1282
1283 let conn = store.lock_conn().unwrap();
1284 let error = conn
1285 .execute(
1286 "DELETE FROM slices WHERE id = ?1",
1287 params![slice.id.0.as_slice()],
1288 )
1289 .unwrap_err();
1290 assert!(error.to_string().contains("slices are append-only"));
1291 }
1292}