Skip to main content

onlyne_store/
server.rs

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
22/// Server store schema revision. The `hop` column set this to 2: the ledger
23/// keeps the hop count of `Causality`, which `onlyne handoff` reads back to
24/// extend a chain. The `expires_at` and `requeued` columns were applied in
25/// place, so they never moved the marker. Version 3 is the tuple rebuild: the
26/// `sessions` row lost its `public_lifecycle` column, which cannot be taken
27/// back from an existing file, so an old layout is refused rather than carried.
28/// Version 4 adds the `ghost_sweeps` audit table: the server settles a `working`
29/// mirror row whose task ledger row already reached a terminal state, and each
30/// settlement writes one row there. A marker-3 file carries no such table, so it
31/// stops at the door on the same string every other mismatch prints.
32/// The ledger's five family-metadata columns — `family`, `hop_budget`,
33/// `origin`, `deadline`, `labels_json` — were applied in place beside
34/// `expires_at` and `requeued`, so they never moved the marker.
35/// Version 5 rekeys the mirror: a session row is addressed by `session_id`
36/// rather than by the delivery it happens to serve, carries `last_seen` beside
37/// `updated_at`, and the deliveries it serves are recorded in `session_tasks`.
38/// A marker-4 file holds rows under the old key and no bindings table, which
39/// cannot be read back as this layout, so it stops at the door on the same
40/// string every other mismatch prints.
41/// Version 6 adds the `hook_cursors` table: an event hook is delivered
42/// at-least-once, so each declared hook records the last event `seq` it
43/// handled successfully and resumes from there after a restart. A marker-5
44/// file holds no such row, so it stops at the door on the same string every
45/// other mismatch prints.
46/// The client store keeps its own revision.
47const 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/// One mirror row as a writer hands it over, and as a reader gets it back.
230///
231/// The address is `session_id`, which is the table's key: a session serves one
232/// delivery at a time, and the deliveries it serves live in `session_tasks`
233/// rather than in this row. `task_id` is the delivery the write is about — the
234/// binding a writer opens — and on a read it is the delivery the session is
235/// serving now, derived from its open `session_tasks` row and absent when the
236/// session is on no delivery.
237#[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    /// When the mirror last saw this session, in unix seconds. Stored as the
252    /// RFC 3339 text the crate's helpers produce.
253    pub last_seen: i64,
254    /// When the projection content last moved, in unix seconds.
255    pub updated_at: i64,
256}
257
258pub type ServerSessionRow = SessionWrite;
259
260/// One `session_tasks` row: which delivery a session serves, and whether it
261/// still serves it.
262#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
263pub struct SessionBindingRow {
264    pub session_id: String,
265    pub task_id: String,
266    /// When this session took the delivery, in unix seconds.
267    pub bound_at: i64,
268    /// When it stopped serving it, in unix seconds. `None` while the delivery
269    /// is the one the session is on.
270    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    /// Causality's hop count from the root task, which is what turns the
291    /// `parent_task` links into a measurable depth.
292    pub hop: i64,
293    /// The persisted expiry deadline of a ttl note, so a restarted server can
294    /// re-arm its sweep.
295    pub expires_at: Option<String>,
296    /// Times this row has moved from in_flight back to queued.
297    #[serde(default)]
298    pub requeued: i64,
299    /// The family's root task id, read off the envelope's causality. Every row
300    /// of one run carries it unchanged, and `onlyne handoff` reads it back to
301    /// mint the next hop into the same family.
302    #[serde(default)]
303    pub family: Option<String>,
304    /// The hops the family may spend, carried on every row of the run.
305    #[serde(default)]
306    pub hop_budget: Option<i64>,
307    /// The role the family reports home to, carried on every row of the run.
308    #[serde(default)]
309    pub origin: Option<String>,
310    /// Wall-clock bound for the whole family, in the RFC 3339 text this
311    /// crate's helpers produce.
312    #[serde(default)]
313    pub deadline: Option<String>,
314    /// The causality's labels as JSON text, so the column keeps whatever map
315    /// the sender attached.
316    #[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    /// The sender principal this row recorded.
363    ///
364    /// The sender column carries the principal beside the admin marker, so a
365    /// row states whether an admin-surface send wrote it (plan §8 line 320)
366    /// without a second column. A column holding a bare principal decodes too.
367    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    /// Whether the send behind this row arrived on the admin surface.
373    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
381/// Encode the sender column: the principal plus the admin marker.
382fn 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/// One settlement the server's ghost sweep recorded in its own audit table.
435///
436/// `seq_before` and `seq_after` are the mirror row's two versions. The sweep
437/// writes through the same settlement path a `repair_*` verb uses, and that path
438/// bumps `seq` by one, so the pair names the exact write an operator is reading
439/// about. `swept_at` holds the kernel's unix seconds and the column holds the
440/// RFC 3339 text this crate's conversion helpers produce.
441#[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    /// The verdict written onto the mirror row, read off the task's ledger row.
451    pub outcome: Outcome,
452    /// What justified the sweep: the evidence tag plus the ledger state it read.
453    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    /// The one connection that writes, behind its own lock.
470    inner: Arc<Mutex<Connection>>,
471    /// The read-only handles every query runs on, so a reader is never queued
472    /// behind the writer above.
473    readers: ReadPool,
474    /// Writes this store made against `sessions` since it opened: one bump per
475    /// `project_session`, `rebind_session`, `publish_mirror_outcome`, and
476    /// `flush_last_seen`. A heartbeat that landed in memory counts as none.
477    session_rows_written: Arc<AtomicU64>,
478    /// The beats this process took and has not written down.
479    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        // The writer above ran the schema gate — the marker check, the DDL, the
493        // in-place column adds — so the readers below only ever open a file
494        // this process has already accepted.
495        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    /// Write one mirror row, and bind the delivery the write carries.
558    ///
559    /// The row is addressed by session id and the `(generation, seq)` gate is
560    /// the row's own, as it was when the row answered for a task. A write that
561    /// lands binds its delivery in the same transaction: a session serves one
562    /// delivery at a time, so taking this one releases whatever the session was
563    /// on before. A write the gate refuses changes nothing, the binding
564    /// included — the row it was refused by is the newer word.
565    ///
566    /// A write that lands carries the session's `last_seen` forward, so the
567    /// beats it has already outlived are spent: [`crate::liveness`] is told
568    /// what the row now holds and drops the ones that are no longer newer.
569    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    /// Move one mirror row to the session the write names, and write it there.
582    ///
583    /// This is what `repair rebind` means once a row is addressed by its
584    /// session: the operator says the delivery is now carried by another
585    /// session, so the row moves to that address instead of a second row
586    /// appearing beside it. The bindings move with it, since they name the same
587    /// session. A row already sitting at the new address is the one being
588    /// replaced, so it goes first: the operator's word is the newer fact.
589    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            // The row moved, so a beat taken under the old address is no longer
615            // this row's. The write's own clock is what the row now holds.
616            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    /// Take one delivery for a session: a `session_tasks` row with `bound_at`.
623    ///
624    /// Answers whether the binding moved, which a pair already open does not.
625    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    /// Stop serving one delivery: its `session_tasks` row gets `released_at`.
636    ///
637    /// Only an open binding is released, so the first release is the one the
638    /// row keeps and a second call changes nothing.
639    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    /// The delivery a session is on now, when it is on one.
650    ///
651    /// At most one binding of a session is open, so the order here only decides
652    /// which row a hand-written pair of open bindings answers with.
653    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    /// Publish a late mirror verdict when the stored projection bytes still
665    /// match the bytes the server compared. The stored tuple remains the
666    /// session version while `observed_json` carries the task verdict.
667    ///
668    /// `last_seen` is left where it stands: this write says the projection
669    /// gained a verdict, not that the session was heard from, and the beats
670    /// that were newer than the row keep answering for it
671    /// ([`ServerLedger::beat_session`]).
672    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    /// Take one beat for a session whose projection did not change.
694    ///
695    /// This is the whole of v2's liveness rule, and it is deliberately not a
696    /// projection write: `(generation, seq)`, `updated_at`, and the event
697    /// stream are all left exactly as they were, and the beat reaches the row
698    /// only when the row is at least `flush_after_secs` behind — the interval
699    /// the server states beside its call, which is also the reader's worst-case
700    /// staleness. Between those flushes the beat lives in memory, where every
701    /// read of the row picks it up ([`crate::liveness`]).
702    ///
703    /// The caller supplies the interval rather than this store: it is a promise
704    /// about what a reader of `last_seen` is owed, and the server is the layer
705    /// that knows the cluster's presence window. The interval is floored at one
706    /// second, so a spec that asks for a window below it gets one flush per
707    /// second rather than one per beat — the v1 write rate this slice removes.
708    ///
709    /// Answers whether the beat reached the table.
710    pub fn beat_session(
711        &self,
712        session_id: &str,
713        at: i64,
714        flush_after_secs: i64,
715    ) -> StoreResult<bool> {
716        // The memory entry moves first. A reader arriving between the two steps
717        // below must be handed this beat, never the value it replaced.
718        self.live.note(session_id, at);
719        let Some(stored) = self.stored_session_row(session_id)? else {
720            // Nothing to refresh: a beat for a session with no row is not a row
721            // this process can make fresher. The entry stays for the row that
722            // may yet be written, and a row that never appears is read by
723            // nobody.
724            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    /// Write one mirror row's `last_seen`, and nothing else.
737    ///
738    /// The row is the durable half of a beat: what a restarted server, and any
739    /// reader of the file rather than of the cluster, has to judge freshness
740    /// by. A mirror whose `last_seen` froze at the last content change is the
741    /// v1 defect that column exists to fix, and it is why the row carries it at
742    /// all — [`ServerLedger::beat_session`] is the only caller, and it decides
743    /// when the row is worth reaching.
744    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    /// One mirror row, addressed by its session id: the live `last_seen` this
755    /// process holds when it has a newer one, the persisted value otherwise.
756    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    /// One mirror row exactly as the file holds it.
763    ///
764    /// The liveness layer needs the persisted value on its own — the interval
765    /// is measured from it — and a caller that wants what a reader is handed
766    /// wants [`ServerLedger::get_session_row`] instead.
767    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    /// The freshest `last_seen` this process can answer for one stored row.
783    ///
784    /// Every read of the table passes through here. A reader must never be
785    /// handed a value the server has already outlived — that is v1's frozen
786    /// mirror, and it is the failure this column was added to end — so the
787    /// answer is the newer of the row's value and the beat held in memory.
788    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    /// The mirror row serving one delivery, read through its binding.
794    ///
795    /// A reader asks "the session serving this task"; the binding is where that
796    /// fact lives, and the row's own `task_id` is derived from it either way.
797    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            // The delivery filter asks which session served the task, which is
818            // the bindings table's answer. Both spellings of a binding count: a
819            // delivery that ended still names the session that carried it, and
820            // a reader of a settled delivery still asks for it.
821            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            // The mirror holds one copy of the published projection, in
832            // `observed_json`, and the lifecycle is a key inside it. The row's
833            // own dimensions stay columns; this one stays derived.
834            //
835            // Bytes that do not parse, or that never carried the key, read back
836            // as `created` on the row's own read path, which decodes the mirror
837            // and falls back to the columns. The filter has to agree: JSON1
838            // answers NULL for a key it cannot find and fails outright on bytes
839            // that are not JSON, and a scan that dropped such a row would leave
840            // a stale session without ever earning its fault row.
841            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    /// How many deliveries are queued for one role's inbox, counted exactly.
907    ///
908    /// [`ServerLedger::queued_for`] takes a limit, so a caller that counted its
909    /// rows would report the cap as the depth: an operator reading "512" when
910    /// the truth is 900 reads a saturated role as a full one. The rows are the
911    /// ones `pull` would hand this role, so `note` rows are left out exactly as
912    /// `relay::pull` leaves them out (`crates/onlyne-server/src/relay.rs`).
913    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    /// `msg_id` plus parsed deadline for every queued or in-flight row whose
931    /// `expires_at` is set, oldest deadline first.
932    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    /// Move one in-flight row back to `queued` and publish its `ledger_state`
961    /// event in the same transaction. A row already `queued` is returned
962    /// unchanged, so a disconnect path that fires twice settles once.
963    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    /// Move one queued row to `expired` and publish its `ledger_state` event in
983    /// the same transaction.
984    /// Settle one row whose deadline passed.
985    ///
986    /// A row reaches this from `queued` while it waits for its role, and from
987    /// `in_flight` when the role held it past the deadline: the sender asked for
988    /// a deadline, so the sweep answers with `expired` in both states.
989    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    /// Move one row to `rejected` and publish its `ledger_state` event in the
1008    /// same transaction. The automatic requeue gate uses this when the row has
1009    /// used its requeue budget.
1010    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    /// Move every `open` fault of a task to `next_state`, publishing one `fault`
1027    /// event per moved row in the same transaction. Returns the moved rows.
1028    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    /// The plan's task-keyed ledger read: every row of one task across roles, in
1096    /// insertion order. `docs/v1-PLAN.md` line 501 reads three rows back in order
1097    /// after a reconnect and line 508 reads one task's rows after a relocation.
1098    /// The ledger carries no monotonic column of its own, so the order is
1099    /// `enqueued_at` with SQLite's `rowid` breaking ties between rows written in
1100    /// one second: `rowid` is the insertion counter, and it makes the order the
1101    /// ledger's write order without a second copy of that fact.
1102    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    /// Persist one ghost-sweep audit row and return its id.
1219    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    /// The recorded sweeps, newest first.
1239    ///
1240    /// The order is the sessions listing's order: `swept_at DESC, rowid DESC`
1241    /// against an ascending index, which SQLite walks backwards. `rowid` breaks
1242    /// the ties inside one second, and it breaks them in the order the pass
1243    /// wrote them.
1244    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    /// The last event `seq` one hook handled successfully, `None` for a hook
1275    /// this cluster has never run.
1276    ///
1277    /// This is the resume point of the at-least-once rule: a hook that
1278    /// restarted picks up after the last event it handled, so nothing between
1279    /// the cursor and the head is skipped (`docs/v2-CONTRACT.md` §"Slice 7").
1280    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    /// Record the last event `seq` one hook handled successfully. A `seq` at or
1292    /// below the recorded one is left alone: a worker that resumed from an
1293    /// older row after a failure must not walk the cursor backwards.
1294    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    /// One connection from the read-only pool.
1345    ///
1346    /// Every statement that only reads goes here rather than through
1347    /// [`ServerLedger::conn`]: the writer's lock is held for the whole of a
1348    /// write, including the part where SQLite waits on the file lock, and a
1349    /// board refresh that shares it waits out the delivery path.
1350    fn read(&self) -> StoreResult<MutexGuard<'_, Connection>> {
1351        self.readers.get()
1352    }
1353
1354    /// Record one statement this store executed against `sessions`.
1355    fn note_session_row_write(&self) {
1356        self.session_rows_written.fetch_add(1, Ordering::Relaxed);
1357    }
1358
1359    /// Statements this store has executed against the `sessions` table.
1360    ///
1361    /// The projection writes, the mirror-outcome publishes, the liveness
1362    /// flushes, and the address moves a rebind makes all count here. It is the
1363    /// observation a liveness claim is made against: "a beat that changed
1364    /// nothing wrote no row" is answered by this number rather than by reading
1365    /// the code, and a test that reads it does not have to be a party to the
1366    /// write path.
1367    pub fn session_row_writes(&self) -> u64 {
1368        self.session_rows_written.load(Ordering::Relaxed)
1369    }
1370}
1371
1372/// How many read-only connections the query paths share.
1373///
1374/// WAL lets readers run beside the writer, so the pool is about *width*: a
1375/// board, a TUI, and an operator's `onlyne sessions` are three reads that may
1376/// be in flight at once, and four leaves room for the server's own scans
1377/// without any of them waiting on another. Each of these is one file handle;
1378/// the pool never grows, because every read behind it is a point lookup or an
1379/// indexed listing.
1380const READ_POOL_SIZE: usize = 4;
1381
1382/// The read-only handles a query runs on.
1383///
1384/// These are opened `SQLITE_OPEN_READ_ONLY`, so a statement that tried to write
1385/// here would fail instead of quietly becoming a second writer: the one writer
1386/// is the connection above, and this pool is the answer to "the reader should
1387/// not queue behind it".
1388#[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    /// One connection, round-robin, so two reads in flight land on two.
1409    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
1420/// The DDL revision travels beside the marker: the two databases evolve apart,
1421/// so the server store states its own number instead of sharing one constant
1422/// that a change to either DDL would invalidate for both.
1423pub(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            // The two revisions are reported apart: a file from an older build
1471            // and a file from a build that speaks another protocol are different
1472            // problems, and an operator who is told only "unsupported" has to
1473            // guess which one they have.
1474            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
1516/// Open one read-only handle on a server database that already exists.
1517///
1518/// No pragmas and no DDL: the writer set `journal_mode=WAL` on the file and ran
1519/// the schema gate before any of these open, and a read-only connection cannot
1520/// change either. The busy timeout is the writer's, so a reader that meets a
1521/// checkpoint waits it out instead of failing a board refresh.
1522fn 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
1531/// Add `ledger.expires_at` to a file whose ledger lacks it.
1532///
1533/// The column landed in place, so an existing database keeps its rows and a
1534/// file at the current marker keeps opening.
1535fn 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
1547/// Add `ledger.requeued` to a file whose ledger lacks it.
1548///
1549/// The column landed in place, so an existing database keeps its rows and a
1550/// file at the current marker keeps opening.
1551fn 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
1566/// Add the five ledger columns the envelope's causality feeds to a file whose
1567/// ledger lacks them: `family`, `hop_budget`, `origin`, `deadline`,
1568/// `labels_json`.
1569///
1570/// The columns landed in place, so an existing database keeps its rows and a
1571/// file at the current marker keeps opening.
1572fn 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
1603// Rejection-path markers only: these names identify pre-v1 tables and the `swarm` prefix that the schema gate refuses (plan §10 line 370).
1604fn is_legacy_table(name: &str) -> bool {
1605    matches!(
1606        name,
1607        "io_cursors" | "loopback_idempotency" | "pending_replies"
1608    ) || name.starts_with("swarm")
1609}
1610
1611/// One moment as RFC 3339 text.
1612///
1613/// Milliseconds rather than whole seconds: three acks inside one second are a
1614/// normal run, and plan §Verification case 4 reads the stamps of a reconnect
1615/// burst as the order they were settled in, which whole seconds collapse.
1616pub 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
1717/// The delivery one session is serving right now, as a subquery over its
1718/// bindings: an unbound session answers with nothing, which is the state a
1719/// claim reports as a session with no delivery.
1720///
1721/// No ordering: at most one binding of a session is open, which is what
1722/// `open_binding_conn` maintains, and an ordered subquery would have SQLite
1723/// spill a sorter into a temporary file on every listing read.
1724pub(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
1728/// One mirror row's columns, as every read of the table selects them.
1729///
1730/// `task_id` is not a column of `sessions`: it is the delivery the row is
1731/// labelled with — the one its session is on now, selected in the second
1732/// position `session_row` reads.
1733pub(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
1741/// The session serving one task, as a subquery over the bindings.
1742///
1743/// The open binding is the session on that delivery now; the fallback to the
1744/// last binding taken keeps a settled delivery readable, exactly as the row it
1745/// used to be keyed by stayed readable.
1746pub(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
1749/// Apply one mirror write, and bind its delivery when the write lands.
1750///
1751/// Shared by the plain write and the operator's rebind, which differ only in
1752/// where the row is addressed before the write.
1753fn 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
1782/// Take one delivery for a session, releasing whatever it was on before.
1783///
1784/// A session serves one delivery at a time, so this is the one place that rule
1785/// is enforced: the session's other open binding is released at the clock this
1786/// take carries. A pair already open stays as it stands — its `bound_at` is
1787/// when this session took that delivery, not when it was last written about —
1788/// and a pair whose binding was released is taken again under this clock.
1789pub(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
1808/// Stop serving one delivery. Only an open binding is released, so the first
1809/// release is the one the row keeps.
1810pub(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
1822/// The sessions listing read, as SQL. Named so the caller and the plan test run
1823/// one text: the order's tie key is the primary key because an explicit `rowid`
1824/// cannot be an index column, and an order SQLite cannot satisfy from an index
1825/// makes it spill a sorter to a temporary file on every read.
1826fn 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
1833/// The ledger listing read, as SQL, on the same terms as the sessions listing.
1834fn 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
1840/// The ghost-sweep listing read, as SQL, on the same terms as the two listings
1841/// above. Named for the same reason: the plan test runs this one text.
1842fn 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
1892/// One [`session_columns`] row. `task_id` is the derived binding, so it is the
1893/// second column rather than one of the table's own.
1894fn 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
1992/// The exact `Event::LedgerState` payload `State::emit` stores, so a store-side
1993/// transition and a server-side one produce the same event row.
1994fn 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
2013/// The exact `Event::Fault` payload for one stored fault row.
2014fn 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
2057/// Grapheme clusters kept by [`head_preview`].
2058pub const OUT_HEAD_CLUSTERS: usize = 200;
2059
2060/// The one-line preview of a message body: the first [`OUT_HEAD_CLUSTERS`]
2061/// grapheme clusters, cut on a cluster boundary, so a combining sequence, a ZWJ
2062/// emoji, and a flag survive whole. `docs/v1-PLAN.md` line 358 keeps the first
2063/// 200 字 of the body in `ledger.out_head` for the operator-facing ledger view,
2064/// and the unit is grapheme clusters: a code-point cut lands inside a cluster
2065/// and the preview renders as a broken glyph.
2066pub 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    // An explicit `head` on the body wins: it is the one display line the
2072    // sender named, and a body that carries its full result in `text` would
2073    // otherwise show the first clusters of the result rather than the line the
2074    // caller wrote. A body with no `head` of its own — the CLI's own
2075    // completion, a plain delivery — falls back to its text, which is the only
2076    // content it has.
2077    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
2084/// Encode a kernel timestamp, unix seconds, as the RFC 3339 text every other
2085/// column in this schema uses. A value outside the representable range encodes
2086/// as the epoch.
2087pub(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
2093/// Decode an RFC 3339 column back to the kernel's unix seconds.
2094pub(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}