Skip to main content

lash_sqlite_store/
persistence.rs

1//! The [`RuntimePersistence`] capability-segment implementations for
2//! [`Store`]: [`SessionCommitStore`], [`SessionExecutionLeaseStore`],
3//! [`QueuedWorkStore`], [`TurnInputStore`], and [`StoreMaintenance`].
4//!
5//! This is the tokio-rusqlite port of the prior store's `persistence.rs`. The
6//! public surface is byte-for-byte the prior store async trait: identical method
7//! names and signatures, so consumers swap backends with a path rename only.
8//!
9//! The translation rules (see `conn.rs`, `lifecycle.rs`, `blobs.rs`):
10//!
11//! * Pure reads run through `self.conn.call(move |conn| { ... })`.
12//! * Read-then-write paths run through `self.conn.write(move |tx| { ... })`
13//!   (`BEGIN IMMEDIATE`, commit on `Ok`, rollback on `Err`) — this is the
14//!   cross-process write-lock guard.
15//! * Paths that may abandon partially-applied writes (the queued-work claim)
16//!   run through `self.conn.write_flow`, deciding commit vs rollback via
17//!   [`TxOutcome`].
18//! * The shared `*_conn` helpers (`try_load_session_head_meta_from_conn`,
19//!   `Self::put_checkpoint_conn`, `Self::load_usage_deltas_conn`,
20//!   `Self::load_session_graph_from_conn`, the queued-work helpers, …) are
21//!   synchronous and take a `&rusqlite::Connection`, so they are reused from
22//!   inside these closures (a `&Transaction` derefs to `&Connection`).
23//! * Closures must be `'static` + `Send`: every borrow of `self`/caller data is
24//!   cloned into an owned value before being moved in.
25
26use 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                // Roll back on a `StoreError` so a failure after the first
518                // write (e.g. a head-revision conflict surfaced mid-commit, or a
519                // backend write error) does not leave the partial transaction
520                // committed, while still carrying the typed error to the caller.
521                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                // Roll back the partially-inserted batch/items on a
835                // `StoreError` while still returning the typed error.
836                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                    // The fence is validated live, so its fencing token is the
865                    // currently-live session-lease generation; claims pin it and
866                    // are claimable only across a different generation (ADR 0029).
867                    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                        // Under `BEGIN IMMEDIATE` this connection already holds
1079                        // the write lock, but the row could still have been
1080                        // claimed by an earlier committed writer (its
1081                        // `claim_token` set and not yet expired). The `WHERE`
1082                        // clause filters those out, so a 0-row update means we
1083                        // lost the race for this batch: treat the whole claim as
1084                        // not-won rather than returning a claim that doesn't
1085                        // actually own the row.
1086                        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                            // Lost the race for this batch. Roll back any sibling
1116                            // rows we already claimed in this transaction so we
1117                            // never return a half-owned claim.
1118                            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                // Lower a `StoreError` into the rollback arm so the closure body
1133                // can keep using `?` while still propagating the error to the
1134                // caller. Encode it as a `Result` carried out of the flow.
1135                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            // A claim is live only while the session-execution-lease generation it
2014            // pins still holds the session lease (ADR 0029).
2015            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}