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