1use std::path::{Path, PathBuf};
2use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
3use std::sync::{Arc, Mutex, MutexGuard};
4use std::time::Duration;
5
6use chrono::{DateTime, SecondsFormat, Utc};
7use onlyne_proto::text::SchemaMismatch;
8use onlyne_proto::{
9 Envelope, Event, FaultEvent, LedgerQuery, LedgerState, LedgerStateEvent, Lifecycle, MsgKind,
10 Outcome, Principal, QueryFaultsArgs, QuerySessionsArgs,
11};
12use rusqlite::types::{Type, Value as SqlValue};
13use rusqlite::{Connection, OpenFlags, OptionalExtension, Row, params, params_from_iter};
14use serde::{Deserialize, Serialize};
15use serde_json::Value;
16use unicode_segmentation::UnicodeSegmentation;
17
18use crate::error::{StoreError, StoreResult};
19use crate::liveness::LiveSessions;
20use crate::transition_allowed;
21
22const SERVER_SCHEMA_VERSION: i64 = 6;
48const PROTOCOL_VERSION: i64 = 1;
49const SERVER_MARKER: &str = "onlyne-server";
50const DEFAULT_LIMIT: i64 = 100;
51
52pub const SERVER_DDL: &str = r#"CREATE TABLE IF NOT EXISTS roles(
53 name TEXT PRIMARY KEY,
54 key TEXT NOT NULL,
55 admin INTEGER NOT NULL,
56 max_sessions INTEGER NOT NULL,
57 spec_hash TEXT NOT NULL,
58 updated_at TEXT NOT NULL
59);
60CREATE TABLE IF NOT EXISTS sessions(
61 session_id TEXT PRIMARY KEY,
62 role TEXT NOT NULL,
63 generation INTEGER NOT NULL,
64 seq INTEGER NOT NULL,
65 agent_state TEXT NOT NULL,
66 delivery_state TEXT NOT NULL,
67 resource_state TEXT NOT NULL,
68 recovery_substate TEXT NOT NULL,
69 -- The (generation, seq) gate. The session reducer's isolate-after-N and
70 -- terminate-after-N policy needs a persisted counter, so this pair carries
71 -- DEFAULT_ISOLATE_AFTER and DEFAULT_TERMINATE_AFTER.
72 desired_json TEXT NOT NULL,
73 -- The client's published projection whole, lifecycle included: there is no
74 -- column beside it to fall back to, and every reader of the mirror parses the
75 -- lifecycle out of these bytes.
76 observed_json TEXT NOT NULL,
77 mismatch_count INTEGER NOT NULL,
78 -- The mirror's own freshness: what a reader judges a stale row by. v1's
79 -- mirror could be hours old and read as current.
80 last_seen TEXT NOT NULL,
81 -- When the projection content last moved. A beat that only refreshes
82 -- `last_seen` leaves it alone.
83 updated_at TEXT NOT NULL
84);
85-- The role-addressed reads: one role's sessions, which the note rule and the
86-- stale scan ask for.
87CREATE INDEX IF NOT EXISTS sessions_role_idx ON sessions(role);
88-- The order `list_sessions` reads in, ascending so the read walks it backwards:
89-- a descending index does not satisfy `updated_at DESC, rowid DESC` — SQLite
90-- leaves a `TEMP B-TREE FOR LAST TERM` and spills it to a temporary file on
91-- every read, and a reader that polls once a second turns that into megabytes
92-- per second of writes nothing asked for.
93CREATE INDEX IF NOT EXISTS sessions_updated_idx ON sessions(updated_at);
94-- Which delivery a session serves, and which sessions have served a delivery.
95-- The mirror row answers for a session; this table is the only place a
96-- delivery binding lives, so a session that served one delivery and then
97-- another is one row here twice rather than two mirror rows.
98CREATE TABLE IF NOT EXISTS session_tasks(
99 session_id TEXT NOT NULL,
100 task_id TEXT NOT NULL,
101 bound_at TEXT NOT NULL,
102 released_at TEXT,
103 PRIMARY KEY (session_id, task_id)
104);
105-- The reverse read: the session serving a task.
106CREATE INDEX IF NOT EXISTS session_tasks_task_idx ON session_tasks(task_id);
107-- The open binding of one session: at most one row per session is unreleased,
108-- because a session serves one delivery at a time.
109CREATE INDEX IF NOT EXISTS session_tasks_open_idx ON session_tasks(session_id, released_at);
110CREATE TABLE IF NOT EXISTS ledger(
111 msg_id TEXT PRIMARY KEY,
112 op_id TEXT UNIQUE,
113 fingerprint TEXT,
114 kind TEXT NOT NULL,
115 from_json TEXT NOT NULL,
116 to_json TEXT NOT NULL,
117 task TEXT,
118 parent_task TEXT,
119 attempt INTEGER NOT NULL,
120 state TEXT NOT NULL,
121 out_head TEXT,
122 reason TEXT,
123 enqueued_at TEXT NOT NULL,
124 acked_at TEXT,
125 -- Nullable because retention pruning clears an acknowledged body after the
126 -- cutoff.
127 body_json TEXT,
128 -- Causality's hop count from the root task. `parent_task` alone gives the
129 -- chain's shape, not its depth; `onlyne handoff` reads both back to extend
130 -- the chain, so the row stores the counter beside the link.
131 hop INTEGER NOT NULL DEFAULT 0,
132 -- Persisted expiry deadline of a ttl note, so a restarted server can re-arm
133 -- its sweep.
134 expires_at TEXT,
135 -- Times this row has moved from in_flight back to queued.
136 requeued INTEGER NOT NULL DEFAULT 0,
137 -- The family's root task id, read off the envelope's causality. `parent_task`
138 -- gives the chain's shape, and this names the arc every hop of one run
139 -- belongs to, which `onlyne ledger` prints beside the hop.
140 family TEXT,
141 -- The hops the family may spend, set by whoever started the run.
142 hop_budget INTEGER,
143 -- The role the family reports home to, carried on every row of the run.
144 origin TEXT,
145 -- Wall-clock bound for the whole family, in the RFC 3339 shape `expires_at`
146 -- uses.
147 deadline TEXT,
148 -- The causality's free-form labels as JSON text.
149 labels_json TEXT
150);
151CREATE INDEX IF NOT EXISTS ledger_state_enqueued_idx ON ledger(state,enqueued_at);
152CREATE INDEX IF NOT EXISTS ledger_task_idx ON ledger(task);
153CREATE INDEX IF NOT EXISTS ledger_kind_state_idx ON ledger(kind,state);
154-- The ledger listing's own reading order, on the same terms as the sessions
155-- index above. A listing filtered by `state` keeps using `ledger_state_enqueued_idx`
156-- and reads that backwards.
157CREATE INDEX IF NOT EXISTS ledger_enqueued_idx ON ledger(enqueued_at);
158CREATE TABLE IF NOT EXISTS events(
159 seq INTEGER PRIMARY KEY,
160 type TEXT NOT NULL,
161 data_json TEXT NOT NULL,
162 created_at TEXT NOT NULL
163);
164CREATE INDEX IF NOT EXISTS events_type_idx ON events(type);
165CREATE TABLE IF NOT EXISTS faults(
166 id INTEGER PRIMARY KEY AUTOINCREMENT,
167 task_id TEXT,
168 role TEXT,
169 session_id TEXT,
170 generation INTEGER,
171 seq INTEGER,
172 desired_json TEXT,
173 observed_json TEXT,
174 intent TEXT,
175 attempt INTEGER,
176 backend_ref TEXT,
177 kind TEXT NOT NULL,
178 reason TEXT NOT NULL,
179 state TEXT NOT NULL,
180 -- Encoded from the kernel's unix seconds through this crate's own helper on
181 -- every write.
182 created_at TEXT NOT NULL
183);
184CREATE INDEX IF NOT EXISTS faults_task_kind_generation_idx ON faults(task_id,kind,generation);
185CREATE INDEX IF NOT EXISTS faults_state_idx ON faults(state);
186-- The ghost sweep's own audit trail, one row per `working` mirror row the server
187-- settled because the task's ledger row had already reached a terminal state.
188-- `seq_before` and `seq_after` are the mirror row's two versions, so an operator
189-- reads the exact write the pass made straight off this row.
190CREATE TABLE IF NOT EXISTS ghost_sweeps(
191 id INTEGER PRIMARY KEY AUTOINCREMENT,
192 task_id TEXT NOT NULL,
193 role TEXT NOT NULL,
194 session_id TEXT NOT NULL,
195 generation INTEGER NOT NULL,
196 seq_before INTEGER NOT NULL,
197 seq_after INTEGER NOT NULL,
198 outcome TEXT NOT NULL,
199 evidence TEXT NOT NULL,
200 swept_at TEXT NOT NULL
201);
202-- The listing's own reading order, on the same terms as the sessions index
203-- above: ascending, so `swept_at DESC, rowid DESC` walks it backwards.
204CREATE INDEX IF NOT EXISTS ghost_sweeps_swept_at_idx ON ghost_sweeps(swept_at);
205CREATE TABLE IF NOT EXISTS inbox_cursors(
206 role TEXT PRIMARY KEY,
207 last_msg_id TEXT,
208 last_seq INTEGER NOT NULL,
209 updated_at TEXT NOT NULL
210);
211CREATE TABLE IF NOT EXISTS hook_cursors(
212 hook TEXT PRIMARY KEY,
213 last_seq INTEGER NOT NULL,
214 updated_at TEXT NOT NULL
215);"#;
216
217const SCHEMA_MARKER_DDL: &str = "CREATE TABLE IF NOT EXISTS schema_marker(name TEXT PRIMARY KEY, version INTEGER NOT NULL, protocol_version INTEGER NOT NULL);";
218
219#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
220pub struct RoleRow {
221 pub name: String,
222 pub key: String,
223 pub admin: bool,
224 pub max_sessions: i64,
225 pub spec_hash: String,
226 pub updated_at: String,
227}
228
229#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
238pub struct SessionWrite {
239 pub session_id: String,
240 pub task_id: Option<String>,
241 pub role: String,
242 pub generation: i64,
243 pub seq: i64,
244 pub agent_state: String,
245 pub delivery_state: String,
246 pub resource_state: String,
247 pub recovery_substate: String,
248 pub desired_json: String,
249 pub observed_json: String,
250 pub mismatch_count: i64,
251 pub last_seen: i64,
254 pub updated_at: i64,
256}
257
258pub type ServerSessionRow = SessionWrite;
259
260#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
263pub struct SessionBindingRow {
264 pub session_id: String,
265 pub task_id: String,
266 pub bound_at: i64,
268 pub released_at: Option<i64>,
271}
272
273#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
274pub struct LedgerRow {
275 pub msg_id: String,
276 pub op_id: Option<String>,
277 pub fingerprint: Option<String>,
278 pub kind: MsgKind,
279 pub from_json: String,
280 pub to_json: String,
281 pub task: Option<String>,
282 pub parent_task: Option<String>,
283 pub attempt: i64,
284 pub state: LedgerState,
285 pub out_head: Option<String>,
286 pub reason: Option<String>,
287 pub enqueued_at: String,
288 pub acked_at: Option<String>,
289 pub body_json: Option<String>,
290 pub hop: i64,
293 pub expires_at: Option<String>,
296 #[serde(default)]
298 pub requeued: i64,
299 #[serde(default)]
303 pub family: Option<String>,
304 #[serde(default)]
306 pub hop_budget: Option<i64>,
307 #[serde(default)]
309 pub origin: Option<String>,
310 #[serde(default)]
313 pub deadline: Option<String>,
314 #[serde(default)]
317 pub labels_json: Option<String>,
318}
319
320impl LedgerRow {
321 pub fn from_envelope(envelope: &Envelope, fingerprint: &str) -> StoreResult<Self> {
322 let body_json = serde_json::to_string(&envelope.body)?;
323 let causality = envelope.causality.as_ref();
324 let expires_at = match (envelope.kind, envelope.ttl_ms) {
325 (MsgKind::Note, Some(ttl)) => Some(rfc3339(
326 envelope.ts + chrono::Duration::milliseconds(ttl as i64),
327 )),
328 _ => None,
329 };
330 Ok(Self {
331 msg_id: envelope.id.clone(),
332 op_id: envelope.op_id.clone(),
333 fingerprint: Some(fingerprint.to_string()),
334 kind: envelope.kind,
335 from_json: sender_column(&envelope.from, envelope.admin)?,
336 to_json: serde_json::to_string(&envelope.to)?,
337 task: causality.map(|c| c.task.clone()),
338 parent_task: causality.and_then(|c| c.parent_task.clone()),
339 attempt: causality.map(|c| i64::from(c.attempt)).unwrap_or(0),
340 hop: causality.map(|c| i64::from(c.hop)).unwrap_or(0),
341 state: LedgerState::Queued,
342 out_head: Some(body_head(envelope, &body_json)),
343 reason: None,
344 enqueued_at: rfc3339(envelope.ts),
345 acked_at: None,
346 body_json: Some(body_json),
347 expires_at,
348 requeued: 0,
349 family: causality.and_then(|c| c.family.clone()),
350 hop_budget: causality.and_then(|c| c.hop_budget).map(i64::from),
351 origin: causality.and_then(|c| c.origin.clone()),
352 deadline: causality.and_then(|c| c.deadline).map(rfc3339),
353 labels_json: match causality.and_then(|c| c.labels.as_ref()) {
354 Some(labels) => Some(serde_json::to_string(labels)?),
355 None => None,
356 },
357 })
358 }
359}
360
361impl LedgerRow {
362 pub fn sender(&self) -> Result<Principal, serde_json::Error> {
368 let value: Value = serde_json::from_str(&self.from_json)?;
369 serde_json::from_value(value.get("principal").cloned().unwrap_or(value))
370 }
371
372 pub fn sender_is_admin(&self) -> bool {
374 serde_json::from_str::<Value>(&self.from_json)
375 .ok()
376 .and_then(|value| value.get("admin").and_then(Value::as_bool))
377 .unwrap_or(false)
378 }
379}
380
381fn sender_column(from: &Principal, admin: bool) -> StoreResult<String> {
383 Ok(serde_json::to_string(&serde_json::json!({
384 "admin": admin,
385 "principal": from,
386 }))?)
387}
388
389#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
390pub enum Append {
391 Accepted(LedgerRow),
392 Duplicate {
393 existing: LedgerRow,
394 fingerprint_matches: bool,
395 },
396}
397
398#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
399pub struct EventRecord {
400 pub seq: i64,
401 pub kind: String,
402 pub data: Value,
403 pub created_at: String,
404}
405
406#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
407pub struct ServerFaultRow {
408 pub id: i64,
409 pub task_id: Option<String>,
410 pub role: Option<String>,
411 pub session_id: Option<String>,
412 pub generation: Option<i64>,
413 pub seq: Option<i64>,
414 pub desired_json: Option<String>,
415 pub observed_json: Option<String>,
416 pub intent: Option<String>,
417 pub attempt: Option<i64>,
418 pub backend_ref: Option<String>,
419 pub kind: String,
420 pub reason: String,
421 pub state: String,
422 pub created_at: i64,
423}
424
425#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
426pub struct FaultQuery {
427 pub task_id: Option<String>,
428 pub role: Option<String>,
429 pub kind: Option<String>,
430 pub open_only: bool,
431 pub limit: u32,
432}
433
434#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
442pub struct GhostSweepRow {
443 pub id: i64,
444 pub task_id: String,
445 pub role: String,
446 pub session_id: String,
447 pub generation: i64,
448 pub seq_before: i64,
449 pub seq_after: i64,
450 pub outcome: Outcome,
452 pub evidence: String,
454 pub swept_at: i64,
455}
456
457#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
458pub struct CursorRow {
459 pub role: String,
460 pub last_msg_id: Option<String>,
461 pub last_seq: i64,
462 pub updated_at: String,
463}
464
465#[derive(Clone, Debug)]
466pub struct ServerLedger {
467 path: PathBuf,
468 retention_days: i64,
469 inner: Arc<Mutex<Connection>>,
471 readers: ReadPool,
474 session_rows_written: Arc<AtomicU64>,
478 live: Arc<LiveSessions>,
480}
481
482impl ServerLedger {
483 pub fn open(path: impl AsRef<Path>, retention_days: u32) -> StoreResult<Self> {
484 let path = path.as_ref().to_path_buf();
485 let conn = open_connection(
486 &path,
487 "server",
488 SERVER_MARKER,
489 SERVER_DDL,
490 SERVER_SCHEMA_VERSION,
491 )?;
492 let readers = ReadPool::open(&path)?;
496 Ok(Self {
497 path,
498 retention_days: i64::from(retention_days).max(1),
499 inner: Arc::new(Mutex::new(conn)),
500 readers,
501 session_rows_written: Arc::new(AtomicU64::new(0)),
502 live: Arc::new(LiveSessions::default()),
503 })
504 }
505
506 pub fn path(&self) -> &Path {
507 &self.path
508 }
509
510 pub fn retention_days(&self) -> i64 {
511 self.retention_days
512 }
513
514 pub fn upsert_role(&self, role: &RoleRow) -> StoreResult<bool> {
515 let conn = self.conn()?;
516 let changed = conn.execute(
517 "INSERT INTO roles(name,key,admin,max_sessions,spec_hash,updated_at) VALUES(?,?,?,?,?,?)
518 ON CONFLICT(name) DO UPDATE SET key=excluded.key,admin=excluded.admin,max_sessions=excluded.max_sessions,spec_hash=excluded.spec_hash,updated_at=excluded.updated_at",
519 params![
520 role.name,
521 role.key,
522 bool_int(role.admin),
523 role.max_sessions,
524 role.spec_hash,
525 unix_to_rfc3339(rfc3339_to_unix(&role.updated_at))
526 ],
527 )?;
528 Ok(changed == 1)
529 }
530
531 pub fn list_roles(&self) -> StoreResult<Vec<RoleRow>> {
532 let conn = self.read()?;
533 let rows = conn
534 .prepare(
535 "SELECT name,key,admin,max_sessions,spec_hash,updated_at FROM roles ORDER BY name",
536 )?
537 .query_map([], role_row)?
538 .collect::<Result<Vec<_>, _>>()?;
539 Ok(rows)
540 }
541
542 pub fn remove_role_missing_from(&self, names: &[String]) -> StoreResult<usize> {
543 let conn = self.conn()?;
544 if names.is_empty() {
545 return Ok(conn.execute("DELETE FROM roles", [])?);
546 }
547 let placeholders = repeat_placeholders(names.len());
548 let sql = format!("DELETE FROM roles WHERE name NOT IN ({placeholders})");
549 let args = names
550 .iter()
551 .cloned()
552 .map(SqlValue::Text)
553 .collect::<Vec<_>>();
554 Ok(conn.execute(&sql, params_from_iter(args))?)
555 }
556
557 pub fn project_session(&self, write: &SessionWrite) -> StoreResult<bool> {
570 let conn = self.conn()?;
571 self.note_session_row_write();
572 let tx = conn.unchecked_transaction()?;
573 let changed = project_session_conn(&tx, write)?;
574 tx.commit()?;
575 if changed {
576 self.live.settle(&write.session_id, write.last_seen);
577 }
578 Ok(changed)
579 }
580
581 pub fn rebind_session(&self, from_session_id: &str, write: &SessionWrite) -> StoreResult<bool> {
590 let conn = self.conn()?;
591 self.note_session_row_write();
592 let tx = conn.unchecked_transaction()?;
593 if from_session_id != write.session_id {
594 tx.execute(
595 "DELETE FROM session_tasks WHERE session_id=?",
596 params![write.session_id],
597 )?;
598 tx.execute(
599 "DELETE FROM sessions WHERE session_id=?",
600 params![write.session_id],
601 )?;
602 tx.execute(
603 "UPDATE session_tasks SET session_id=? WHERE session_id=?",
604 params![write.session_id, from_session_id],
605 )?;
606 tx.execute(
607 "UPDATE sessions SET session_id=? WHERE session_id=?",
608 params![write.session_id, from_session_id],
609 )?;
610 }
611 let changed = project_session_conn(&tx, write)?;
612 tx.commit()?;
613 if changed {
614 self.live.settle(from_session_id, write.last_seen);
617 self.live.settle(&write.session_id, write.last_seen);
618 }
619 Ok(changed)
620 }
621
622 pub fn open_binding(
626 &self,
627 session_id: &str,
628 task_id: &str,
629 bound_at: i64,
630 ) -> StoreResult<bool> {
631 let conn = self.conn()?;
632 Ok(open_binding_conn(&conn, session_id, task_id, bound_at)? > 0)
633 }
634
635 pub fn release_binding(
640 &self,
641 session_id: &str,
642 task_id: &str,
643 released_at: i64,
644 ) -> StoreResult<bool> {
645 let conn = self.conn()?;
646 Ok(release_binding_conn(&conn, session_id, task_id, released_at)? == 1)
647 }
648
649 pub fn open_binding_of(&self, session_id: &str) -> StoreResult<Option<SessionBindingRow>> {
654 let conn = self.read()?;
655 Ok(conn
656 .query_row(
657 "SELECT session_id,task_id,bound_at,released_at FROM session_tasks WHERE session_id=? AND released_at IS NULL ORDER BY bound_at DESC,task_id DESC LIMIT 1",
658 params![session_id],
659 session_binding_row,
660 )
661 .optional()?)
662 }
663
664 pub fn publish_mirror_outcome(
673 &self,
674 session_id: &str,
675 observed_json: &str,
676 expected_observed_json: &str,
677 updated_at: i64,
678 ) -> StoreResult<bool> {
679 let conn = self.conn()?;
680 self.note_session_row_write();
681 let changed = conn.execute(
682 "UPDATE sessions SET observed_json=?,updated_at=? WHERE session_id=? AND observed_json=?",
683 params![
684 observed_json,
685 unix_to_rfc3339(updated_at),
686 session_id,
687 expected_observed_json
688 ],
689 )?;
690 Ok(changed == 1)
691 }
692
693 pub fn beat_session(
711 &self,
712 session_id: &str,
713 at: i64,
714 flush_after_secs: i64,
715 ) -> StoreResult<bool> {
716 self.live.note(session_id, at);
719 let Some(stored) = self.stored_session_row(session_id)? else {
720 return Ok(false);
725 };
726 if at.saturating_sub(stored.last_seen) < flush_after_secs.max(1) {
727 return Ok(false);
728 }
729 let flushed = self.flush_last_seen(session_id, at)?;
730 if flushed {
731 self.live.settle(session_id, at);
732 }
733 Ok(flushed)
734 }
735
736 fn flush_last_seen(&self, session_id: &str, last_seen: i64) -> StoreResult<bool> {
745 let conn = self.conn()?;
746 self.note_session_row_write();
747 let changed = conn.execute(
748 "UPDATE sessions SET last_seen=? WHERE session_id=?",
749 params![unix_to_rfc3339(last_seen), session_id],
750 )?;
751 Ok(changed == 1)
752 }
753
754 pub fn get_session_row(&self, session_id: &str) -> StoreResult<Option<ServerSessionRow>> {
757 Ok(self
758 .stored_session_row(session_id)?
759 .map(|row| self.with_live_seen(row)))
760 }
761
762 fn stored_session_row(&self, session_id: &str) -> StoreResult<Option<ServerSessionRow>> {
768 let conn = self.read()?;
769 let row = conn
770 .query_row(
771 &format!(
772 "SELECT {} FROM sessions WHERE session_id=?",
773 session_columns()
774 ),
775 params![session_id],
776 session_row,
777 )
778 .optional()?;
779 Ok(row)
780 }
781
782 fn with_live_seen(&self, mut row: ServerSessionRow) -> ServerSessionRow {
789 row.last_seen = self.live.freshest(&row.session_id, row.last_seen);
790 row
791 }
792
793 pub fn session_row_for_task(&self, task_id: &str) -> StoreResult<Option<ServerSessionRow>> {
798 let conn = self.read()?;
799 let row = conn
800 .query_row(
801 &format!(
802 "SELECT {} FROM sessions WHERE session_id={SESSION_ID_FOR_TASK}",
803 session_columns()
804 ),
805 params![task_id],
806 session_row,
807 )
808 .optional()?;
809 Ok(row.map(|row| self.with_live_seen(row)))
810 }
811
812 pub fn list_sessions(&self, filter: QuerySessionsArgs) -> StoreResult<Vec<ServerSessionRow>> {
813 let conn = self.read()?;
814 let mut clauses = Vec::new();
815 let mut args = Vec::new();
816 if let Some(task_id) = filter.task_id {
817 clauses.push(
822 "session_id IN (SELECT session_id FROM session_tasks WHERE task_id=?)".to_string(),
823 );
824 args.push(SqlValue::Text(task_id));
825 }
826 if let Some(role) = filter.role {
827 clauses.push("role=?".to_string());
828 args.push(SqlValue::Text(role));
829 }
830 if let Some(lifecycle) = filter.lifecycle {
831 clauses.push(
842 "IFNULL(CASE WHEN json_valid(observed_json) THEN json_extract(observed_json,'$.lifecycle') END,?)=?".to_string(),
843 );
844 args.push(SqlValue::Text(string_tag(&Lifecycle::Created)?));
845 args.push(SqlValue::Text(string_tag(&lifecycle)?));
846 }
847 let where_sql = where_sql(&clauses);
848 let limit = sql_limit(filter.limit);
849 args.push(SqlValue::Integer(limit));
850 let sql = sessions_list_sql(&where_sql);
851 let rows = conn
852 .prepare(&sql)?
853 .query_map(params_from_iter(args), session_row)?
854 .collect::<Result<Vec<_>, _>>()?;
855 Ok(rows
856 .into_iter()
857 .map(|row| self.with_live_seen(row))
858 .collect())
859 }
860
861 pub fn append_ledger(&self, row: &LedgerRow) -> StoreResult<Append> {
862 let conn = self.conn()?;
863 if let Some(op_id) = row.op_id.as_deref() {
864 if let Some(existing) = ledger_by_op_id(&conn, op_id)? {
865 return Ok(Append::Duplicate {
866 fingerprint_matches: existing.fingerprint == row.fingerprint,
867 existing,
868 });
869 }
870 }
871 insert_ledger_row(&conn, row)?;
872 Ok(Append::Accepted(row.clone()))
873 }
874
875 pub fn mark_in_flight(&self, msg_id: &str) -> StoreResult<bool> {
876 self.transition_msg(msg_id, LedgerState::InFlight, None, None)
877 }
878
879 pub fn mark_acked(&self, msg_id: &str, at: DateTime<Utc>) -> StoreResult<bool> {
880 self.transition_msg(msg_id, LedgerState::Acked, Some(rfc3339(at)), None)
881 }
882
883 pub fn mark_rejected(&self, msg_id: &str, reason: &str) -> StoreResult<bool> {
884 self.transition_msg(
885 msg_id,
886 LedgerState::Rejected,
887 None,
888 Some(reason.to_string()),
889 )
890 }
891
892 pub fn expire_queued_before(&self, now: DateTime<Utc>) -> StoreResult<usize> {
893 ensure_transition_allowed(LedgerState::Queued, LedgerState::Expired)?;
894 let conn = self.conn()?;
895 let changed = conn.execute(
896 "UPDATE ledger SET state='expired',reason='expired' WHERE state='queued' AND enqueued_at < ?",
897 params![rfc3339(now)],
898 )?;
899 Ok(changed)
900 }
901
902 pub fn queued_for(&self, role: &str, limit: u32) -> StoreResult<Vec<LedgerRow>> {
903 self.ledger_for_role(role, LedgerState::Queued, limit)
904 }
905
906 pub fn queued_count_for(&self, role: &str) -> StoreResult<u32> {
914 let conn = self.read()?;
915 let count = conn
916 .query_row(
917 "SELECT COUNT(*) FROM ledger WHERE state=? AND json_extract(to_json,'$.role.role')=? AND kind<>?",
918 params![LedgerState::Queued.as_str(), role, MsgKind::Note.as_str()],
919 |row| row.get::<_, u32>(0),
920 )
921 .optional()?
922 .unwrap_or(0);
923 Ok(count)
924 }
925
926 pub fn in_flight_for(&self, role: &str) -> StoreResult<Vec<LedgerRow>> {
927 self.ledger_for_role(role, LedgerState::InFlight, 500)
928 }
929
930 pub fn pending_expiries(&self) -> StoreResult<Vec<(String, DateTime<Utc>)>> {
933 let conn = self.read()?;
934 let rows = conn
935 .prepare(
936 "SELECT msg_id, expires_at FROM ledger WHERE state IN ('queued','in_flight') AND expires_at IS NOT NULL ORDER BY expires_at,rowid",
937 )?
938 .query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?
939 .collect::<Result<Vec<_>, _>>()?;
940 Ok(rows
941 .into_iter()
942 .filter_map(|(msg_id, expires_at)| {
943 DateTime::parse_from_rfc3339(&expires_at)
944 .ok()
945 .map(|deadline| (msg_id, deadline.with_timezone(&Utc)))
946 })
947 .collect())
948 }
949
950 pub fn requeue_in_flight(&self, role: &str) -> StoreResult<usize> {
951 ensure_transition_allowed(LedgerState::InFlight, LedgerState::Queued)?;
952 let conn = self.conn()?;
953 let changed = conn.execute(
954 "UPDATE ledger SET state='queued' WHERE state='in_flight' AND json_extract(to_json,'$.role.role')=?",
955 params![role],
956 )?;
957 Ok(changed)
958 }
959
960 pub fn requeue_one(&self, msg_id: &str) -> StoreResult<LedgerRow> {
964 let conn = self.conn()?;
965 let current = ledger_by_msg_id(&conn, msg_id)?;
966 if current.state == LedgerState::Queued {
967 return Ok(current);
968 }
969 ensure_transition_allowed(current.state, LedgerState::Queued)?;
970 let tx = conn.unchecked_transaction()?;
971 tx.execute(
972 "UPDATE ledger SET state='queued',requeued=requeued+1 WHERE msg_id=?",
973 params![msg_id],
974 )?;
975 let updated = ledger_by_msg_id(&tx, msg_id)?;
976 let event = ledger_state_event(&updated, LedgerState::Queued, None)?;
977 append_event_conn(&tx, "ledger_state", &event)?;
978 tx.commit()?;
979 Ok(updated)
980 }
981
982 pub fn expire_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow> {
990 let conn = self.conn()?;
991 let current = ledger_by_msg_id(&conn, msg_id)?;
992 if !matches!(current.state, LedgerState::Queued | LedgerState::InFlight) {
993 ensure_transition_allowed(current.state, LedgerState::Expired)?;
994 }
995 let tx = conn.unchecked_transaction()?;
996 tx.execute(
997 "UPDATE ledger SET state='expired',reason=? WHERE msg_id=?",
998 params![reason, msg_id],
999 )?;
1000 let updated = ledger_by_msg_id(&tx, msg_id)?;
1001 let event = ledger_state_event(&updated, LedgerState::Expired, Some(reason.to_string()))?;
1002 append_event_conn(&tx, "ledger_state", &event)?;
1003 tx.commit()?;
1004 Ok(updated)
1005 }
1006
1007 pub fn fail_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow> {
1011 let conn = self.conn()?;
1012 let current = ledger_by_msg_id(&conn, msg_id)?;
1013 ensure_transition_allowed(current.state, LedgerState::Rejected)?;
1014 let tx = conn.unchecked_transaction()?;
1015 tx.execute(
1016 "UPDATE ledger SET state='rejected',reason=? WHERE msg_id=?",
1017 params![reason, msg_id],
1018 )?;
1019 let updated = ledger_by_msg_id(&tx, msg_id)?;
1020 let event = ledger_state_event(&updated, LedgerState::Rejected, Some(reason.to_string()))?;
1021 append_event_conn(&tx, "ledger_state", &event)?;
1022 tx.commit()?;
1023 Ok(updated)
1024 }
1025
1026 pub fn update_fault_state(
1029 &self,
1030 task_id: &str,
1031 next_state: &str,
1032 reason: &str,
1033 ) -> StoreResult<Vec<ServerFaultRow>> {
1034 let conn = self.conn()?;
1035 let tx = conn.unchecked_transaction()?;
1036 let ids = {
1037 let mut stmt =
1038 tx.prepare("SELECT id FROM faults WHERE task_id=? AND state='open' ORDER BY id")?;
1039 stmt.query_map(params![task_id], |r| r.get::<_, i64>(0))?
1040 .collect::<Result<Vec<_>, _>>()?
1041 };
1042 let mut moved = Vec::with_capacity(ids.len());
1043 for id in ids {
1044 tx.execute(
1045 "UPDATE faults SET state=?,reason=? WHERE id=?",
1046 params![next_state, reason, id],
1047 )?;
1048 let row = fault_row_by_id(&tx, id)?;
1049 let event = fault_event(&row)?;
1050 append_event_conn(&tx, "fault", &event)?;
1051 moved.push(row);
1052 }
1053 tx.commit()?;
1054 Ok(moved)
1055 }
1056
1057 pub fn ledger_query(&self, query: LedgerQuery) -> StoreResult<Vec<LedgerRow>> {
1058 let conn = self.read()?;
1059 let mut clauses = Vec::new();
1060 let mut args = Vec::new();
1061 if let Some(task) = query.task {
1062 clauses.push("task=?".to_string());
1063 args.push(SqlValue::Text(task));
1064 }
1065 if let Some(op_id) = query.op_id {
1066 clauses.push("op_id=?".to_string());
1067 args.push(SqlValue::Text(op_id));
1068 }
1069 if let Some(msg_id) = query.msg_id {
1070 clauses.push("msg_id=?".to_string());
1071 args.push(SqlValue::Text(msg_id));
1072 }
1073 if let Some(role) = query.role {
1074 clauses.push("(json_extract(from_json,'$.principal.role.role')=? OR json_extract(to_json,'$.role.role')=?)".to_string());
1075 args.push(SqlValue::Text(role.clone()));
1076 args.push(SqlValue::Text(role));
1077 }
1078 if let Some(state) = query.state {
1079 clauses.push("state=?".to_string());
1080 args.push(SqlValue::Text(state.as_str().to_string()));
1081 }
1082 if let Some(kind) = query.kind {
1083 clauses.push("kind=?".to_string());
1084 args.push(SqlValue::Text(kind.as_str().to_string()));
1085 }
1086 args.push(SqlValue::Integer(sql_limit(query.limit)));
1087 let sql = ledger_list_sql(&where_sql(&clauses));
1088 let rows = conn
1089 .prepare(&sql)?
1090 .query_map(params_from_iter(args), ledger_row)?
1091 .collect::<Result<Vec<_>, _>>()?;
1092 Ok(rows)
1093 }
1094
1095 pub fn ledger_task(&self, task: &str, limit: u32) -> StoreResult<Vec<LedgerRow>> {
1103 let conn = self.read()?;
1104 let rows = conn
1105 .prepare(&format!(
1106 "SELECT {LEDGER_COLUMNS} FROM ledger WHERE task=? ORDER BY enqueued_at,rowid LIMIT ?"
1107 ))?
1108 .query_map(params![task, sql_limit(limit)], ledger_row)?
1109 .collect::<Result<Vec<_>, _>>()?;
1110 Ok(rows)
1111 }
1112
1113 pub fn append_event(&self, kind: &str, data: &Value) -> StoreResult<i64> {
1114 let conn = self.conn()?;
1115 append_event_conn(&conn, kind, data)
1116 }
1117
1118 pub fn events_since(&self, seq: i64, limit: u32) -> StoreResult<Vec<EventRecord>> {
1119 let conn = self.read()?;
1120 events_since_conn(&conn, seq, limit)
1121 }
1122
1123 pub fn event_head(&self) -> StoreResult<i64> {
1124 let conn = self.read()?;
1125 event_head_conn(&conn)
1126 }
1127
1128 pub fn prune(&self, older_than: DateTime<Utc>) -> StoreResult<usize> {
1129 let conn = self.conn()?;
1130 let changed = conn.execute(
1131 "UPDATE ledger SET body_json=NULL WHERE state='acked' AND acked_at IS NOT NULL AND acked_at < ?",
1132 params![rfc3339(older_than)],
1133 )?;
1134 Ok(changed)
1135 }
1136
1137 pub fn record_fault(&self, fault: &ServerFaultRow) -> StoreResult<i64> {
1138 let conn = self.conn()?;
1139 conn.execute(
1140 "INSERT INTO faults(task_id,role,session_id,generation,seq,desired_json,observed_json,intent,attempt,backend_ref,kind,reason,state,created_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
1141 params![
1142 fault.task_id,
1143 fault.role,
1144 fault.session_id,
1145 fault.generation,
1146 fault.seq,
1147 fault.desired_json,
1148 fault.observed_json,
1149 fault.intent,
1150 fault.attempt,
1151 fault.backend_ref,
1152 fault.kind,
1153 fault.reason,
1154 fault.state,
1155 unix_to_rfc3339(fault.created_at)
1156 ],
1157 )?;
1158 Ok(conn.last_insert_rowid())
1159 }
1160
1161 pub fn open_faults(&self) -> StoreResult<Vec<ServerFaultRow>> {
1162 self.faults_query(FaultQuery {
1163 open_only: true,
1164 limit: 500,
1165 ..FaultQuery::default()
1166 })
1167 }
1168
1169 pub fn ack_fault(&self, fault_id: i64) -> StoreResult<bool> {
1170 let conn = self.conn()?;
1171 Ok(conn.execute(
1172 "UPDATE faults SET state='acked' WHERE id=? AND state<>'acked'",
1173 params![fault_id],
1174 )? == 1)
1175 }
1176
1177 pub fn faults_query(&self, query: FaultQuery) -> StoreResult<Vec<ServerFaultRow>> {
1178 let conn = self.read()?;
1179 let mut clauses = Vec::new();
1180 let mut args = Vec::new();
1181 if let Some(task_id) = query.task_id {
1182 clauses.push("task_id=?".to_string());
1183 args.push(SqlValue::Text(task_id));
1184 }
1185 if let Some(role) = query.role {
1186 clauses.push("role=?".to_string());
1187 args.push(SqlValue::Text(role));
1188 }
1189 if let Some(kind) = query.kind {
1190 clauses.push("kind=?".to_string());
1191 args.push(SqlValue::Text(kind));
1192 }
1193 if query.open_only {
1194 clauses.push("state='open'".to_string());
1195 }
1196 args.push(SqlValue::Integer(sql_limit(query.limit)));
1197 let sql = format!(
1198 "SELECT {FAULT_COLUMNS} FROM faults{} ORDER BY id LIMIT ?",
1199 where_sql(&clauses)
1200 );
1201 let rows = conn
1202 .prepare(&sql)?
1203 .query_map(params_from_iter(args), server_fault_row)?
1204 .collect::<Result<Vec<_>, _>>()?;
1205 Ok(rows)
1206 }
1207
1208 pub fn faults_query_proto(&self, query: QueryFaultsArgs) -> StoreResult<Vec<ServerFaultRow>> {
1209 self.faults_query(FaultQuery {
1210 task_id: query.task_id,
1211 role: query.role,
1212 kind: query.kind,
1213 open_only: query.open_only,
1214 limit: query.limit,
1215 })
1216 }
1217
1218 pub fn record_ghost_sweep(&self, sweep: &GhostSweepRow) -> StoreResult<i64> {
1220 let conn = self.conn()?;
1221 conn.execute(
1222 "INSERT INTO ghost_sweeps(task_id,role,session_id,generation,seq_before,seq_after,outcome,evidence,swept_at) VALUES(?,?,?,?,?,?,?,?,?)",
1223 params![
1224 sweep.task_id,
1225 sweep.role,
1226 sweep.session_id,
1227 sweep.generation,
1228 sweep.seq_before,
1229 sweep.seq_after,
1230 string_tag(&sweep.outcome)?,
1231 sweep.evidence,
1232 unix_to_rfc3339(sweep.swept_at)
1233 ],
1234 )?;
1235 Ok(conn.last_insert_rowid())
1236 }
1237
1238 pub fn list_ghost_sweeps(&self, limit: u32) -> StoreResult<Vec<GhostSweepRow>> {
1245 let conn = self.read()?;
1246 let rows = conn
1247 .prepare(&ghost_sweeps_list_sql())?
1248 .query_map(params![sql_limit(limit)], ghost_sweep_row)?
1249 .collect::<Result<Vec<_>, _>>()?;
1250 Ok(rows)
1251 }
1252
1253 pub fn cursor_for(&self, role: &str) -> StoreResult<Option<CursorRow>> {
1254 let conn = self.read()?;
1255 Ok(conn
1256 .query_row(
1257 "SELECT role,last_msg_id,last_seq,updated_at FROM inbox_cursors WHERE role=?",
1258 params![role],
1259 cursor_row,
1260 )
1261 .optional()?)
1262 }
1263
1264 pub fn set_cursor(&self, role: &str, msg_id: Option<&str>, seq: i64) -> StoreResult<bool> {
1265 let conn = self.conn()?;
1266 let updated_at = rfc3339(Utc::now());
1267 Ok(conn.execute(
1268 "INSERT INTO inbox_cursors(role,last_msg_id,last_seq,updated_at) VALUES(?,?,?,?)
1269 ON CONFLICT(role) DO UPDATE SET last_msg_id=excluded.last_msg_id,last_seq=excluded.last_seq,updated_at=excluded.updated_at",
1270 params![role, msg_id, seq, updated_at],
1271 )? == 1)
1272 }
1273
1274 pub fn hook_cursor(&self, hook: &str) -> StoreResult<Option<i64>> {
1281 let conn = self.read()?;
1282 Ok(conn
1283 .query_row(
1284 "SELECT last_seq FROM hook_cursors WHERE hook=?",
1285 params![hook],
1286 |row| row.get(0),
1287 )
1288 .optional()?)
1289 }
1290
1291 pub fn set_hook_cursor(&self, hook: &str, seq: i64) -> StoreResult<bool> {
1295 let conn = self.conn()?;
1296 let updated_at = rfc3339(Utc::now());
1297 Ok(conn.execute(
1298 "INSERT INTO hook_cursors(hook,last_seq,updated_at) VALUES(?,?,?)
1299 ON CONFLICT(hook) DO UPDATE SET last_seq=excluded.last_seq,updated_at=excluded.updated_at
1300 WHERE excluded.last_seq > hook_cursors.last_seq",
1301 params![hook, seq, updated_at],
1302 )? == 1)
1303 }
1304
1305 fn transition_msg(
1306 &self,
1307 msg_id: &str,
1308 to: LedgerState,
1309 acked_at: Option<String>,
1310 reason: Option<String>,
1311 ) -> StoreResult<bool> {
1312 let conn = self.conn()?;
1313 let from = current_ledger_state(&conn, msg_id)?;
1314 ensure_transition_allowed(from, to)?;
1315 let changed = conn.execute(
1316 "UPDATE ledger SET state=?,acked_at=COALESCE(?,acked_at),reason=COALESCE(?,reason) WHERE msg_id=?",
1317 params![to.as_str(), acked_at, reason, msg_id],
1318 )?;
1319 Ok(changed == 1)
1320 }
1321
1322 fn ledger_for_role(
1323 &self,
1324 role: &str,
1325 state: LedgerState,
1326 limit: u32,
1327 ) -> StoreResult<Vec<LedgerRow>> {
1328 let conn = self.read()?;
1329 let rows = conn
1330 .prepare(&format!(
1331 "SELECT {LEDGER_COLUMNS} FROM ledger WHERE state=? AND json_extract(to_json,'$.role.role')=? ORDER BY enqueued_at,rowid LIMIT ?"
1332 ))?
1333 .query_map(params![state.as_str(), role, sql_limit(limit)], ledger_row)?
1334 .collect::<Result<Vec<_>, _>>()?;
1335 Ok(rows)
1336 }
1337
1338 fn conn(&self) -> StoreResult<MutexGuard<'_, Connection>> {
1339 self.inner
1340 .lock()
1341 .map_err(|_| StoreError::Sqlite("database mutex poisoned".to_string()))
1342 }
1343
1344 fn read(&self) -> StoreResult<MutexGuard<'_, Connection>> {
1351 self.readers.get()
1352 }
1353
1354 fn note_session_row_write(&self) {
1356 self.session_rows_written.fetch_add(1, Ordering::Relaxed);
1357 }
1358
1359 pub fn session_row_writes(&self) -> u64 {
1368 self.session_rows_written.load(Ordering::Relaxed)
1369 }
1370}
1371
1372const READ_POOL_SIZE: usize = 4;
1381
1382#[derive(Clone, Debug)]
1389struct ReadPool {
1390 path: PathBuf,
1391 conns: Arc<Vec<Mutex<Connection>>>,
1392 next: Arc<AtomicUsize>,
1393}
1394
1395impl ReadPool {
1396 fn open(path: &Path) -> StoreResult<Self> {
1397 let mut conns = Vec::with_capacity(READ_POOL_SIZE);
1398 for _ in 0..READ_POOL_SIZE {
1399 conns.push(Mutex::new(open_reader(path)?));
1400 }
1401 Ok(Self {
1402 path: path.to_path_buf(),
1403 conns: Arc::new(conns),
1404 next: Arc::new(AtomicUsize::new(0)),
1405 })
1406 }
1407
1408 fn get(&self) -> StoreResult<MutexGuard<'_, Connection>> {
1410 let index = self.next.fetch_add(1, Ordering::Relaxed) % self.conns.len();
1411 self.conns[index].lock().map_err(|_| {
1412 StoreError::Sqlite(format!(
1413 "read connection mutex poisoned for {}",
1414 self.path.display()
1415 ))
1416 })
1417 }
1418}
1419
1420pub(crate) fn open_connection(
1424 path: &Path,
1425 which: &'static str,
1426 marker: &str,
1427 ddl: &str,
1428 schema_version: i64,
1429) -> StoreResult<Connection> {
1430 if let Some(parent) = path.parent() {
1431 std::fs::create_dir_all(parent).map_err(|e| StoreError::Sqlite(e.to_string()))?;
1432 }
1433 let conn = Connection::open(path)?;
1434 conn.busy_timeout(Duration::from_millis(5000))?;
1435 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
1436 ensure_schema(&conn, which, marker, ddl, schema_version)?;
1437 Ok(conn)
1438}
1439
1440fn ensure_schema(
1441 conn: &Connection,
1442 which: &'static str,
1443 marker: &str,
1444 ddl: &str,
1445 schema_version: i64,
1446) -> StoreResult<()> {
1447 let tables = user_tables(conn)?;
1448 let legacy: Vec<String> = tables
1449 .iter()
1450 .filter(|name| is_legacy_table(name))
1451 .cloned()
1452 .collect();
1453 if !legacy.is_empty() {
1454 return Err(StoreError::unsupported_schema(
1455 which,
1456 SchemaMismatch::LegacyTables { names: legacy },
1457 ));
1458 }
1459 conn.execute_batch(SCHEMA_MARKER_DDL)?;
1460 let marker_row = conn
1461 .query_row(
1462 "SELECT version,protocol_version FROM schema_marker WHERE name=?",
1463 params![marker],
1464 |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)),
1465 )
1466 .optional()?;
1467 match marker_row {
1468 Some((version, protocol)) if version == schema_version && protocol == PROTOCOL_VERSION => {}
1469 Some((version, protocol)) => {
1470 return Err(if version != schema_version {
1475 StoreError::unsupported_schema(
1476 which,
1477 SchemaMismatch::Version {
1478 found: version,
1479 expected: schema_version,
1480 },
1481 )
1482 } else {
1483 StoreError::unsupported_schema(
1484 which,
1485 SchemaMismatch::Protocol {
1486 found: protocol,
1487 expected: PROTOCOL_VERSION,
1488 },
1489 )
1490 });
1491 }
1492 None => {
1493 let marker_count: i64 =
1494 conn.query_row("SELECT COUNT(*) FROM schema_marker", [], |r| r.get(0))?;
1495 if marker_count > 0 || !tables.is_empty() {
1496 return Err(StoreError::unsupported_schema(
1497 which,
1498 SchemaMismatch::NotEmpty {
1499 tables: tables.len(),
1500 },
1501 ));
1502 }
1503 conn.execute(
1504 "INSERT INTO schema_marker(name,version,protocol_version) VALUES(?,?,?)",
1505 params![marker, schema_version, PROTOCOL_VERSION],
1506 )?;
1507 }
1508 }
1509 conn.execute_batch(ddl)?;
1510 ensure_ledger_expires_at(conn)?;
1511 ensure_ledger_requeued(conn)?;
1512 ensure_ledger_causality_columns(conn)?;
1513 Ok(())
1514}
1515
1516fn open_reader(path: &Path) -> StoreResult<Connection> {
1523 let conn = Connection::open_with_flags(
1524 path,
1525 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
1526 )?;
1527 conn.busy_timeout(Duration::from_millis(5000))?;
1528 Ok(conn)
1529}
1530
1531fn ensure_ledger_expires_at(conn: &Connection) -> StoreResult<()> {
1536 let columns = conn
1537 .prepare("PRAGMA table_info(ledger)")?
1538 .query_map([], |row| row.get::<_, String>(1))?
1539 .collect::<Result<Vec<_>, _>>()?;
1540 if columns.is_empty() || columns.iter().any(|name| name == "expires_at") {
1541 return Ok(());
1542 }
1543 conn.execute("ALTER TABLE ledger ADD COLUMN expires_at TEXT", [])?;
1544 Ok(())
1545}
1546
1547fn ensure_ledger_requeued(conn: &Connection) -> StoreResult<()> {
1552 let columns = conn
1553 .prepare("PRAGMA table_info(ledger)")?
1554 .query_map([], |row| row.get::<_, String>(1))?
1555 .collect::<Result<Vec<_>, _>>()?;
1556 if columns.is_empty() || columns.iter().any(|name| name == "requeued") {
1557 return Ok(());
1558 }
1559 conn.execute(
1560 "ALTER TABLE ledger ADD COLUMN requeued INTEGER NOT NULL DEFAULT 0",
1561 [],
1562 )?;
1563 Ok(())
1564}
1565
1566fn ensure_ledger_causality_columns(conn: &Connection) -> StoreResult<()> {
1573 let columns = conn
1574 .prepare("PRAGMA table_info(ledger)")?
1575 .query_map([], |row| row.get::<_, String>(1))?
1576 .collect::<Result<Vec<_>, _>>()?;
1577 if columns.is_empty() {
1578 return Ok(());
1579 }
1580 for (name, kind) in [
1581 ("family", "TEXT"),
1582 ("hop_budget", "INTEGER"),
1583 ("origin", "TEXT"),
1584 ("deadline", "TEXT"),
1585 ("labels_json", "TEXT"),
1586 ] {
1587 if columns.iter().any(|column| column == name) {
1588 continue;
1589 }
1590 conn.execute(&format!("ALTER TABLE ledger ADD COLUMN {name} {kind}"), [])?;
1591 }
1592 Ok(())
1593}
1594
1595fn user_tables(conn: &Connection) -> StoreResult<Vec<String>> {
1596 let rows = conn
1597 .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name")?
1598 .query_map([], |row| row.get::<_, String>(0))?
1599 .collect::<Result<Vec<_>, _>>()?;
1600 Ok(rows)
1601}
1602
1603fn is_legacy_table(name: &str) -> bool {
1605 matches!(
1606 name,
1607 "io_cursors" | "loopback_idempotency" | "pending_replies"
1608 ) || name.starts_with("swarm")
1609}
1610
1611pub fn rfc3339(now: DateTime<Utc>) -> String {
1617 now.to_rfc3339_opts(SecondsFormat::Millis, true)
1618}
1619
1620pub(crate) fn append_event_conn(conn: &Connection, kind: &str, data: &Value) -> StoreResult<i64> {
1621 conn.execute(
1622 "INSERT INTO events(type,data_json,created_at) VALUES(?,?,?)",
1623 params![kind, serde_json::to_string(data)?, rfc3339(Utc::now())],
1624 )?;
1625 Ok(conn.last_insert_rowid())
1626}
1627
1628pub(crate) fn events_since_conn(
1629 conn: &Connection,
1630 seq: i64,
1631 limit: u32,
1632) -> StoreResult<Vec<EventRecord>> {
1633 let rows = conn
1634 .prepare(
1635 "SELECT seq,type,data_json,created_at FROM events WHERE seq>? ORDER BY seq LIMIT ?",
1636 )?
1637 .query_map(params![seq, sql_limit(limit)], event_row)?
1638 .collect::<Result<Vec<_>, _>>()?;
1639 Ok(rows)
1640}
1641
1642pub(crate) fn event_head_conn(conn: &Connection) -> StoreResult<i64> {
1643 Ok(conn.query_row("SELECT COALESCE(MAX(seq),0) FROM events", [], |r| r.get(0))?)
1644}
1645
1646fn current_ledger_state(conn: &Connection, msg_id: &str) -> StoreResult<LedgerState> {
1647 let state = conn
1648 .query_row(
1649 "SELECT state FROM ledger WHERE msg_id=?",
1650 params![msg_id],
1651 |r| r.get::<_, String>(0),
1652 )
1653 .optional()?
1654 .ok_or(StoreError::NotFound)?;
1655 parse_ledger_state(&state).ok_or(StoreError::Serialization(format!(
1656 "invalid ledger state {state}"
1657 )))
1658}
1659
1660fn ensure_transition_allowed(from: LedgerState, to: LedgerState) -> StoreResult<()> {
1661 if transition_allowed(from, to) {
1662 return Ok(());
1663 }
1664 Err(StoreError::InvalidState {
1665 from: from.as_str().to_string(),
1666 to: to.as_str().to_string(),
1667 })
1668}
1669
1670fn insert_ledger_row(conn: &Connection, row: &LedgerRow) -> StoreResult<()> {
1671 conn.execute(
1672 "INSERT INTO ledger(msg_id,op_id,fingerprint,kind,from_json,to_json,task,parent_task,attempt,hop,state,out_head,reason,enqueued_at,acked_at,body_json,expires_at,requeued,family,hop_budget,origin,deadline,labels_json) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
1673 params![
1674 row.msg_id,
1675 row.op_id,
1676 row.fingerprint,
1677 row.kind.as_str(),
1678 row.from_json,
1679 row.to_json,
1680 row.task,
1681 row.parent_task,
1682 row.attempt,
1683 row.hop,
1684 row.state.as_str(),
1685 row.out_head,
1686 row.reason,
1687 row.enqueued_at,
1688 row.acked_at,
1689 row.body_json,
1690 row.expires_at,
1691 row.requeued,
1692 row.family,
1693 row.hop_budget,
1694 row.origin,
1695 row.deadline,
1696 row.labels_json
1697 ],
1698 )?;
1699 Ok(())
1700}
1701
1702fn ledger_by_op_id(conn: &Connection, op_id: &str) -> StoreResult<Option<LedgerRow>> {
1703 Ok(conn
1704 .query_row(
1705 &format!("SELECT {LEDGER_COLUMNS} FROM ledger WHERE op_id=?"),
1706 params![op_id],
1707 ledger_row,
1708 )
1709 .optional()?)
1710}
1711
1712const LEDGER_COLUMNS: &str = "msg_id,op_id,fingerprint,kind,from_json,to_json,task,parent_task,attempt,state,out_head,reason,enqueued_at,acked_at,body_json,hop,expires_at,requeued,family,hop_budget,origin,deadline,labels_json";
1713const FAULT_COLUMNS: &str = "id,task_id,role,session_id,generation,seq,desired_json,observed_json,intent,attempt,backend_ref,kind,reason,state,created_at";
1714const GHOST_SWEEP_COLUMNS: &str =
1715 "id,task_id,role,session_id,generation,seq_before,seq_after,outcome,evidence,swept_at";
1716
1717pub(crate) const OPEN_BINDING_TASK: &str =
1725 "(SELECT st.task_id FROM session_tasks st WHERE st.session_id=sessions.session_id
1726 AND st.released_at IS NULL LIMIT 1)";
1727
1728pub(crate) fn session_columns() -> String {
1734 format!(
1735 "session_id,{OPEN_BINDING_TASK},\
1736 role,generation,seq,agent_state,delivery_state,resource_state,recovery_substate,\
1737 desired_json,observed_json,mismatch_count,last_seen,updated_at"
1738 )
1739}
1740
1741pub(crate) const SESSION_ID_FOR_TASK: &str = "(SELECT session_id FROM session_tasks WHERE task_id=?
1747 ORDER BY (released_at IS NULL) DESC, bound_at DESC, session_id DESC LIMIT 1)";
1748
1749fn project_session_conn(conn: &Connection, write: &SessionWrite) -> StoreResult<bool> {
1754 let changed = conn.execute(
1755 "INSERT INTO sessions(session_id,role,generation,seq,agent_state,delivery_state,resource_state,recovery_substate,desired_json,observed_json,mismatch_count,last_seen,updated_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)
1756 ON CONFLICT(session_id) DO UPDATE SET role=excluded.role,generation=excluded.generation,seq=excluded.seq,agent_state=excluded.agent_state,delivery_state=excluded.delivery_state,resource_state=excluded.resource_state,recovery_substate=excluded.recovery_substate,desired_json=excluded.desired_json,observed_json=excluded.observed_json,mismatch_count=excluded.mismatch_count,last_seen=excluded.last_seen,updated_at=excluded.updated_at
1757 WHERE excluded.generation > sessions.generation OR (excluded.generation = sessions.generation AND excluded.seq > sessions.seq)",
1758 params![
1759 write.session_id,
1760 write.role,
1761 write.generation,
1762 write.seq,
1763 write.agent_state,
1764 write.delivery_state,
1765 write.resource_state,
1766 write.recovery_substate,
1767 write.desired_json,
1768 write.observed_json,
1769 write.mismatch_count,
1770 unix_to_rfc3339(write.last_seen),
1771 unix_to_rfc3339(write.updated_at)
1772 ],
1773 )?;
1774 if changed == 1 {
1775 if let Some(task_id) = write.task_id.as_deref() {
1776 open_binding_conn(conn, &write.session_id, task_id, write.last_seen)?;
1777 }
1778 }
1779 Ok(changed == 1)
1780}
1781
1782pub(crate) fn open_binding_conn(
1790 conn: &Connection,
1791 session_id: &str,
1792 task_id: &str,
1793 bound_at: i64,
1794) -> StoreResult<usize> {
1795 let at = unix_to_rfc3339(bound_at);
1796 conn.execute(
1797 "UPDATE session_tasks SET released_at=? WHERE session_id=? AND released_at IS NULL AND task_id<>?",
1798 params![at, session_id, task_id],
1799 )?;
1800 Ok(conn.execute(
1801 "INSERT INTO session_tasks(session_id,task_id,bound_at,released_at) VALUES(?,?,?,NULL)
1802 ON CONFLICT(session_id,task_id) DO UPDATE SET bound_at=excluded.bound_at,released_at=NULL
1803 WHERE session_tasks.released_at IS NOT NULL",
1804 params![session_id, task_id, at],
1805 )?)
1806}
1807
1808pub(crate) fn release_binding_conn(
1811 conn: &Connection,
1812 session_id: &str,
1813 task_id: &str,
1814 released_at: i64,
1815) -> StoreResult<usize> {
1816 Ok(conn.execute(
1817 "UPDATE session_tasks SET released_at=? WHERE session_id=? AND task_id=? AND released_at IS NULL",
1818 params![unix_to_rfc3339(released_at), session_id, task_id],
1819 )?)
1820}
1821
1822fn sessions_list_sql(where_sql: &str) -> String {
1827 format!(
1828 "SELECT {} FROM sessions{where_sql} ORDER BY updated_at DESC,rowid DESC LIMIT ?",
1829 session_columns()
1830 )
1831}
1832
1833fn ledger_list_sql(where_sql: &str) -> String {
1835 format!(
1836 "SELECT {LEDGER_COLUMNS} FROM ledger{where_sql} ORDER BY enqueued_at DESC,rowid DESC LIMIT ?"
1837 )
1838}
1839
1840fn ghost_sweeps_list_sql() -> String {
1843 format!(
1844 "SELECT {GHOST_SWEEP_COLUMNS} FROM ghost_sweeps ORDER BY swept_at DESC,rowid DESC LIMIT ?"
1845 )
1846}
1847
1848fn ledger_row(r: &Row<'_>) -> rusqlite::Result<LedgerRow> {
1849 let kind: String = r.get(3)?;
1850 let state: String = r.get(9)?;
1851 Ok(LedgerRow {
1852 msg_id: r.get(0)?,
1853 op_id: r.get(1)?,
1854 fingerprint: r.get(2)?,
1855 kind: parse_msg_kind(&kind)
1856 .ok_or_else(|| conversion_error(3, format!("invalid message kind {kind}")))?,
1857 from_json: r.get(4)?,
1858 to_json: r.get(5)?,
1859 task: r.get(6)?,
1860 parent_task: r.get(7)?,
1861 attempt: r.get(8)?,
1862 state: parse_ledger_state(&state)
1863 .ok_or_else(|| conversion_error(9, format!("invalid ledger state {state}")))?,
1864 out_head: r.get(10)?,
1865 reason: r.get(11)?,
1866 enqueued_at: r.get(12)?,
1867 acked_at: r.get(13)?,
1868 body_json: r.get(14)?,
1869 hop: r.get(15)?,
1870 expires_at: r.get(16)?,
1871 requeued: r.get(17)?,
1872 family: r.get(18)?,
1873 hop_budget: r.get(19)?,
1874 origin: r.get(20)?,
1875 deadline: r.get(21)?,
1876 labels_json: r.get(22)?,
1877 })
1878}
1879
1880fn role_row(r: &Row<'_>) -> rusqlite::Result<RoleRow> {
1881 let admin: i64 = r.get(2)?;
1882 Ok(RoleRow {
1883 name: r.get(0)?,
1884 key: r.get(1)?,
1885 admin: admin != 0,
1886 max_sessions: r.get(3)?,
1887 spec_hash: r.get(4)?,
1888 updated_at: r.get(5)?,
1889 })
1890}
1891
1892fn session_row(r: &Row<'_>) -> rusqlite::Result<ServerSessionRow> {
1895 Ok(ServerSessionRow {
1896 session_id: r.get(0)?,
1897 task_id: r.get(1)?,
1898 role: r.get(2)?,
1899 generation: r.get(3)?,
1900 seq: r.get(4)?,
1901 agent_state: r.get(5)?,
1902 delivery_state: r.get(6)?,
1903 resource_state: r.get(7)?,
1904 recovery_substate: r.get(8)?,
1905 desired_json: r.get(9)?,
1906 observed_json: r.get(10)?,
1907 mismatch_count: r.get(11)?,
1908 last_seen: rfc3339_to_unix(&r.get::<_, String>(12)?),
1909 updated_at: rfc3339_to_unix(&r.get::<_, String>(13)?),
1910 })
1911}
1912
1913fn session_binding_row(r: &Row<'_>) -> rusqlite::Result<SessionBindingRow> {
1914 Ok(SessionBindingRow {
1915 session_id: r.get(0)?,
1916 task_id: r.get(1)?,
1917 bound_at: rfc3339_to_unix(&r.get::<_, String>(2)?),
1918 released_at: r
1919 .get::<_, Option<String>>(3)?
1920 .map(|text| rfc3339_to_unix(&text)),
1921 })
1922}
1923
1924fn event_row(r: &Row<'_>) -> rusqlite::Result<EventRecord> {
1925 let text: String = r.get(2)?;
1926 let data = serde_json::from_str(&text).map_err(|e| conversion_error(2, e.to_string()))?;
1927 Ok(EventRecord {
1928 seq: r.get(0)?,
1929 kind: r.get(1)?,
1930 data,
1931 created_at: r.get(3)?,
1932 })
1933}
1934
1935fn server_fault_row(r: &Row<'_>) -> rusqlite::Result<ServerFaultRow> {
1936 Ok(ServerFaultRow {
1937 id: r.get(0)?,
1938 task_id: r.get(1)?,
1939 role: r.get(2)?,
1940 session_id: r.get(3)?,
1941 generation: r.get(4)?,
1942 seq: r.get(5)?,
1943 desired_json: r.get(6)?,
1944 observed_json: r.get(7)?,
1945 intent: r.get(8)?,
1946 attempt: r.get(9)?,
1947 backend_ref: r.get(10)?,
1948 kind: r.get(11)?,
1949 reason: r.get(12)?,
1950 state: r.get(13)?,
1951 created_at: rfc3339_to_unix(&r.get::<_, String>(14)?),
1952 })
1953}
1954
1955fn ghost_sweep_row(r: &Row<'_>) -> rusqlite::Result<GhostSweepRow> {
1956 let outcome: String = r.get(7)?;
1957 Ok(GhostSweepRow {
1958 id: r.get(0)?,
1959 task_id: r.get(1)?,
1960 role: r.get(2)?,
1961 session_id: r.get(3)?,
1962 generation: r.get(4)?,
1963 seq_before: r.get(5)?,
1964 seq_after: r.get(6)?,
1965 outcome: parse_outcome(&outcome)
1966 .ok_or_else(|| conversion_error(7, format!("invalid outcome {outcome}")))?,
1967 evidence: r.get(8)?,
1968 swept_at: rfc3339_to_unix(&r.get::<_, String>(9)?),
1969 })
1970}
1971
1972fn ledger_by_msg_id(conn: &Connection, msg_id: &str) -> StoreResult<LedgerRow> {
1973 conn.query_row(
1974 &format!("SELECT {LEDGER_COLUMNS} FROM ledger WHERE msg_id=?"),
1975 params![msg_id],
1976 ledger_row,
1977 )
1978 .optional()?
1979 .ok_or(StoreError::NotFound)
1980}
1981
1982fn fault_row_by_id(conn: &Connection, id: i64) -> StoreResult<ServerFaultRow> {
1983 conn.query_row(
1984 &format!("SELECT {FAULT_COLUMNS} FROM faults WHERE id=?"),
1985 params![id],
1986 server_fault_row,
1987 )
1988 .optional()?
1989 .ok_or(StoreError::NotFound)
1990}
1991
1992fn ledger_state_event(
1995 row: &LedgerRow,
1996 state: LedgerState,
1997 reason: Option<String>,
1998) -> StoreResult<Value> {
1999 let event = Event::LedgerState(LedgerStateEvent {
2000 msg_id: row.msg_id.clone(),
2001 op_id: row.op_id.clone(),
2002 kind: row.kind,
2003 from: row.sender().unwrap_or_else(|_| Principal::role("unknown")),
2004 to: serde_json::from_str(&row.to_json).unwrap_or_else(|_| Principal::role("unknown")),
2005 task: row.task.clone(),
2006 state,
2007 outcome: None,
2008 reason,
2009 });
2010 Ok(serde_json::to_value(&event)?)
2011}
2012
2013fn fault_event(row: &ServerFaultRow) -> StoreResult<Value> {
2015 let event = Event::Fault(FaultEvent {
2016 id: row.id,
2017 task_id: row.task_id.clone(),
2018 role: row.role.clone(),
2019 session_id: row.session_id.clone(),
2020 generation: row.generation.map(|value| value.max(0) as u64),
2021 seq: row.seq.map(|value| value.max(0) as u64),
2022 kind: row.kind.clone(),
2023 reason: row.reason.clone(),
2024 desired: row
2025 .desired_json
2026 .as_deref()
2027 .and_then(|text| serde_json::from_str(text).ok()),
2028 observed: row
2029 .observed_json
2030 .as_deref()
2031 .and_then(|text| serde_json::from_str(text).ok()),
2032 intent: row.intent.clone(),
2033 attempt: row.attempt.map(|value| value.max(0) as u64),
2034 backend_ref: row
2035 .backend_ref
2036 .as_deref()
2037 .and_then(|text| serde_json::from_str(text).ok()),
2038 state: Some(row.state.clone()),
2039 created_at: Some(row.created_at),
2040 });
2041 Ok(serde_json::to_value(&event)?)
2042}
2043
2044fn cursor_row(r: &Row<'_>) -> rusqlite::Result<CursorRow> {
2045 Ok(CursorRow {
2046 role: r.get(0)?,
2047 last_msg_id: r.get(1)?,
2048 last_seq: r.get(2)?,
2049 updated_at: r.get(3)?,
2050 })
2051}
2052
2053fn bool_int(value: bool) -> i64 {
2054 if value { 1 } else { 0 }
2055}
2056
2057pub const OUT_HEAD_CLUSTERS: usize = 200;
2059
2060pub fn head_preview(text: &str) -> String {
2067 text.graphemes(true).take(OUT_HEAD_CLUSTERS).collect()
2068}
2069
2070fn body_head(envelope: &Envelope, body_json: &str) -> String {
2071 if let Some(head) = envelope.body.head.as_deref().filter(|h| !h.is_empty()) {
2078 return head_preview(head);
2079 }
2080 let text = envelope.body.text.as_deref().unwrap_or(body_json);
2081 head_preview(text)
2082}
2083
2084pub(crate) fn unix_to_rfc3339(seconds: i64) -> String {
2088 DateTime::from_timestamp(seconds, 0)
2089 .map(rfc3339)
2090 .unwrap_or_else(|| rfc3339(DateTime::UNIX_EPOCH))
2091}
2092
2093pub(crate) fn rfc3339_to_unix(text: &str) -> i64 {
2095 DateTime::parse_from_rfc3339(text)
2096 .map(|value| value.timestamp())
2097 .unwrap_or(0)
2098}
2099
2100fn parse_msg_kind(value: &str) -> Option<MsgKind> {
2101 match value {
2102 "task" => Some(MsgKind::Task),
2103 "completion" => Some(MsgKind::Completion),
2104 "note" => Some(MsgKind::Note),
2105 "control" => Some(MsgKind::Control),
2106 _ => None,
2107 }
2108}
2109
2110fn parse_ledger_state(value: &str) -> Option<LedgerState> {
2111 match value {
2112 "queued" => Some(LedgerState::Queued),
2113 "in_flight" => Some(LedgerState::InFlight),
2114 "acked" => Some(LedgerState::Acked),
2115 "rejected" => Some(LedgerState::Rejected),
2116 "expired" => Some(LedgerState::Expired),
2117 _ => None,
2118 }
2119}
2120
2121fn parse_outcome(value: &str) -> Option<Outcome> {
2122 match value {
2123 "done" => Some(Outcome::Done),
2124 "failed" => Some(Outcome::Failed),
2125 "cancelled" => Some(Outcome::Cancelled),
2126 "blocked" => Some(Outcome::Blocked),
2127 _ => None,
2128 }
2129}
2130
2131pub(crate) fn string_tag<T: Serialize>(value: &T) -> StoreResult<String> {
2132 match serde_json::to_value(value)? {
2133 Value::String(text) => Ok(text),
2134 other => Err(StoreError::Serialization(format!(
2135 "state serialized as {other}"
2136 ))),
2137 }
2138}
2139
2140fn sql_limit(limit: u32) -> i64 {
2141 if limit == 0 {
2142 DEFAULT_LIMIT
2143 } else {
2144 i64::from(limit.min(500))
2145 }
2146}
2147
2148fn repeat_placeholders(len: usize) -> String {
2149 std::iter::repeat_n("?", len).collect::<Vec<_>>().join(",")
2150}
2151
2152fn where_sql(clauses: &[String]) -> String {
2153 if clauses.is_empty() {
2154 String::new()
2155 } else {
2156 format!(" WHERE {}", clauses.join(" AND "))
2157 }
2158}
2159
2160fn conversion_error(index: usize, message: String) -> rusqlite::Error {
2161 rusqlite::Error::FromSqlConversionFailure(
2162 index,
2163 Type::Text,
2164 Box::new(StoreError::Serialization(message)),
2165 )
2166}