use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, MutexGuard};
use crate::session::{FaultRecord, SessionLedger, SessionRecord, VersionedSession};
use chrono::{DateTime, Utc};
use onlyne_proto::Causality;
use onlyne_proto::LiveSession;
use onlyne_proto::TaskState;
use rusqlite::types::Type;
use rusqlite::{Connection, OptionalExtension, Row, params};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::error::{StoreError, StoreResult};
use crate::server::{
EventRecord, OPEN_BINDING_TASK, SESSION_ID_FOR_TASK, append_event_conn, event_head_conn,
events_since_conn, open_binding_conn, open_connection, release_binding_conn, rfc3339,
string_tag,
};
const CLIENT_MARKER: &str = "onlyne-client";
const CLIENT_SCHEMA_VERSION: i64 = 3;
const DEFAULT_LIMIT: i64 = 100;
pub const INTENT_FLUSH_BATCH_SIZE: u32 = 100;
pub const CLIENT_DDL: &str = r#"-- The client's own mirror of the sessions it holds, keyed the way the server's
-- mirror is: a session row answers for a session, and which delivery that
-- session serves lives in `session_tasks` below.
--
-- Two columns are the client's alone, because only the process that holds a
-- session can answer them: `backend` names the backend that ran it (`acp`,
-- `exec`, `orca`, `external`, … — which one comes from the role's drive under
-- the placement this machine resolved) and `backend_ref` is the reference that
-- backend answers to. There is no `role` column: one client serves one role,
-- and the workspace config owns that fact.
CREATE TABLE IF NOT EXISTS sessions(
session_id TEXT PRIMARY KEY,
generation INTEGER NOT NULL,
seq INTEGER NOT NULL,
agent_state TEXT NOT NULL,
delivery_state TEXT NOT NULL,
resource_state TEXT NOT NULL,
recovery_substate TEXT NOT NULL,
observed_json TEXT NOT NULL,
backend TEXT,
backend_ref TEXT,
-- The (generation, seq) gate. The reducer's isolate-after-N and
-- terminate-after-N policy needs a persisted counter, so this pair carries
-- DEFAULT_ISOLATE_AFTER and DEFAULT_TERMINATE_AFTER.
desired_json TEXT NOT NULL,
mismatch_count INTEGER NOT NULL DEFAULT 0,
-- When this client last wrote about the session. A reader judges a row's
-- freshness by it, and the tuple's own version lives in the pair above.
last_seen TEXT NOT NULL,
-- When the tuple last moved. Both clocks are the kernel's unix seconds,
-- encoded through this crate's own helper on every write.
updated_at TEXT NOT NULL
);
-- Which delivery a session serves: the same table, columns and keys as the
-- server's, and the only place a binding lives on this side too.
CREATE TABLE IF NOT EXISTS session_tasks(
session_id TEXT NOT NULL,
task_id TEXT NOT NULL,
bound_at TEXT NOT NULL,
released_at TEXT,
PRIMARY KEY (session_id, task_id)
);
CREATE INDEX IF NOT EXISTS session_tasks_task_idx ON session_tasks(task_id);
CREATE INDEX IF NOT EXISTS session_tasks_open_idx ON session_tasks(session_id, released_at);
-- The task's own record. How its work ended is not a session dimension: the
-- session tuple says what this session can prove about its agent, intent,
-- resource, and recovery line, and `project` needs the task's verdict handed
-- in before it can say whether the session is over. A row is opened when a
-- delivery becomes a session, from the envelope's causality, and settled once;
-- a verdict that arrives for a task this process never opened writes its own
-- row, so a completion is never dropped because of who opened what. `kind` is
-- how the delivery reached this role: `root` for work given here, `relay` for
-- work handed down from a parent task.
CREATE TABLE IF NOT EXISTS task(
task_id TEXT PRIMARY KEY,
kind TEXT,
parent_task TEXT,
hop INTEGER NOT NULL DEFAULT 0,
attempt INTEGER NOT NULL DEFAULT 0,
-- TaskState's own snake_case tag from onlyne-proto's lifecycle vocabulary.
-- `pending` is the open state, and the settle write is the only thing that
-- leaves it, so `settled_at IS NULL` and `task_state = 'pending'` answer the
-- same question.
task_state TEXT NOT NULL,
opened_at TEXT NOT NULL,
settled_at TEXT
);
CREATE TABLE IF NOT EXISTS intents(
op_id TEXT PRIMARY KEY,
env_json TEXT NOT NULL,
attempt INTEGER NOT NULL,
state TEXT NOT NULL,
next_attempt_at TEXT NOT NULL,
receipt_json TEXT,
last_error TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS intents_state_due_idx ON intents(state,next_attempt_at);
CREATE TABLE IF NOT EXISTS out_head_cache(
task_id TEXT PRIMARY KEY,
head TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS prose_cache(
role TEXT PRIMARY KEY,
prose TEXT NOT NULL,
spec_hash TEXT NOT NULL,
cached_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS config_cache(
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
-- The client keeps its own fault rows: the bridge records them locally, so a
-- restart still shows what the process found wrong.
CREATE TABLE IF NOT EXISTS faults(
id INTEGER PRIMARY KEY AUTOINCREMENT,
task_id TEXT,
role TEXT,
session_id TEXT,
generation INTEGER,
seq INTEGER,
desired_json TEXT,
observed_json TEXT,
intent TEXT,
-- An exhausted intent is observable through the local fault and the fault
-- report; the attempt count is what an operator reads before intervening.
attempt INTEGER,
backend_ref TEXT,
kind TEXT NOT NULL,
reason TEXT NOT NULL,
state TEXT NOT NULL,
-- Encoded from the kernel's unix seconds through this crate's own helper on
-- every write.
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS faults_task_kind_generation_idx ON faults(task_id,kind,generation);
CREATE TABLE IF NOT EXISTS events(
seq INTEGER PRIMARY KEY,
type TEXT NOT NULL,
data_json TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS events_type_idx ON events(type);"#;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IntentRow {
pub op_id: String,
pub env_json: Value,
pub attempt: i64,
pub state: String,
pub next_attempt_at: String,
pub receipt_json: Option<Value>,
pub last_error: Option<String>,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TaskRow {
pub task_id: String,
pub kind: Option<String>,
pub parent_task: Option<String>,
pub hop: i64,
pub attempt: i64,
pub task_state: TaskState,
pub opened_at: String,
pub settled_at: Option<String>,
}
#[derive(Clone, Debug)]
pub struct ClientStore {
path: PathBuf,
inner: Arc<Mutex<Connection>>,
}
impl ClientStore {
pub fn open(path: impl AsRef<Path>) -> StoreResult<Self> {
let path = path.as_ref().to_path_buf();
let conn = open_connection(
&path,
"client",
CLIENT_MARKER,
CLIENT_DDL,
CLIENT_SCHEMA_VERSION,
)?;
Ok(Self {
path,
inner: Arc::new(Mutex::new(conn)),
})
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn enqueue_intent(&self, op_id: &str, env_json: &Value) -> StoreResult<bool> {
let conn = self.conn()?;
let now = rfc3339(Utc::now());
let changed = conn.execute(
"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,?,?)",
params![op_id, serde_json::to_string(env_json)?, now, now, now],
)?;
Ok(changed == 1)
}
pub fn due_intents(&self, now: DateTime<Utc>, limit: u32) -> StoreResult<Vec<IntentRow>> {
let conn = self.conn()?;
let rows = conn
.prepare(
"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 ?",
)?
.query_map(params![rfc3339(now), sql_limit(limit)], intent_row)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub fn bump_intent(
&self,
op_id: &str,
next_attempt_at: DateTime<Utc>,
error: &str,
) -> StoreResult<bool> {
let conn = self.conn()?;
let changed = conn.execute(
"UPDATE intents SET state='retrying',attempt=attempt+1,next_attempt_at=?,last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
params![rfc3339(next_attempt_at), error, rfc3339(Utc::now()), op_id],
)?;
Ok(changed == 1)
}
pub fn defer_intent(
&self,
op_id: &str,
next_attempt_at: DateTime<Utc>,
error: &str,
) -> StoreResult<bool> {
let conn = self.conn()?;
let changed = conn.execute(
"UPDATE intents SET state='retrying',next_attempt_at=?,last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
params![rfc3339(next_attempt_at), error, rfc3339(Utc::now()), op_id],
)?;
Ok(changed == 1)
}
pub fn accept_intent(&self, op_id: &str, receipt_json: &Value) -> StoreResult<bool> {
let conn = self.conn()?;
let changed = conn.execute(
"UPDATE intents SET state='accepted',receipt_json=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
params![serde_json::to_string(receipt_json)?, rfc3339(Utc::now()), op_id],
)?;
Ok(changed == 1)
}
pub fn exhaust_intent(&self, op_id: &str, error: &str) -> StoreResult<bool> {
let conn = self.conn()?;
let changed = conn.execute(
"UPDATE intents SET state='exhausted',last_error=?,updated_at=? WHERE op_id=? AND state IN ('pending','retrying')",
params![error, rfc3339(Utc::now()), op_id],
)?;
Ok(changed == 1)
}
pub fn pending_intent_count(&self) -> StoreResult<i64> {
let conn = self.conn()?;
Ok(conn.query_row(
"SELECT COUNT(*) FROM intents WHERE state IN ('pending','retrying')",
[],
|r| r.get(0),
)?)
}
pub fn delete_intent(&self, op_id: &str) -> StoreResult<bool> {
let conn = self.conn()?;
Ok(conn.execute("DELETE FROM intents WHERE op_id=?", params![op_id])? == 1)
}
pub fn flush_order(&self) -> StoreResult<Vec<IntentRow>> {
self.due_intents(Utc::now(), INTENT_FLUSH_BATCH_SIZE)
}
pub fn append_event(&self, kind: &str, data: &Value) -> StoreResult<i64> {
let conn = self.conn()?;
append_event_conn(&conn, kind, data)
}
pub fn events_since(&self, seq: i64, limit: u32) -> StoreResult<Vec<EventRecord>> {
let conn = self.conn()?;
events_since_conn(&conn, seq, limit)
}
pub fn event_head(&self) -> StoreResult<i64> {
let conn = self.conn()?;
event_head_conn(&conn)
}
pub fn put_out_head(&self, task_id: &str, head: &str) -> StoreResult<bool> {
let conn = self.conn()?;
Ok(conn.execute(
"INSERT INTO out_head_cache(task_id,head) VALUES(?,?) ON CONFLICT(task_id) DO UPDATE SET head=excluded.head",
params![task_id, crate::server::head_preview(head)],
)? == 1)
}
pub fn out_head(&self, task_id: &str) -> StoreResult<Option<String>> {
let conn = self.conn()?;
Ok(conn
.query_row(
"SELECT head FROM out_head_cache WHERE task_id=?",
params![task_id],
|r| r.get(0),
)
.optional()?)
}
pub fn put_prose(&self, role: &str, prose: &str, spec_hash: &str) -> StoreResult<bool> {
let conn = self.conn()?;
Ok(conn.execute(
"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",
params![role, prose, spec_hash, rfc3339(Utc::now())],
)? == 1)
}
pub fn prose(&self, role: &str) -> StoreResult<Option<(String, String)>> {
let conn = self.conn()?;
Ok(conn
.query_row(
"SELECT prose,spec_hash FROM prose_cache WHERE role=?",
params![role],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?)
}
pub fn put_config(&self, key: &str, value: &str) -> StoreResult<bool> {
let conn = self.conn()?;
Ok(conn.execute(
"INSERT INTO config_cache(key,value) VALUES(?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value",
params![key, value],
)? == 1)
}
pub fn config(&self, key: &str) -> StoreResult<Option<String>> {
let conn = self.conn()?;
Ok(conn
.query_row(
"SELECT value FROM config_cache WHERE key=?",
params![key],
|r| r.get(0),
)
.optional()?)
}
pub fn bump_session_version(
&self,
task_id: &str,
generation: u64,
seq: u64,
) -> StoreResult<bool> {
let conn = self.conn()?;
let generation = generation as i64;
let seq = seq as i64;
let changed = conn.execute(
&format!(
"UPDATE sessions SET generation=?,seq=?,updated_at=? WHERE session_id={SESSION_ID_FOR_TASK} AND (? > generation OR (? = generation AND ? > seq))"
),
params![
generation,
seq,
rfc3339(Utc::now()),
task_id,
generation,
generation,
seq
],
)?;
Ok(changed == 1)
}
pub fn open_task(&self, causality: &Causality, kind: &str) -> StoreResult<bool> {
let conn = self.conn()?;
let now = rfc3339(Utc::now());
let changed = conn.execute(
"INSERT INTO task(task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at) VALUES(?,?,?,?,?,'pending',?,NULL)
ON CONFLICT(task_id) DO UPDATE SET kind=excluded.kind,parent_task=excluded.parent_task,hop=excluded.hop,attempt=MAX(excluded.attempt,task.attempt)",
params![
causality.task,
kind,
causality.parent_task,
i64::from(causality.hop),
i64::from(causality.attempt),
now
],
)?;
Ok(changed == 1)
}
pub fn settle_task(&self, task_id: &str, task_state: TaskState) -> StoreResult<bool> {
let conn = self.conn()?;
let now = rfc3339(Utc::now());
let changed = conn.execute(
"INSERT INTO task(task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at) VALUES(?,NULL,NULL,0,0,?,?,?)
ON CONFLICT(task_id) DO UPDATE SET task_state=excluded.task_state,settled_at=excluded.settled_at WHERE task.settled_at IS NULL",
params![task_id, string_tag(&task_state)?, now, now],
)?;
Ok(changed == 1)
}
pub fn task(&self, task_id: &str) -> StoreResult<Option<TaskRow>> {
let conn = self.conn()?;
let row = conn
.query_row(
"SELECT task_id,kind,parent_task,hop,attempt,task_state,opened_at,settled_at FROM task WHERE task_id=?",
params![task_id],
task_row,
)
.optional()?;
Ok(row)
}
pub fn open_tasks(&self, limit: u32) -> StoreResult<Vec<TaskRow>> {
let conn = self.conn()?;
let rows = conn
.prepare(
"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 ?",
)?
.query_map(params![sql_limit(limit)], task_row)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub fn active_sessions(&self) -> StoreResult<Vec<LiveSession>> {
let conn = self.conn()?;
let mut stmt = conn.prepare(&format!(
"SELECT session_id,{OPEN_BINDING_TASK},resource_state FROM sessions \
WHERE agent_state IN ('booting', 'ready', 'idle') \
AND (resource_state != 'closed' OR agent_state IN ('ready', 'idle')) \
ORDER BY session_id"
))?;
let rows = stmt
.query_map(params![], |row| {
let resource_state: String = row.get(2)?;
Ok(LiveSession {
session_id: row.get(0)?,
task_id: row.get(1)?,
suspended: resource_state == "closed",
})
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub fn release_binding(&self, session_id: &str, task_id: &str) -> StoreResult<usize> {
let conn = self.conn()?;
release_binding_conn(&conn, session_id, task_id, Utc::now().timestamp())
}
pub fn bind_task(&self, session_id: &str, task_id: &str) -> StoreResult<usize> {
let conn = self.conn()?;
open_binding_conn(&conn, session_id, task_id, Utc::now().timestamp())
}
fn conn(&self) -> StoreResult<MutexGuard<'_, Connection>> {
self.inner
.lock()
.map_err(|_| StoreError::Sqlite("database mutex poisoned".to_string()))
}
}
impl SessionLedger for ClientStore {
fn get_session(&self, task_id: &str) -> anyhow::Result<Option<SessionRecord>> {
let conn = self.conn()?;
let row = conn
.query_row(
&format!(
"SELECT {} FROM sessions WHERE session_id={SESSION_ID_FOR_TASK}",
CLIENT_SESSION_COLUMNS
),
params![task_id],
|row| session_record_row(row, task_id),
)
.optional()?;
Ok(row)
}
fn upsert_session(&self, task_id: &str, version: &VersionedSession) -> anyhow::Result<bool> {
let conn = self.conn()?;
let (backend, session_id) = backend_parts(task_id, &version.backend_ref);
let tx = conn.unchecked_transaction()?;
let changed = tx.execute(
"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(?,?,?,?,?,?,?,?,?,?,?,?,?,?)
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
WHERE excluded.generation > sessions.generation OR (excluded.generation = sessions.generation AND excluded.seq > sessions.seq)",
params![
&session_id,
version.generation,
version.seq,
version.agent_state,
version.delivery_state,
version.resource_state,
version.recovery_substate,
version.observed_json,
backend,
version.backend_ref,
version.desired_json,
version.mismatch_count,
crate::server::unix_to_rfc3339(version.updated_at),
crate::server::unix_to_rfc3339(version.updated_at)
],
)?;
if changed == 1 {
crate::server::open_binding_conn(&tx, &session_id, task_id, version.updated_at)?;
}
tx.commit()?;
Ok(changed == 1)
}
fn task_is_known(&self, task_id: &str) -> anyhow::Result<bool> {
let conn = self.conn()?;
let session_count: i64 = conn.query_row(
"SELECT COUNT(*) FROM session_tasks WHERE task_id=?",
params![task_id],
|r| r.get(0),
)?;
if session_count > 0 {
return Ok(true);
}
Ok(intent_stats(&conn, task_id)?.is_some())
}
fn task_attempt(&self, task_id: &str) -> anyhow::Result<i64> {
let conn = self.conn()?;
Ok(intent_stats(&conn, task_id)?.unwrap_or(0))
}
fn list_faults(&self, task_id: &str) -> anyhow::Result<Vec<FaultRecord>> {
let conn = self.conn()?;
let rows = conn
.prepare(
"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",
)?
.query_map(params![task_id], fault_record_row)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
fn insert_fault(&self, fault: &FaultRecord) -> anyhow::Result<i64> {
let conn = self.conn()?;
conn.execute(
"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,?,?,?,?,?,?,?,?,?,?,?,?)",
params![
fault.task_id,
fault.session_id,
fault.generation,
fault.seq,
fault.desired_json,
fault.observed_json,
fault.intent,
fault.attempt,
fault.backend_ref,
fault.kind,
fault.reason,
fault.state,
crate::server::unix_to_rfc3339(fault.created_at)
],
)?;
Ok(conn.last_insert_rowid())
}
fn emit(&self, kind: &str, data: Value) {
if let Err(err) = self.append_event(kind, &data) {
tracing::warn!(kind, error = %err, "session ledger event was not stored");
}
}
fn note_alert(&self, line: String) {
tracing::warn!(alert = %line, "session ledger alert");
}
}
fn intent_stats(conn: &Connection, task_id: &str) -> StoreResult<Option<i64>> {
let mut stmt = conn.prepare("SELECT env_json,attempt FROM intents")?;
let rows = stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?)))?;
let mut max_attempt: Option<i64> = None;
for row in rows {
let (env_json, attempt) = row?;
if intent_task_matches(&env_json, task_id) {
max_attempt = Some(max_attempt.map_or(attempt, |current| current.max(attempt)));
}
}
Ok(max_attempt)
}
fn intent_task_matches(env_json: &str, task_id: &str) -> bool {
let Ok(value) = serde_json::from_str::<Value>(env_json) else {
return false;
};
value
.get("causality")
.and_then(|v| v.get("task"))
.and_then(Value::as_str)
== Some(task_id)
}
fn backend_parts(task_id: &str, backend_ref: &str) -> (Option<String>, String) {
let parsed = serde_json::from_str::<Value>(backend_ref).ok();
let backend = parsed
.as_ref()
.and_then(|v| v.get("backend"))
.and_then(Value::as_str)
.map(str::to_string);
let session_id = parsed
.as_ref()
.and_then(|v| v.get("task_id"))
.and_then(Value::as_str)
.unwrap_or(task_id)
.to_string();
(backend, session_id)
}
const 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";
fn session_record_row(r: &Row<'_>, task_id: &str) -> rusqlite::Result<SessionRecord> {
Ok(SessionRecord {
session_id: r.get(0)?,
task_id: task_id.to_string(),
agent_state: r.get(1)?,
delivery_state: r.get(2)?,
resource_state: r.get(3)?,
recovery_substate: r.get(4)?,
desired_json: r.get(5)?,
observed_json: r.get(6)?,
generation: r.get(7)?,
seq: r.get(8)?,
backend_ref: r.get(9)?,
mismatch_count: r.get(10)?,
updated_at: crate::server::rfc3339_to_unix(&r.get::<_, String>(11)?),
})
}
fn task_row(r: &Row<'_>) -> rusqlite::Result<TaskRow> {
let state_word: String = r.get(5)?;
let task_state = serde_json::from_value(Value::String(state_word.clone()))
.map_err(|e| conversion_error(5, format!("task_state column holds {state_word:?}: {e}")))?;
Ok(TaskRow {
task_id: r.get(0)?,
kind: r.get(1)?,
parent_task: r.get(2)?,
hop: r.get(3)?,
attempt: r.get(4)?,
task_state,
opened_at: r.get(6)?,
settled_at: r.get(7)?,
})
}
fn fault_record_row(r: &Row<'_>) -> rusqlite::Result<FaultRecord> {
Ok(FaultRecord {
id: r.get(0)?,
task_id: r.get::<_, Option<String>>(1)?.unwrap_or_default(),
session_id: r.get::<_, Option<String>>(2)?.unwrap_or_default(),
generation: r.get::<_, Option<i64>>(3)?.unwrap_or_default(),
seq: r.get::<_, Option<i64>>(4)?.unwrap_or_default(),
desired_json: r
.get::<_, Option<String>>(5)?
.unwrap_or_else(|| "{}".to_string()),
observed_json: r
.get::<_, Option<String>>(6)?
.unwrap_or_else(|| "{}".to_string()),
intent: r.get::<_, Option<String>>(7)?.unwrap_or_default(),
attempt: r.get::<_, Option<i64>>(8)?.unwrap_or_default(),
backend_ref: r
.get::<_, Option<String>>(9)?
.unwrap_or_else(|| "{}".to_string()),
kind: r.get(10)?,
reason: r.get(11)?,
state: r.get(12)?,
created_at: r
.get::<_, Option<String>>(13)?
.map(|text| crate::server::rfc3339_to_unix(&text))
.unwrap_or_default(),
})
}
fn intent_row(r: &Row<'_>) -> rusqlite::Result<IntentRow> {
let env_text: String = r.get(1)?;
let receipt_text: Option<String> = r.get(5)?;
let env_json =
serde_json::from_str(&env_text).map_err(|e| conversion_error(1, e.to_string()))?;
let receipt_json = receipt_text
.map(|text| serde_json::from_str(&text).map_err(|e| conversion_error(5, e.to_string())))
.transpose()?;
Ok(IntentRow {
op_id: r.get(0)?,
env_json,
attempt: r.get(2)?,
state: r.get(3)?,
next_attempt_at: r.get(4)?,
receipt_json,
last_error: r.get(6)?,
created_at: r.get(7)?,
updated_at: r.get(8)?,
})
}
fn sql_limit(limit: u32) -> i64 {
if limit == 0 {
DEFAULT_LIMIT
} else {
i64::from(limit.min(500))
}
}
fn conversion_error(index: usize, message: String) -> rusqlite::Error {
rusqlite::Error::FromSqlConversionFailure(
index,
Type::Text,
Box::new(StoreError::Serialization(message)),
)
}