Skip to main content

onlyne_store/
client.rs

1use std::path::{Path, PathBuf};
2use std::sync::{Arc, Mutex, MutexGuard};
3
4use crate::session::{FaultRecord, SessionLedger, SessionRecord, VersionedSession};
5use chrono::{DateTime, Utc};
6use onlyne_proto::Causality;
7use onlyne_proto::LiveSession;
8use onlyne_proto::TaskState;
9use rusqlite::types::Type;
10use rusqlite::{Connection, OptionalExtension, Row, params};
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13
14use crate::error::{StoreError, StoreResult};
15use crate::server::{
16    EventRecord, OPEN_BINDING_TASK, SESSION_ID_FOR_TASK, append_event_conn, event_head_conn,
17    events_since_conn, open_binding_conn, open_connection, release_binding_conn, rfc3339,
18    string_tag,
19};
20
21const CLIENT_MARKER: &str = "onlyne-client";
22/// The client DDL's own revision; the server store carries a separate one.
23/// Version 2 is the tuple rebuild: the `sessions` row lost its
24/// `public_lifecycle` column, and the task moved into a table of its own.
25/// Version 3 rekeys the client's mirror the way the server's is keyed: the
26/// `sessions` row is addressed by `session_id`, carries `last_seen`, and the
27/// deliveries the client's sessions serve live in `session_tasks`. A marker-2
28/// file holds rows under the old key and no bindings table, which cannot be
29/// read back as this layout, so it stops at the door on the same string every
30/// other mismatch prints.
31const CLIENT_SCHEMA_VERSION: i64 = 3;
32const DEFAULT_LIMIT: i64 = 100;
33/// Rows one flush pass takes from the intent queue.
34///
35/// The queue is durable, so a client that stayed offline through a long outage
36/// can hold tens of thousands of pending rows. A pass that read them all held
37/// the store's single connection for the length of the queue and starved every
38/// other caller of the same database, so one pass takes this many due rows and
39/// the next pass takes the rest.
40pub const INTENT_FLUSH_BATCH_SIZE: u32 = 100;
41
42pub const CLIENT_DDL: &str = r#"-- The client's own mirror of the sessions it holds, keyed the way the server's
43-- mirror is: a session row answers for a session, and which delivery that
44-- session serves lives in `session_tasks` below.
45--
46-- Two columns are the client's alone, because only the process that holds a
47-- session can answer them: `backend` names the backend that ran it (`acp`,
48-- `exec`, `orca`, `external`, … — which one comes from the role's drive under
49-- the placement this machine resolved) and `backend_ref` is the reference that
50-- backend answers to. There is no `role` column: one client serves one role,
51-- and the workspace config owns that fact.
52CREATE TABLE IF NOT EXISTS sessions(
53  session_id TEXT PRIMARY KEY,
54  generation INTEGER NOT NULL,
55  seq INTEGER NOT NULL,
56  agent_state TEXT NOT NULL,
57  delivery_state TEXT NOT NULL,
58  resource_state TEXT NOT NULL,
59  recovery_substate TEXT NOT NULL,
60  observed_json TEXT NOT NULL,
61  backend TEXT,
62  backend_ref TEXT,
63  -- The (generation, seq) gate. The reducer's isolate-after-N and
64  -- terminate-after-N policy needs a persisted counter, so this pair carries
65  -- DEFAULT_ISOLATE_AFTER and DEFAULT_TERMINATE_AFTER.
66  desired_json TEXT NOT NULL,
67  mismatch_count INTEGER NOT NULL DEFAULT 0,
68  -- When this client last wrote about the session. A reader judges a row's
69  -- freshness by it, and the tuple's own version lives in the pair above.
70  last_seen TEXT NOT NULL,
71  -- When the tuple last moved. Both clocks are the kernel's unix seconds,
72  -- encoded through this crate's own helper on every write.
73  updated_at TEXT NOT NULL
74);
75-- Which delivery a session serves: the same table, columns and keys as the
76-- server's, and the only place a binding lives on this side too.
77CREATE TABLE IF NOT EXISTS session_tasks(
78  session_id TEXT NOT NULL,
79  task_id TEXT NOT NULL,
80  bound_at TEXT NOT NULL,
81  released_at TEXT,
82  PRIMARY KEY (session_id, task_id)
83);
84CREATE INDEX IF NOT EXISTS session_tasks_task_idx ON session_tasks(task_id);
85CREATE INDEX IF NOT EXISTS session_tasks_open_idx ON session_tasks(session_id, released_at);
86-- The task's own record. How its work ended is not a session dimension: the
87-- session tuple says what this session can prove about its agent, intent,
88-- resource, and recovery line, and `project` needs the task's verdict handed
89-- in before it can say whether the session is over. A row is opened when a
90-- delivery becomes a session, from the envelope's causality, and settled once;
91-- a verdict that arrives for a task this process never opened writes its own
92-- row, so a completion is never dropped because of who opened what. `kind` is
93-- how the delivery reached this role: `root` for work given here, `relay` for
94-- work handed down from a parent task.
95CREATE TABLE IF NOT EXISTS task(
96  task_id TEXT PRIMARY KEY,
97  kind TEXT,
98  parent_task TEXT,
99  hop INTEGER NOT NULL DEFAULT 0,
100  attempt INTEGER NOT NULL DEFAULT 0,
101  -- TaskState's own snake_case tag from onlyne-proto's lifecycle vocabulary.
102  -- `pending` is the open state, and the settle write is the only thing that
103  -- leaves it, so `settled_at IS NULL` and `task_state = 'pending'` answer the
104  -- same question.
105  task_state TEXT NOT NULL,
106  opened_at TEXT NOT NULL,
107  settled_at TEXT
108);
109CREATE TABLE IF NOT EXISTS intents(
110  op_id TEXT PRIMARY KEY,
111  env_json TEXT NOT NULL,
112  attempt INTEGER NOT NULL,
113  state TEXT NOT NULL,
114  next_attempt_at TEXT NOT NULL,
115  receipt_json TEXT,
116  last_error TEXT,
117  created_at TEXT NOT NULL,
118  updated_at TEXT NOT NULL
119);
120CREATE INDEX IF NOT EXISTS intents_state_due_idx ON intents(state,next_attempt_at);
121CREATE TABLE IF NOT EXISTS out_head_cache(
122  task_id TEXT PRIMARY KEY,
123  head TEXT NOT NULL
124);
125CREATE TABLE IF NOT EXISTS prose_cache(
126  role TEXT PRIMARY KEY,
127  prose TEXT NOT NULL,
128  spec_hash TEXT NOT NULL,
129  cached_at TEXT NOT NULL
130);
131CREATE TABLE IF NOT EXISTS config_cache(
132  key TEXT PRIMARY KEY,
133  value TEXT NOT NULL
134);
135-- The client keeps its own fault rows: the bridge records them locally, so a
136-- restart still shows what the process found wrong.
137CREATE TABLE IF NOT EXISTS faults(
138  id INTEGER PRIMARY KEY AUTOINCREMENT,
139  task_id TEXT,
140  role TEXT,
141  session_id TEXT,
142  generation INTEGER,
143  seq INTEGER,
144  desired_json TEXT,
145  observed_json TEXT,
146  intent TEXT,
147  -- An exhausted intent is observable through the local fault and the fault
148  -- report; the attempt count is what an operator reads before intervening.
149  attempt INTEGER,
150  backend_ref TEXT,
151  kind TEXT NOT NULL,
152  reason TEXT NOT NULL,
153  state TEXT NOT NULL,
154  -- Encoded from the kernel's unix seconds through this crate's own helper on
155  -- every write.
156  created_at TEXT NOT NULL
157);
158CREATE INDEX IF NOT EXISTS faults_task_kind_generation_idx ON faults(task_id,kind,generation);
159CREATE TABLE IF NOT EXISTS events(
160  seq INTEGER PRIMARY KEY,
161  type TEXT NOT NULL,
162  data_json TEXT NOT NULL,
163  created_at TEXT NOT NULL
164);
165CREATE INDEX IF NOT EXISTS events_type_idx ON events(type);"#;
166
167#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
168pub struct IntentRow {
169    pub op_id: String,
170    pub env_json: Value,
171    pub attempt: i64,
172    pub state: String,
173    pub next_attempt_at: String,
174    pub receipt_json: Option<Value>,
175    pub last_error: Option<String>,
176    pub created_at: String,
177    pub updated_at: String,
178}
179
180/// One task record, as the `task` table holds it. The chain columns come from
181/// the delivery envelope's causality, so a task says who caused it and how deep
182/// it sits without its owner being asked. `task_state` is the whole answer to
183/// "how did this end": no session row carries it, and `None` from [`ClientStore::task`]
184/// means this role never opened the task, not that it is in flight.
185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
186pub struct TaskRow {
187    pub task_id: String,
188    pub kind: Option<String>,
189    pub parent_task: Option<String>,
190    pub hop: i64,
191    pub attempt: i64,
192    pub task_state: TaskState,
193    pub opened_at: String,
194    pub settled_at: Option<String>,
195}
196
197#[derive(Clone, Debug)]
198pub struct ClientStore {
199    path: PathBuf,
200    inner: Arc<Mutex<Connection>>,
201}
202
203impl ClientStore {
204    pub fn open(path: impl AsRef<Path>) -> StoreResult<Self> {
205        let path = path.as_ref().to_path_buf();
206        let conn = open_connection(
207            &path,
208            "client",
209            CLIENT_MARKER,
210            CLIENT_DDL,
211            CLIENT_SCHEMA_VERSION,
212        )?;
213        Ok(Self {
214            path,
215            inner: Arc::new(Mutex::new(conn)),
216        })
217    }
218
219    pub fn path(&self) -> &Path {
220        &self.path
221    }
222
223    pub fn enqueue_intent(&self, op_id: &str, env_json: &Value) -> StoreResult<bool> {
224        let conn = self.conn()?;
225        let now = rfc3339(Utc::now());
226        let changed = conn.execute(
227            "INSERT OR IGNORE INTO intents(op_id,env_json,attempt,state,next_attempt_at,receipt_json,last_error,created_at,updated_at) VALUES(?,?,0,'pending',?,NULL,NULL,?,?)",
228            params![op_id, serde_json::to_string(env_json)?, now, now, now],
229        )?;
230        Ok(changed == 1)
231    }
232
233    /// The rows due by `now`, oldest deadline first, at most `limit` of them.
234    ///
235    /// This is [`ClientStore::flush_order`] with the clock supplied by the
236    /// caller, which is how a test reads a retry the machine has parked in the
237    /// future without waiting on it.
238    pub fn due_intents(&self, now: DateTime<Utc>, limit: u32) -> StoreResult<Vec<IntentRow>> {
239        let conn = self.conn()?;
240        let rows = conn
241            .prepare(
242                "SELECT op_id,env_json,attempt,state,next_attempt_at,receipt_json,last_error,created_at,updated_at FROM intents WHERE state IN ('pending','retrying') AND next_attempt_at<=? ORDER BY next_attempt_at,created_at,rowid LIMIT ?",
243            )?
244            .query_map(params![rfc3339(now), sql_limit(limit)], intent_row)?
245            .collect::<Result<Vec<_>, _>>()?;
246        Ok(rows)
247    }
248
249    /// Store a caller-supplied retry time and move the intent to `retrying`.
250    /// The caller owns the attempt ceiling; the client's `IntentMachine` checks `attempts` and calls `exhaust_intent` when the ceiling is reached.
251    pub fn bump_intent(
252        &self,
253        op_id: &str,
254        next_attempt_at: DateTime<Utc>,
255        error: &str,
256    ) -> StoreResult<bool> {
257        let conn = self.conn()?;
258        let changed = conn.execute(
259            "UPDATE intents SET state='retrying',attempt=attempt+1,next_attempt_at=?,last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
260            params![rfc3339(next_attempt_at), error, rfc3339(Utc::now()), op_id],
261        )?;
262        Ok(changed == 1)
263    }
264
265    /// Push an intent's next attempt forward without consuming its budget.
266    ///
267    /// A transport failure is not the server refusing the message, so the plan's
268    /// disconnect rule keeps the queue intact and the reconnect flushes it
269    /// (plan §6 line 289). Counting those against `intent.attempts` would drop a
270    /// completion that was written while the link was down.
271    pub fn defer_intent(
272        &self,
273        op_id: &str,
274        next_attempt_at: DateTime<Utc>,
275        error: &str,
276    ) -> StoreResult<bool> {
277        let conn = self.conn()?;
278        let changed = conn.execute(
279            "UPDATE intents SET state='retrying',next_attempt_at=?,last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
280            params![rfc3339(next_attempt_at), error, rfc3339(Utc::now()), op_id],
281        )?;
282        Ok(changed == 1)
283    }
284
285    pub fn accept_intent(&self, op_id: &str, receipt_json: &Value) -> StoreResult<bool> {
286        let conn = self.conn()?;
287        let changed = conn.execute(
288            "UPDATE intents SET state='accepted',receipt_json=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
289            params![serde_json::to_string(receipt_json)?, rfc3339(Utc::now()), op_id],
290        )?;
291        Ok(changed == 1)
292    }
293
294    pub fn exhaust_intent(&self, op_id: &str, error: &str) -> StoreResult<bool> {
295        let conn = self.conn()?;
296        let changed = conn.execute(
297            "UPDATE intents SET state='exhausted',last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
298            params![error, rfc3339(Utc::now()), op_id],
299        )?;
300        Ok(changed == 1)
301    }
302
303    pub fn pending_intent_count(&self) -> StoreResult<i64> {
304        let conn = self.conn()?;
305        Ok(conn.query_row(
306            "SELECT COUNT(*) FROM intents WHERE state IN ('pending','retrying')",
307            [],
308            |r| r.get(0),
309        )?)
310    }
311
312    /// Remove a row the server refused for good.
313    ///
314    /// The row is gone rather than parked in a terminal state because a
315    /// permanent refusal will never turn into an acceptance, and the queue must
316    /// not carry it forward on every pass. Like every write here it runs on the
317    /// store's one connection, under the store's lock and its busy timeout: a
318    /// caller that opened a second handle to the file to drop this row wrote
319    /// around that serialization, and whether the row went depended on how that
320    /// handle happened to treat a busy database.
321    pub fn delete_intent(&self, op_id: &str) -> StoreResult<bool> {
322        let conn = self.conn()?;
323        Ok(conn.execute("DELETE FROM intents WHERE op_id=?", params![op_id])? == 1)
324    }
325
326    /// The rows one flush pass should send, in the order it should send them.
327    ///
328    /// Three bounds, each of which the last pass lacked. Only a row whose
329    /// `next_attempt_at` has arrived is returned, so a retry the machine pushed
330    /// into the future stays put until it is due instead of being sent again on
331    /// the strength of being queued; the deadline orders the batch, so the row
332    /// that has waited longest is first and cannot be starved by later arrivals;
333    /// and the batch is capped at [`INTENT_FLUSH_BATCH_SIZE`] rows, so reading a
334    /// queue an offline client filled costs one bounded read.
335    ///
336    /// A caller that means to look past the current deadline — a test asking what
337    /// a run left in the queue — reads [`ClientStore::due_intents`] with its own
338    /// horizon instead, which is the same query with the clock it supplies.
339    pub fn flush_order(&self) -> StoreResult<Vec<IntentRow>> {
340        self.due_intents(Utc::now(), INTENT_FLUSH_BATCH_SIZE)
341    }
342
343    pub fn append_event(&self, kind: &str, data: &Value) -> StoreResult<i64> {
344        let conn = self.conn()?;
345        append_event_conn(&conn, kind, data)
346    }
347
348    pub fn events_since(&self, seq: i64, limit: u32) -> StoreResult<Vec<EventRecord>> {
349        let conn = self.conn()?;
350        events_since_conn(&conn, seq, limit)
351    }
352
353    pub fn event_head(&self) -> StoreResult<i64> {
354        let conn = self.conn()?;
355        event_head_conn(&conn)
356    }
357
358    pub fn put_out_head(&self, task_id: &str, head: &str) -> StoreResult<bool> {
359        let conn = self.conn()?;
360        Ok(conn.execute(
361            "INSERT INTO out_head_cache(task_id,head) VALUES(?,?) ON CONFLICT(task_id) DO UPDATE SET head=excluded.head",
362            params![task_id, crate::server::head_preview(head)],
363        )? == 1)
364    }
365
366    pub fn out_head(&self, task_id: &str) -> StoreResult<Option<String>> {
367        let conn = self.conn()?;
368        Ok(conn
369            .query_row(
370                "SELECT head FROM out_head_cache WHERE task_id=?",
371                params![task_id],
372                |r| r.get(0),
373            )
374            .optional()?)
375    }
376
377    pub fn put_prose(&self, role: &str, prose: &str, spec_hash: &str) -> StoreResult<bool> {
378        let conn = self.conn()?;
379        Ok(conn.execute(
380            "INSERT INTO prose_cache(role,prose,spec_hash,cached_at) VALUES(?,?,?,?) ON CONFLICT(role) DO UPDATE SET prose=excluded.prose,spec_hash=excluded.spec_hash,cached_at=excluded.cached_at",
381            params![role, prose, spec_hash, rfc3339(Utc::now())],
382        )? == 1)
383    }
384
385    pub fn prose(&self, role: &str) -> StoreResult<Option<(String, String)>> {
386        let conn = self.conn()?;
387        Ok(conn
388            .query_row(
389                "SELECT prose,spec_hash FROM prose_cache WHERE role=?",
390                params![role],
391                |r| Ok((r.get(0)?, r.get(1)?)),
392            )
393            .optional()?)
394    }
395
396    pub fn put_config(&self, key: &str, value: &str) -> StoreResult<bool> {
397        let conn = self.conn()?;
398        Ok(conn.execute(
399            "INSERT INTO config_cache(key,value) VALUES(?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value",
400            params![key, value],
401        )? == 1)
402    }
403
404    pub fn config(&self, key: &str) -> StoreResult<Option<String>> {
405        let conn = self.conn()?;
406        Ok(conn
407            .query_row(
408                "SELECT value FROM config_cache WHERE key=?",
409                params![key],
410                |r| r.get(0),
411            )
412            .optional()?)
413    }
414
415    /// Write a heartbeat's (generation, seq) onto the stored row when that
416    /// pair is strictly newer. The observation tuple stays as stored. Returns
417    /// true when the row changed.
418    pub fn bump_session_version(
419        &self,
420        task_id: &str,
421        generation: u64,
422        seq: u64,
423    ) -> StoreResult<bool> {
424        let conn = self.conn()?;
425        let generation = generation as i64;
426        let seq = seq as i64;
427        let changed = conn.execute(
428            &format!(
429                "UPDATE sessions SET generation=?,seq=?,updated_at=? WHERE session_id={SESSION_ID_FOR_TASK} AND (? > generation OR (? = generation AND ? > seq))"
430            ),
431            params![
432                generation,
433                seq,
434                rfc3339(Utc::now()),
435                task_id,
436                generation,
437                generation,
438                seq
439            ],
440        )?;
441        Ok(changed == 1)
442    }
443
444    /// Open the task record for one delivery.
445    ///
446    /// The causality of the envelope that became a session is the record: the
447    /// task id, its parent, its depth, and the redelivery count it arrived on.
448    /// An already-open task keeps its settle, its `opened_at`, and its verdict —
449    /// a re-dispatch refreshes only the chain it arrived on, and takes the
450    /// larger attempt so the count never runs backwards. `kind` is how the
451    /// delivery reached this role — `root` or `relay` — never the envelope's
452    /// message kind, which says nothing a reader can act on. Answers true when
453    /// the record was written, which covers both the fresh open and a chain
454    /// refresh on a task that is already there.
455    pub fn open_task(&self, causality: &Causality, kind: &str) -> StoreResult<bool> {
456        let conn = self.conn()?;
457        let now = rfc3339(Utc::now());
458        let changed = conn.execute(
459            "INSERT INTO task(task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at) VALUES(?,?,?,?,?,'pending',?,NULL)
460             ON CONFLICT(task_id) DO UPDATE SET kind=excluded.kind,parent_task=excluded.parent_task,hop=excluded.hop,attempt=MAX(excluded.attempt,task.attempt)",
461            params![
462                causality.task,
463                kind,
464                causality.parent_task,
465                i64::from(causality.hop),
466                i64::from(causality.attempt),
467                now
468            ],
469        )?;
470        Ok(changed == 1)
471    }
472
473    /// Settle one task with the verdict its completion carries.
474    ///
475    /// The first terminal verdict wins. A task sitting open takes the verdict
476    /// and stamps `settled_at`; one already settled keeps its record and answers
477    /// `false`, which is what stops a zombie session that came back after a
478    /// newer one answered from rewriting the answer it already gave.
479    ///
480    /// The write never depends on the open having run. A task this process never
481    /// dispatched — a completion reported for work an earlier process opened, or
482    /// a delivery that found its slot already serving and returned before the
483    /// open — gets its record here, with no kind, no parent, and both clocks at
484    /// this verdict. Without that, the verdict would be dropped on the floor and
485    /// the session could never project its way out of `working`.
486    pub fn settle_task(&self, task_id: &str, task_state: TaskState) -> StoreResult<bool> {
487        let conn = self.conn()?;
488        let now = rfc3339(Utc::now());
489        let changed = conn.execute(
490            "INSERT INTO task(task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at) VALUES(?,NULL,NULL,0,0,?,?,?)
491             ON CONFLICT(task_id) DO UPDATE SET task_state=excluded.task_state,settled_at=excluded.settled_at WHERE task.settled_at IS NULL",
492            params![task_id, string_tag(&task_state)?, now, now],
493        )?;
494        Ok(changed == 1)
495    }
496
497    /// The record of one task, when this role opened it.
498    pub fn task(&self, task_id: &str) -> StoreResult<Option<TaskRow>> {
499        let conn = self.conn()?;
500        let row = conn
501            .query_row(
502                "SELECT task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at FROM task WHERE task_id=?",
503                params![task_id],
504                task_row,
505            )
506            .optional()?;
507        Ok(row)
508    }
509
510    /// Tasks this role has opened and not settled, oldest first.
511    pub fn open_tasks(&self, limit: u32) -> StoreResult<Vec<TaskRow>> {
512        let conn = self.conn()?;
513        let rows = conn
514            .prepare(
515                "SELECT task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at FROM task WHERE settled_at IS NULL ORDER BY opened_at,rowid LIMIT ?",
516            )?
517            .query_map(params![sql_limit(limit)], task_row)?
518            .collect::<Result<Vec<_>, _>>()?;
519        Ok(rows)
520    }
521
522    /// The sessions this client holds and can recover, for the `hello` claim:
523    /// each one's id, the delivery it is on now, and whether its process has
524    /// been released.
525    ///
526    /// A fresh process after a crash has empty slots and this table is the only
527    /// witness left that the sessions exist, so it must declare them to prevent
528    /// duplicate dispatch.
529    ///
530    /// Only claims sessions whose runtime may still be there: the ones that have
531    /// not mounted yet or are between turns (`booting`, `ready`, `idle`). A
532    /// session that was mid-turn when this process died is not claimed — the
533    /// plugin that was running it is gone — so the server requeues its
534    /// deliveries instead.
535    ///
536    /// A suspended session — one whose resource is `closed` while its agent has
537    /// not gone — is claimed with `suspended: true`. The work it holds is still
538    /// owed, so the server holds its row rather than requeueing it, and the
539    /// session resumes that delivery instead of opening a second one. An exited
540    /// session (`agent_state` `gone`) is never claimed: it holds nothing.
541    pub fn active_sessions(&self) -> StoreResult<Vec<LiveSession>> {
542        let conn = self.conn()?;
543        let mut stmt = conn.prepare(&format!(
544            "SELECT session_id,{OPEN_BINDING_TASK},resource_state FROM sessions \
545             WHERE agent_state IN ('booting', 'ready', 'idle') \
546             AND (resource_state != 'closed' OR agent_state IN ('ready', 'idle')) \
547             ORDER BY session_id"
548        ))?;
549        let rows = stmt
550            .query_map(params![], |row| {
551                let resource_state: String = row.get(2)?;
552                Ok(LiveSession {
553                    session_id: row.get(0)?,
554                    task_id: row.get(1)?,
555                    suspended: resource_state == "closed",
556                })
557            })?
558            .collect::<Result<Vec<_>, _>>()?;
559        Ok(rows)
560    }
561
562    /// Stop serving one delivery: its `session_tasks` row gets `released_at`.
563    ///
564    /// A `task` or `role` scope session outlives the delivery it served, so
565    /// settling that delivery is not the end of the binding. Releasing it here
566    /// is what stops the session from reading as bound to settled work: the next
567    /// delivery's binding would close the other open one anyway
568    /// ([`crate::server::open_binding_conn`]), but a session that goes idle with
569    /// no next delivery keeps answering `OPEN_BINDING_TASK` with a task that is
570    /// over, and `hello` then claims it for work nobody owes. Only an open
571    /// binding is released, so the first call is the one the row keeps and a
572    /// second changes nothing; the count is how many rows moved.
573    pub fn release_binding(&self, session_id: &str, task_id: &str) -> StoreResult<usize> {
574        let conn = self.conn()?;
575        release_binding_conn(&conn, session_id, task_id, Utc::now().timestamp())
576    }
577
578    /// Take one delivery for a session: its `session_tasks` row opens now.
579    ///
580    /// A scoped session serves its next delivery without its tuple moving, and
581    /// the binding has to open first: `SESSION_ID_FOR_TASK` reads it, so a write
582    /// for a delivery whose binding is not open yet lands on no row at all.
583    /// Opening a pair also closes the session's other open binding, so the
584    /// session is on one delivery at a time either way. The count is how many
585    /// rows the insert moved — zero for a pair already open.
586    pub fn bind_task(&self, session_id: &str, task_id: &str) -> StoreResult<usize> {
587        let conn = self.conn()?;
588        open_binding_conn(&conn, session_id, task_id, Utc::now().timestamp())
589    }
590
591    fn conn(&self) -> StoreResult<MutexGuard<'_, Connection>> {
592        self.inner
593            .lock()
594            .map_err(|_| StoreError::Sqlite("database mutex poisoned".to_string()))
595    }
596}
597
598impl SessionLedger for ClientStore {
599    /// The session serving one delivery, read through its binding.
600    ///
601    /// The record answers for the delivery the caller named: a session that
602    /// served one delivery and then another still answers for the first, which
603    /// is what a caller holding that delivery's tuple asks for.
604    fn get_session(&self, task_id: &str) -> anyhow::Result<Option<SessionRecord>> {
605        let conn = self.conn()?;
606        let row = conn
607            .query_row(
608                &format!(
609                    "SELECT {} FROM sessions WHERE session_id={SESSION_ID_FOR_TASK}",
610                    CLIENT_SESSION_COLUMNS
611                ),
612                params![task_id],
613                |row| session_record_row(row, task_id),
614            )
615            .optional()?;
616        Ok(row)
617    }
618
619    /// Write one session tuple, and bind the delivery it is about.
620    ///
621    /// The row is addressed by the session's own id, which the backend
622    /// reference names; the delivery the tuple serves is the binding this write
623    /// opens, and a session serves one delivery at a time.
624    fn upsert_session(&self, task_id: &str, version: &VersionedSession) -> anyhow::Result<bool> {
625        let conn = self.conn()?;
626        let (backend, session_id) = backend_parts(task_id, &version.backend_ref);
627        let tx = conn.unchecked_transaction()?;
628        let changed = tx.execute(
629            "INSERT INTO sessions(session_id,generation,seq,agent_state,delivery_state,resource_state,recovery_substate,observed_json,backend,backend_ref,desired_json,mismatch_count,last_seen,updated_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)
630             ON CONFLICT(session_id) DO UPDATE SET 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,observed_json=excluded.observed_json,backend=COALESCE(excluded.backend,sessions.backend),backend_ref=excluded.backend_ref,desired_json=excluded.desired_json,mismatch_count=excluded.mismatch_count,last_seen=excluded.last_seen,updated_at=excluded.updated_at
631             WHERE excluded.generation > sessions.generation OR (excluded.generation = sessions.generation AND excluded.seq > sessions.seq)",
632            params![
633                &session_id,
634                version.generation,
635                version.seq,
636                version.agent_state,
637                version.delivery_state,
638                version.resource_state,
639                version.recovery_substate,
640                version.observed_json,
641                backend,
642                version.backend_ref,
643                version.desired_json,
644                version.mismatch_count,
645                crate::server::unix_to_rfc3339(version.updated_at),
646                crate::server::unix_to_rfc3339(version.updated_at)
647            ],
648        )?;
649        if changed == 1 {
650            crate::server::open_binding_conn(&tx, &session_id, task_id, version.updated_at)?;
651        }
652        tx.commit()?;
653        Ok(changed == 1)
654    }
655
656    fn task_is_known(&self, task_id: &str) -> anyhow::Result<bool> {
657        let conn = self.conn()?;
658        // The binding is the record of a delivery having become a session here,
659        // which is exactly what this asks.
660        let session_count: i64 = conn.query_row(
661            "SELECT COUNT(*) FROM session_tasks WHERE task_id=?",
662            params![task_id],
663            |r| r.get(0),
664        )?;
665        if session_count > 0 {
666            return Ok(true);
667        }
668        Ok(intent_stats(&conn, task_id)?.is_some())
669    }
670
671    fn task_attempt(&self, task_id: &str) -> anyhow::Result<i64> {
672        let conn = self.conn()?;
673        Ok(intent_stats(&conn, task_id)?.unwrap_or(0))
674    }
675
676    fn list_faults(&self, task_id: &str) -> anyhow::Result<Vec<FaultRecord>> {
677        let conn = self.conn()?;
678        let rows = conn
679            .prepare(
680                "SELECT id,task_id,session_id,generation,seq,desired_json,observed_json,intent,attempt,backend_ref,kind,reason,state,created_at FROM faults WHERE task_id=? ORDER BY id",
681            )?
682            .query_map(params![task_id], fault_record_row)?
683            .collect::<Result<Vec<_>, _>>()?;
684        Ok(rows)
685    }
686
687    fn insert_fault(&self, fault: &FaultRecord) -> anyhow::Result<i64> {
688        let conn = self.conn()?;
689        conn.execute(
690            "INSERT INTO faults(task_id,role,session_id,generation,seq,desired_json,observed_json,intent,attempt,backend_ref,kind,reason,state,created_at) VALUES(?,NULL,?,?,?,?,?,?,?,?,?,?,?,?)",
691            params![
692                fault.task_id,
693                fault.session_id,
694                fault.generation,
695                fault.seq,
696                fault.desired_json,
697                fault.observed_json,
698                fault.intent,
699                fault.attempt,
700                fault.backend_ref,
701                fault.kind,
702                fault.reason,
703                fault.state,
704                crate::server::unix_to_rfc3339(fault.created_at)
705            ],
706        )?;
707        Ok(conn.last_insert_rowid())
708    }
709
710    fn emit(&self, kind: &str, data: Value) {
711        if let Err(err) = self.append_event(kind, &data) {
712            tracing::warn!(kind, error = %err, "session ledger event was not stored");
713        }
714    }
715
716    fn note_alert(&self, line: String) {
717        tracing::warn!(alert = %line, "session ledger alert");
718    }
719}
720
721fn intent_stats(conn: &Connection, task_id: &str) -> StoreResult<Option<i64>> {
722    let mut stmt = conn.prepare("SELECT env_json,attempt FROM intents")?;
723    let rows = stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?)))?;
724    let mut max_attempt: Option<i64> = None;
725    for row in rows {
726        let (env_json, attempt) = row?;
727        if intent_task_matches(&env_json, task_id) {
728            max_attempt = Some(max_attempt.map_or(attempt, |current| current.max(attempt)));
729        }
730    }
731    Ok(max_attempt)
732}
733
734fn intent_task_matches(env_json: &str, task_id: &str) -> bool {
735    let Ok(value) = serde_json::from_str::<Value>(env_json) else {
736        return false;
737    };
738    value
739        .get("causality")
740        .and_then(|v| v.get("task"))
741        .and_then(Value::as_str)
742        == Some(task_id)
743}
744
745/// What the stored backend reference says about the session itself: the backend
746/// that ran it, and the id this client files the session under.
747///
748/// `backend_ref` is the whole `SessionRef` a backend handed back, serialized,
749/// so its `task_id` is the session's own key on this side — a client-held
750/// session shares it with the delivery that opened it, which is also the id
751/// that delivery is reported under. A reference carrying neither leaves the
752/// session under the delivery it opened, which is the spelling this store has
753/// always used.
754fn backend_parts(task_id: &str, backend_ref: &str) -> (Option<String>, String) {
755    let parsed = serde_json::from_str::<Value>(backend_ref).ok();
756    let backend = parsed
757        .as_ref()
758        .and_then(|v| v.get("backend"))
759        .and_then(Value::as_str)
760        .map(str::to_string);
761    let session_id = parsed
762        .as_ref()
763        .and_then(|v| v.get("task_id"))
764        .and_then(Value::as_str)
765        .unwrap_or(task_id)
766        .to_string();
767    (backend, session_id)
768}
769
770/// One stored session tuple's columns, in the order [`session_record_row`]
771/// reads them.
772const CLIENT_SESSION_COLUMNS: &str = "session_id,agent_state,delivery_state,resource_state,recovery_substate,desired_json,observed_json,generation,seq,backend_ref,mismatch_count,updated_at";
773
774/// One stored tuple, carrying its own key beside the delivery the caller asked
775/// about: the row answers for a session, and the delivery it serves is the
776/// binding, which the caller holds.
777fn session_record_row(r: &Row<'_>, task_id: &str) -> rusqlite::Result<SessionRecord> {
778    Ok(SessionRecord {
779        session_id: r.get(0)?,
780        task_id: task_id.to_string(),
781        agent_state: r.get(1)?,
782        delivery_state: r.get(2)?,
783        resource_state: r.get(3)?,
784        recovery_substate: r.get(4)?,
785        desired_json: r.get(5)?,
786        observed_json: r.get(6)?,
787        generation: r.get(7)?,
788        seq: r.get(8)?,
789        backend_ref: r.get(9)?,
790        mismatch_count: r.get(10)?,
791        updated_at: crate::server::rfc3339_to_unix(&r.get::<_, String>(11)?),
792    })
793}
794
795fn task_row(r: &Row<'_>) -> rusqlite::Result<TaskRow> {
796    let state_word: String = r.get(5)?;
797    let task_state = serde_json::from_value(Value::String(state_word.clone()))
798        .map_err(|e| conversion_error(5, format!("task_state column holds {state_word:?}: {e}")))?;
799    Ok(TaskRow {
800        task_id: r.get(0)?,
801        kind: r.get(1)?,
802        parent_task: r.get(2)?,
803        hop: r.get(3)?,
804        attempt: r.get(4)?,
805        task_state,
806        opened_at: r.get(6)?,
807        settled_at: r.get(7)?,
808    })
809}
810
811fn fault_record_row(r: &Row<'_>) -> rusqlite::Result<FaultRecord> {
812    Ok(FaultRecord {
813        id: r.get(0)?,
814        task_id: r.get::<_, Option<String>>(1)?.unwrap_or_default(),
815        session_id: r.get::<_, Option<String>>(2)?.unwrap_or_default(),
816        generation: r.get::<_, Option<i64>>(3)?.unwrap_or_default(),
817        seq: r.get::<_, Option<i64>>(4)?.unwrap_or_default(),
818        desired_json: r
819            .get::<_, Option<String>>(5)?
820            .unwrap_or_else(|| "{}".to_string()),
821        observed_json: r
822            .get::<_, Option<String>>(6)?
823            .unwrap_or_else(|| "{}".to_string()),
824        intent: r.get::<_, Option<String>>(7)?.unwrap_or_default(),
825        attempt: r.get::<_, Option<i64>>(8)?.unwrap_or_default(),
826        backend_ref: r
827            .get::<_, Option<String>>(9)?
828            .unwrap_or_else(|| "{}".to_string()),
829        kind: r.get(10)?,
830        reason: r.get(11)?,
831        state: r.get(12)?,
832        created_at: r
833            .get::<_, Option<String>>(13)?
834            .map(|text| crate::server::rfc3339_to_unix(&text))
835            .unwrap_or_default(),
836    })
837}
838
839fn intent_row(r: &Row<'_>) -> rusqlite::Result<IntentRow> {
840    let env_text: String = r.get(1)?;
841    let receipt_text: Option<String> = r.get(5)?;
842    let env_json =
843        serde_json::from_str(&env_text).map_err(|e| conversion_error(1, e.to_string()))?;
844    let receipt_json = receipt_text
845        .map(|text| serde_json::from_str(&text).map_err(|e| conversion_error(5, e.to_string())))
846        .transpose()?;
847    Ok(IntentRow {
848        op_id: r.get(0)?,
849        env_json,
850        attempt: r.get(2)?,
851        state: r.get(3)?,
852        next_attempt_at: r.get(4)?,
853        receipt_json,
854        last_error: r.get(6)?,
855        created_at: r.get(7)?,
856        updated_at: r.get(8)?,
857    })
858}
859
860fn sql_limit(limit: u32) -> i64 {
861    if limit == 0 {
862        DEFAULT_LIMIT
863    } else {
864        i64::from(limit.min(500))
865    }
866}
867
868fn conversion_error(index: usize, message: String) -> rusqlite::Error {
869    rusqlite::Error::FromSqlConversionFailure(
870        index,
871        Type::Text,
872        Box::new(StoreError::Serialization(message)),
873    )
874}