lash_sqlite_store/process_registry/
support.rs1use super::*;
2
3pub(super) fn process_status_label(record: &ProcessRecord) -> &'static str {
4 record.status.label()
5}
6
7impl SqliteProcessRegistry {
8 pub async fn open(
13 path: &Path,
14 session_store_root: impl Into<PathBuf>,
15 ) -> tokio_rusqlite::Result<Self> {
16 Self::open_with_clock(path, Arc::new(lash_core::SystemClock), session_store_root).await
17 }
18
19 pub async fn open_with_clock(
20 path: &Path,
21 clock: Arc<dyn lash_core::Clock>,
22 session_store_root: impl Into<PathBuf>,
23 ) -> tokio_rusqlite::Result<Self> {
24 Self::open_configured(path, clock, session_store_root.into()).await
25 }
26
27 async fn open_configured(
28 path: &Path,
29 clock: Arc<dyn lash_core::Clock>,
30 process_session_store_root: PathBuf,
31 ) -> tokio_rusqlite::Result<Self> {
32 let conn = SqliteConnection::open(path).await?;
33 ensure_process_schema(&conn).await?;
34 apply_pragmas(&conn, StoreBacking::File).await?;
35 Ok(Self {
36 conn,
37 clock,
38 process_session_store_root: Some(process_session_store_root),
39 })
40 }
41
42 pub async fn memory() -> tokio_rusqlite::Result<Self> {
43 Self::memory_with_clock(Arc::new(lash_core::SystemClock)).await
44 }
45
46 pub async fn memory_with_clock(
47 clock: Arc<dyn lash_core::Clock>,
48 ) -> tokio_rusqlite::Result<Self> {
49 let conn = SqliteConnection::open_in_memory().await?;
50 ensure_process_schema(&conn).await?;
51 apply_pragmas(&conn, StoreBacking::Memory).await?;
52 Ok(Self {
53 conn,
54 clock,
55 process_session_store_root: None,
56 })
57 }
58
59 pub(crate) fn load_process_conn(
60 conn: &Connection,
61 process_id: &str,
62 ) -> Result<Option<ProcessRecord>, lash_core::PluginError> {
63 let json: Option<String> = conn
64 .query_row(
65 "SELECT record_json FROM processes WHERE process_id = ?1",
66 params![process_id],
67 |row| row.get(0),
68 )
69 .optional()
70 .map_err(process_sqlite_error)?;
71 json.map(|json| serde_json::from_str(&json).map_err(process_decode_error))
72 .transpose()
73 }
74
75 pub(crate) fn save_process_conn(
76 conn: &Connection,
77 record: &ProcessRecord,
78 ) -> Result<(), lash_core::PluginError> {
79 let change_seq = Self::next_change_seq_conn(conn)?;
80 conn.execute(
81 "UPDATE processes
82 SET updated_at_ms = ?2, change_seq = ?3, status = ?4, record_json = ?5
83 WHERE process_id = ?1",
84 params![
85 record.id.as_str(),
86 record.updated_at_ms as i64,
87 change_seq as i64,
88 process_status_label(record),
89 process_encode_json(record)?
90 ],
91 )
92 .map_err(process_sqlite_error)?;
93 Ok(())
94 }
95
96 pub(super) fn next_change_seq_conn(conn: &Connection) -> Result<u64, lash_core::PluginError> {
97 conn.execute(
98 "UPDATE process_change_clock
99 SET current_seq = current_seq + 1
100 WHERE singleton = 1",
101 [],
102 )
103 .map_err(process_sqlite_error)?;
104 conn.query_row(
105 "SELECT current_seq FROM process_change_clock WHERE singleton = 1",
106 [],
107 |row| row.get::<_, i64>(0),
108 )
109 .map(|seq| seq as u64)
110 .map_err(process_sqlite_error)
111 }
112
113 pub(crate) fn load_event_by_key_conn(
114 conn: &Connection,
115 process_id: &str,
116 replay_key: &str,
117 ) -> Result<Option<(String, ProcessEvent)>, lash_core::PluginError> {
118 let row: Option<(String, String)> = conn
119 .query_row(
120 "SELECT payload_hash, event_json
121 FROM process_events
122 WHERE process_id = ?1 AND idempotency_key = ?2",
123 params![process_id, replay_key],
124 |row| Ok((row.get(0)?, row.get(1)?)),
125 )
126 .optional()
127 .map_err(process_sqlite_error)?;
128 row.map(|(hash, json)| {
129 serde_json::from_str(&json)
130 .map(|event| (hash, event))
131 .map_err(process_decode_error)
132 })
133 .transpose()
134 }
135
136 pub(crate) fn load_process_lease_conn(
137 conn: &Connection,
138 process_id: &str,
139 ) -> Result<Option<ProcessLease>, lash_core::PluginError> {
140 conn.query_row(
141 "SELECT lease_owner_id, lease_token, lease_fencing_token,
142 lease_claimed_at_ms, lease_expires_at_ms,
143 lease_owner_incarnation_id, lease_owner_liveness_json
144 FROM process_leases
145 WHERE process_id = ?1",
146 params![process_id],
147 |row| {
148 let owner_id: Option<String> = row.get(0)?;
149 let lease_token: Option<String> = row.get(1)?;
150 let incarnation_id: Option<String> = row.get(5)?;
151 let liveness_json: Option<String> = row.get(6)?;
152 let (Some(owner_id), Some(lease_token)) = (owner_id, lease_token) else {
153 return Ok(None);
154 };
155 Ok(Some(ProcessLease {
156 schema_version: PROCESS_LEASE_SCHEMA_VERSION,
157 process_id: process_id.to_string(),
158 owner: process_lease_owner_from_columns(
159 owner_id,
160 incarnation_id,
161 liveness_json,
162 ),
163 lease_token,
164 fencing_token: row.get::<_, i64>(2)? as u64,
165 claimed_at_epoch_ms: row.get::<_, i64>(3)? as u64,
166 expires_at_epoch_ms: row.get::<_, i64>(4)? as u64,
167 }))
168 },
169 )
170 .optional()
171 .map(|lease| lease.flatten())
172 .map_err(process_sqlite_error)
173 }
174
175 pub(super) fn acquire_process_lease_conn(
178 conn: &Connection,
179 process_id: &str,
180 owner: &LeaseOwnerIdentity,
181 fencing_token: u64,
182 now: u64,
183 lease_ttl_ms: u64,
184 ) -> Result<ProcessLease, lash_core::PluginError> {
185 let lease = ProcessLease {
186 schema_version: PROCESS_LEASE_SCHEMA_VERSION,
187 process_id: process_id.to_string(),
188 owner: owner.clone(),
189 lease_token: format!(
190 "{:x}",
191 Sha256::digest(
192 format!(
193 "{process_id}:{}:{}:{now}:{fencing_token}",
194 owner.owner_id, owner.incarnation_id
195 )
196 .as_bytes()
197 )
198 ),
199 fencing_token,
200 claimed_at_epoch_ms: now,
201 expires_at_epoch_ms: now.saturating_add(lease_ttl_ms),
202 };
203 conn.execute(
204 "INSERT INTO process_leases (
205 process_id, lease_owner_id, lease_owner_incarnation_id,
206 lease_owner_liveness_json, lease_token, lease_fencing_token,
207 lease_claimed_at_ms, lease_expires_at_ms
208 )
209 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
210 ON CONFLICT(process_id) DO UPDATE SET
211 lease_owner_id = excluded.lease_owner_id,
212 lease_owner_incarnation_id = excluded.lease_owner_incarnation_id,
213 lease_owner_liveness_json = excluded.lease_owner_liveness_json,
214 lease_token = excluded.lease_token,
215 lease_fencing_token = excluded.lease_fencing_token,
216 lease_claimed_at_ms = excluded.lease_claimed_at_ms,
217 lease_expires_at_ms = excluded.lease_expires_at_ms",
218 params![
219 lease.process_id.as_str(),
220 lease.owner.owner_id.as_str(),
221 lease.owner.incarnation_id.as_str(),
222 encode_process_lease_liveness(&lease.owner.liveness)?,
223 lease.lease_token.as_str(),
224 lease.fencing_token as i64,
225 lease.claimed_at_epoch_ms as i64,
226 lease.expires_at_epoch_ms as i64,
227 ],
228 )
229 .map_err(process_sqlite_error)?;
230 Ok(lease)
231 }
232
233 pub(super) fn list_grants_for_scope_conn(
234 conn: &Connection,
235 session_scope: &SessionScope,
236 live_only: bool,
237 ) -> Result<Vec<ProcessHandleGrantEntry>, lash_core::PluginError> {
238 let session_scope_id = session_scope.id();
239 let status_clause = if live_only {
240 "AND p.status = 'running'"
241 } else {
242 ""
243 };
244 let mut stmt = conn
245 .prepare(&format!(
246 "SELECT g.process_id, g.descriptor_json, p.record_json
247 FROM process_handle_grants g
248 JOIN processes p ON p.process_id = g.process_id
249 WHERE g.scope_id = ?1 {status_clause}
250 ORDER BY g.process_id ASC"
251 ))
252 .map_err(process_sqlite_error)?;
253 let rows = stmt
254 .query_map(params![session_scope_id.as_str()], |row| {
255 Ok((
256 row.get::<_, String>(0)?,
257 row.get::<_, String>(1)?,
258 row.get::<_, String>(2)?,
259 ))
260 })
261 .map_err(process_sqlite_error)?;
262 let mut entries = Vec::new();
263 for row in rows {
264 let (process_id, descriptor_json, record_json) = row.map_err(process_sqlite_error)?;
265 let descriptor: ProcessHandleDescriptor =
266 serde_json::from_str(&descriptor_json).map_err(process_decode_error)?;
267 let record: ProcessRecord =
268 serde_json::from_str(&record_json).map_err(process_decode_error)?;
269 entries.push((
270 ProcessHandleGrant {
271 session_id: session_scope.session_id.clone(),
272 process_id,
273 descriptor,
274 },
275 record,
276 ));
277 }
278 Ok(entries)
279 }
280}
281
282pub(crate) fn tx_outcome<T>(
287 result: Result<T, lash_core::PluginError>,
288) -> TxOutcome<Result<T, lash_core::PluginError>> {
289 match result {
290 Ok(value) => TxOutcome::Commit(Ok(value)),
291 Err(err) => TxOutcome::Rollback(Err(err)),
292 }
293}