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 let mut enqueued_queue_batches = Vec::new();
476 for (index, batch) in commit.enqueued_queue_batches.iter().enumerate() {
477 if batch.session_id != commit.session_id {
478 return Err(StoreError::SessionBindingMismatch {
479 bound_session_id: commit.session_id.clone(),
480 attempted_session_id: batch.session_id.clone(),
481 });
482 }
483 enqueued_queue_batches.push(enqueue_queued_work_conn(
484 tx,
485 batch,
486 now,
487 enqueue_nonce_start.saturating_add(index as u64),
488 )?);
489 }
490 let result = RuntimeCommitResult {
491 head_revision: next_revision,
492 checkpoint_ref: stored_checkpoint.checkpoint_ref,
493 manifest: stored_checkpoint.manifest,
494 enqueued_queue_batches,
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 cancel_pending_turn_inputs(
1711 &self,
1712 session_id: &str,
1713 targets: &[lash_core::PendingTurnInputCancelTarget],
1714 ) -> Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> {
1715 let session_id = session_id.to_string();
1716 let targets = targets.to_vec();
1717 let now = self.clock.timestamp_ms();
1718 self.conn
1719 .write_flow(move |tx| {
1720 let outcome: Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> =
1721 (|| {
1722 let mut results = Vec::with_capacity(targets.len());
1723 for target in targets {
1724 let outcome = match load_pending_turn_input_row_by_target_conn(
1725 tx,
1726 &session_id,
1727 &target,
1728 )? {
1729 Some(row) => cancel_pending_turn_input_row_conn(tx, row, now)?,
1730 None => lash_core::PendingTurnInputCancelOutcome::NotFound,
1731 };
1732 results
1733 .push(lash_core::PendingTurnInputCancelResult { target, outcome });
1734 }
1735 Ok(results)
1736 })();
1737 match outcome {
1738 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1739 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1740 }
1741 })
1742 .await
1743 .map_err(sqlite_error)?
1744 }
1745
1746 async fn cancel_pending_turn_input_suffix(
1747 &self,
1748 session_id: &str,
1749 anchor: &lash_core::PendingTurnInputCancelTarget,
1750 ) -> Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> {
1751 let session_id = session_id.to_string();
1752 let anchor = anchor.clone();
1753 let now = self.clock.timestamp_ms();
1754 self.conn
1755 .write_flow(move |tx| {
1756 let outcome: Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> =
1757 (|| {
1758 let Some(anchor_row) =
1759 load_pending_turn_input_row_by_target_conn(tx, &session_id, &anchor)?
1760 else {
1761 return Ok(
1762 lash_core::PendingTurnInputSuffixCancelOutcome::AnchorNotFound {
1763 anchor,
1764 },
1765 );
1766 };
1767 let rows = {
1768 let mut stmt = tx
1769 .prepare(
1770 "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
1771 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
1772 claim_owner_id, claim_owner_incarnation_id,
1773 claim_owner_liveness_json, claim_token, claim_session_lease_generation
1774 FROM pending_turn_inputs
1775 WHERE session_id = ?1 AND enqueue_seq >= ?2
1776 ORDER BY enqueue_seq ASC",
1777 )
1778 .map_err(sqlite_error)?;
1779 let rows = stmt
1780 .query_map(
1781 params![session_id, anchor_row.enqueue_seq as i64],
1782 pending_turn_input_row_from_sql,
1783 )
1784 .map_err(sqlite_error)?;
1785 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1786 };
1787 let mut outcomes = Vec::with_capacity(rows.len());
1788 for row in rows {
1789 outcomes.push(cancel_pending_turn_input_row_conn(tx, row, now)?);
1790 }
1791 Ok(lash_core::PendingTurnInputSuffixCancelOutcome::Outcomes {
1792 anchor,
1793 outcomes,
1794 })
1795 })();
1796 match outcome {
1797 Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1798 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1799 }
1800 })
1801 .await
1802 .map_err(sqlite_error)?
1803 }
1804
1805 async fn claim_active_turn_inputs(
1806 &self,
1807 session_id: &str,
1808 session_execution_lease: &SessionExecutionLeaseFence,
1809 owner: &LeaseOwnerIdentity,
1810 turn_id: &str,
1811 checkpoint: lash_core::CheckpointKind,
1812 max_inputs: usize,
1813 ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1814 claim_pending_turn_inputs_sqlite(
1815 &self.conn,
1816 self.clock.timestamp_ms(),
1817 session_id,
1818 session_execution_lease,
1819 owner,
1820 max_inputs,
1821 lash_core::TurnInputClaimMode::ActiveTurn {
1822 turn_id: turn_id.to_string(),
1823 checkpoint,
1824 },
1825 )
1826 .await
1827 }
1828
1829 async fn claim_next_turn_inputs(
1830 &self,
1831 session_id: &str,
1832 session_execution_lease: &SessionExecutionLeaseFence,
1833 owner: &LeaseOwnerIdentity,
1834 max_inputs: usize,
1835 ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1836 claim_pending_turn_inputs_sqlite(
1837 &self.conn,
1838 self.clock.timestamp_ms(),
1839 session_id,
1840 session_execution_lease,
1841 owner,
1842 max_inputs,
1843 lash_core::TurnInputClaimMode::NextTurn,
1844 )
1845 .await
1846 }
1847
1848 async fn abandon_turn_input_claim(
1849 &self,
1850 claim: &lash_core::TurnInputClaim,
1851 ) -> Result<(), StoreError> {
1852 let session_id = claim.session_id.clone();
1853 let claim_id = claim.claim_id.clone();
1854 let lease_token = claim.lease_token.clone();
1855 let restored_state = match claim.mode {
1856 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
1857 lash_core::TurnInputState::PendingActive
1858 }
1859 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
1860 };
1861 self.conn
1862 .write(move |tx| {
1863 tx.execute(
1864 "UPDATE pending_turn_inputs
1865 SET state = CASE
1866 WHEN state = ?4 THEN ?5
1867 ELSE state
1868 END,
1869 claim_id = NULL,
1870 claim_owner_id = NULL,
1871 claim_owner_incarnation_id = NULL,
1872 claim_owner_liveness_json = NULL,
1873 claim_token = NULL,
1874 claim_session_lease_generation = 0
1875 WHERE session_id = ?1 AND claim_id = ?2 AND claim_token = ?3",
1876 params![
1877 session_id,
1878 claim_id,
1879 lease_token,
1880 lash_core::TurnInputState::Accepted.as_str(),
1881 restored_state.as_str(),
1882 ],
1883 )
1884 })
1885 .await
1886 .map_err(sqlite_error)?;
1887 Ok(())
1888 }
1889
1890 async fn abandon_turn_input_claims(
1891 &self,
1892 claims: &[lash_core::TurnInputClaim],
1893 ) -> Result<(), StoreError> {
1894 if claims.is_empty() {
1895 return Ok(());
1896 }
1897 let mut sql = "UPDATE pending_turn_inputs
1898 SET state = CASE
1899 WHEN state = 'accepted' THEN 'pending_active'
1900 ELSE state
1901 END,
1902 claim_id = NULL,
1903 claim_owner_id = NULL,
1904 claim_owner_incarnation_id = NULL,
1905 claim_owner_liveness_json = NULL,
1906 claim_token = NULL,
1907 claim_session_lease_generation = 0
1908 WHERE (session_id, claim_id, claim_token) IN ("
1909 .to_string();
1910 let mut values: Vec<rusqlite::types::Value> = Vec::with_capacity(claims.len() * 3);
1911 for (index, claim) in claims.iter().enumerate() {
1912 if index > 0 {
1913 sql.push_str(", ");
1914 }
1915 sql.push_str("(?, ?, ?)");
1916 values.push(claim.session_id.clone().into());
1917 values.push(claim.claim_id.clone().into());
1918 values.push(claim.lease_token.clone().into());
1919 }
1920 sql.push(')');
1921 self.conn
1922 .write(move |tx| tx.execute(&sql, rusqlite::params_from_iter(values.iter())))
1923 .await
1924 .map_err(sqlite_error)?;
1925 Ok(())
1926 }
1927}
1928
1929#[async_trait::async_trait]
1930impl StoreMaintenance for Store {
1931 async fn tombstone_nodes(&self, ids: &[String]) -> Result<(), StoreError> {
1932 if ids.is_empty() {
1933 return Ok(());
1934 }
1935 let ids = ids.to_vec();
1936 self.conn
1937 .write(move |tx| {
1938 let mut stmt =
1939 tx.prepare("UPDATE graph_nodes SET tombstoned = 1 WHERE node_id = ?1")?;
1940 for id in &ids {
1941 stmt.execute(params![id])?;
1942 }
1943 Ok(())
1944 })
1945 .await
1946 .map_err(sqlite_error)
1947 }
1948
1949 async fn vacuum(&self) -> Result<VacuumReport, StoreError> {
1950 let (removed_node_count, removed_pending_turn_input_tombstone_count) = self
1951 .conn
1952 .write(move |tx| {
1953 let removed_node_count =
1954 tx.execute("DELETE FROM graph_nodes WHERE tombstoned = 1", [])?;
1955 let removed_pending_turn_input_tombstone_count = tx.execute(
1956 "DELETE FROM pending_turn_inputs
1957 WHERE state IN (?1, ?2)",
1958 params![
1959 lash_core::TurnInputState::Cancelled.as_str(),
1960 lash_core::TurnInputState::Completed.as_str()
1961 ],
1962 )?;
1963 Ok((
1964 removed_node_count,
1965 removed_pending_turn_input_tombstone_count,
1966 ))
1967 })
1968 .await
1969 .map_err(sqlite_error)?;
1970 Ok(VacuumReport {
1971 removed_node_count,
1972 removed_pending_turn_input_tombstone_count,
1973 })
1974 }
1975
1976 async fn gc_unreachable(&self) -> Result<GcReport, StoreError> {
1977 Ok(Store::gc_unreachable(self).await)
1978 }
1979}
1980
1981fn derive_pending_turn_input_id(
1982 session_id: &str,
1983 source_key: Option<&str>,
1984 now_epoch_ms: u64,
1985 nonce: u64,
1986) -> String {
1987 format!(
1988 "ti:{:x}",
1989 Sha256::digest(format!("{session_id}:{source_key:?}:{now_epoch_ms}:{nonce}").as_bytes())
1990 )
1991}
1992
1993fn cancel_pending_turn_input_row_conn(
1994 conn: &Connection,
1995 row: PendingTurnInputRow,
1996 now_epoch_ms: u64,
1997) -> Result<lash_core::PendingTurnInputCancelOutcome, StoreError> {
1998 let mut input = pending_turn_input_from_row(row.clone())?;
1999 match input.state {
2000 lash_core::TurnInputState::Cancelled => Ok(
2001 lash_core::PendingTurnInputCancelOutcome::AlreadyCancelled(input),
2002 ),
2003 lash_core::TurnInputState::Completed => Ok(
2004 lash_core::PendingTurnInputCancelOutcome::AlreadyCompleted(input),
2005 ),
2006 lash_core::TurnInputState::Accepted => {
2007 Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2008 claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2009 input,
2010 })
2011 }
2012 lash_core::TurnInputState::PendingActive | lash_core::TurnInputState::DeferredNextTurn => {
2013 let live_claim = row.claim_token.is_some()
2016 && load_session_execution_lease_row_conn(conn, &row.session_id)?.is_some_and(
2017 |lease| {
2018 lease.lease_token.is_some()
2019 && lease.expires_at_ms > now_epoch_ms
2020 && lease.fencing_token == row.claim_session_lease_generation
2021 },
2022 );
2023 if live_claim {
2024 return Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2025 claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2026 input,
2027 });
2028 }
2029 conn.execute(
2030 "UPDATE pending_turn_inputs
2031 SET state = ?3,
2032 claim_id = NULL,
2033 claim_owner_id = NULL,
2034 claim_owner_incarnation_id = NULL,
2035 claim_owner_liveness_json = NULL,
2036 claim_token = NULL,
2037 claim_session_lease_generation = 0
2038 WHERE session_id = ?1 AND input_id = ?2",
2039 params![
2040 row.session_id,
2041 row.input_id,
2042 lash_core::TurnInputState::Cancelled.as_str(),
2043 ],
2044 )
2045 .map_err(sqlite_error)?;
2046 input.state = lash_core::TurnInputState::Cancelled;
2047 Ok(lash_core::PendingTurnInputCancelOutcome::Cancelled(input))
2048 }
2049 }
2050}
2051
2052#[allow(clippy::too_many_arguments)]
2053async fn checkpoint_work_pending_sqlite(
2054 conn: &SqliteConnection,
2055 now: u64,
2056 session_id: &str,
2057 generation: u64,
2058 turn_id: &str,
2059 checkpoint: lash_core::CheckpointKind,
2060 max_inputs: usize,
2061 max_batches: usize,
2062) -> Result<bool, StoreError> {
2063 if max_inputs == 0 && max_batches == 0 {
2064 return Ok(false);
2065 }
2066 let session_id = session_id.to_string();
2067 let turn_id = turn_id.to_string();
2068 conn.call(move |conn| {
2069 let outcome: Result<bool, StoreError> = (|| {
2070 let head_candidate = sqlite_queued_work_head_candidate_cte(
2071 QueuedWorkClaimBoundary::ActiveTurnCheckpoint,
2072 );
2073 let sql = format!(
2074 "WITH {head_candidate}
2075 SELECT (
2076 ?7 > 0 AND EXISTS (
2077 SELECT 1
2078 FROM pending_turn_inputs
2079 WHERE session_id = ?1
2080 AND state = ?4
2081 AND (claim_token IS NULL OR claim_session_lease_generation <> ?3)
2082 AND json_extract(ingress_json, '$.scope') = 'active_turn'
2083 AND json_extract(ingress_json, '$.turn_id') = ?5
2084 AND (
2085 ?6 = 'before_completion'
2086 OR COALESCE(
2087 json_extract(ingress_json, '$.min_boundary'),
2088 'after_work'
2089 ) = 'after_work'
2090 )
2091 LIMIT 1
2092 )
2093 ) OR (
2094 ?8 > 0 AND EXISTS (
2095 SELECT 1
2096 FROM queued_work_head_candidate AS head
2097 JOIN queued_work_items AS item
2098 ON item.batch_id = head.head_batch_id
2099 WHERE json_extract(item.payload_json, '$.type') <> 'session_command'
2100 LIMIT 1
2101 )
2102 )"
2103 );
2104 let pending: i64 = conn
2105 .query_row(
2106 &sql,
2107 params![
2108 session_id,
2109 now as i64,
2110 generation as i64,
2111 lash_core::TurnInputState::PendingActive.as_str(),
2112 turn_id,
2113 match checkpoint {
2114 lash_core::CheckpointKind::AfterWork => "after_work",
2115 lash_core::CheckpointKind::BeforeCompletion => "before_completion",
2116 },
2117 max_inputs as i64,
2118 max_batches as i64,
2119 ],
2120 |row| row.get(0),
2121 )
2122 .map_err(sqlite_error)?;
2123 Ok(pending != 0)
2124 })();
2125 Ok(outcome)
2126 })
2127 .await
2128 .map_err(sqlite_error)?
2129}
2130
2131#[allow(clippy::too_many_arguments)]
2132fn claim_ready_queued_work_sqlite_conn(
2133 tx: &Connection,
2134 now: u64,
2135 session_id: &str,
2136 session_execution_lease: &SessionExecutionLeaseFence,
2137 owner: &LeaseOwnerIdentity,
2138 boundary: QueuedWorkClaimBoundary,
2139 max_batches: usize,
2140) -> Result<TxOutcome<Option<QueuedWorkClaim>>, StoreError> {
2141 if max_batches == 0 {
2142 return Ok(TxOutcome::Commit(None));
2143 }
2144 let generation = session_execution_lease.fencing_token;
2145 let candidate_rows = {
2146 let mut stmt = tx
2147 .prepare(&sqlite_queued_work_claim_candidates_sql(boundary))
2148 .map_err(sqlite_error)?;
2149 let rows = stmt
2150 .query_map(
2151 params![
2152 session_id,
2153 now as i64,
2154 generation as i64,
2155 claim_scan_limit(max_batches)
2156 ],
2157 queued_batch_row_from_sql,
2158 )
2159 .map_err(sqlite_error)?;
2160 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2161 };
2162 let candidate_rows = candidate_rows
2163 .into_iter()
2164 .filter(|row| row.claim_token.is_none() || row.claim_session_lease_generation != generation)
2165 .collect::<Vec<_>>();
2166 let candidate_batches = candidate_rows
2167 .iter()
2168 .map(|row| queued_work_batch_from_conn(tx, row.clone()))
2169 .collect::<Result<Vec<_>, StoreError>>()?;
2170 let candidates = candidate_rows
2171 .iter()
2172 .zip(candidate_batches.iter())
2173 .map(|(row, batch)| {
2174 Ok(ClaimCandidate {
2175 enqueue_seq: row.enqueue_seq,
2176 claim_fencing_token: row.claim_fencing_token,
2177 work_class: batch.work_class().ok_or_else(|| {
2178 StoreError::Backend(format!(
2179 "queued-work batch `{}` has mixed or empty payload classes",
2180 batch.batch_id
2181 ))
2182 })?,
2183 delivery_policy: decode_delivery_policy(row.delivery_policy.clone())?,
2184 slot_policy: decode_slot_policy(row.slot_policy.clone())?,
2185 merge_key: decode_merge_key(row.merge_key_json.clone())?,
2186 })
2187 })
2188 .collect::<Result<Vec<_>, StoreError>>()?;
2189 let selected_len = select_turn_work_claim_prefix(&candidates, boundary, max_batches);
2190 if selected_len == 0 {
2191 return Ok(TxOutcome::Commit(None));
2192 }
2193 let mut selected = candidate_rows;
2194 selected.truncate(selected_len);
2195 let mut selected_batches = candidate_batches;
2196 selected_batches.truncate(selected_len);
2197 let lease = QueuedWorkClaimLease::derive(&candidates[0], session_id, owner, now, generation);
2198 let liveness_json = encode_liveness(&owner.liveness)?;
2199 for row in &selected {
2200 let claimed = tx
2201 .execute(
2202 "UPDATE queued_work_batches
2203 SET claim_id = ?3,
2204 claim_owner_id = ?4,
2205 claim_owner_incarnation_id = ?5,
2206 claim_owner_liveness_json = ?6,
2207 claim_token = ?7,
2208 claim_fencing_token = claim_fencing_token + 1,
2209 claim_session_lease_generation = ?8
2210 WHERE session_id = ?1
2211 AND batch_id = ?2
2212 AND (
2213 claim_token IS NULL
2214 OR claim_session_lease_generation <> ?8
2215 )",
2216 params![
2217 session_id,
2218 row.batch_id,
2219 lease.claim_id,
2220 owner.owner_id.as_str(),
2221 owner.incarnation_id.as_str(),
2222 liveness_json.as_str(),
2223 lease.lease_token,
2224 lease.session_lease_generation as i64,
2225 ],
2226 )
2227 .map_err(sqlite_error)?;
2228 if claimed == 0 {
2229 return Ok(TxOutcome::Rollback(None));
2230 }
2231 }
2232 Ok(TxOutcome::Commit(Some(QueuedWorkClaim {
2233 session_id: session_id.to_string(),
2234 claim_id: lease.claim_id,
2235 owner: owner.clone(),
2236 lease_token: lease.lease_token,
2237 fencing_token: lease.fencing_token,
2238 session_lease_generation: lease.session_lease_generation,
2239 batches: selected_batches,
2240 })))
2241}
2242
2243#[allow(clippy::too_many_arguments)]
2244fn claim_pending_turn_inputs_sqlite_conn(
2245 tx: &Connection,
2246 now: u64,
2247 session_id: &str,
2248 session_execution_lease: &SessionExecutionLeaseFence,
2249 owner: &LeaseOwnerIdentity,
2250 max_inputs: usize,
2251 mode: lash_core::TurnInputClaimMode,
2252) -> Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> {
2253 if max_inputs == 0 {
2254 return Ok(TxOutcome::Commit(None));
2255 }
2256 let generation = session_execution_lease.fencing_token;
2257 let wanted_state = match &mode {
2258 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2259 lash_core::TurnInputState::PendingActive
2260 }
2261 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2262 };
2263 let candidate_rows = {
2264 let mut sql = "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2265 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2266 claim_owner_id, claim_owner_incarnation_id,
2267 claim_owner_liveness_json, claim_token, claim_session_lease_generation
2268 FROM pending_turn_inputs
2269 WHERE session_id = ? AND state = ?
2270 AND (
2271 claim_token IS NULL
2272 OR claim_session_lease_generation <> ?
2273 )"
2274 .to_string();
2275 let mut values: Vec<rusqlite::types::Value> = vec![
2276 session_id.to_string().into(),
2277 wanted_state.as_str().to_string().into(),
2278 (generation as i64).into(),
2279 ];
2280 if let lash_core::TurnInputClaimMode::ActiveTurn {
2281 turn_id,
2282 checkpoint,
2283 } = &mode
2284 {
2285 sql.push_str(
2286 " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2287 AND json_extract(ingress_json, '$.turn_id') = ?",
2288 );
2289 values.push(turn_id.clone().into());
2290 if *checkpoint == lash_core::CheckpointKind::AfterWork {
2291 sql.push_str(
2292 " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2293 );
2294 }
2295 }
2296 sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2297 values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2298 let mut stmt = tx.prepare(&sql).map_err(sqlite_error)?;
2299 let rows = stmt
2300 .query_map(
2301 rusqlite::params_from_iter(values.iter()),
2302 pending_turn_input_row_from_sql,
2303 )
2304 .map_err(sqlite_error)?;
2305 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2306 };
2307 let selected = candidate_rows
2308 .into_iter()
2309 .take(max_inputs)
2310 .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2311 .collect::<Result<Vec<_>, StoreError>>()?;
2312 let Some((head, _)) = selected.first() else {
2313 return Ok(TxOutcome::Commit(None));
2314 };
2315 let lease = TurnInputClaimLease::derive(head, session_id, owner, now, generation);
2316 let liveness_json = encode_liveness(&owner.liveness)?;
2317 let state_after_claim = match &mode {
2318 lash_core::TurnInputClaimMode::ActiveTurn { .. } => lash_core::TurnInputState::Accepted,
2319 lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2320 };
2321 let mut inputs = Vec::new();
2322 for (row, mut input) in selected {
2323 let claimed = tx
2324 .execute(
2325 "UPDATE pending_turn_inputs
2326 SET state = ?3,
2327 claim_id = ?4,
2328 claim_owner_id = ?5,
2329 claim_owner_incarnation_id = ?6,
2330 claim_owner_liveness_json = ?7,
2331 claim_token = ?8,
2332 claim_fencing_token = claim_fencing_token + 1,
2333 claim_session_lease_generation = ?9
2334 WHERE session_id = ?1
2335 AND input_id = ?2
2336 AND (
2337 claim_token IS NULL
2338 OR claim_session_lease_generation <> ?9
2339 )",
2340 params![
2341 session_id,
2342 row.input_id,
2343 state_after_claim.as_str(),
2344 lease.claim_id,
2345 owner.owner_id.as_str(),
2346 owner.incarnation_id.as_str(),
2347 liveness_json.as_str(),
2348 lease.lease_token,
2349 lease.session_lease_generation as i64,
2350 ],
2351 )
2352 .map_err(sqlite_error)?;
2353 if claimed == 0 {
2354 return Ok(TxOutcome::Rollback(None));
2355 }
2356 input.state = state_after_claim;
2357 inputs.push(input);
2358 }
2359 Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2360 session_id: session_id.to_string(),
2361 claim_id: lease.claim_id,
2362 owner: owner.clone(),
2363 lease_token: lease.lease_token,
2364 fencing_token: lease.fencing_token,
2365 session_lease_generation: lease.session_lease_generation,
2366 mode,
2367 inputs,
2368 })))
2369}
2370
2371async fn claim_pending_turn_inputs_sqlite(
2372 conn: &SqliteConnection,
2373 now: u64,
2374 session_id: &str,
2375 session_execution_lease: &SessionExecutionLeaseFence,
2376 owner: &LeaseOwnerIdentity,
2377 max_inputs: usize,
2378 mode: lash_core::TurnInputClaimMode,
2379) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
2380 if max_inputs == 0 {
2381 return Ok(None);
2382 }
2383 let session_id = session_id.to_string();
2384 let session_execution_lease = session_execution_lease.clone();
2385 let owner = owner.clone();
2386 conn.write_flow(move |tx| {
2387 let outcome: Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> = (|| {
2388 ensure_session_execution_lease_conn(
2389 tx,
2390 &session_id,
2391 &session_execution_lease,
2392 now,
2393 )?;
2394 let generation = session_execution_lease.fencing_token;
2395 let wanted_state = match &mode {
2396 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2397 lash_core::TurnInputState::PendingActive
2398 }
2399 lash_core::TurnInputClaimMode::NextTurn => {
2400 lash_core::TurnInputState::DeferredNextTurn
2401 }
2402 };
2403 let candidate_rows = {
2404 let mut sql =
2405 "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2406 state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2407 claim_owner_id, claim_owner_incarnation_id,
2408 claim_owner_liveness_json, claim_token, claim_session_lease_generation
2409 FROM pending_turn_inputs
2410 WHERE session_id = ? AND state = ?
2411 AND (
2412 claim_token IS NULL
2413 OR claim_session_lease_generation <> ?
2414 )"
2415 .to_string();
2416 let mut values: Vec<rusqlite::types::Value> = vec![
2417 session_id.clone().into(),
2418 wanted_state.as_str().to_string().into(),
2419 (generation as i64).into(),
2420 ];
2421 if let lash_core::TurnInputClaimMode::ActiveTurn {
2422 turn_id,
2423 checkpoint,
2424 } = &mode
2425 {
2426 sql.push_str(
2427 " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2428 AND json_extract(ingress_json, '$.turn_id') = ?",
2429 );
2430 values.push(turn_id.clone().into());
2431 if *checkpoint == lash_core::CheckpointKind::AfterWork {
2432 sql.push_str(
2433 " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2434 );
2435 }
2436 }
2437 sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2438 values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2439 let mut stmt = tx
2440 .prepare(&sql)
2441 .map_err(sqlite_error)?;
2442 let rows = stmt
2443 .query_map(
2444 rusqlite::params_from_iter(values.iter()),
2445 pending_turn_input_row_from_sql,
2446 )
2447 .map_err(sqlite_error)?;
2448 rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2449 };
2450 let selected = candidate_rows
2451 .into_iter()
2452 .take(max_inputs)
2453 .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2454 .collect::<Result<Vec<_>, StoreError>>()?;
2455 let Some((head, _)) = selected.first() else {
2456 return Ok(TxOutcome::Commit(None));
2457 };
2458 let lease = TurnInputClaimLease::derive(head, &session_id, &owner, now, generation);
2459 let liveness_json = encode_liveness(&owner.liveness)?;
2460 let state_after_claim = match &mode {
2461 lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2462 lash_core::TurnInputState::Accepted
2463 }
2464 lash_core::TurnInputClaimMode::NextTurn => {
2465 lash_core::TurnInputState::DeferredNextTurn
2466 }
2467 };
2468 let mut inputs = Vec::new();
2469 for (row, mut input) in selected {
2470 let claimed = tx
2471 .execute(
2472 "UPDATE pending_turn_inputs
2473 SET state = ?3,
2474 claim_id = ?4,
2475 claim_owner_id = ?5,
2476 claim_owner_incarnation_id = ?6,
2477 claim_owner_liveness_json = ?7,
2478 claim_token = ?8,
2479 claim_fencing_token = claim_fencing_token + 1,
2480 claim_session_lease_generation = ?9
2481 WHERE session_id = ?1
2482 AND input_id = ?2
2483 AND (
2484 claim_token IS NULL
2485 OR claim_session_lease_generation <> ?9
2486 )",
2487 params![
2488 session_id,
2489 row.input_id,
2490 state_after_claim.as_str(),
2491 lease.claim_id,
2492 owner.owner_id.as_str(),
2493 owner.incarnation_id.as_str(),
2494 liveness_json.as_str(),
2495 lease.lease_token,
2496 lease.session_lease_generation as i64,
2497 ],
2498 )
2499 .map_err(sqlite_error)?;
2500 if claimed == 0 {
2501 return Ok(TxOutcome::Rollback(None));
2502 }
2503 input.state = state_after_claim;
2504 inputs.push(input);
2505 }
2506 Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2507 session_id: session_id.clone(),
2508 claim_id: lease.claim_id,
2509 owner: owner.clone(),
2510 lease_token: lease.lease_token,
2511 fencing_token: lease.fencing_token,
2512 session_lease_generation: lease.session_lease_generation,
2513 mode,
2514 inputs,
2515 })))
2516 })(
2517 );
2518 match outcome {
2519 Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
2520 Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
2521 Err(err) => Ok(TxOutcome::Rollback(Err(err))),
2522 }
2523 })
2524 .await
2525 .map_err(sqlite_error)?
2526}
2527
2528struct SessionExecutionLeaseRow {
2529 owner: Option<LeaseOwnerIdentity>,
2530 lease_token: Option<String>,
2531 fencing_token: u64,
2532 claimed_at_ms: u64,
2533 expires_at_ms: u64,
2534}
2535
2536fn load_session_execution_lease_row_conn(
2537 conn: &Connection,
2538 session_id: &str,
2539) -> Result<Option<SessionExecutionLeaseRow>, StoreError> {
2540 let row = conn
2541 .query_row(
2542 "SELECT lease_owner_id, lease_token, lease_fencing_token,
2543 lease_claimed_at_ms, lease_expires_at_ms,
2544 lease_owner_incarnation_id, lease_owner_liveness_json
2545 FROM session_execution_leases
2546 WHERE session_id = ?1",
2547 params![session_id],
2548 |row| {
2549 let owner_id: Option<String> = row.get(0)?;
2550 let incarnation_id: Option<String> = row.get(5)?;
2551 let liveness_json: Option<String> = row.get(6)?;
2552 Ok(SessionExecutionLeaseRow {
2553 owner: lease_owner_from_columns(owner_id, incarnation_id, liveness_json),
2554 lease_token: row.get(1)?,
2555 fencing_token: row.get::<_, i64>(2)? as u64,
2556 claimed_at_ms: row.get::<_, i64>(3)? as u64,
2557 expires_at_ms: row.get::<_, i64>(4)? as u64,
2558 })
2559 },
2560 )
2561 .optional()
2562 .map_err(sqlite_error)?;
2563 Ok(row)
2564}
2565
2566fn lease_owner_from_columns(
2567 owner_id: Option<String>,
2568 incarnation_id: Option<String>,
2569 liveness_json: Option<String>,
2570) -> Option<LeaseOwnerIdentity> {
2571 owner_id.map(|owner_id| LeaseOwnerIdentity {
2572 incarnation_id: incarnation_id.unwrap_or_else(|| owner_id.clone()),
2573 owner_id,
2574 liveness: liveness_json
2575 .as_deref()
2576 .and_then(|json| serde_json::from_str(json).ok())
2577 .unwrap_or(LeaseOwnerLiveness::Opaque),
2578 })
2579}
2580
2581fn encode_liveness(liveness: &LeaseOwnerLiveness) -> Result<String, StoreError> {
2582 serde_json::to_string(liveness)
2583 .map_err(|err| StoreError::Backend(format!("failed to encode lease liveness: {err}")))
2584}
2585
2586fn row_to_session_execution_lease(
2587 session_id: &str,
2588 row: SessionExecutionLeaseRow,
2589) -> Result<SessionExecutionLease, StoreError> {
2590 Ok(SessionExecutionLease {
2591 session_id: session_id.to_string(),
2592 owner: row
2593 .owner
2594 .ok_or_else(|| StoreError::Backend("live session lease missing owner".to_string()))?,
2595 lease_token: row.lease_token.ok_or_else(|| {
2596 StoreError::Backend("live session lease missing lease token".to_string())
2597 })?,
2598 fencing_token: row.fencing_token,
2599 claimed_at_epoch_ms: row.claimed_at_ms,
2600 expires_at_epoch_ms: row.expires_at_ms,
2601 })
2602}
2603
2604fn acquire_session_execution_lease_conn(
2605 conn: &Connection,
2606 session_id: &str,
2607 owner: &LeaseOwnerIdentity,
2608 previous_fencing_token: u64,
2609 now: u64,
2610 lease_ttl_ms: u64,
2611) -> Result<SessionExecutionLease, StoreError> {
2612 let fencing_token = previous_fencing_token.saturating_add(1);
2613 let lease_token = format!(
2614 "{}:{}:{}:{now}:{fencing_token}",
2615 session_id, owner.owner_id, owner.incarnation_id
2616 );
2617 let expires_at = now.saturating_add(lease_ttl_ms);
2618 let liveness_json = encode_liveness(&owner.liveness)?;
2619 conn.execute(
2620 "INSERT INTO session_execution_leases (
2621 session_id, lease_owner_id, lease_owner_incarnation_id, lease_owner_liveness_json,
2622 lease_token, lease_fencing_token, lease_claimed_at_ms, lease_expires_at_ms
2623 )
2624 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
2625 ON CONFLICT(session_id) DO UPDATE SET
2626 lease_owner_id = excluded.lease_owner_id,
2627 lease_owner_incarnation_id = excluded.lease_owner_incarnation_id,
2628 lease_owner_liveness_json = excluded.lease_owner_liveness_json,
2629 lease_token = excluded.lease_token,
2630 lease_fencing_token = excluded.lease_fencing_token,
2631 lease_claimed_at_ms = excluded.lease_claimed_at_ms,
2632 lease_expires_at_ms = excluded.lease_expires_at_ms",
2633 params![
2634 session_id,
2635 owner.owner_id,
2636 owner.incarnation_id,
2637 liveness_json,
2638 lease_token,
2639 fencing_token as i64,
2640 now as i64,
2641 expires_at as i64
2642 ],
2643 )
2644 .map_err(sqlite_error)?;
2645 Ok(SessionExecutionLease {
2646 session_id: session_id.to_string(),
2647 owner: owner.clone(),
2648 lease_token,
2649 fencing_token,
2650 claimed_at_epoch_ms: now,
2651 expires_at_epoch_ms: expires_at,
2652 })
2653}
2654
2655fn ensure_session_execution_lease_conn(
2656 conn: &Connection,
2657 session_id: &str,
2658 fence: &SessionExecutionLeaseFence,
2659 now: u64,
2660) -> Result<(), StoreError> {
2661 if fence.session_id != session_id {
2662 return Err(StoreError::SessionExecutionLeaseExpired {
2663 session_id: session_id.to_string(),
2664 });
2665 }
2666 let current = load_session_execution_lease_row_conn(conn, session_id)?;
2667 let Some(current) = current else {
2668 return Err(StoreError::SessionExecutionLeaseExpired {
2669 session_id: session_id.to_string(),
2670 });
2671 };
2672 if current
2673 .owner
2674 .as_ref()
2675 .is_some_and(|owner| owner.same_incarnation(&fence.owner))
2676 && current.lease_token.as_deref() == Some(fence.lease_token.as_str())
2677 && current.fencing_token == fence.fencing_token
2678 && current.expires_at_ms > now
2679 {
2680 Ok(())
2681 } else {
2682 Err(StoreError::SessionExecutionLeaseExpired {
2683 session_id: session_id.to_string(),
2684 })
2685 }
2686}
2687
2688fn release_session_execution_lease_conn(
2689 conn: &Connection,
2690 completion: &SessionExecutionLeaseCompletion,
2691) -> Result<(), StoreError> {
2692 conn.execute(
2693 "UPDATE session_execution_leases
2694 SET lease_owner_id = NULL,
2695 lease_owner_incarnation_id = NULL,
2696 lease_owner_liveness_json = NULL,
2697 lease_token = NULL,
2698 lease_claimed_at_ms = 0,
2699 lease_expires_at_ms = 0
2700 WHERE session_id = ?1
2701 AND lease_owner_id = ?2
2702 AND lease_owner_incarnation_id = ?3
2703 AND lease_token = ?4
2704 AND lease_fencing_token = ?5",
2705 params![
2706 completion.session_id,
2707 completion.owner.owner_id,
2708 completion.owner.incarnation_id,
2709 completion.lease_token,
2710 completion.fencing_token as i64
2711 ],
2712 )
2713 .map_err(sqlite_error)?;
2714 Ok(())
2715}