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";
22const CLIENT_SCHEMA_VERSION: i64 = 3;
32const DEFAULT_LIMIT: i64 = 100;
33pub 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#[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
745fn 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
770const 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
774fn 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}