1use super::*;
27
28const SQLITE_QUEUED_WORK_HEAD_CANDIDATE_PREDICATE: &str = "session_id = ?1
29 AND available_at_ms <= ?2
30 AND (
31 claim_token IS NULL
32 OR claim_session_lease_generation <> ?3
33 )";
34
35fn sqlite_queued_work_head_candidate_cte(boundary: QueuedWorkClaimBoundary) -> String {
36 let delivery_gate = match boundary {
37 QueuedWorkClaimBoundary::Idle => "",
38 QueuedWorkClaimBoundary::ActiveTurnCheckpoint => {
39 "WHERE head_delivery_policy = 'earliest_safe_boundary'"
40 }
41 };
42 format!(
43 "queued_work_head_candidate AS (
44 SELECT head_enqueue_seq, head_batch_id, head_delivery_policy
45 FROM (
46 SELECT enqueue_seq AS head_enqueue_seq,
47 batch_id AS head_batch_id,
48 delivery_policy AS head_delivery_policy
49 FROM queued_work_batches
50 WHERE {SQLITE_QUEUED_WORK_HEAD_CANDIDATE_PREDICATE}
51 ORDER BY enqueue_seq ASC
52 LIMIT 1
53 ) AS unfiltered_head
54 {delivery_gate}
55 )"
56 )
57}
58
59fn sqlite_queued_work_claim_candidates_sql(boundary: QueuedWorkClaimBoundary) -> String {
60 let head_candidate = sqlite_queued_work_head_candidate_cte(boundary);
61 format!(
62 "WITH {head_candidate}
63 SELECT enqueue_seq, batch_id, session_id, source_key, delivery_policy,
64 slot_policy, merge_key_json, available_at_ms, enqueued_at_ms,
65 claim_fencing_token, claim_owner_id, claim_owner_incarnation_id,
66 claim_owner_liveness_json, claim_token, claim_session_lease_generation
67 FROM queued_work_batches
68 CROSS JOIN queued_work_head_candidate
69 WHERE {SQLITE_QUEUED_WORK_HEAD_CANDIDATE_PREDICATE}
70 ORDER BY enqueue_seq ASC
71 LIMIT ?4"
72 )
73}
74
75#[async_trait::async_trait]
76impl SessionCommitStore for Store {
77 fn durability_tier(&self) -> DurabilityTier {
78 DurabilityTier::Durable
79 }
80
81 async fn load_session(
82 &self,
83 scope: SessionReadScope,
84 ) -> Result<Option<PersistedSessionRead>, StoreError> {
85 self.conn
86 .call(move |conn| {
87 let outcome: Result<Option<PersistedSessionRead>, StoreError> = (|| {
88 let Some(meta) = try_load_session_head_meta_from_conn(conn)? else {
89 return Ok(None);
90 };
91 let leaf_node_id = match &scope {
92 SessionReadScope::FullGraph => meta.leaf_node_id.clone(),
93 SessionReadScope::ActivePath { leaf_node_id } => {
94 leaf_node_id.clone().or_else(|| meta.leaf_node_id.clone())
95 }
96 };
97 let mut graph = match scope {
98 SessionReadScope::FullGraph => {
99 Self::load_session_graph_from_conn(conn, meta.leaf_node_id.clone())
100 }
101 SessionReadScope::ActivePath { .. } => {
102 Self::load_active_path_session_graph_from_conn(
103 conn,
104 leaf_node_id.clone(),
105 )
106 .map_err(sqlite_error)?
107 }
108 };
109 graph.set_leaf_node_id(leaf_node_id);
110 let checkpoint = meta
111 .checkpoint_ref
112 .as_ref()
113 .map(|blob_ref| Self::get_checkpoint_conn(conn, blob_ref))
114 .transpose()?
115 .flatten();
116 Ok(Some(PersistedSessionRead {
117 session_id: meta.session_id,
118 head_revision: meta.head_revision,
119 config: meta.config,
120 agent_frames: meta.agent_frames,
121 current_agent_frame_id: meta.current_agent_frame_id,
122 graph,
123 checkpoint_ref: meta.checkpoint_ref,
124 checkpoint,
125 token_ledger: merge_token_ledger_entries(Self::load_usage_deltas_conn(
126 conn,
127 )),
128 }))
129 })(
130 );
131 Ok(outcome)
132 })
133 .await
134 .map_err(sqlite_error)?
135 }
136
137 async fn load_node(
138 &self,
139 node_id: &str,
140 ) -> Result<Option<lash_core::SessionNodeRecord>, StoreError> {
141 let node_id = node_id.to_string();
142 let row: Option<String> = self
143 .conn
144 .call(move |conn| {
145 conn.query_row(
146 "SELECT node_json FROM graph_nodes WHERE node_id = ?1 AND tombstoned = 0",
147 params![node_id],
148 |row| row.get(0),
149 )
150 .optional()
151 })
152 .await
153 .map_err(sqlite_error)?;
154 Ok(row.and_then(|json| serde_json::from_str(&json).ok()))
155 }
156
157 async fn commit_runtime_state(
158 &self,
159 commit: RuntimeCommit,
160 ) -> Result<RuntimeCommitResult, StoreError> {
161 let blob_profile = self.options.blob_profile;
162 let now = self.clock.timestamp_ms();
163 let enqueue_nonce_start = self.commit_count.fetch_add(
164 commit.enqueued_queue_batches.len() as u64,
165 AtomicOrdering::Relaxed,
166 );
167 let result = self
168 .conn
169 .write_flow(move |tx| {
170 let outcome: Result<RuntimeCommitResult, StoreError> = (|| {
171 let existing = try_load_session_head_meta_from_conn(tx)?;
172 if let Some(bound_session_id) =
173 existing.as_ref().map(|meta| meta.session_id.as_str())
174 && bound_session_id != commit.session_id
175 {
176 return Err(StoreError::SessionBindingMismatch {
177 bound_session_id: bound_session_id.to_string(),
178 attempted_session_id: commit.session_id.clone(),
179 });
180 }
181 if let Some(completed) = &commit.turn_commit {
182 if completed.session_id != commit.session_id {
183 return Err(StoreError::RuntimeTurnCommitConflict {
184 session_id: completed.session_id.clone(),
185 turn_id: completed.turn_id.clone(),
186 });
187 }
188 let prior: Option<(String, String)> = tx
189 .query_row(
190 "SELECT turn_commit_hash, result_json FROM runtime_turn_commits
191 WHERE session_id = ?1 AND turn_id = ?2",
192 params![completed.session_id, completed.turn_id],
193 |row| Ok((row.get(0)?, row.get(1)?)),
194 )
195 .optional()
196 .map_err(sqlite_error)?;
197 if let Some((turn_commit_hash, result_json)) = prior {
198 if turn_commit_hash == completed.turn_commit_hash {
199 let result: RuntimeCommitResult =
200 serde_json::from_str(&result_json).map_err(|err| {
201 StoreError::Backend(format!(
202 "failed to decode runtime turn commit result: {err}"
203 ))
204 })?;
205 if let Some(completion) =
206 commit.release_session_execution_lease.as_ref()
207 {
208 release_session_execution_lease_conn(tx, completion)?;
209 }
210 return Ok(result);
211 }
212 return Err(StoreError::RuntimeTurnCommitConflict {
213 session_id: completed.session_id.clone(),
214 turn_id: completed.turn_id.clone(),
215 });
216 }
217 }
218 let actual_revision = existing.as_ref().map_or(0, |meta| meta.head_revision);
219 let expected_revision = commit.expected_head_revision.unwrap_or(0);
220 if expected_revision != actual_revision {
221 return Err(StoreError::HeadRevisionConflict {
222 expected: commit.expected_head_revision,
223 actual: actual_revision,
224 });
225 }
226 for completed in &commit.completed_queue_claims {
227 if completed.session_id != commit.session_id {
228 return Err(StoreError::QueuedWorkClaimSuperseded {
229 session_id: completed.session_id.clone(),
230 claim_id: completed.claim_id.clone(),
231 });
232 }
233 ensure_queued_work_completion_conn(tx, completed)?;
234 }
235 for completed in &commit.completed_turn_input_claims {
236 if completed.session_id != commit.session_id {
237 return Err(StoreError::TurnInputClaimSuperseded {
238 session_id: completed.session_id.clone(),
239 claim_id: completed.claim_id.clone(),
240 });
241 }
242 let owned_rows: usize = tx
243 .query_row(
244 "SELECT COUNT(*)
245 FROM pending_turn_inputs
246 WHERE session_id = ?1
247 AND claim_id = ?2
248 AND claim_token = ?3",
249 params![
250 completed.session_id,
251 completed.claim_id,
252 completed.lease_token
253 ],
254 |row| row.get::<_, i64>(0),
255 )
256 .map_err(sqlite_error)? as usize;
257 ensure_turn_input_completion_owns_all_inputs(completed, owned_rows)?;
258 }
259
260 let stored_checkpoint =
261 Self::put_checkpoint_conn(tx, &commit.checkpoint, blob_profile)
262 .map_err(sqlite_error)?;
263
264 if !commit.usage_deltas.is_empty() {
265 let mut stmt = tx
266 .prepare(
267 "INSERT INTO usage_deltas (
268 source, model, input_tokens, output_tokens, cache_read_input_tokens, cache_write_input_tokens, reasoning_output_tokens
269 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
270 )
271 .map_err(sqlite_error)?;
272 for entry in &commit.usage_deltas {
273 stmt.execute(params![
274 entry.source,
275 entry.model,
276 entry.usage.input_tokens,
277 entry.usage.output_tokens,
278 entry.usage.cache_read_input_tokens,
279 entry.usage.cache_write_input_tokens,
280 entry.usage.reasoning_output_tokens,
281 ])
282 .map_err(sqlite_error)?;
283 }
284 }
285
286 let leaf_node_id = match &commit.graph {
287 GraphCommitDelta::Unchanged { leaf_node_id } => leaf_node_id.clone(),
288 GraphCommitDelta::Append {
289 nodes,
290 leaf_node_id,
291 } => {
292 for node in nodes {
293 let node_json = encode_json(node);
294 tx.execute(
295 "INSERT INTO graph_nodes (node_id, node_json) VALUES (?1, ?2)",
296 params![node.node_id, node_json],
297 )
298 .map_err(sqlite_error)?;
299 }
300 leaf_node_id.clone()
301 }
302 GraphCommitDelta::ReplaceFull(graph) => {
303 tx.execute("DELETE FROM graph_nodes", [])
304 .map_err(sqlite_error)?;
305 for node in &graph.nodes {
306 let node_json = encode_json(node);
307 tx.execute(
308 "INSERT INTO graph_nodes (node_id, node_json) VALUES (?1, ?2)",
309 params![node.node_id, node_json],
310 )
311 .map_err(sqlite_error)?;
312 }
313 graph.leaf_node_id.clone()
314 }
315 };
316 let graph_node_count: usize = tx
317 .query_row(
318 "SELECT COUNT(*) FROM graph_nodes WHERE tombstoned = 0",
319 [],
320 |row| row.get::<_, i64>(0),
321 )
322 .map_err(sqlite_error)? as usize;
323 let next_revision = actual_revision + 1;
324 let meta = SessionHeadMeta {
325 schema_version: lash_core::store::SESSION_HEAD_META_SCHEMA_VERSION,
326 session_id: commit.session_id.clone(),
327 head_revision: next_revision,
328 config: commit.config.clone(),
329 agent_frames: commit.agent_frames.clone(),
330 current_agent_frame_id: commit.current_agent_frame_id.clone(),
331 checkpoint_ref: Some(stored_checkpoint.checkpoint_ref.clone()),
332 leaf_node_id,
333 graph_node_count,
334 token_ledger: Vec::new(),
335 };
336 tx.execute(
337 "INSERT OR REPLACE INTO session_head (singleton, session_id, head_json, head_revision)
338 VALUES (1, ?1, ?2, ?3)",
339 params![
340 meta.session_id,
341 encode_json(&meta),
342 meta.head_revision as i64
343 ],
344 )
345 .map_err(sqlite_error)?;
346 for completed in &commit.completed_queue_claims {
347 for batch_id in &completed.batch_ids {
348 tx.execute(
349 "DELETE FROM queued_work_batches
350 WHERE session_id = ?1
351 AND batch_id = ?2
352 AND claim_id = ?3
353 AND claim_token = ?4",
354 params![
355 completed.session_id,
356 batch_id,
357 completed.claim_id,
358 completed.lease_token
359 ],
360 )
361 .map_err(sqlite_error)?;
362 }
363 }
364 for completed in &commit.completed_turn_input_claims {
365 for input_id in &completed.input_ids {
366 tx.execute(
367 "UPDATE pending_turn_inputs
368 SET state = ?5,
369 claim_id = NULL,
370 claim_owner_id = NULL,
371 claim_owner_incarnation_id = NULL,
372 claim_owner_liveness_json = NULL,
373 claim_token = NULL,
374 claim_session_lease_generation = 0
375 WHERE session_id = ?1
376 AND input_id = ?2
377 AND claim_id = ?3
378 AND claim_token = ?4",
379 params![
380 completed.session_id,
381 input_id,
382 completed.claim_id,
383 completed.lease_token,
384 lash_core::TurnInputState::Completed.as_str(),
385 ],
386 )
387 .map_err(sqlite_error)?;
388 }
389 }
390 if let Some(turn_id) = commit.interrupted_turn_input_turn_id.as_deref() {
391 let input_ids = {
392 let mut stmt = tx
393 .prepare(
394 "SELECT input_id, ingress_json
395 FROM pending_turn_inputs
396 WHERE session_id = ?1 AND state = ?2",
397 )
398 .map_err(sqlite_error)?;
399 let rows = stmt
400 .query_map(
401 params![
402 commit.session_id,
403 lash_core::TurnInputState::PendingActive.as_str()
404 ],
405 |row| {
406 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
407 },
408 )
409 .map_err(sqlite_error)?;
410 let mut input_ids = Vec::new();
411 for row in rows {
412 let (input_id, ingress_json) = row.map_err(sqlite_error)?;
413 let ingress = decode_turn_input_ingress(ingress_json)?;
414 if ingress
415 .active_turn_id()
416 .is_some_and(|active| active == turn_id)
417 {
418 input_ids.push(input_id);
419 }
420 }
421 input_ids
422 };
423 let next_turn_ingress = encode_json(&lash_core::TurnInputIngress::NextTurn);
424 let mut stmt = tx
425 .prepare(
426 "UPDATE pending_turn_inputs
427 SET state = ?3,
428 ingress_json = ?4,
429 claim_id = NULL,
430 claim_owner_id = NULL,
431 claim_owner_incarnation_id = NULL,
432 claim_owner_liveness_json = NULL,
433 claim_token = NULL,
434 claim_session_lease_generation = 0
435 WHERE session_id = ?1 AND input_id = ?2",
436 )
437 .map_err(sqlite_error)?;
438 for input_id in input_ids {
439 stmt.execute(params![
440 commit.session_id,
441 input_id,
442 lash_core::TurnInputState::DeferredNextTurn.as_str(),
443 next_turn_ingress
444 ])
445 .map_err(sqlite_error)?;
446 }
447 }
448 if !commit.committed_attachment_ids.is_empty() {
449 let now = now as i64;
450 let mut stmt = tx
451 .prepare(
452 "UPDATE attachment_manifest
453 SET committed_at_ms = COALESCE(committed_at_ms, ?1)
454 WHERE attachment_id = ?2 AND session_id = ?3",
455 )
456 .map_err(sqlite_error)?;
457 for id in &commit.committed_attachment_ids {
458 stmt.execute(params![now, id.as_str(), commit.session_id])
459 .map_err(sqlite_error)?;
460 }
461 }
462 if let Some(turn_commit) = &commit.turn_commit {
463 tx.execute(
464 "UPDATE attachment_manifest
465 SET committed_at_ms = COALESCE(committed_at_ms, ?1)
466 WHERE session_id = ?2
467 AND owner_kind = 'turn'
468 AND owner_id = ?3
469 AND committed_at_ms IS NULL",
470 params![now as i64, commit.session_id, turn_commit.turn_id],
471 )
472 .map_err(sqlite_error)?;
473 }
474 let mut enqueued_queue_batches = Vec::new();
475 for (index, batch) in commit.enqueued_queue_batches.iter().enumerate() {
476 if batch.session_id != commit.session_id {
477 return Err(StoreError::SessionBindingMismatch {
478 bound_session_id: commit.session_id.clone(),
479 attempted_session_id: batch.session_id.clone(),
480 });
481 }
482 enqueued_queue_batches.push(enqueue_queued_work_conn(
483 tx,
484 batch,
485 now,
486 enqueue_nonce_start.saturating_add(index as u64),
487 )?);
488 }
489 let result = RuntimeCommitResult {
490 head_revision: next_revision,
491 checkpoint_ref: stored_checkpoint.checkpoint_ref,
492 manifest: stored_checkpoint.manifest,
493 enqueued_queue_batches,
494 turn_input_applications: commit.turn_input_applications(),
495 };
496 if let Some(completed) = &commit.turn_commit {
497 tx.execute(
498 "INSERT INTO runtime_turn_commits (
499 session_id, turn_id, turn_commit_hash, result_json, committed_at_ms
500 )
501 VALUES (?1, ?2, ?3, ?4, ?5)",
502 params![
503 completed.session_id,
504 completed.turn_id,
505 completed.turn_commit_hash,
506 encode_json(&result),
507 now as i64
508 ],
509 )
510 .map_err(sqlite_error)?;
511 }
512 if let Some(completion) = commit.release_session_execution_lease.as_ref() {
513 release_session_execution_lease_conn(tx, completion)?;
514 }
515 Ok(result)
516 })();
517 match outcome {
522 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
523 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
524 }
525 })
526 .await
527 .map_err(sqlite_error)??;
528 self.maybe_auto_gc().await;
529 Ok(result)
530 }
531
532 async fn save_session_meta(&self, meta: SessionMeta) -> Result<(), StoreError> {
533 Store::save_session_meta(self, meta).await;
534 Ok(())
535 }
536
537 async fn load_session_meta(&self) -> Result<Option<SessionMeta>, StoreError> {
538 Ok(Store::load_session_meta(self).await)
539 }
540}
541
542#[async_trait::async_trait]
543impl SessionExecutionLeaseStore for Store {
544 async fn try_claim_session_execution_lease(
545 &self,
546 session_id: &str,
547 owner: &LeaseOwnerIdentity,
548 lease_ttl_ms: u64,
549 ) -> Result<SessionExecutionLeaseClaimOutcome, StoreError> {
550 let session_id = session_id.to_string();
551 let owner = owner.clone();
552 let now = self.clock.timestamp_ms();
553 self.conn
554 .write_flow(move |tx| {
555 let outcome: Result<SessionExecutionLeaseClaimOutcome, StoreError> = (|| {
556 let current = load_session_execution_lease_row_conn(tx, &session_id)?;
557 if current.as_ref().is_some_and(|lease| {
558 lease.lease_token.is_some() && lease.expires_at_ms > now
559 }) {
560 let current = current.expect("checked current lease is present");
561 if current
562 .owner
563 .as_ref()
564 .is_some_and(|current_owner| current_owner.same_incarnation(&owner))
565 {
566 let expires_at = now.saturating_add(lease_ttl_ms);
567 tx.execute(
568 "UPDATE session_execution_leases
569 SET lease_expires_at_ms = ?2
570 WHERE session_id = ?1",
571 params![session_id, expires_at as i64],
572 )
573 .map_err(sqlite_error)?;
574 return Ok(SessionExecutionLeaseClaimOutcome::Acquired(
575 SessionExecutionLease {
576 session_id,
577 owner,
578 lease_token: current.lease_token.expect("live lease token set"),
579 fencing_token: current.fencing_token,
580 claimed_at_epoch_ms: current.claimed_at_ms,
581 expires_at_epoch_ms: expires_at,
582 },
583 ));
584 }
585 return Ok(SessionExecutionLeaseClaimOutcome::Busy {
586 holder: row_to_session_execution_lease(&session_id, current)?,
587 });
588 }
589 Ok(SessionExecutionLeaseClaimOutcome::Acquired(
590 acquire_session_execution_lease_conn(
591 tx,
592 &session_id,
593 &owner,
594 current.as_ref().map_or(0, |lease| lease.fencing_token),
595 now,
596 lease_ttl_ms,
597 )?,
598 ))
599 })(
600 );
601 match outcome {
602 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
603 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
604 }
605 })
606 .await
607 .map_err(sqlite_error)?
608 }
609
610 async fn reclaim_session_execution_lease(
611 &self,
612 session_id: &str,
613 owner: &LeaseOwnerIdentity,
614 observed_holder: &SessionExecutionLeaseFence,
615 lease_ttl_ms: u64,
616 ) -> Result<SessionExecutionLeaseClaimOutcome, StoreError> {
617 let session_id = session_id.to_string();
618 let owner = owner.clone();
619 let observed_holder = observed_holder.clone();
620 let now = self.clock.timestamp_ms();
621 self.conn
622 .write_flow(move |tx| {
623 let outcome: Result<SessionExecutionLeaseClaimOutcome, StoreError> = (|| {
624 let current = load_session_execution_lease_row_conn(tx, &session_id)?;
625 let Some(current) = current else {
626 return Ok(SessionExecutionLeaseClaimOutcome::Acquired(
627 acquire_session_execution_lease_conn(
628 tx,
629 &session_id,
630 &owner,
631 0,
632 now,
633 lease_ttl_ms,
634 )?,
635 ));
636 };
637 if current.lease_token.is_none() || current.expires_at_ms <= now {
638 return Ok(SessionExecutionLeaseClaimOutcome::Acquired(
639 acquire_session_execution_lease_conn(
640 tx,
641 &session_id,
642 &owner,
643 current.fencing_token,
644 now,
645 lease_ttl_ms,
646 )?,
647 ));
648 }
649 let holder = row_to_session_execution_lease(&session_id, current)?;
650 if observed_holder.session_id == session_id
651 && holder.owner.same_incarnation(&observed_holder.owner)
652 && holder.lease_token == observed_holder.lease_token
653 && holder.fencing_token == observed_holder.fencing_token
654 && holder.owner.is_definitely_dead_for_claimant(&owner)
655 {
656 let fencing_token = holder.fencing_token.saturating_add(1);
657 let lease_token = format!(
658 "{}:{}:{}:{now}:{fencing_token}",
659 session_id, owner.owner_id, owner.incarnation_id
660 );
661 let expires_at = now.saturating_add(lease_ttl_ms);
662 let liveness_json = encode_liveness(&owner.liveness)?;
663 let changed = tx
664 .execute(
665 "UPDATE session_execution_leases
666 SET lease_owner_id = ?1,
667 lease_owner_incarnation_id = ?2,
668 lease_owner_liveness_json = ?3,
669 lease_token = ?4,
670 lease_fencing_token = ?5,
671 lease_claimed_at_ms = ?6,
672 lease_expires_at_ms = ?7
673 WHERE session_id = ?8
674 AND lease_owner_id = ?9
675 AND lease_owner_incarnation_id = ?10
676 AND lease_token = ?11
677 AND lease_fencing_token = ?12",
678 params![
679 owner.owner_id,
680 owner.incarnation_id,
681 liveness_json,
682 lease_token,
683 fencing_token as i64,
684 now as i64,
685 expires_at as i64,
686 session_id,
687 observed_holder.owner.owner_id,
688 observed_holder.owner.incarnation_id,
689 observed_holder.lease_token,
690 observed_holder.fencing_token as i64,
691 ],
692 )
693 .map_err(sqlite_error)?;
694 if changed == 1 {
695 return Ok(SessionExecutionLeaseClaimOutcome::Acquired(
696 SessionExecutionLease {
697 session_id,
698 owner,
699 lease_token,
700 fencing_token,
701 claimed_at_epoch_ms: now,
702 expires_at_epoch_ms: expires_at,
703 },
704 ));
705 }
706 let current = load_session_execution_lease_row_conn(tx, &session_id)?;
707 if current.as_ref().is_some_and(|lease| {
708 lease.lease_token.is_some() && lease.expires_at_ms > now
709 }) {
710 let current = current.expect("checked current lease is present");
711 return Ok(SessionExecutionLeaseClaimOutcome::Busy {
712 holder: row_to_session_execution_lease(&session_id, current)?,
713 });
714 }
715 let previous_fencing_token =
716 current.as_ref().map_or(0, |lease| lease.fencing_token);
717 return Ok(SessionExecutionLeaseClaimOutcome::Acquired(
718 acquire_session_execution_lease_conn(
719 tx,
720 &session_id,
721 &owner,
722 previous_fencing_token,
723 now,
724 lease_ttl_ms,
725 )?,
726 ));
727 }
728 Ok(SessionExecutionLeaseClaimOutcome::Busy { holder })
729 })(
730 );
731 match outcome {
732 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
733 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
734 }
735 })
736 .await
737 .map_err(sqlite_error)?
738 }
739
740 async fn renew_session_execution_lease(
741 &self,
742 fence: &SessionExecutionLeaseFence,
743 lease_ttl_ms: u64,
744 ) -> Result<SessionExecutionLease, StoreError> {
745 let fence = fence.clone();
746 let now = self.clock.timestamp_ms();
747 self.conn
748 .write_flow(move |tx| {
749 let outcome: Result<SessionExecutionLease, StoreError> = (|| {
750 let current = load_session_execution_lease_row_conn(tx, &fence.session_id)?;
751 let Some(current) = current else {
752 return Err(StoreError::SessionExecutionLeaseExpired {
753 session_id: fence.session_id.clone(),
754 });
755 };
756 if !current
757 .owner
758 .as_ref()
759 .is_some_and(|owner| owner.same_incarnation(&fence.owner))
760 || current.lease_token.as_deref() != Some(fence.lease_token.as_str())
761 || current.fencing_token != fence.fencing_token
762 || current.expires_at_ms <= now
763 {
764 return Err(StoreError::SessionExecutionLeaseExpired {
765 session_id: fence.session_id.clone(),
766 });
767 }
768 let expires_at = now.saturating_add(lease_ttl_ms);
769 tx.execute(
770 "UPDATE session_execution_leases
771 SET lease_expires_at_ms = ?5
772 WHERE session_id = ?1
773 AND lease_owner_id = ?2
774 AND lease_owner_incarnation_id = ?3
775 AND lease_token = ?4
776 AND lease_fencing_token = ?6",
777 params![
778 fence.session_id,
779 fence.owner.owner_id,
780 fence.owner.incarnation_id,
781 fence.lease_token,
782 expires_at as i64,
783 fence.fencing_token as i64
784 ],
785 )
786 .map_err(sqlite_error)?;
787 Ok(SessionExecutionLease {
788 session_id: fence.session_id,
789 owner: fence.owner,
790 lease_token: fence.lease_token,
791 fencing_token: fence.fencing_token,
792 claimed_at_epoch_ms: current.claimed_at_ms,
793 expires_at_epoch_ms: expires_at,
794 })
795 })();
796 match outcome {
797 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
798 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
799 }
800 })
801 .await
802 .map_err(sqlite_error)?
803 }
804
805 async fn release_session_execution_lease(
806 &self,
807 completion: &SessionExecutionLeaseCompletion,
808 ) -> Result<(), StoreError> {
809 let completion = completion.clone();
810 self.conn
811 .write_flow(move |tx| {
812 let outcome = release_session_execution_lease_conn(tx, &completion);
813 match outcome {
814 Ok(()) => Ok(TxOutcome::Commit(Ok(()))),
815 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
816 }
817 })
818 .await
819 .map_err(sqlite_error)?
820 }
821}
822
823#[async_trait::async_trait]
824impl QueuedWorkStore for Store {
825 async fn enqueue_queued_work(
826 &self,
827 batch: QueuedWorkBatchDraft,
828 ) -> Result<QueuedWorkBatch, StoreError> {
829 let nonce = self.commit_count.fetch_add(1, AtomicOrdering::Relaxed);
830 let now = self.clock.timestamp_ms();
831 self.conn
832 .write_flow(move |tx| {
833 let outcome = enqueue_queued_work_conn(tx, &batch, now, nonce);
834 match outcome {
837 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
838 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
839 }
840 })
841 .await
842 .map_err(sqlite_error)?
843 }
844
845 async fn claim_leading_ready_session_command(
846 &self,
847 session_id: &str,
848 session_execution_lease: &SessionExecutionLeaseFence,
849 owner: &LeaseOwnerIdentity,
850 ) -> Result<Option<QueuedWorkClaim>, StoreError> {
851 let session_id = session_id.to_string();
852 let session_execution_lease = session_execution_lease.clone();
853 let owner = owner.clone();
854 let now = self.clock.timestamp_ms();
855 self.conn
856 .write_flow(move |tx| {
857 let outcome: Result<TxOutcome<Option<QueuedWorkClaim>>, StoreError> = (|| {
858 ensure_session_execution_lease_conn(
859 tx,
860 &session_id,
861 &session_execution_lease,
862 now,
863 )?;
864 let generation = session_execution_lease.fencing_token;
868 let candidate_rows = {
869 let mut stmt = tx
870 .prepare(&sqlite_queued_work_claim_candidates_sql(
871 QueuedWorkClaimBoundary::Idle,
872 ))
873 .map_err(sqlite_error)?;
874 let rows = stmt
875 .query_map(
876 params![
877 session_id,
878 now as i64,
879 generation as i64,
880 claim_scan_limit(1)
881 ],
882 queued_batch_row_from_sql,
883 )
884 .map_err(sqlite_error)?;
885 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
886 };
887 let candidate_rows = candidate_rows
888 .into_iter()
889 .filter(|row| {
890 row.claim_token.is_none()
891 || row.claim_session_lease_generation != generation
892 })
893 .collect::<Vec<_>>();
894 let candidate_batches = candidate_rows
895 .iter()
896 .map(|row| queued_work_batch_from_conn(tx, row.clone()))
897 .collect::<Result<Vec<_>, StoreError>>()?;
898 let candidates = candidate_rows
899 .iter()
900 .zip(candidate_batches.iter())
901 .map(|(row, batch)| {
902 Ok(ClaimCandidate {
903 enqueue_seq: row.enqueue_seq,
904 claim_fencing_token: row.claim_fencing_token,
905 work_class: batch.work_class().ok_or_else(|| {
906 StoreError::Backend(format!(
907 "queued-work batch `{}` has mixed or empty payload classes",
908 batch.batch_id
909 ))
910 })?,
911 delivery_policy: decode_delivery_policy(
912 row.delivery_policy.clone(),
913 )?,
914 slot_policy: decode_slot_policy(row.slot_policy.clone())?,
915 merge_key: decode_merge_key(row.merge_key_json.clone())?,
916 })
917 })
918 .collect::<Result<Vec<_>, StoreError>>()?;
919 let selected_len = select_leading_session_command(&candidates);
920 if selected_len == 0 {
921 return Ok(TxOutcome::Commit(None));
922 }
923 let mut selected = candidate_rows;
924 selected.truncate(selected_len);
925 let mut selected_batches = candidate_batches;
926 selected_batches.truncate(selected_len);
927 let lease = QueuedWorkClaimLease::derive(
928 &candidates[0],
929 &session_id,
930 &owner,
931 now,
932 generation,
933 );
934 let liveness_json = encode_liveness(&owner.liveness)?;
935 for row in &selected {
936 let claimed = tx
937 .execute(
938 "UPDATE queued_work_batches
939 SET claim_id = ?3,
940 claim_owner_id = ?4,
941 claim_owner_incarnation_id = ?5,
942 claim_owner_liveness_json = ?6,
943 claim_token = ?7,
944 claim_fencing_token = claim_fencing_token + 1,
945 claim_session_lease_generation = ?8
946 WHERE session_id = ?1
947 AND batch_id = ?2
948 AND (
949 claim_token IS NULL
950 OR claim_session_lease_generation <> ?8
951 )",
952 params![
953 session_id,
954 row.batch_id,
955 lease.claim_id,
956 owner.owner_id.as_str(),
957 owner.incarnation_id.as_str(),
958 liveness_json.as_str(),
959 lease.lease_token,
960 lease.session_lease_generation as i64,
961 ],
962 )
963 .map_err(sqlite_error)?;
964 if claimed == 0 {
965 return Ok(TxOutcome::Rollback(None));
966 }
967 }
968 Ok(TxOutcome::Commit(Some(QueuedWorkClaim {
969 session_id: session_id.clone(),
970 claim_id: lease.claim_id,
971 owner: owner.clone(),
972 lease_token: lease.lease_token,
973 fencing_token: lease.fencing_token,
974 session_lease_generation: lease.session_lease_generation,
975 batches: selected_batches,
976 })))
977 })(
978 );
979 match outcome {
980 Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
981 Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
982 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
983 }
984 })
985 .await
986 .map_err(sqlite_error)?
987 }
988
989 async fn claim_ready_queued_work(
990 &self,
991 session_id: &str,
992 session_execution_lease: &SessionExecutionLeaseFence,
993 owner: &LeaseOwnerIdentity,
994 boundary: QueuedWorkClaimBoundary,
995 max_batches: usize,
996 ) -> Result<Option<QueuedWorkClaim>, StoreError> {
997 if max_batches == 0 {
998 return Ok(None);
999 }
1000 let session_id = session_id.to_string();
1001 let session_execution_lease = session_execution_lease.clone();
1002 let owner = owner.clone();
1003 let now = self.clock.timestamp_ms();
1004 self.conn
1005 .write_flow(move |tx| {
1006 let outcome: Result<TxOutcome<Option<QueuedWorkClaim>>, StoreError> = (|| {
1007 ensure_session_execution_lease_conn(
1008 tx,
1009 &session_id,
1010 &session_execution_lease,
1011 now,
1012 )?;
1013 let generation = session_execution_lease.fencing_token;
1014 let candidate_rows = {
1015 let mut stmt = tx
1016 .prepare(&sqlite_queued_work_claim_candidates_sql(boundary))
1017 .map_err(sqlite_error)?;
1018 let rows = stmt
1019 .query_map(
1020 params![
1021 session_id,
1022 now as i64,
1023 generation as i64,
1024 claim_scan_limit(max_batches)
1025 ],
1026 queued_batch_row_from_sql,
1027 )
1028 .map_err(sqlite_error)?;
1029 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1030 };
1031 let candidate_rows = candidate_rows
1032 .into_iter()
1033 .filter(|row| {
1034 row.claim_token.is_none()
1035 || row.claim_session_lease_generation != generation
1036 })
1037 .collect::<Vec<_>>();
1038 let candidate_batches = queued_work_batches_from_conn(tx, &candidate_rows)?;
1039 let candidates = candidate_rows
1040 .iter()
1041 .zip(candidate_batches.iter())
1042 .map(|(row, batch)| {
1043 Ok(ClaimCandidate {
1044 enqueue_seq: row.enqueue_seq,
1045 claim_fencing_token: row.claim_fencing_token,
1046 work_class: batch.work_class().ok_or_else(|| {
1047 StoreError::Backend(format!(
1048 "queued-work batch `{}` has mixed or empty payload classes",
1049 batch.batch_id
1050 ))
1051 })?,
1052 delivery_policy: decode_delivery_policy(
1053 row.delivery_policy.clone(),
1054 )?,
1055 slot_policy: decode_slot_policy(row.slot_policy.clone())?,
1056 merge_key: decode_merge_key(row.merge_key_json.clone())?,
1057 })
1058 })
1059 .collect::<Result<Vec<_>, StoreError>>()?;
1060 let selected_len =
1061 select_turn_work_claim_prefix(&candidates, boundary, max_batches);
1062 if selected_len == 0 {
1063 return Ok(TxOutcome::Commit(None));
1064 }
1065 let mut selected = candidate_rows;
1066 selected.truncate(selected_len);
1067 let mut selected_batches = candidate_batches;
1068 selected_batches.truncate(selected_len);
1069 let lease = QueuedWorkClaimLease::derive(
1070 &candidates[0],
1071 &session_id,
1072 &owner,
1073 now,
1074 generation,
1075 );
1076 let liveness_json = encode_liveness(&owner.liveness)?;
1077 for row in &selected {
1078 let claimed = tx
1087 .execute(
1088 "UPDATE queued_work_batches
1089 SET claim_id = ?3,
1090 claim_owner_id = ?4,
1091 claim_owner_incarnation_id = ?5,
1092 claim_owner_liveness_json = ?6,
1093 claim_token = ?7,
1094 claim_fencing_token = claim_fencing_token + 1,
1095 claim_session_lease_generation = ?8
1096 WHERE session_id = ?1
1097 AND batch_id = ?2
1098 AND (
1099 claim_token IS NULL
1100 OR claim_session_lease_generation <> ?8
1101 )",
1102 params![
1103 session_id,
1104 row.batch_id,
1105 lease.claim_id,
1106 owner.owner_id.as_str(),
1107 owner.incarnation_id.as_str(),
1108 liveness_json.as_str(),
1109 lease.lease_token,
1110 lease.session_lease_generation as i64,
1111 ],
1112 )
1113 .map_err(sqlite_error)?;
1114 if claimed == 0 {
1115 return Ok(TxOutcome::Rollback(None));
1119 }
1120 }
1121 Ok(TxOutcome::Commit(Some(QueuedWorkClaim {
1122 session_id: session_id.clone(),
1123 claim_id: lease.claim_id,
1124 owner: owner.clone(),
1125 lease_token: lease.lease_token,
1126 fencing_token: lease.fencing_token,
1127 session_lease_generation: lease.session_lease_generation,
1128 batches: selected_batches,
1129 })))
1130 })(
1131 );
1132 match outcome {
1136 Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
1137 Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
1138 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1139 }
1140 })
1141 .await
1142 .map_err(sqlite_error)?
1143 }
1144
1145 async fn claim_checkpoint_work(
1146 &self,
1147 session_id: &str,
1148 session_execution_lease: &SessionExecutionLeaseFence,
1149 owner: &LeaseOwnerIdentity,
1150 turn_id: &str,
1151 checkpoint: lash_core::CheckpointKind,
1152 max_inputs: usize,
1153 max_batches: usize,
1154 ) -> Result<(Option<lash_core::TurnInputClaim>, Option<QueuedWorkClaim>), StoreError> {
1155 #[cfg(test)]
1156 self.checkpoint_probe_count
1157 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1158 let now = self.clock.timestamp_ms();
1159 if !checkpoint_work_pending_sqlite(
1160 &self.conn,
1161 now,
1162 session_id,
1163 session_execution_lease.fencing_token,
1164 turn_id,
1165 checkpoint,
1166 max_inputs,
1167 max_batches,
1168 )
1169 .await?
1170 {
1171 return Ok((None, None));
1172 }
1173
1174 #[cfg(test)]
1175 self.checkpoint_write_transaction_count
1176 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1177 let session_id = session_id.to_string();
1178 let session_execution_lease = session_execution_lease.clone();
1179 let owner = owner.clone();
1180 let turn_id = turn_id.to_string();
1181 self.conn
1182 .write_flow(move |tx| {
1183 let outcome: Result<
1184 TxOutcome<(Option<lash_core::TurnInputClaim>, Option<QueuedWorkClaim>)>,
1185 StoreError,
1186 > = (|| {
1187 ensure_session_execution_lease_conn(
1188 tx,
1189 &session_id,
1190 &session_execution_lease,
1191 now,
1192 )?;
1193 let input = claim_pending_turn_inputs_sqlite_conn(
1194 tx,
1195 now,
1196 &session_id,
1197 &session_execution_lease,
1198 &owner,
1199 max_inputs,
1200 lash_core::TurnInputClaimMode::ActiveTurn {
1201 turn_id,
1202 checkpoint,
1203 },
1204 )?;
1205 let input = match input {
1206 TxOutcome::Commit(input) => input,
1207 TxOutcome::Rollback(input) => {
1208 return Ok(TxOutcome::Rollback((input, None)));
1209 }
1210 };
1211 let queued = claim_ready_queued_work_sqlite_conn(
1212 tx,
1213 now,
1214 &session_id,
1215 &session_execution_lease,
1216 &owner,
1217 QueuedWorkClaimBoundary::ActiveTurnCheckpoint,
1218 max_batches,
1219 )?;
1220 match queued {
1221 TxOutcome::Commit(queued) => Ok(TxOutcome::Commit((input, queued))),
1222 TxOutcome::Rollback(queued) => Ok(TxOutcome::Rollback((None, queued))),
1223 }
1224 })();
1225 match outcome {
1226 Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
1227 Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
1228 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1229 }
1230 })
1231 .await
1232 .map_err(sqlite_error)?
1233 }
1234
1235 async fn claim_ready_queued_work_by_batch_ids(
1236 &self,
1237 session_id: &str,
1238 session_execution_lease: &SessionExecutionLeaseFence,
1239 owner: &LeaseOwnerIdentity,
1240 boundary: QueuedWorkClaimBoundary,
1241 batch_ids: &[String],
1242 ) -> Result<Option<QueuedWorkClaim>, StoreError> {
1243 if batch_ids.is_empty() {
1244 return Ok(None);
1245 }
1246 let session_id = session_id.to_string();
1247 let fence = session_execution_lease.clone();
1248 let owner = owner.clone();
1249 let batch_ids = batch_ids.to_vec();
1250 let now = self.clock.timestamp_ms();
1251 self.conn
1252 .write_flow(move |tx| {
1253 let outcome: Result<Option<QueuedWorkClaim>, StoreError> = (|| {
1254 ensure_session_execution_lease_conn(tx, &session_id, &fence, now)?;
1255 let generation = fence.fencing_token;
1256 let mut rows = Vec::new();
1257 let mut batches = Vec::new();
1258 for batch_id in &batch_ids {
1259 let row = tx
1260 .query_row(
1261 "SELECT enqueue_seq, batch_id, session_id, source_key,
1262 delivery_policy, slot_policy, merge_key_json,
1263 available_at_ms, enqueued_at_ms, claim_fencing_token,
1264 claim_owner_id, claim_owner_incarnation_id,
1265 claim_owner_liveness_json, claim_token,
1266 claim_session_lease_generation
1267 FROM queued_work_batches
1268 WHERE session_id = ?1 AND batch_id = ?2
1269 AND available_at_ms <= ?3
1270 AND (claim_token IS NULL
1271 OR claim_session_lease_generation <> ?4)",
1272 params![session_id, batch_id, now as i64, generation as i64],
1273 queued_batch_row_from_sql,
1274 )
1275 .optional()
1276 .map_err(sqlite_error)?;
1277 let Some(row) = row else {
1278 return Ok(None);
1279 };
1280 let batch = queued_work_batch_from_conn(tx, row.clone())?;
1281 if batch.work_class() != Some(lash_core::runtime::QueuedWorkClass::TurnWork)
1282 {
1283 return Ok(None);
1284 }
1285 rows.push(row);
1286 batches.push(batch);
1287 }
1288 let candidates = rows
1289 .iter()
1290 .map(|row| {
1291 Ok(ClaimCandidate {
1292 enqueue_seq: row.enqueue_seq,
1293 claim_fencing_token: row.claim_fencing_token,
1294 work_class: lash_core::runtime::QueuedWorkClass::TurnWork,
1295 delivery_policy: decode_delivery_policy(
1296 row.delivery_policy.clone(),
1297 )?,
1298 slot_policy: decode_slot_policy(row.slot_policy.clone())?,
1299 merge_key: decode_merge_key(row.merge_key_json.clone())?,
1300 })
1301 })
1302 .collect::<Result<Vec<_>, StoreError>>()?;
1303 if select_turn_work_claim_prefix(&candidates, boundary, candidates.len())
1304 != candidates.len()
1305 {
1306 return Ok(None);
1307 }
1308 let lease = QueuedWorkClaimLease::derive(
1309 &candidates[0],
1310 &session_id,
1311 &owner,
1312 now,
1313 generation,
1314 );
1315 let owner_liveness_json = encode_liveness(&owner.liveness)?;
1316 for row in &rows {
1317 let changed = tx
1318 .execute(
1319 "UPDATE queued_work_batches
1320 SET claim_id = ?3, claim_owner_id = ?4,
1321 claim_owner_incarnation_id = ?5,
1322 claim_owner_liveness_json = ?6, claim_token = ?7,
1323 claim_fencing_token = claim_fencing_token + 1,
1324 claim_session_lease_generation = ?8
1325 WHERE session_id = ?1 AND batch_id = ?2
1326 AND (claim_token IS NULL
1327 OR claim_session_lease_generation <> ?8)",
1328 params![
1329 session_id,
1330 row.batch_id,
1331 lease.claim_id,
1332 owner.owner_id,
1333 owner.incarnation_id,
1334 owner_liveness_json,
1335 lease.lease_token,
1336 generation as i64,
1337 ],
1338 )
1339 .map_err(sqlite_error)?;
1340 if changed != 1 {
1341 return Ok(None);
1342 }
1343 }
1344 Ok(Some(QueuedWorkClaim {
1345 session_id,
1346 claim_id: lease.claim_id,
1347 owner,
1348 lease_token: lease.lease_token,
1349 fencing_token: lease.fencing_token,
1350 session_lease_generation: lease.session_lease_generation,
1351 batches,
1352 }))
1353 })();
1354 match outcome {
1355 Ok(Some(value)) => Ok(TxOutcome::Commit(Ok(Some(value)))),
1356 Ok(None) => Ok(TxOutcome::Rollback(Ok(None))),
1357 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1358 }
1359 })
1360 .await
1361 .map_err(sqlite_error)?
1362 }
1363
1364 async fn abandon_queued_work_claim(&self, claim: &QueuedWorkClaim) -> Result<(), StoreError> {
1365 let session_id = claim.session_id.clone();
1366 let claim_id = claim.claim_id.clone();
1367 let lease_token = claim.lease_token.clone();
1368 self.conn
1369 .write(move |tx| {
1370 tx.execute(
1371 "UPDATE queued_work_batches
1372 SET claim_id = NULL,
1373 claim_owner_id = NULL,
1374 claim_owner_incarnation_id = NULL,
1375 claim_owner_liveness_json = NULL,
1376 claim_token = NULL,
1377 claim_session_lease_generation = 0
1378 WHERE session_id = ?1 AND claim_id = ?2 AND claim_token = ?3",
1379 params![session_id, claim_id, lease_token],
1380 )
1381 })
1382 .await
1383 .map_err(sqlite_error)?;
1384 Ok(())
1385 }
1386
1387 async fn abandon_queued_work_claims(
1388 &self,
1389 claims: &[QueuedWorkClaim],
1390 ) -> Result<(), StoreError> {
1391 if claims.is_empty() {
1392 return Ok(());
1393 }
1394 let mut sql = "UPDATE queued_work_batches
1395 SET claim_id = NULL,
1396 claim_owner_id = NULL,
1397 claim_owner_incarnation_id = NULL,
1398 claim_owner_liveness_json = NULL,
1399 claim_token = NULL,
1400 claim_session_lease_generation = 0
1401 WHERE (session_id, claim_id, claim_token) IN ("
1402 .to_string();
1403 let mut values: Vec<rusqlite::types::Value> = Vec::with_capacity(claims.len() * 3);
1404 for (index, claim) in claims.iter().enumerate() {
1405 if index > 0 {
1406 sql.push_str(", ");
1407 }
1408 sql.push_str("(?, ?, ?)");
1409 values.push(claim.session_id.clone().into());
1410 values.push(claim.claim_id.clone().into());
1411 values.push(claim.lease_token.clone().into());
1412 }
1413 sql.push(')');
1414 self.conn
1415 .write(move |tx| tx.execute(&sql, rusqlite::params_from_iter(values.iter())))
1416 .await
1417 .map_err(sqlite_error)?;
1418 Ok(())
1419 }
1420
1421 async fn cancel_queued_work_batch(
1422 &self,
1423 session_id: &str,
1424 batch_id: &str,
1425 ) -> Result<Option<QueuedWorkBatch>, StoreError> {
1426 let session_id = session_id.to_string();
1427 let batch_id = batch_id.to_string();
1428 let now = self.clock.timestamp_ms() as i64;
1429 self.conn
1430 .write_flow(move |tx| {
1431 let outcome: Result<Option<QueuedWorkBatch>, StoreError> = (|| {
1432 let row = tx
1433 .query_row(
1434 "SELECT enqueue_seq, batch_id, session_id, source_key, delivery_policy,
1435 slot_policy, merge_key_json, available_at_ms, enqueued_at_ms,
1436 claim_fencing_token, claim_owner_id, claim_owner_incarnation_id,
1437 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1438 FROM queued_work_batches
1439 WHERE session_id = ?1
1440 AND batch_id = ?2
1441 AND (claim_token IS NULL OR NOT EXISTS (
1442 SELECT 1 FROM session_execution_leases sel
1443 WHERE sel.session_id = ?1
1444 AND sel.lease_token IS NOT NULL
1445 AND sel.lease_expires_at_ms > ?3
1446 AND sel.lease_fencing_token
1447 = queued_work_batches.claim_session_lease_generation
1448 ))",
1449 params![session_id, batch_id, now],
1450 queued_batch_row_from_sql,
1451 )
1452 .optional()
1453 .map_err(sqlite_error)?;
1454 let Some(row) = row else {
1455 return Ok(None);
1456 };
1457 let batch = queued_work_batch_from_conn(tx, row)?;
1458 tx.execute(
1459 "DELETE FROM queued_work_batches
1460 WHERE session_id = ?1
1461 AND batch_id = ?2
1462 AND (claim_token IS NULL OR NOT EXISTS (
1463 SELECT 1 FROM session_execution_leases sel
1464 WHERE sel.session_id = ?1
1465 AND sel.lease_token IS NOT NULL
1466 AND sel.lease_expires_at_ms > ?3
1467 AND sel.lease_fencing_token
1468 = queued_work_batches.claim_session_lease_generation
1469 ))",
1470 params![session_id, batch_id, now],
1471 )
1472 .map_err(sqlite_error)?;
1473 Ok(Some(batch))
1474 })();
1475 match outcome {
1476 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1477 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1478 }
1479 })
1480 .await
1481 .map_err(sqlite_error)?
1482 }
1483
1484 async fn list_queued_work(&self, session_id: &str) -> Result<Vec<QueuedWorkBatch>, StoreError> {
1485 let session_id = session_id.to_string();
1486 self.conn
1487 .call(move |conn| {
1488 let outcome: Result<Vec<QueuedWorkBatch>, StoreError> = (|| {
1489 let rows = {
1490 let mut stmt = conn
1491 .prepare(
1492 "SELECT enqueue_seq, batch_id, session_id, source_key, delivery_policy,
1493 slot_policy, merge_key_json, available_at_ms, enqueued_at_ms,
1494 claim_fencing_token, claim_owner_id, claim_owner_incarnation_id,
1495 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1496 FROM queued_work_batches
1497 WHERE session_id = ?1
1498 ORDER BY enqueue_seq ASC",
1499 )
1500 .map_err(sqlite_error)?;
1501 let rows = stmt
1502 .query_map(params![session_id], queued_batch_row_from_sql)
1503 .map_err(sqlite_error)?;
1504 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1505 };
1506 rows.into_iter()
1507 .map(|row| queued_work_batch_from_conn(conn, row))
1508 .collect()
1509 })();
1510 Ok(outcome)
1511 })
1512 .await
1513 .map_err(sqlite_error)?
1514 }
1515
1516 async fn list_pending_queued_work(
1517 &self,
1518 session_id: &str,
1519 ) -> Result<Vec<QueuedWorkBatch>, StoreError> {
1520 let session_id = session_id.to_string();
1521 let now = self.clock.timestamp_ms();
1522 self.conn
1523 .call(move |conn| {
1524 let outcome: Result<Vec<QueuedWorkBatch>, StoreError> = (|| {
1525 let rows = {
1526 let mut stmt = conn
1527 .prepare(
1528 "SELECT enqueue_seq, batch_id, session_id, source_key, delivery_policy,
1529 slot_policy, merge_key_json, available_at_ms, enqueued_at_ms,
1530 claim_fencing_token, claim_owner_id, claim_owner_incarnation_id,
1531 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1532 FROM queued_work_batches
1533 WHERE session_id = ?1
1534 AND (claim_token IS NULL OR NOT EXISTS (
1535 SELECT 1 FROM session_execution_leases sel
1536 WHERE sel.session_id = ?1
1537 AND sel.lease_token IS NOT NULL
1538 AND sel.lease_expires_at_ms > ?2
1539 AND sel.lease_fencing_token
1540 = queued_work_batches.claim_session_lease_generation
1541 ))
1542 ORDER BY enqueue_seq ASC",
1543 )
1544 .map_err(sqlite_error)?;
1545 let rows = stmt
1546 .query_map(
1547 params![session_id, now as i64],
1548 queued_batch_row_from_sql,
1549 )
1550 .map_err(sqlite_error)?;
1551 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1552 };
1553 rows.into_iter()
1554 .map(|row| queued_work_batch_from_conn(conn, row))
1555 .collect()
1556 })();
1557 Ok(outcome)
1558 })
1559 .await
1560 .map_err(sqlite_error)?
1561 }
1562}
1563
1564#[async_trait::async_trait]
1565impl TurnInputStore for Store {
1566 async fn enqueue_pending_turn_input(
1567 &self,
1568 draft: lash_core::PendingTurnInputDraft,
1569 ) -> Result<lash_core::PendingTurnInput, StoreError> {
1570 let nonce = self.commit_count.fetch_add(1, AtomicOrdering::Relaxed);
1571 let now = self.clock.timestamp_ms();
1572 self.conn
1573 .write_flow(move |tx| {
1574 let outcome: Result<lash_core::PendingTurnInput, StoreError> = (|| {
1575 if let Some(source_key) = draft.source_key.as_deref() {
1576 let existing_id: Option<String> = tx
1577 .query_row(
1578 "SELECT input_id
1579 FROM pending_turn_inputs
1580 WHERE session_id = ?1 AND source_key = ?2",
1581 params![draft.session_id, source_key],
1582 |row| row.get(0),
1583 )
1584 .optional()
1585 .map_err(sqlite_error)?;
1586 if let Some(input_id) = existing_id {
1587 let existing = load_pending_turn_input_by_id_conn(
1588 tx,
1589 &draft.session_id,
1590 &input_id,
1591 )?
1592 .ok_or_else(|| {
1593 StoreError::Backend(
1594 "pending turn input source row disappeared".to_string(),
1595 )
1596 })?;
1597 if !draft.submitted_content_matches(&existing).map_err(|err| {
1598 StoreError::Backend(format!(
1599 "failed to compare pending turn input submission: {err}"
1600 ))
1601 })? {
1602 return Err(StoreError::PendingTurnInputSourceKeyConflict {
1603 session_id: draft.session_id.clone(),
1604 source_key: source_key.to_string(),
1605 existing_input_id: existing.input_id.clone(),
1606 });
1607 }
1608 return Ok(existing);
1609 }
1610 }
1611 let input_id = draft.input_id.clone().unwrap_or_else(|| {
1612 derive_pending_turn_input_id(
1613 &draft.session_id,
1614 draft.source_key.as_deref(),
1615 now,
1616 nonce,
1617 )
1618 });
1619 let state = match draft.ingress {
1620 lash_core::TurnInputIngress::ActiveTurn { .. } => {
1621 lash_core::TurnInputState::PendingActive
1622 }
1623 lash_core::TurnInputIngress::NextTurn => {
1624 lash_core::TurnInputState::DeferredNextTurn
1625 }
1626 };
1627 tx.execute(
1628 "INSERT INTO pending_turn_inputs (
1629 input_id, session_id, source_key, ingress_json, state,
1630 input_json, enqueued_at_ms
1631 )
1632 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1633 params![
1634 input_id,
1635 draft.session_id,
1636 draft.source_key.as_deref(),
1637 encode_json(&draft.ingress),
1638 state.as_str(),
1639 encode_json(&draft.input),
1640 now as i64,
1641 ],
1642 )
1643 .map_err(sqlite_error)?;
1644 load_pending_turn_input_by_id_conn(tx, &draft.session_id, &input_id)?
1645 .ok_or_else(|| {
1646 StoreError::Backend("pending turn input insert disappeared".to_string())
1647 })
1648 })();
1649 match outcome {
1650 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1651 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1652 }
1653 })
1654 .await
1655 .map_err(sqlite_error)?
1656 }
1657
1658 async fn list_pending_turn_inputs(
1659 &self,
1660 session_id: &str,
1661 ) -> Result<Vec<lash_core::PendingTurnInput>, StoreError> {
1662 let session_id = session_id.to_string();
1663 let now = self.clock.timestamp_ms();
1664 self.conn
1665 .call(move |conn| {
1666 let outcome: Result<Vec<lash_core::PendingTurnInput>, StoreError> = (|| {
1667 let rows = {
1668 let mut stmt = conn
1669 .prepare(
1670 "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
1671 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
1672 claim_owner_id, claim_owner_incarnation_id,
1673 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1674 FROM pending_turn_inputs
1675 WHERE session_id = ?1
1676 AND state IN (?2, ?3)
1677 AND (claim_token IS NULL OR NOT EXISTS (
1678 SELECT 1 FROM session_execution_leases sel
1679 WHERE sel.session_id = ?1
1680 AND sel.lease_token IS NOT NULL
1681 AND sel.lease_expires_at_ms > ?4
1682 AND sel.lease_fencing_token
1683 = pending_turn_inputs.claim_session_lease_generation
1684 ))
1685 ORDER BY enqueue_seq ASC",
1686 )
1687 .map_err(sqlite_error)?;
1688 let rows = stmt
1689 .query_map(
1690 params![
1691 session_id,
1692 lash_core::TurnInputState::PendingActive.as_str(),
1693 lash_core::TurnInputState::DeferredNextTurn.as_str(),
1694 now as i64
1695 ],
1696 pending_turn_input_row_from_sql,
1697 )
1698 .map_err(sqlite_error)?;
1699 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1700 };
1701 rows.into_iter().map(pending_turn_input_from_row).collect()
1702 })(
1703 );
1704 Ok(outcome)
1705 })
1706 .await
1707 .map_err(sqlite_error)?
1708 }
1709
1710 async fn list_turn_input_applications(
1711 &self,
1712 session_id: &str,
1713 ) -> Result<Vec<lash_core::TurnInputApplication>, StoreError> {
1714 let session_id = session_id.to_string();
1715 self.conn
1716 .call(move |conn| {
1717 let outcome = (|| {
1718 let mut stmt = conn
1719 .prepare(
1720 "SELECT turn_id, result_json
1721 FROM runtime_turn_commits
1722 WHERE session_id = ?1",
1723 )
1724 .map_err(sqlite_error)?;
1725 let rows = stmt
1726 .query_map(params![session_id], |row| {
1727 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1728 })
1729 .map_err(sqlite_error)?;
1730 let mut commits = Vec::new();
1731 for row in rows {
1732 let (turn_id, result_json) = row.map_err(sqlite_error)?;
1733 let result: RuntimeCommitResult = serde_json::from_str(&result_json)
1734 .map_err(|err| {
1735 StoreError::Backend(format!(
1736 "failed to decode runtime turn commit result: {err}"
1737 ))
1738 })?;
1739 commits.push((
1740 result.head_revision,
1741 turn_id,
1742 result.turn_input_applications,
1743 ));
1744 }
1745 commits.sort_by(|left, right| {
1746 (left.0, left.1.as_str()).cmp(&(right.0, right.1.as_str()))
1747 });
1748 Ok(commits
1749 .into_iter()
1750 .flat_map(|(_, _, applications)| applications)
1751 .collect())
1752 })();
1753 Ok(outcome)
1754 })
1755 .await
1756 .map_err(sqlite_error)?
1757 }
1758
1759 async fn cancel_pending_turn_inputs(
1760 &self,
1761 session_id: &str,
1762 targets: &[lash_core::PendingTurnInputCancelTarget],
1763 ) -> Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> {
1764 let session_id = session_id.to_string();
1765 let targets = targets.to_vec();
1766 let now = self.clock.timestamp_ms();
1767 self.conn
1768 .write_flow(move |tx| {
1769 let outcome: Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> =
1770 (|| {
1771 let mut results = Vec::with_capacity(targets.len());
1772 for target in targets {
1773 let outcome = match load_pending_turn_input_row_by_target_conn(
1774 tx,
1775 &session_id,
1776 &target,
1777 )? {
1778 Some(row) => cancel_pending_turn_input_row_conn(tx, row, now)?,
1779 None => lash_core::PendingTurnInputCancelOutcome::NotFound,
1780 };
1781 results
1782 .push(lash_core::PendingTurnInputCancelResult { target, outcome });
1783 }
1784 Ok(results)
1785 })();
1786 match outcome {
1787 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1788 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1789 }
1790 })
1791 .await
1792 .map_err(sqlite_error)?
1793 }
1794
1795 async fn cancel_pending_turn_input_suffix(
1796 &self,
1797 session_id: &str,
1798 anchor: &lash_core::PendingTurnInputCancelTarget,
1799 ) -> Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> {
1800 let session_id = session_id.to_string();
1801 let anchor = anchor.clone();
1802 let now = self.clock.timestamp_ms();
1803 self.conn
1804 .write_flow(move |tx| {
1805 let outcome: Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> =
1806 (|| {
1807 let Some(anchor_row) =
1808 load_pending_turn_input_row_by_target_conn(tx, &session_id, &anchor)?
1809 else {
1810 return Ok(
1811 lash_core::PendingTurnInputSuffixCancelOutcome::AnchorNotFound {
1812 anchor,
1813 },
1814 );
1815 };
1816 let rows = {
1817 let mut stmt = tx
1818 .prepare(
1819 "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
1820 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
1821 claim_owner_id, claim_owner_incarnation_id,
1822 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1823 FROM pending_turn_inputs
1824 WHERE session_id = ?1 AND enqueue_seq >= ?2
1825 ORDER BY enqueue_seq ASC",
1826 )
1827 .map_err(sqlite_error)?;
1828 let rows = stmt
1829 .query_map(
1830 params![session_id, anchor_row.enqueue_seq as i64],
1831 pending_turn_input_row_from_sql,
1832 )
1833 .map_err(sqlite_error)?;
1834 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1835 };
1836 let mut outcomes = Vec::with_capacity(rows.len());
1837 for row in rows {
1838 outcomes.push(cancel_pending_turn_input_row_conn(tx, row, now)?);
1839 }
1840 Ok(lash_core::PendingTurnInputSuffixCancelOutcome::Outcomes {
1841 anchor,
1842 outcomes,
1843 })
1844 })();
1845 match outcome {
1846 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1847 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1848 }
1849 })
1850 .await
1851 .map_err(sqlite_error)?
1852 }
1853
1854 async fn claim_active_turn_inputs(
1855 &self,
1856 session_id: &str,
1857 session_execution_lease: &SessionExecutionLeaseFence,
1858 owner: &LeaseOwnerIdentity,
1859 turn_id: &str,
1860 checkpoint: lash_core::CheckpointKind,
1861 max_inputs: usize,
1862 ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1863 claim_pending_turn_inputs_sqlite(
1864 &self.conn,
1865 self.clock.timestamp_ms(),
1866 session_id,
1867 session_execution_lease,
1868 owner,
1869 max_inputs,
1870 lash_core::TurnInputClaimMode::ActiveTurn {
1871 turn_id: turn_id.to_string(),
1872 checkpoint,
1873 },
1874 )
1875 .await
1876 }
1877
1878 async fn claim_next_turn_inputs(
1879 &self,
1880 session_id: &str,
1881 session_execution_lease: &SessionExecutionLeaseFence,
1882 owner: &LeaseOwnerIdentity,
1883 max_inputs: usize,
1884 ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1885 claim_pending_turn_inputs_sqlite(
1886 &self.conn,
1887 self.clock.timestamp_ms(),
1888 session_id,
1889 session_execution_lease,
1890 owner,
1891 max_inputs,
1892 lash_core::TurnInputClaimMode::NextTurn,
1893 )
1894 .await
1895 }
1896
1897 async fn abandon_turn_input_claim(
1898 &self,
1899 claim: &lash_core::TurnInputClaim,
1900 ) -> Result<(), StoreError> {
1901 let session_id = claim.session_id.clone();
1902 let claim_id = claim.claim_id.clone();
1903 let lease_token = claim.lease_token.clone();
1904 let restored_state = match claim.mode {
1905 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
1906 lash_core::TurnInputState::PendingActive
1907 }
1908 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
1909 };
1910 self.conn
1911 .write(move |tx| {
1912 tx.execute(
1913 "UPDATE pending_turn_inputs
1914 SET state = CASE
1915 WHEN state = ?4 THEN ?5
1916 ELSE state
1917 END,
1918 claim_id = NULL,
1919 claim_owner_id = NULL,
1920 claim_owner_incarnation_id = NULL,
1921 claim_owner_liveness_json = NULL,
1922 claim_token = NULL,
1923 claim_session_lease_generation = 0
1924 WHERE session_id = ?1 AND claim_id = ?2 AND claim_token = ?3",
1925 params![
1926 session_id,
1927 claim_id,
1928 lease_token,
1929 lash_core::TurnInputState::Accepted.as_str(),
1930 restored_state.as_str(),
1931 ],
1932 )
1933 })
1934 .await
1935 .map_err(sqlite_error)?;
1936 Ok(())
1937 }
1938
1939 async fn abandon_turn_input_claims(
1940 &self,
1941 claims: &[lash_core::TurnInputClaim],
1942 ) -> Result<(), StoreError> {
1943 if claims.is_empty() {
1944 return Ok(());
1945 }
1946 let mut sql = "UPDATE pending_turn_inputs
1947 SET state = CASE
1948 WHEN state = 'accepted' THEN 'pending_active'
1949 ELSE state
1950 END,
1951 claim_id = NULL,
1952 claim_owner_id = NULL,
1953 claim_owner_incarnation_id = NULL,
1954 claim_owner_liveness_json = NULL,
1955 claim_token = NULL,
1956 claim_session_lease_generation = 0
1957 WHERE (session_id, claim_id, claim_token) IN ("
1958 .to_string();
1959 let mut values: Vec<rusqlite::types::Value> = Vec::with_capacity(claims.len() * 3);
1960 for (index, claim) in claims.iter().enumerate() {
1961 if index > 0 {
1962 sql.push_str(", ");
1963 }
1964 sql.push_str("(?, ?, ?)");
1965 values.push(claim.session_id.clone().into());
1966 values.push(claim.claim_id.clone().into());
1967 values.push(claim.lease_token.clone().into());
1968 }
1969 sql.push(')');
1970 self.conn
1971 .write(move |tx| tx.execute(&sql, rusqlite::params_from_iter(values.iter())))
1972 .await
1973 .map_err(sqlite_error)?;
1974 Ok(())
1975 }
1976}
1977
1978#[async_trait::async_trait]
1979impl StoreMaintenance for Store {
1980 async fn tombstone_nodes(&self, ids: &[String]) -> Result<(), StoreError> {
1981 if ids.is_empty() {
1982 return Ok(());
1983 }
1984 let ids = ids.to_vec();
1985 self.conn
1986 .write(move |tx| {
1987 let mut stmt =
1988 tx.prepare("UPDATE graph_nodes SET tombstoned = 1 WHERE node_id = ?1")?;
1989 for id in &ids {
1990 stmt.execute(params![id])?;
1991 }
1992 Ok(())
1993 })
1994 .await
1995 .map_err(sqlite_error)
1996 }
1997
1998 async fn vacuum(&self) -> Result<VacuumReport, StoreError> {
1999 let (removed_node_count, removed_pending_turn_input_tombstone_count) = self
2000 .conn
2001 .write(move |tx| {
2002 let removed_node_count =
2003 tx.execute("DELETE FROM graph_nodes WHERE tombstoned = 1", [])?;
2004 let removed_pending_turn_input_tombstone_count = tx.execute(
2005 "DELETE FROM pending_turn_inputs
2006 WHERE state IN (?1, ?2)",
2007 params![
2008 lash_core::TurnInputState::Cancelled.as_str(),
2009 lash_core::TurnInputState::Completed.as_str()
2010 ],
2011 )?;
2012 Ok((
2013 removed_node_count,
2014 removed_pending_turn_input_tombstone_count,
2015 ))
2016 })
2017 .await
2018 .map_err(sqlite_error)?;
2019 Ok(VacuumReport {
2020 removed_node_count,
2021 removed_pending_turn_input_tombstone_count,
2022 })
2023 }
2024
2025 async fn gc_unreachable(&self) -> Result<GcReport, StoreError> {
2026 Ok(Store::gc_unreachable(self).await)
2027 }
2028}
2029
2030fn derive_pending_turn_input_id(
2031 session_id: &str,
2032 source_key: Option<&str>,
2033 now_epoch_ms: u64,
2034 nonce: u64,
2035) -> String {
2036 format!(
2037 "ti:{:x}",
2038 Sha256::digest(format!("{session_id}:{source_key:?}:{now_epoch_ms}:{nonce}").as_bytes())
2039 )
2040}
2041
2042fn cancel_pending_turn_input_row_conn(
2043 conn: &Connection,
2044 row: PendingTurnInputRow,
2045 now_epoch_ms: u64,
2046) -> Result<lash_core::PendingTurnInputCancelOutcome, StoreError> {
2047 let mut input = pending_turn_input_from_row(row.clone())?;
2048 match input.state {
2049 lash_core::TurnInputState::Cancelled => Ok(
2050 lash_core::PendingTurnInputCancelOutcome::AlreadyCancelled(input),
2051 ),
2052 lash_core::TurnInputState::Completed => Ok(
2053 lash_core::PendingTurnInputCancelOutcome::AlreadyCompleted(input),
2054 ),
2055 lash_core::TurnInputState::Accepted => {
2056 Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2057 claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2058 input,
2059 })
2060 }
2061 lash_core::TurnInputState::PendingActive | lash_core::TurnInputState::DeferredNextTurn => {
2062 let live_claim = row.claim_token.is_some()
2065 && load_session_execution_lease_row_conn(conn, &row.session_id)?.is_some_and(
2066 |lease| {
2067 lease.lease_token.is_some()
2068 && lease.expires_at_ms > now_epoch_ms
2069 && lease.fencing_token == row.claim_session_lease_generation
2070 },
2071 );
2072 if live_claim {
2073 return Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2074 claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2075 input,
2076 });
2077 }
2078 conn.execute(
2079 "UPDATE pending_turn_inputs
2080 SET state = ?3,
2081 claim_id = NULL,
2082 claim_owner_id = NULL,
2083 claim_owner_incarnation_id = NULL,
2084 claim_owner_liveness_json = NULL,
2085 claim_token = NULL,
2086 claim_session_lease_generation = 0
2087 WHERE session_id = ?1 AND input_id = ?2",
2088 params![
2089 row.session_id,
2090 row.input_id,
2091 lash_core::TurnInputState::Cancelled.as_str(),
2092 ],
2093 )
2094 .map_err(sqlite_error)?;
2095 input.state = lash_core::TurnInputState::Cancelled;
2096 Ok(lash_core::PendingTurnInputCancelOutcome::Cancelled(input))
2097 }
2098 }
2099}
2100
2101#[allow(clippy::too_many_arguments)]
2102async fn checkpoint_work_pending_sqlite(
2103 conn: &SqliteConnection,
2104 now: u64,
2105 session_id: &str,
2106 generation: u64,
2107 turn_id: &str,
2108 checkpoint: lash_core::CheckpointKind,
2109 max_inputs: usize,
2110 max_batches: usize,
2111) -> Result<bool, StoreError> {
2112 if max_inputs == 0 && max_batches == 0 {
2113 return Ok(false);
2114 }
2115 let session_id = session_id.to_string();
2116 let turn_id = turn_id.to_string();
2117 conn.call(move |conn| {
2118 let outcome: Result<bool, StoreError> = (|| {
2119 let head_candidate = sqlite_queued_work_head_candidate_cte(
2120 QueuedWorkClaimBoundary::ActiveTurnCheckpoint,
2121 );
2122 let sql = format!(
2123 "WITH {head_candidate}
2124 SELECT (
2125 ?7 > 0 AND EXISTS (
2126 SELECT 1
2127 FROM pending_turn_inputs
2128 WHERE session_id = ?1
2129 AND state = ?4
2130 AND (claim_token IS NULL OR claim_session_lease_generation <> ?3)
2131 AND json_extract(ingress_json, '$.scope') = 'active_turn'
2132 AND json_extract(ingress_json, '$.turn_id') = ?5
2133 AND (
2134 ?6 = 'before_completion'
2135 OR COALESCE(
2136 json_extract(ingress_json, '$.min_boundary'),
2137 'after_work'
2138 ) = 'after_work'
2139 )
2140 LIMIT 1
2141 )
2142 ) OR (
2143 ?8 > 0 AND EXISTS (
2144 SELECT 1
2145 FROM queued_work_head_candidate AS head
2146 JOIN queued_work_items AS item
2147 ON item.batch_id = head.head_batch_id
2148 WHERE json_extract(item.payload_json, '$.type') <> 'session_command'
2149 LIMIT 1
2150 )
2151 )"
2152 );
2153 let pending: i64 = conn
2154 .query_row(
2155 &sql,
2156 params![
2157 session_id,
2158 now as i64,
2159 generation as i64,
2160 lash_core::TurnInputState::PendingActive.as_str(),
2161 turn_id,
2162 match checkpoint {
2163 lash_core::CheckpointKind::AfterWork => "after_work",
2164 lash_core::CheckpointKind::BeforeCompletion => "before_completion",
2165 },
2166 max_inputs as i64,
2167 max_batches as i64,
2168 ],
2169 |row| row.get(0),
2170 )
2171 .map_err(sqlite_error)?;
2172 Ok(pending != 0)
2173 })();
2174 Ok(outcome)
2175 })
2176 .await
2177 .map_err(sqlite_error)?
2178}
2179
2180#[allow(clippy::too_many_arguments)]
2181fn claim_ready_queued_work_sqlite_conn(
2182 tx: &Connection,
2183 now: u64,
2184 session_id: &str,
2185 session_execution_lease: &SessionExecutionLeaseFence,
2186 owner: &LeaseOwnerIdentity,
2187 boundary: QueuedWorkClaimBoundary,
2188 max_batches: usize,
2189) -> Result<TxOutcome<Option<QueuedWorkClaim>>, StoreError> {
2190 if max_batches == 0 {
2191 return Ok(TxOutcome::Commit(None));
2192 }
2193 let generation = session_execution_lease.fencing_token;
2194 let candidate_rows = {
2195 let mut stmt = tx
2196 .prepare(&sqlite_queued_work_claim_candidates_sql(boundary))
2197 .map_err(sqlite_error)?;
2198 let rows = stmt
2199 .query_map(
2200 params![
2201 session_id,
2202 now as i64,
2203 generation as i64,
2204 claim_scan_limit(max_batches)
2205 ],
2206 queued_batch_row_from_sql,
2207 )
2208 .map_err(sqlite_error)?;
2209 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2210 };
2211 let candidate_rows = candidate_rows
2212 .into_iter()
2213 .filter(|row| row.claim_token.is_none() || row.claim_session_lease_generation != generation)
2214 .collect::<Vec<_>>();
2215 let candidate_batches = candidate_rows
2216 .iter()
2217 .map(|row| queued_work_batch_from_conn(tx, row.clone()))
2218 .collect::<Result<Vec<_>, StoreError>>()?;
2219 let candidates = candidate_rows
2220 .iter()
2221 .zip(candidate_batches.iter())
2222 .map(|(row, batch)| {
2223 Ok(ClaimCandidate {
2224 enqueue_seq: row.enqueue_seq,
2225 claim_fencing_token: row.claim_fencing_token,
2226 work_class: batch.work_class().ok_or_else(|| {
2227 StoreError::Backend(format!(
2228 "queued-work batch `{}` has mixed or empty payload classes",
2229 batch.batch_id
2230 ))
2231 })?,
2232 delivery_policy: decode_delivery_policy(row.delivery_policy.clone())?,
2233 slot_policy: decode_slot_policy(row.slot_policy.clone())?,
2234 merge_key: decode_merge_key(row.merge_key_json.clone())?,
2235 })
2236 })
2237 .collect::<Result<Vec<_>, StoreError>>()?;
2238 let selected_len = select_turn_work_claim_prefix(&candidates, boundary, max_batches);
2239 if selected_len == 0 {
2240 return Ok(TxOutcome::Commit(None));
2241 }
2242 let mut selected = candidate_rows;
2243 selected.truncate(selected_len);
2244 let mut selected_batches = candidate_batches;
2245 selected_batches.truncate(selected_len);
2246 let lease = QueuedWorkClaimLease::derive(&candidates[0], session_id, owner, now, generation);
2247 let liveness_json = encode_liveness(&owner.liveness)?;
2248 for row in &selected {
2249 let claimed = tx
2250 .execute(
2251 "UPDATE queued_work_batches
2252 SET claim_id = ?3,
2253 claim_owner_id = ?4,
2254 claim_owner_incarnation_id = ?5,
2255 claim_owner_liveness_json = ?6,
2256 claim_token = ?7,
2257 claim_fencing_token = claim_fencing_token + 1,
2258 claim_session_lease_generation = ?8
2259 WHERE session_id = ?1
2260 AND batch_id = ?2
2261 AND (
2262 claim_token IS NULL
2263 OR claim_session_lease_generation <> ?8
2264 )",
2265 params![
2266 session_id,
2267 row.batch_id,
2268 lease.claim_id,
2269 owner.owner_id.as_str(),
2270 owner.incarnation_id.as_str(),
2271 liveness_json.as_str(),
2272 lease.lease_token,
2273 lease.session_lease_generation as i64,
2274 ],
2275 )
2276 .map_err(sqlite_error)?;
2277 if claimed == 0 {
2278 return Ok(TxOutcome::Rollback(None));
2279 }
2280 }
2281 Ok(TxOutcome::Commit(Some(QueuedWorkClaim {
2282 session_id: session_id.to_string(),
2283 claim_id: lease.claim_id,
2284 owner: owner.clone(),
2285 lease_token: lease.lease_token,
2286 fencing_token: lease.fencing_token,
2287 session_lease_generation: lease.session_lease_generation,
2288 batches: selected_batches,
2289 })))
2290}
2291
2292#[allow(clippy::too_many_arguments)]
2293fn claim_pending_turn_inputs_sqlite_conn(
2294 tx: &Connection,
2295 now: u64,
2296 session_id: &str,
2297 session_execution_lease: &SessionExecutionLeaseFence,
2298 owner: &LeaseOwnerIdentity,
2299 max_inputs: usize,
2300 mode: lash_core::TurnInputClaimMode,
2301) -> Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> {
2302 if max_inputs == 0 {
2303 return Ok(TxOutcome::Commit(None));
2304 }
2305 let generation = session_execution_lease.fencing_token;
2306 let wanted_state = match &mode {
2307 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2308 lash_core::TurnInputState::PendingActive
2309 }
2310 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2311 };
2312 let candidate_rows = {
2313 let mut sql = "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2314 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2315 claim_owner_id, claim_owner_incarnation_id,
2316 claim_owner_liveness_json, claim_token, claim_session_lease_generation
2317 FROM pending_turn_inputs
2318 WHERE session_id = ? AND state = ?
2319 AND (
2320 claim_token IS NULL
2321 OR claim_session_lease_generation <> ?
2322 )"
2323 .to_string();
2324 let mut values: Vec<rusqlite::types::Value> = vec![
2325 session_id.to_string().into(),
2326 wanted_state.as_str().to_string().into(),
2327 (generation as i64).into(),
2328 ];
2329 if let lash_core::TurnInputClaimMode::ActiveTurn {
2330 turn_id,
2331 checkpoint,
2332 } = &mode
2333 {
2334 sql.push_str(
2335 " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2336 AND json_extract(ingress_json, '$.turn_id') = ?",
2337 );
2338 values.push(turn_id.clone().into());
2339 if *checkpoint == lash_core::CheckpointKind::AfterWork {
2340 sql.push_str(
2341 " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2342 );
2343 }
2344 }
2345 sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2346 values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2347 let mut stmt = tx.prepare(&sql).map_err(sqlite_error)?;
2348 let rows = stmt
2349 .query_map(
2350 rusqlite::params_from_iter(values.iter()),
2351 pending_turn_input_row_from_sql,
2352 )
2353 .map_err(sqlite_error)?;
2354 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2355 };
2356 let selected = candidate_rows
2357 .into_iter()
2358 .take(max_inputs)
2359 .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2360 .collect::<Result<Vec<_>, StoreError>>()?;
2361 let Some((head, _)) = selected.first() else {
2362 return Ok(TxOutcome::Commit(None));
2363 };
2364 let lease = TurnInputClaimLease::derive(head, session_id, owner, now, generation);
2365 let liveness_json = encode_liveness(&owner.liveness)?;
2366 let state_after_claim = match &mode {
2367 lash_core::TurnInputClaimMode::ActiveTurn { .. } => lash_core::TurnInputState::Accepted,
2368 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2369 };
2370 let mut inputs = Vec::new();
2371 for (row, mut input) in selected {
2372 let claimed = tx
2373 .execute(
2374 "UPDATE pending_turn_inputs
2375 SET state = ?3,
2376 claim_id = ?4,
2377 claim_owner_id = ?5,
2378 claim_owner_incarnation_id = ?6,
2379 claim_owner_liveness_json = ?7,
2380 claim_token = ?8,
2381 claim_fencing_token = claim_fencing_token + 1,
2382 claim_session_lease_generation = ?9
2383 WHERE session_id = ?1
2384 AND input_id = ?2
2385 AND (
2386 claim_token IS NULL
2387 OR claim_session_lease_generation <> ?9
2388 )",
2389 params![
2390 session_id,
2391 row.input_id,
2392 state_after_claim.as_str(),
2393 lease.claim_id,
2394 owner.owner_id.as_str(),
2395 owner.incarnation_id.as_str(),
2396 liveness_json.as_str(),
2397 lease.lease_token,
2398 lease.session_lease_generation as i64,
2399 ],
2400 )
2401 .map_err(sqlite_error)?;
2402 if claimed == 0 {
2403 return Ok(TxOutcome::Rollback(None));
2404 }
2405 input.state = state_after_claim;
2406 inputs.push(input);
2407 }
2408 Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2409 session_id: session_id.to_string(),
2410 claim_id: lease.claim_id,
2411 owner: owner.clone(),
2412 lease_token: lease.lease_token,
2413 fencing_token: lease.fencing_token,
2414 session_lease_generation: lease.session_lease_generation,
2415 mode,
2416 inputs,
2417 applications: Vec::new(),
2418 })))
2419}
2420
2421async fn claim_pending_turn_inputs_sqlite(
2422 conn: &SqliteConnection,
2423 now: u64,
2424 session_id: &str,
2425 session_execution_lease: &SessionExecutionLeaseFence,
2426 owner: &LeaseOwnerIdentity,
2427 max_inputs: usize,
2428 mode: lash_core::TurnInputClaimMode,
2429) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
2430 if max_inputs == 0 {
2431 return Ok(None);
2432 }
2433 let session_id = session_id.to_string();
2434 let session_execution_lease = session_execution_lease.clone();
2435 let owner = owner.clone();
2436 conn.write_flow(move |tx| {
2437 let outcome: Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> = (|| {
2438 ensure_session_execution_lease_conn(
2439 tx,
2440 &session_id,
2441 &session_execution_lease,
2442 now,
2443 )?;
2444 let generation = session_execution_lease.fencing_token;
2445 let wanted_state = match &mode {
2446 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2447 lash_core::TurnInputState::PendingActive
2448 }
2449 lash_core::TurnInputClaimMode::NextTurn => {
2450 lash_core::TurnInputState::DeferredNextTurn
2451 }
2452 };
2453 let candidate_rows = {
2454 let mut sql =
2455 "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2456 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2457 claim_owner_id, claim_owner_incarnation_id,
2458 claim_owner_liveness_json, claim_token, claim_session_lease_generation
2459 FROM pending_turn_inputs
2460 WHERE session_id = ? AND state = ?
2461 AND (
2462 claim_token IS NULL
2463 OR claim_session_lease_generation <> ?
2464 )"
2465 .to_string();
2466 let mut values: Vec<rusqlite::types::Value> = vec![
2467 session_id.clone().into(),
2468 wanted_state.as_str().to_string().into(),
2469 (generation as i64).into(),
2470 ];
2471 if let lash_core::TurnInputClaimMode::ActiveTurn {
2472 turn_id,
2473 checkpoint,
2474 } = &mode
2475 {
2476 sql.push_str(
2477 " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2478 AND json_extract(ingress_json, '$.turn_id') = ?",
2479 );
2480 values.push(turn_id.clone().into());
2481 if *checkpoint == lash_core::CheckpointKind::AfterWork {
2482 sql.push_str(
2483 " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2484 );
2485 }
2486 }
2487 sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2488 values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2489 let mut stmt = tx
2490 .prepare(&sql)
2491 .map_err(sqlite_error)?;
2492 let rows = stmt
2493 .query_map(
2494 rusqlite::params_from_iter(values.iter()),
2495 pending_turn_input_row_from_sql,
2496 )
2497 .map_err(sqlite_error)?;
2498 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2499 };
2500 let selected = candidate_rows
2501 .into_iter()
2502 .take(max_inputs)
2503 .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2504 .collect::<Result<Vec<_>, StoreError>>()?;
2505 let Some((head, _)) = selected.first() else {
2506 return Ok(TxOutcome::Commit(None));
2507 };
2508 let lease = TurnInputClaimLease::derive(head, &session_id, &owner, now, generation);
2509 let liveness_json = encode_liveness(&owner.liveness)?;
2510 let state_after_claim = match &mode {
2511 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2512 lash_core::TurnInputState::Accepted
2513 }
2514 lash_core::TurnInputClaimMode::NextTurn => {
2515 lash_core::TurnInputState::DeferredNextTurn
2516 }
2517 };
2518 let mut inputs = Vec::new();
2519 for (row, mut input) in selected {
2520 let claimed = tx
2521 .execute(
2522 "UPDATE pending_turn_inputs
2523 SET state = ?3,
2524 claim_id = ?4,
2525 claim_owner_id = ?5,
2526 claim_owner_incarnation_id = ?6,
2527 claim_owner_liveness_json = ?7,
2528 claim_token = ?8,
2529 claim_fencing_token = claim_fencing_token + 1,
2530 claim_session_lease_generation = ?9
2531 WHERE session_id = ?1
2532 AND input_id = ?2
2533 AND (
2534 claim_token IS NULL
2535 OR claim_session_lease_generation <> ?9
2536 )",
2537 params![
2538 session_id,
2539 row.input_id,
2540 state_after_claim.as_str(),
2541 lease.claim_id,
2542 owner.owner_id.as_str(),
2543 owner.incarnation_id.as_str(),
2544 liveness_json.as_str(),
2545 lease.lease_token,
2546 lease.session_lease_generation as i64,
2547 ],
2548 )
2549 .map_err(sqlite_error)?;
2550 if claimed == 0 {
2551 return Ok(TxOutcome::Rollback(None));
2552 }
2553 input.state = state_after_claim;
2554 inputs.push(input);
2555 }
2556 Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2557 session_id: session_id.clone(),
2558 claim_id: lease.claim_id,
2559 owner: owner.clone(),
2560 lease_token: lease.lease_token,
2561 fencing_token: lease.fencing_token,
2562 session_lease_generation: lease.session_lease_generation,
2563 mode,
2564 inputs,
2565 applications: Vec::new(),
2566 })))
2567 })(
2568 );
2569 match outcome {
2570 Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
2571 Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
2572 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
2573 }
2574 })
2575 .await
2576 .map_err(sqlite_error)?
2577}
2578
2579struct SessionExecutionLeaseRow {
2580 owner: Option<LeaseOwnerIdentity>,
2581 lease_token: Option<String>,
2582 fencing_token: u64,
2583 claimed_at_ms: u64,
2584 expires_at_ms: u64,
2585}
2586
2587fn load_session_execution_lease_row_conn(
2588 conn: &Connection,
2589 session_id: &str,
2590) -> Result<Option<SessionExecutionLeaseRow>, StoreError> {
2591 let row = conn
2592 .query_row(
2593 "SELECT lease_owner_id, lease_token, lease_fencing_token,
2594 lease_claimed_at_ms, lease_expires_at_ms,
2595 lease_owner_incarnation_id, lease_owner_liveness_json
2596 FROM session_execution_leases
2597 WHERE session_id = ?1",
2598 params![session_id],
2599 |row| {
2600 let owner_id: Option<String> = row.get(0)?;
2601 let incarnation_id: Option<String> = row.get(5)?;
2602 let liveness_json: Option<String> = row.get(6)?;
2603 Ok(SessionExecutionLeaseRow {
2604 owner: lease_owner_from_columns(owner_id, incarnation_id, liveness_json),
2605 lease_token: row.get(1)?,
2606 fencing_token: row.get::<_, i64>(2)? as u64,
2607 claimed_at_ms: row.get::<_, i64>(3)? as u64,
2608 expires_at_ms: row.get::<_, i64>(4)? as u64,
2609 })
2610 },
2611 )
2612 .optional()
2613 .map_err(sqlite_error)?;
2614 Ok(row)
2615}
2616
2617fn lease_owner_from_columns(
2618 owner_id: Option<String>,
2619 incarnation_id: Option<String>,
2620 liveness_json: Option<String>,
2621) -> Option<LeaseOwnerIdentity> {
2622 owner_id.map(|owner_id| LeaseOwnerIdentity {
2623 incarnation_id: incarnation_id.unwrap_or_else(|| owner_id.clone()),
2624 owner_id,
2625 liveness: liveness_json
2626 .as_deref()
2627 .and_then(|json| serde_json::from_str(json).ok())
2628 .unwrap_or(LeaseOwnerLiveness::Opaque),
2629 })
2630}
2631
2632fn encode_liveness(liveness: &LeaseOwnerLiveness) -> Result<String, StoreError> {
2633 serde_json::to_string(liveness)
2634 .map_err(|err| StoreError::Backend(format!("failed to encode lease liveness: {err}")))
2635}
2636
2637fn row_to_session_execution_lease(
2638 session_id: &str,
2639 row: SessionExecutionLeaseRow,
2640) -> Result<SessionExecutionLease, StoreError> {
2641 Ok(SessionExecutionLease {
2642 session_id: session_id.to_string(),
2643 owner: row
2644 .owner
2645 .ok_or_else(|| StoreError::Backend("live session lease missing owner".to_string()))?,
2646 lease_token: row.lease_token.ok_or_else(|| {
2647 StoreError::Backend("live session lease missing lease token".to_string())
2648 })?,
2649 fencing_token: row.fencing_token,
2650 claimed_at_epoch_ms: row.claimed_at_ms,
2651 expires_at_epoch_ms: row.expires_at_ms,
2652 })
2653}
2654
2655fn acquire_session_execution_lease_conn(
2656 conn: &Connection,
2657 session_id: &str,
2658 owner: &LeaseOwnerIdentity,
2659 previous_fencing_token: u64,
2660 now: u64,
2661 lease_ttl_ms: u64,
2662) -> Result<SessionExecutionLease, StoreError> {
2663 let fencing_token = previous_fencing_token.saturating_add(1);
2664 let lease_token = format!(
2665 "{}:{}:{}:{now}:{fencing_token}",
2666 session_id, owner.owner_id, owner.incarnation_id
2667 );
2668 let expires_at = now.saturating_add(lease_ttl_ms);
2669 let liveness_json = encode_liveness(&owner.liveness)?;
2670 conn.execute(
2671 "INSERT INTO session_execution_leases (
2672 session_id, lease_owner_id, lease_owner_incarnation_id, lease_owner_liveness_json,
2673 lease_token, lease_fencing_token, lease_claimed_at_ms, lease_expires_at_ms
2674 )
2675 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
2676 ON CONFLICT(session_id) DO UPDATE SET
2677 lease_owner_id = excluded.lease_owner_id,
2678 lease_owner_incarnation_id = excluded.lease_owner_incarnation_id,
2679 lease_owner_liveness_json = excluded.lease_owner_liveness_json,
2680 lease_token = excluded.lease_token,
2681 lease_fencing_token = excluded.lease_fencing_token,
2682 lease_claimed_at_ms = excluded.lease_claimed_at_ms,
2683 lease_expires_at_ms = excluded.lease_expires_at_ms",
2684 params![
2685 session_id,
2686 owner.owner_id,
2687 owner.incarnation_id,
2688 liveness_json,
2689 lease_token,
2690 fencing_token as i64,
2691 now as i64,
2692 expires_at as i64
2693 ],
2694 )
2695 .map_err(sqlite_error)?;
2696 Ok(SessionExecutionLease {
2697 session_id: session_id.to_string(),
2698 owner: owner.clone(),
2699 lease_token,
2700 fencing_token,
2701 claimed_at_epoch_ms: now,
2702 expires_at_epoch_ms: expires_at,
2703 })
2704}
2705
2706fn ensure_session_execution_lease_conn(
2707 conn: &Connection,
2708 session_id: &str,
2709 fence: &SessionExecutionLeaseFence,
2710 now: u64,
2711) -> Result<(), StoreError> {
2712 if fence.session_id != session_id {
2713 return Err(StoreError::SessionExecutionLeaseExpired {
2714 session_id: session_id.to_string(),
2715 });
2716 }
2717 let current = load_session_execution_lease_row_conn(conn, session_id)?;
2718 let Some(current) = current else {
2719 return Err(StoreError::SessionExecutionLeaseExpired {
2720 session_id: session_id.to_string(),
2721 });
2722 };
2723 if current
2724 .owner
2725 .as_ref()
2726 .is_some_and(|owner| owner.same_incarnation(&fence.owner))
2727 && current.lease_token.as_deref() == Some(fence.lease_token.as_str())
2728 && current.fencing_token == fence.fencing_token
2729 && current.expires_at_ms > now
2730 {
2731 Ok(())
2732 } else {
2733 Err(StoreError::SessionExecutionLeaseExpired {
2734 session_id: session_id.to_string(),
2735 })
2736 }
2737}
2738
2739fn release_session_execution_lease_conn(
2740 conn: &Connection,
2741 completion: &SessionExecutionLeaseCompletion,
2742) -> Result<(), StoreError> {
2743 conn.execute(
2744 "UPDATE session_execution_leases
2745 SET lease_owner_id = NULL,
2746 lease_owner_incarnation_id = NULL,
2747 lease_owner_liveness_json = NULL,
2748 lease_token = NULL,
2749 lease_claimed_at_ms = 0,
2750 lease_expires_at_ms = 0
2751 WHERE session_id = ?1
2752 AND lease_owner_id = ?2
2753 AND lease_owner_incarnation_id = ?3
2754 AND lease_token = ?4
2755 AND lease_fencing_token = ?5",
2756 params![
2757 completion.session_id,
2758 completion.owner.owner_id,
2759 completion.owner.incarnation_id,
2760 completion.lease_token,
2761 completion.fencing_token as i64
2762 ],
2763 )
2764 .map_err(sqlite_error)?;
2765 Ok(())
2766}