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 actual_revision = existing.as_ref().map_or(0, |meta| meta.head_revision);
219                    let expected_revision = commit.expected_head_revision.unwrap_or(0);
220                    if expected_revision != actual_revision {
221                        return Err(StoreError::HeadRevisionConflict {
222                            expected: commit.expected_head_revision,
223                            actual: actual_revision,
224                        });
225                    }
226                    for completed in &commit.completed_queue_claims {
227                        if completed.session_id != commit.session_id {
228                            return Err(StoreError::QueuedWorkClaimSuperseded {
229                                session_id: completed.session_id.clone(),
230                                claim_id: completed.claim_id.clone(),
231                            });
232                        }
233                        ensure_queued_work_completion_conn(tx, completed)?;
234                    }
235                    for completed in &commit.completed_turn_input_claims {
236                        if completed.session_id != commit.session_id {
237                            return Err(StoreError::TurnInputClaimSuperseded {
238                                session_id: completed.session_id.clone(),
239                                claim_id: completed.claim_id.clone(),
240                            });
241                        }
242                        let owned_rows: usize = tx
243                            .query_row(
244                                "SELECT COUNT(*)
245                                 FROM pending_turn_inputs
246                                 WHERE session_id = ?1
247                                   AND claim_id = ?2
248                                   AND claim_token = ?3",
249                                params![
250                                    completed.session_id,
251                                    completed.claim_id,
252                                    completed.lease_token
253                                ],
254                                |row| row.get::<_, i64>(0),
255                            )
256                            .map_err(sqlite_error)? as usize;
257                        ensure_turn_input_completion_owns_all_inputs(completed, owned_rows)?;
258                    }
259
260                    let stored_checkpoint =
261                        Self::put_checkpoint_conn(tx, &commit.checkpoint, blob_profile)
262                            .map_err(sqlite_error)?;
263
264                    if !commit.usage_deltas.is_empty() {
265                        let mut stmt = tx
266                            .prepare(
267                                "INSERT INTO usage_deltas (
268                                    source, model, input_tokens, output_tokens, cache_read_input_tokens, cache_write_input_tokens, reasoning_output_tokens
269                                ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
270                            )
271                            .map_err(sqlite_error)?;
272                        for entry in &commit.usage_deltas {
273                            stmt.execute(params![
274                                entry.source,
275                                entry.model,
276                                entry.usage.input_tokens,
277                                entry.usage.output_tokens,
278                                entry.usage.cache_read_input_tokens,
279                                entry.usage.cache_write_input_tokens,
280                                entry.usage.reasoning_output_tokens,
281                            ])
282                            .map_err(sqlite_error)?;
283                        }
284                    }
285
286                    let leaf_node_id = match &commit.graph {
287                        GraphCommitDelta::Unchanged { leaf_node_id } => leaf_node_id.clone(),
288                        GraphCommitDelta::Append {
289                            nodes,
290                            leaf_node_id,
291                        } => {
292                            for node in nodes {
293                                let node_json = encode_json(node);
294                                tx.execute(
295                                    "INSERT INTO graph_nodes (node_id, node_json) VALUES (?1, ?2)",
296                                    params![node.node_id, node_json],
297                                )
298                                .map_err(sqlite_error)?;
299                            }
300                            leaf_node_id.clone()
301                        }
302                        GraphCommitDelta::ReplaceFull(graph) => {
303                            tx.execute("DELETE FROM graph_nodes", [])
304                                .map_err(sqlite_error)?;
305                            for node in &graph.nodes {
306                                let node_json = encode_json(node);
307                                tx.execute(
308                                    "INSERT INTO graph_nodes (node_id, node_json) VALUES (?1, ?2)",
309                                    params![node.node_id, node_json],
310                                )
311                                .map_err(sqlite_error)?;
312                            }
313                            graph.leaf_node_id.clone()
314                        }
315                    };
316                    let graph_node_count: usize = tx
317                        .query_row(
318                            "SELECT COUNT(*) FROM graph_nodes WHERE tombstoned = 0",
319                            [],
320                            |row| row.get::<_, i64>(0),
321                        )
322                        .map_err(sqlite_error)? as usize;
323                    let next_revision = actual_revision + 1;
324                    let meta = SessionHeadMeta {
325                        schema_version: lash_core::store::SESSION_HEAD_META_SCHEMA_VERSION,
326                        session_id: commit.session_id.clone(),
327                        head_revision: next_revision,
328                        config: commit.config.clone(),
329                        agent_frames: commit.agent_frames.clone(),
330                        current_agent_frame_id: commit.current_agent_frame_id.clone(),
331                        checkpoint_ref: Some(stored_checkpoint.checkpoint_ref.clone()),
332                        leaf_node_id,
333                        graph_node_count,
334                        token_ledger: Vec::new(),
335                    };
336                    tx.execute(
337                        "INSERT OR REPLACE INTO session_head (singleton, session_id, head_json, head_revision)
338                         VALUES (1, ?1, ?2, ?3)",
339                        params![
340                            meta.session_id,
341                            encode_json(&meta),
342                            meta.head_revision as i64
343                        ],
344                    )
345                    .map_err(sqlite_error)?;
346                    for completed in &commit.completed_queue_claims {
347                        for batch_id in &completed.batch_ids {
348                            tx.execute(
349                                "DELETE FROM queued_work_batches
350                                 WHERE session_id = ?1
351                                   AND batch_id = ?2
352                                   AND claim_id = ?3
353                                   AND claim_token = ?4",
354                                params![
355                                    completed.session_id,
356                                    batch_id,
357                                    completed.claim_id,
358                                    completed.lease_token
359                                ],
360                            )
361                            .map_err(sqlite_error)?;
362                        }
363                    }
364                    for completed in &commit.completed_turn_input_claims {
365                        for input_id in &completed.input_ids {
366                            tx.execute(
367                                "UPDATE pending_turn_inputs
368                                 SET state = ?5,
369                                     claim_id = NULL,
370                                     claim_owner_id = NULL,
371                                     claim_owner_incarnation_id = NULL,
372                                     claim_owner_liveness_json = NULL,
373                                     claim_token = NULL,
374                                     claim_session_lease_generation = 0
375                                 WHERE session_id = ?1
376                                   AND input_id = ?2
377                                   AND claim_id = ?3
378                                   AND claim_token = ?4",
379                                params![
380                                    completed.session_id,
381                                    input_id,
382                                    completed.claim_id,
383                                    completed.lease_token,
384                                    lash_core::TurnInputState::Completed.as_str(),
385                                ],
386                            )
387                            .map_err(sqlite_error)?;
388                        }
389                    }
390                    if let Some(turn_id) = commit.interrupted_turn_input_turn_id.as_deref() {
391                        let input_ids = {
392                            let mut stmt = tx
393                                .prepare(
394                                    "SELECT input_id, ingress_json
395                                     FROM pending_turn_inputs
396                                     WHERE session_id = ?1 AND state = ?2",
397                                )
398                                .map_err(sqlite_error)?;
399                            let rows = stmt
400                                .query_map(
401                                    params![
402                                        commit.session_id,
403                                        lash_core::TurnInputState::PendingActive.as_str()
404                                    ],
405                                    |row| {
406                                        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
407                                    },
408                                )
409                                .map_err(sqlite_error)?;
410                            let mut input_ids = Vec::new();
411                            for row in rows {
412                                let (input_id, ingress_json) = row.map_err(sqlite_error)?;
413                                let ingress = decode_turn_input_ingress(ingress_json)?;
414                                if ingress
415                                    .active_turn_id()
416                                    .is_some_and(|active| active == turn_id)
417                                {
418                                    input_ids.push(input_id);
419                                }
420                            }
421                            input_ids
422                        };
423                        let next_turn_ingress = encode_json(&lash_core::TurnInputIngress::NextTurn);
424                        let mut stmt = tx
425                            .prepare(
426                                "UPDATE pending_turn_inputs
427                                 SET state = ?3,
428                                     ingress_json = ?4,
429                                     claim_id = NULL,
430                                     claim_owner_id = NULL,
431                                     claim_owner_incarnation_id = NULL,
432                                     claim_owner_liveness_json = NULL,
433                                     claim_token = NULL,
434                                     claim_session_lease_generation = 0
435                                 WHERE session_id = ?1 AND input_id = ?2",
436                            )
437                            .map_err(sqlite_error)?;
438                        for input_id in input_ids {
439                            stmt.execute(params![
440                                commit.session_id,
441                                input_id,
442                                lash_core::TurnInputState::DeferredNextTurn.as_str(),
443                                next_turn_ingress
444                            ])
445                            .map_err(sqlite_error)?;
446                        }
447                    }
448                    if !commit.committed_attachment_ids.is_empty() {
449                        let now = now as i64;
450                        let mut stmt = tx
451                            .prepare(
452                                "UPDATE attachment_manifest
453                                 SET committed_at_ms = COALESCE(committed_at_ms, ?1)
454                                 WHERE attachment_id = ?2 AND session_id = ?3",
455                            )
456                            .map_err(sqlite_error)?;
457                        for id in &commit.committed_attachment_ids {
458                            stmt.execute(params![now, id.as_str(), commit.session_id])
459                                .map_err(sqlite_error)?;
460                        }
461                    }
462                    if let Some(turn_commit) = &commit.turn_commit {
463                        tx.execute(
464                            "UPDATE attachment_manifest
465                             SET committed_at_ms = COALESCE(committed_at_ms, ?1)
466                             WHERE session_id = ?2
467                               AND owner_kind = 'turn'
468                               AND owner_id = ?3
469                               AND committed_at_ms IS NULL",
470                            params![now as i64, commit.session_id, turn_commit.turn_id],
471                        )
472                        .map_err(sqlite_error)?;
473                    }
474                    let mut enqueued_queue_batches = Vec::new();
475                    for (index, batch) in commit.enqueued_queue_batches.iter().enumerate() {
476                        if batch.session_id != commit.session_id {
477                            return Err(StoreError::SessionBindingMismatch {
478                                bound_session_id: commit.session_id.clone(),
479                                attempted_session_id: batch.session_id.clone(),
480                            });
481                        }
482                        enqueued_queue_batches.push(enqueue_queued_work_conn(
483                            tx,
484                            batch,
485                            now,
486                            enqueue_nonce_start.saturating_add(index as u64),
487                        )?);
488                    }
489                    let result = RuntimeCommitResult {
490                        head_revision: next_revision,
491                        checkpoint_ref: stored_checkpoint.checkpoint_ref,
492                        manifest: stored_checkpoint.manifest,
493                        enqueued_queue_batches,
494                        turn_input_applications: commit.turn_input_applications(),
495                    };
496                    if let Some(completed) = &commit.turn_commit {
497                        tx.execute(
498                            "INSERT INTO runtime_turn_commits (
499                                session_id, turn_id, turn_commit_hash, result_json, committed_at_ms
500                             )
501                             VALUES (?1, ?2, ?3, ?4, ?5)",
502                            params![
503                                completed.session_id,
504                                completed.turn_id,
505                                completed.turn_commit_hash,
506                                encode_json(&result),
507                                now as i64
508                            ],
509                        )
510                        .map_err(sqlite_error)?;
511                    }
512                    if let Some(completion) = commit.release_session_execution_lease.as_ref() {
513                        release_session_execution_lease_conn(tx, completion)?;
514                    }
515                    Ok(result)
516                })();
517                // 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 list_turn_input_applications(
1711        &self,
1712        session_id: &str,
1713    ) -> Result<Vec<lash_core::TurnInputApplication>, StoreError> {
1714        let session_id = session_id.to_string();
1715        self.conn
1716            .call(move |conn| {
1717                let outcome = (|| {
1718                    let mut stmt = conn
1719                        .prepare(
1720                            "SELECT turn_id, result_json
1721                             FROM runtime_turn_commits
1722                             WHERE session_id = ?1",
1723                        )
1724                        .map_err(sqlite_error)?;
1725                    let rows = stmt
1726                        .query_map(params![session_id], |row| {
1727                            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1728                        })
1729                        .map_err(sqlite_error)?;
1730                    let mut commits = Vec::new();
1731                    for row in rows {
1732                        let (turn_id, result_json) = row.map_err(sqlite_error)?;
1733                        let result: RuntimeCommitResult = serde_json::from_str(&result_json)
1734                            .map_err(|err| {
1735                                StoreError::Backend(format!(
1736                                    "failed to decode runtime turn commit result: {err}"
1737                                ))
1738                            })?;
1739                        commits.push((
1740                            result.head_revision,
1741                            turn_id,
1742                            result.turn_input_applications,
1743                        ));
1744                    }
1745                    commits.sort_by(|left, right| {
1746                        (left.0, left.1.as_str()).cmp(&(right.0, right.1.as_str()))
1747                    });
1748                    Ok(commits
1749                        .into_iter()
1750                        .flat_map(|(_, _, applications)| applications)
1751                        .collect())
1752                })();
1753                Ok(outcome)
1754            })
1755            .await
1756            .map_err(sqlite_error)?
1757    }
1758
1759    async fn cancel_pending_turn_inputs(
1760        &self,
1761        session_id: &str,
1762        targets: &[lash_core::PendingTurnInputCancelTarget],
1763    ) -> Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> {
1764        let session_id = session_id.to_string();
1765        let targets = targets.to_vec();
1766        let now = self.clock.timestamp_ms();
1767        self.conn
1768            .write_flow(move |tx| {
1769                let outcome: Result<Vec<lash_core::PendingTurnInputCancelResult>, StoreError> =
1770                    (|| {
1771                        let mut results = Vec::with_capacity(targets.len());
1772                        for target in targets {
1773                            let outcome = match load_pending_turn_input_row_by_target_conn(
1774                                tx,
1775                                &session_id,
1776                                &target,
1777                            )? {
1778                                Some(row) => cancel_pending_turn_input_row_conn(tx, row, now)?,
1779                                None => lash_core::PendingTurnInputCancelOutcome::NotFound,
1780                            };
1781                            results
1782                                .push(lash_core::PendingTurnInputCancelResult { target, outcome });
1783                        }
1784                        Ok(results)
1785                    })();
1786                match outcome {
1787                    Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1788                    Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1789                }
1790            })
1791            .await
1792            .map_err(sqlite_error)?
1793    }
1794
1795    async fn cancel_pending_turn_input_suffix(
1796        &self,
1797        session_id: &str,
1798        anchor: &lash_core::PendingTurnInputCancelTarget,
1799    ) -> Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> {
1800        let session_id = session_id.to_string();
1801        let anchor = anchor.clone();
1802        let now = self.clock.timestamp_ms();
1803        self.conn
1804            .write_flow(move |tx| {
1805                let outcome: Result<lash_core::PendingTurnInputSuffixCancelOutcome, StoreError> =
1806                    (|| {
1807                        let Some(anchor_row) =
1808                            load_pending_turn_input_row_by_target_conn(tx, &session_id, &anchor)?
1809                        else {
1810                            return Ok(
1811                                lash_core::PendingTurnInputSuffixCancelOutcome::AnchorNotFound {
1812                                    anchor,
1813                                },
1814                            );
1815                        };
1816                        let rows = {
1817                            let mut stmt = tx
1818                                .prepare(
1819                                    "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
1820                                            state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
1821                                            claim_owner_id, claim_owner_incarnation_id,
1822                                            claim_owner_liveness_json, claim_token, claim_session_lease_generation
1823                                     FROM pending_turn_inputs
1824                                     WHERE session_id = ?1 AND enqueue_seq >= ?2
1825                                     ORDER BY enqueue_seq ASC",
1826                                )
1827                                .map_err(sqlite_error)?;
1828                            let rows = stmt
1829                                .query_map(
1830                                    params![session_id, anchor_row.enqueue_seq as i64],
1831                                    pending_turn_input_row_from_sql,
1832                                )
1833                                .map_err(sqlite_error)?;
1834                            rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
1835                        };
1836                        let mut outcomes = Vec::with_capacity(rows.len());
1837                        for row in rows {
1838                            outcomes.push(cancel_pending_turn_input_row_conn(tx, row, now)?);
1839                        }
1840                        Ok(lash_core::PendingTurnInputSuffixCancelOutcome::Outcomes {
1841                            anchor,
1842                            outcomes,
1843                        })
1844                    })();
1845                match outcome {
1846                    Ok(value) => Ok(TxOutcome::Commit(Ok(value))),
1847                    Err(err) => Ok(TxOutcome::Rollback(Err(err))),
1848                }
1849            })
1850            .await
1851            .map_err(sqlite_error)?
1852    }
1853
1854    async fn claim_active_turn_inputs(
1855        &self,
1856        session_id: &str,
1857        session_execution_lease: &SessionExecutionLeaseFence,
1858        owner: &LeaseOwnerIdentity,
1859        turn_id: &str,
1860        checkpoint: lash_core::CheckpointKind,
1861        max_inputs: usize,
1862    ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1863        claim_pending_turn_inputs_sqlite(
1864            &self.conn,
1865            self.clock.timestamp_ms(),
1866            session_id,
1867            session_execution_lease,
1868            owner,
1869            max_inputs,
1870            lash_core::TurnInputClaimMode::ActiveTurn {
1871                turn_id: turn_id.to_string(),
1872                checkpoint,
1873            },
1874        )
1875        .await
1876    }
1877
1878    async fn claim_next_turn_inputs(
1879        &self,
1880        session_id: &str,
1881        session_execution_lease: &SessionExecutionLeaseFence,
1882        owner: &LeaseOwnerIdentity,
1883        max_inputs: usize,
1884    ) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
1885        claim_pending_turn_inputs_sqlite(
1886            &self.conn,
1887            self.clock.timestamp_ms(),
1888            session_id,
1889            session_execution_lease,
1890            owner,
1891            max_inputs,
1892            lash_core::TurnInputClaimMode::NextTurn,
1893        )
1894        .await
1895    }
1896
1897    async fn abandon_turn_input_claim(
1898        &self,
1899        claim: &lash_core::TurnInputClaim,
1900    ) -> Result<(), StoreError> {
1901        let session_id = claim.session_id.clone();
1902        let claim_id = claim.claim_id.clone();
1903        let lease_token = claim.lease_token.clone();
1904        let restored_state = match claim.mode {
1905            lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
1906                lash_core::TurnInputState::PendingActive
1907            }
1908            lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
1909        };
1910        self.conn
1911            .write(move |tx| {
1912                tx.execute(
1913                    "UPDATE pending_turn_inputs
1914                     SET state = CASE
1915                             WHEN state = ?4 THEN ?5
1916                             ELSE state
1917                         END,
1918                         claim_id = NULL,
1919                         claim_owner_id = NULL,
1920                         claim_owner_incarnation_id = NULL,
1921                         claim_owner_liveness_json = NULL,
1922                         claim_token = NULL,
1923                         claim_session_lease_generation = 0
1924                     WHERE session_id = ?1 AND claim_id = ?2 AND claim_token = ?3",
1925                    params![
1926                        session_id,
1927                        claim_id,
1928                        lease_token,
1929                        lash_core::TurnInputState::Accepted.as_str(),
1930                        restored_state.as_str(),
1931                    ],
1932                )
1933            })
1934            .await
1935            .map_err(sqlite_error)?;
1936        Ok(())
1937    }
1938
1939    async fn abandon_turn_input_claims(
1940        &self,
1941        claims: &[lash_core::TurnInputClaim],
1942    ) -> Result<(), StoreError> {
1943        if claims.is_empty() {
1944            return Ok(());
1945        }
1946        let mut sql = "UPDATE pending_turn_inputs
1947             SET state = CASE
1948                     WHEN state = 'accepted' THEN 'pending_active'
1949                     ELSE state
1950                 END,
1951                 claim_id = NULL,
1952                 claim_owner_id = NULL,
1953                 claim_owner_incarnation_id = NULL,
1954                 claim_owner_liveness_json = NULL,
1955                 claim_token = NULL,
1956                 claim_session_lease_generation = 0
1957             WHERE (session_id, claim_id, claim_token) IN ("
1958            .to_string();
1959        let mut values: Vec<rusqlite::types::Value> = Vec::with_capacity(claims.len() * 3);
1960        for (index, claim) in claims.iter().enumerate() {
1961            if index > 0 {
1962                sql.push_str(", ");
1963            }
1964            sql.push_str("(?, ?, ?)");
1965            values.push(claim.session_id.clone().into());
1966            values.push(claim.claim_id.clone().into());
1967            values.push(claim.lease_token.clone().into());
1968        }
1969        sql.push(')');
1970        self.conn
1971            .write(move |tx| tx.execute(&sql, rusqlite::params_from_iter(values.iter())))
1972            .await
1973            .map_err(sqlite_error)?;
1974        Ok(())
1975    }
1976}
1977
1978#[async_trait::async_trait]
1979impl StoreMaintenance for Store {
1980    async fn tombstone_nodes(&self, ids: &[String]) -> Result<(), StoreError> {
1981        if ids.is_empty() {
1982            return Ok(());
1983        }
1984        let ids = ids.to_vec();
1985        self.conn
1986            .write(move |tx| {
1987                let mut stmt =
1988                    tx.prepare("UPDATE graph_nodes SET tombstoned = 1 WHERE node_id = ?1")?;
1989                for id in &ids {
1990                    stmt.execute(params![id])?;
1991                }
1992                Ok(())
1993            })
1994            .await
1995            .map_err(sqlite_error)
1996    }
1997
1998    async fn vacuum(&self) -> Result<VacuumReport, StoreError> {
1999        let (removed_node_count, removed_pending_turn_input_tombstone_count) = self
2000            .conn
2001            .write(move |tx| {
2002                let removed_node_count =
2003                    tx.execute("DELETE FROM graph_nodes WHERE tombstoned = 1", [])?;
2004                let removed_pending_turn_input_tombstone_count = tx.execute(
2005                    "DELETE FROM pending_turn_inputs
2006                     WHERE state IN (?1, ?2)",
2007                    params![
2008                        lash_core::TurnInputState::Cancelled.as_str(),
2009                        lash_core::TurnInputState::Completed.as_str()
2010                    ],
2011                )?;
2012                Ok((
2013                    removed_node_count,
2014                    removed_pending_turn_input_tombstone_count,
2015                ))
2016            })
2017            .await
2018            .map_err(sqlite_error)?;
2019        Ok(VacuumReport {
2020            removed_node_count,
2021            removed_pending_turn_input_tombstone_count,
2022        })
2023    }
2024
2025    async fn gc_unreachable(&self) -> Result<GcReport, StoreError> {
2026        Ok(Store::gc_unreachable(self).await)
2027    }
2028}
2029
2030fn derive_pending_turn_input_id(
2031    session_id: &str,
2032    source_key: Option<&str>,
2033    now_epoch_ms: u64,
2034    nonce: u64,
2035) -> String {
2036    format!(
2037        "ti:{:x}",
2038        Sha256::digest(format!("{session_id}:{source_key:?}:{now_epoch_ms}:{nonce}").as_bytes())
2039    )
2040}
2041
2042fn cancel_pending_turn_input_row_conn(
2043    conn: &Connection,
2044    row: PendingTurnInputRow,
2045    now_epoch_ms: u64,
2046) -> Result<lash_core::PendingTurnInputCancelOutcome, StoreError> {
2047    let mut input = pending_turn_input_from_row(row.clone())?;
2048    match input.state {
2049        lash_core::TurnInputState::Cancelled => Ok(
2050            lash_core::PendingTurnInputCancelOutcome::AlreadyCancelled(input),
2051        ),
2052        lash_core::TurnInputState::Completed => Ok(
2053            lash_core::PendingTurnInputCancelOutcome::AlreadyCompleted(input),
2054        ),
2055        lash_core::TurnInputState::Accepted => {
2056            Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2057                claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2058                input,
2059            })
2060        }
2061        lash_core::TurnInputState::PendingActive | lash_core::TurnInputState::DeferredNextTurn => {
2062            // A claim is live only while the session-execution-lease generation it
2063            // pins still holds the session lease (ADR 0029).
2064            let live_claim = row.claim_token.is_some()
2065                && load_session_execution_lease_row_conn(conn, &row.session_id)?.is_some_and(
2066                    |lease| {
2067                        lease.lease_token.is_some()
2068                            && lease.expires_at_ms > now_epoch_ms
2069                            && lease.fencing_token == row.claim_session_lease_generation
2070                    },
2071                );
2072            if live_claim {
2073                return Ok(lash_core::PendingTurnInputCancelOutcome::AlreadyClaimed {
2074                    claim: pending_turn_input_claim_diagnostics_from_row(&row, input.state),
2075                    input,
2076                });
2077            }
2078            conn.execute(
2079                "UPDATE pending_turn_inputs
2080                 SET state = ?3,
2081                     claim_id = NULL,
2082                     claim_owner_id = NULL,
2083                     claim_owner_incarnation_id = NULL,
2084                     claim_owner_liveness_json = NULL,
2085                     claim_token = NULL,
2086                     claim_session_lease_generation = 0
2087                 WHERE session_id = ?1 AND input_id = ?2",
2088                params![
2089                    row.session_id,
2090                    row.input_id,
2091                    lash_core::TurnInputState::Cancelled.as_str(),
2092                ],
2093            )
2094            .map_err(sqlite_error)?;
2095            input.state = lash_core::TurnInputState::Cancelled;
2096            Ok(lash_core::PendingTurnInputCancelOutcome::Cancelled(input))
2097        }
2098    }
2099}
2100
2101#[allow(clippy::too_many_arguments)]
2102async fn checkpoint_work_pending_sqlite(
2103    conn: &SqliteConnection,
2104    now: u64,
2105    session_id: &str,
2106    generation: u64,
2107    turn_id: &str,
2108    checkpoint: lash_core::CheckpointKind,
2109    max_inputs: usize,
2110    max_batches: usize,
2111) -> Result<bool, StoreError> {
2112    if max_inputs == 0 && max_batches == 0 {
2113        return Ok(false);
2114    }
2115    let session_id = session_id.to_string();
2116    let turn_id = turn_id.to_string();
2117    conn.call(move |conn| {
2118        let outcome: Result<bool, StoreError> = (|| {
2119            let head_candidate = sqlite_queued_work_head_candidate_cte(
2120                QueuedWorkClaimBoundary::ActiveTurnCheckpoint,
2121            );
2122            let sql = format!(
2123                "WITH {head_candidate}
2124                 SELECT (
2125                    ?7 > 0 AND EXISTS (
2126                        SELECT 1
2127                        FROM pending_turn_inputs
2128                        WHERE session_id = ?1
2129                          AND state = ?4
2130                          AND (claim_token IS NULL OR claim_session_lease_generation <> ?3)
2131                          AND json_extract(ingress_json, '$.scope') = 'active_turn'
2132                          AND json_extract(ingress_json, '$.turn_id') = ?5
2133                          AND (
2134                              ?6 = 'before_completion'
2135                              OR COALESCE(
2136                                  json_extract(ingress_json, '$.min_boundary'),
2137                                  'after_work'
2138                              ) = 'after_work'
2139                          )
2140                        LIMIT 1
2141                    )
2142                ) OR (
2143                    ?8 > 0 AND EXISTS (
2144                        SELECT 1
2145                        FROM queued_work_head_candidate AS head
2146                        JOIN queued_work_items AS item
2147                          ON item.batch_id = head.head_batch_id
2148                        WHERE json_extract(item.payload_json, '$.type') <> 'session_command'
2149                        LIMIT 1
2150                    )
2151                )"
2152            );
2153            let pending: i64 = conn
2154                .query_row(
2155                    &sql,
2156                    params![
2157                        session_id,
2158                        now as i64,
2159                        generation as i64,
2160                        lash_core::TurnInputState::PendingActive.as_str(),
2161                        turn_id,
2162                        match checkpoint {
2163                            lash_core::CheckpointKind::AfterWork => "after_work",
2164                            lash_core::CheckpointKind::BeforeCompletion => "before_completion",
2165                        },
2166                        max_inputs as i64,
2167                        max_batches as i64,
2168                    ],
2169                    |row| row.get(0),
2170                )
2171                .map_err(sqlite_error)?;
2172            Ok(pending != 0)
2173        })();
2174        Ok(outcome)
2175    })
2176    .await
2177    .map_err(sqlite_error)?
2178}
2179
2180#[allow(clippy::too_many_arguments)]
2181fn claim_ready_queued_work_sqlite_conn(
2182    tx: &Connection,
2183    now: u64,
2184    session_id: &str,
2185    session_execution_lease: &SessionExecutionLeaseFence,
2186    owner: &LeaseOwnerIdentity,
2187    boundary: QueuedWorkClaimBoundary,
2188    max_batches: usize,
2189) -> Result<TxOutcome<Option<QueuedWorkClaim>>, StoreError> {
2190    if max_batches == 0 {
2191        return Ok(TxOutcome::Commit(None));
2192    }
2193    let generation = session_execution_lease.fencing_token;
2194    let candidate_rows = {
2195        let mut stmt = tx
2196            .prepare(&sqlite_queued_work_claim_candidates_sql(boundary))
2197            .map_err(sqlite_error)?;
2198        let rows = stmt
2199            .query_map(
2200                params![
2201                    session_id,
2202                    now as i64,
2203                    generation as i64,
2204                    claim_scan_limit(max_batches)
2205                ],
2206                queued_batch_row_from_sql,
2207            )
2208            .map_err(sqlite_error)?;
2209        rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2210    };
2211    let candidate_rows = candidate_rows
2212        .into_iter()
2213        .filter(|row| row.claim_token.is_none() || row.claim_session_lease_generation != generation)
2214        .collect::<Vec<_>>();
2215    let candidate_batches = candidate_rows
2216        .iter()
2217        .map(|row| queued_work_batch_from_conn(tx, row.clone()))
2218        .collect::<Result<Vec<_>, StoreError>>()?;
2219    let candidates = candidate_rows
2220        .iter()
2221        .zip(candidate_batches.iter())
2222        .map(|(row, batch)| {
2223            Ok(ClaimCandidate {
2224                enqueue_seq: row.enqueue_seq,
2225                claim_fencing_token: row.claim_fencing_token,
2226                work_class: batch.work_class().ok_or_else(|| {
2227                    StoreError::Backend(format!(
2228                        "queued-work batch `{}` has mixed or empty payload classes",
2229                        batch.batch_id
2230                    ))
2231                })?,
2232                delivery_policy: decode_delivery_policy(row.delivery_policy.clone())?,
2233                slot_policy: decode_slot_policy(row.slot_policy.clone())?,
2234                merge_key: decode_merge_key(row.merge_key_json.clone())?,
2235            })
2236        })
2237        .collect::<Result<Vec<_>, StoreError>>()?;
2238    let selected_len = select_turn_work_claim_prefix(&candidates, boundary, max_batches);
2239    if selected_len == 0 {
2240        return Ok(TxOutcome::Commit(None));
2241    }
2242    let mut selected = candidate_rows;
2243    selected.truncate(selected_len);
2244    let mut selected_batches = candidate_batches;
2245    selected_batches.truncate(selected_len);
2246    let lease = QueuedWorkClaimLease::derive(&candidates[0], session_id, owner, now, generation);
2247    let liveness_json = encode_liveness(&owner.liveness)?;
2248    for row in &selected {
2249        let claimed = tx
2250            .execute(
2251                "UPDATE queued_work_batches
2252                 SET claim_id = ?3,
2253                     claim_owner_id = ?4,
2254                     claim_owner_incarnation_id = ?5,
2255                     claim_owner_liveness_json = ?6,
2256                     claim_token = ?7,
2257                     claim_fencing_token = claim_fencing_token + 1,
2258                     claim_session_lease_generation = ?8
2259                 WHERE session_id = ?1
2260                   AND batch_id = ?2
2261                   AND (
2262                        claim_token IS NULL
2263                        OR claim_session_lease_generation <> ?8
2264                   )",
2265                params![
2266                    session_id,
2267                    row.batch_id,
2268                    lease.claim_id,
2269                    owner.owner_id.as_str(),
2270                    owner.incarnation_id.as_str(),
2271                    liveness_json.as_str(),
2272                    lease.lease_token,
2273                    lease.session_lease_generation as i64,
2274                ],
2275            )
2276            .map_err(sqlite_error)?;
2277        if claimed == 0 {
2278            return Ok(TxOutcome::Rollback(None));
2279        }
2280    }
2281    Ok(TxOutcome::Commit(Some(QueuedWorkClaim {
2282        session_id: session_id.to_string(),
2283        claim_id: lease.claim_id,
2284        owner: owner.clone(),
2285        lease_token: lease.lease_token,
2286        fencing_token: lease.fencing_token,
2287        session_lease_generation: lease.session_lease_generation,
2288        batches: selected_batches,
2289    })))
2290}
2291
2292#[allow(clippy::too_many_arguments)]
2293fn claim_pending_turn_inputs_sqlite_conn(
2294    tx: &Connection,
2295    now: u64,
2296    session_id: &str,
2297    session_execution_lease: &SessionExecutionLeaseFence,
2298    owner: &LeaseOwnerIdentity,
2299    max_inputs: usize,
2300    mode: lash_core::TurnInputClaimMode,
2301) -> Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> {
2302    if max_inputs == 0 {
2303        return Ok(TxOutcome::Commit(None));
2304    }
2305    let generation = session_execution_lease.fencing_token;
2306    let wanted_state = match &mode {
2307        lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2308            lash_core::TurnInputState::PendingActive
2309        }
2310        lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2311    };
2312    let candidate_rows = {
2313        let mut sql = "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2314                        state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2315                        claim_owner_id, claim_owner_incarnation_id,
2316                        claim_owner_liveness_json, claim_token, claim_session_lease_generation
2317                 FROM pending_turn_inputs
2318                 WHERE session_id = ? AND state = ?
2319                   AND (
2320                        claim_token IS NULL
2321                        OR claim_session_lease_generation <> ?
2322                   )"
2323        .to_string();
2324        let mut values: Vec<rusqlite::types::Value> = vec![
2325            session_id.to_string().into(),
2326            wanted_state.as_str().to_string().into(),
2327            (generation as i64).into(),
2328        ];
2329        if let lash_core::TurnInputClaimMode::ActiveTurn {
2330            turn_id,
2331            checkpoint,
2332        } = &mode
2333        {
2334            sql.push_str(
2335                " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2336                  AND json_extract(ingress_json, '$.turn_id') = ?",
2337            );
2338            values.push(turn_id.clone().into());
2339            if *checkpoint == lash_core::CheckpointKind::AfterWork {
2340                sql.push_str(
2341                    " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2342                );
2343            }
2344        }
2345        sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2346        values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2347        let mut stmt = tx.prepare(&sql).map_err(sqlite_error)?;
2348        let rows = stmt
2349            .query_map(
2350                rusqlite::params_from_iter(values.iter()),
2351                pending_turn_input_row_from_sql,
2352            )
2353            .map_err(sqlite_error)?;
2354        rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2355    };
2356    let selected = candidate_rows
2357        .into_iter()
2358        .take(max_inputs)
2359        .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2360        .collect::<Result<Vec<_>, StoreError>>()?;
2361    let Some((head, _)) = selected.first() else {
2362        return Ok(TxOutcome::Commit(None));
2363    };
2364    let lease = TurnInputClaimLease::derive(head, session_id, owner, now, generation);
2365    let liveness_json = encode_liveness(&owner.liveness)?;
2366    let state_after_claim = match &mode {
2367        lash_core::TurnInputClaimMode::ActiveTurn { .. } => lash_core::TurnInputState::Accepted,
2368        lash_core::TurnInputClaimMode::NextTurn => lash_core::TurnInputState::DeferredNextTurn,
2369    };
2370    let mut inputs = Vec::new();
2371    for (row, mut input) in selected {
2372        let claimed = tx
2373            .execute(
2374                "UPDATE pending_turn_inputs
2375                 SET state = ?3,
2376                     claim_id = ?4,
2377                     claim_owner_id = ?5,
2378                     claim_owner_incarnation_id = ?6,
2379                     claim_owner_liveness_json = ?7,
2380                     claim_token = ?8,
2381                     claim_fencing_token = claim_fencing_token + 1,
2382                     claim_session_lease_generation = ?9
2383                 WHERE session_id = ?1
2384                   AND input_id = ?2
2385                   AND (
2386                        claim_token IS NULL
2387                        OR claim_session_lease_generation <> ?9
2388                   )",
2389                params![
2390                    session_id,
2391                    row.input_id,
2392                    state_after_claim.as_str(),
2393                    lease.claim_id,
2394                    owner.owner_id.as_str(),
2395                    owner.incarnation_id.as_str(),
2396                    liveness_json.as_str(),
2397                    lease.lease_token,
2398                    lease.session_lease_generation as i64,
2399                ],
2400            )
2401            .map_err(sqlite_error)?;
2402        if claimed == 0 {
2403            return Ok(TxOutcome::Rollback(None));
2404        }
2405        input.state = state_after_claim;
2406        inputs.push(input);
2407    }
2408    Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2409        session_id: session_id.to_string(),
2410        claim_id: lease.claim_id,
2411        owner: owner.clone(),
2412        lease_token: lease.lease_token,
2413        fencing_token: lease.fencing_token,
2414        session_lease_generation: lease.session_lease_generation,
2415        mode,
2416        inputs,
2417        applications: Vec::new(),
2418    })))
2419}
2420
2421async fn claim_pending_turn_inputs_sqlite(
2422    conn: &SqliteConnection,
2423    now: u64,
2424    session_id: &str,
2425    session_execution_lease: &SessionExecutionLeaseFence,
2426    owner: &LeaseOwnerIdentity,
2427    max_inputs: usize,
2428    mode: lash_core::TurnInputClaimMode,
2429) -> Result<Option<lash_core::TurnInputClaim>, StoreError> {
2430    if max_inputs == 0 {
2431        return Ok(None);
2432    }
2433    let session_id = session_id.to_string();
2434    let session_execution_lease = session_execution_lease.clone();
2435    let owner = owner.clone();
2436    conn.write_flow(move |tx| {
2437        let outcome: Result<TxOutcome<Option<lash_core::TurnInputClaim>>, StoreError> = (|| {
2438            ensure_session_execution_lease_conn(
2439                tx,
2440                &session_id,
2441                &session_execution_lease,
2442                now,
2443            )?;
2444            let generation = session_execution_lease.fencing_token;
2445            let wanted_state = match &mode {
2446                lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2447                    lash_core::TurnInputState::PendingActive
2448                }
2449                lash_core::TurnInputClaimMode::NextTurn => {
2450                    lash_core::TurnInputState::DeferredNextTurn
2451                }
2452            };
2453            let candidate_rows = {
2454                let mut sql =
2455                    "SELECT enqueue_seq, input_id, session_id, source_key, ingress_json,
2456                            state, input_json, enqueued_at_ms, claim_id, claim_fencing_token,
2457                            claim_owner_id, claim_owner_incarnation_id,
2458                            claim_owner_liveness_json, claim_token, claim_session_lease_generation
2459                     FROM pending_turn_inputs
2460                     WHERE session_id = ? AND state = ?
2461                       AND (
2462                            claim_token IS NULL
2463                            OR claim_session_lease_generation <> ?
2464                       )"
2465                    .to_string();
2466                let mut values: Vec<rusqlite::types::Value> = vec![
2467                    session_id.clone().into(),
2468                    wanted_state.as_str().to_string().into(),
2469                    (generation as i64).into(),
2470                ];
2471                if let lash_core::TurnInputClaimMode::ActiveTurn {
2472                    turn_id,
2473                    checkpoint,
2474                } = &mode
2475                {
2476                    sql.push_str(
2477                        " AND json_extract(ingress_json, '$.scope') = 'active_turn'
2478                          AND json_extract(ingress_json, '$.turn_id') = ?",
2479                    );
2480                    values.push(turn_id.clone().into());
2481                    if *checkpoint == lash_core::CheckpointKind::AfterWork {
2482                        sql.push_str(
2483                            " AND COALESCE(json_extract(ingress_json, '$.min_boundary'), 'after_work') = 'after_work'",
2484                        );
2485                    }
2486                }
2487                sql.push_str(" ORDER BY enqueue_seq ASC LIMIT ?");
2488                values.push(i64::try_from(max_inputs).unwrap_or(i64::MAX).into());
2489                let mut stmt = tx
2490                    .prepare(&sql)
2491                    .map_err(sqlite_error)?;
2492                let rows = stmt
2493                    .query_map(
2494                        rusqlite::params_from_iter(values.iter()),
2495                        pending_turn_input_row_from_sql,
2496                    )
2497                    .map_err(sqlite_error)?;
2498                rows.collect::<Result<Vec<_>, _>>().map_err(sqlite_error)?
2499            };
2500            let selected = candidate_rows
2501                .into_iter()
2502                .take(max_inputs)
2503                .map(|row| Ok((row.clone(), pending_turn_input_from_row(row)?)))
2504                .collect::<Result<Vec<_>, StoreError>>()?;
2505            let Some((head, _)) = selected.first() else {
2506                return Ok(TxOutcome::Commit(None));
2507            };
2508            let lease = TurnInputClaimLease::derive(head, &session_id, &owner, now, generation);
2509            let liveness_json = encode_liveness(&owner.liveness)?;
2510            let state_after_claim = match &mode {
2511                lash_core::TurnInputClaimMode::ActiveTurn { .. } => {
2512                    lash_core::TurnInputState::Accepted
2513                }
2514                lash_core::TurnInputClaimMode::NextTurn => {
2515                    lash_core::TurnInputState::DeferredNextTurn
2516                }
2517            };
2518            let mut inputs = Vec::new();
2519            for (row, mut input) in selected {
2520                let claimed = tx
2521                    .execute(
2522                        "UPDATE pending_turn_inputs
2523                         SET state = ?3,
2524                             claim_id = ?4,
2525                             claim_owner_id = ?5,
2526                             claim_owner_incarnation_id = ?6,
2527                             claim_owner_liveness_json = ?7,
2528                             claim_token = ?8,
2529                             claim_fencing_token = claim_fencing_token + 1,
2530                             claim_session_lease_generation = ?9
2531                         WHERE session_id = ?1
2532                           AND input_id = ?2
2533                           AND (
2534                                claim_token IS NULL
2535                                OR claim_session_lease_generation <> ?9
2536                           )",
2537                        params![
2538                            session_id,
2539                            row.input_id,
2540                            state_after_claim.as_str(),
2541                            lease.claim_id,
2542                            owner.owner_id.as_str(),
2543                            owner.incarnation_id.as_str(),
2544                            liveness_json.as_str(),
2545                            lease.lease_token,
2546                            lease.session_lease_generation as i64,
2547                        ],
2548                    )
2549                    .map_err(sqlite_error)?;
2550                if claimed == 0 {
2551                    return Ok(TxOutcome::Rollback(None));
2552                }
2553                input.state = state_after_claim;
2554                inputs.push(input);
2555            }
2556            Ok(TxOutcome::Commit(Some(lash_core::TurnInputClaim {
2557                session_id: session_id.clone(),
2558                claim_id: lease.claim_id,
2559                owner: owner.clone(),
2560                lease_token: lease.lease_token,
2561                fencing_token: lease.fencing_token,
2562                session_lease_generation: lease.session_lease_generation,
2563                mode,
2564                inputs,
2565                applications: Vec::new(),
2566            })))
2567        })(
2568        );
2569        match outcome {
2570            Ok(TxOutcome::Commit(value)) => Ok(TxOutcome::Commit(Ok(value))),
2571            Ok(TxOutcome::Rollback(value)) => Ok(TxOutcome::Rollback(Ok(value))),
2572            Err(err) => Ok(TxOutcome::Rollback(Err(err))),
2573        }
2574    })
2575    .await
2576    .map_err(sqlite_error)?
2577}
2578
2579struct SessionExecutionLeaseRow {
2580    owner: Option<LeaseOwnerIdentity>,
2581    lease_token: Option<String>,
2582    fencing_token: u64,
2583    claimed_at_ms: u64,
2584    expires_at_ms: u64,
2585}
2586
2587fn load_session_execution_lease_row_conn(
2588    conn: &Connection,
2589    session_id: &str,
2590) -> Result<Option<SessionExecutionLeaseRow>, StoreError> {
2591    let row = conn
2592        .query_row(
2593            "SELECT lease_owner_id, lease_token, lease_fencing_token,
2594                    lease_claimed_at_ms, lease_expires_at_ms,
2595                    lease_owner_incarnation_id, lease_owner_liveness_json
2596             FROM session_execution_leases
2597             WHERE session_id = ?1",
2598            params![session_id],
2599            |row| {
2600                let owner_id: Option<String> = row.get(0)?;
2601                let incarnation_id: Option<String> = row.get(5)?;
2602                let liveness_json: Option<String> = row.get(6)?;
2603                Ok(SessionExecutionLeaseRow {
2604                    owner: lease_owner_from_columns(owner_id, incarnation_id, liveness_json),
2605                    lease_token: row.get(1)?,
2606                    fencing_token: row.get::<_, i64>(2)? as u64,
2607                    claimed_at_ms: row.get::<_, i64>(3)? as u64,
2608                    expires_at_ms: row.get::<_, i64>(4)? as u64,
2609                })
2610            },
2611        )
2612        .optional()
2613        .map_err(sqlite_error)?;
2614    Ok(row)
2615}
2616
2617fn lease_owner_from_columns(
2618    owner_id: Option<String>,
2619    incarnation_id: Option<String>,
2620    liveness_json: Option<String>,
2621) -> Option<LeaseOwnerIdentity> {
2622    owner_id.map(|owner_id| LeaseOwnerIdentity {
2623        incarnation_id: incarnation_id.unwrap_or_else(|| owner_id.clone()),
2624        owner_id,
2625        liveness: liveness_json
2626            .as_deref()
2627            .and_then(|json| serde_json::from_str(json).ok())
2628            .unwrap_or(LeaseOwnerLiveness::Opaque),
2629    })
2630}
2631
2632fn encode_liveness(liveness: &LeaseOwnerLiveness) -> Result<String, StoreError> {
2633    serde_json::to_string(liveness)
2634        .map_err(|err| StoreError::Backend(format!("failed to encode lease liveness: {err}")))
2635}
2636
2637fn row_to_session_execution_lease(
2638    session_id: &str,
2639    row: SessionExecutionLeaseRow,
2640) -> Result<SessionExecutionLease, StoreError> {
2641    Ok(SessionExecutionLease {
2642        session_id: session_id.to_string(),
2643        owner: row
2644            .owner
2645            .ok_or_else(|| StoreError::Backend("live session lease missing owner".to_string()))?,
2646        lease_token: row.lease_token.ok_or_else(|| {
2647            StoreError::Backend("live session lease missing lease token".to_string())
2648        })?,
2649        fencing_token: row.fencing_token,
2650        claimed_at_epoch_ms: row.claimed_at_ms,
2651        expires_at_epoch_ms: row.expires_at_ms,
2652    })
2653}
2654
2655fn acquire_session_execution_lease_conn(
2656    conn: &Connection,
2657    session_id: &str,
2658    owner: &LeaseOwnerIdentity,
2659    previous_fencing_token: u64,
2660    now: u64,
2661    lease_ttl_ms: u64,
2662) -> Result<SessionExecutionLease, StoreError> {
2663    let fencing_token = previous_fencing_token.saturating_add(1);
2664    let lease_token = format!(
2665        "{}:{}:{}:{now}:{fencing_token}",
2666        session_id, owner.owner_id, owner.incarnation_id
2667    );
2668    let expires_at = now.saturating_add(lease_ttl_ms);
2669    let liveness_json = encode_liveness(&owner.liveness)?;
2670    conn.execute(
2671        "INSERT INTO session_execution_leases (
2672            session_id, lease_owner_id, lease_owner_incarnation_id, lease_owner_liveness_json,
2673            lease_token, lease_fencing_token, lease_claimed_at_ms, lease_expires_at_ms
2674         )
2675         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
2676         ON CONFLICT(session_id) DO UPDATE SET
2677            lease_owner_id = excluded.lease_owner_id,
2678            lease_owner_incarnation_id = excluded.lease_owner_incarnation_id,
2679            lease_owner_liveness_json = excluded.lease_owner_liveness_json,
2680            lease_token = excluded.lease_token,
2681            lease_fencing_token = excluded.lease_fencing_token,
2682            lease_claimed_at_ms = excluded.lease_claimed_at_ms,
2683            lease_expires_at_ms = excluded.lease_expires_at_ms",
2684        params![
2685            session_id,
2686            owner.owner_id,
2687            owner.incarnation_id,
2688            liveness_json,
2689            lease_token,
2690            fencing_token as i64,
2691            now as i64,
2692            expires_at as i64
2693        ],
2694    )
2695    .map_err(sqlite_error)?;
2696    Ok(SessionExecutionLease {
2697        session_id: session_id.to_string(),
2698        owner: owner.clone(),
2699        lease_token,
2700        fencing_token,
2701        claimed_at_epoch_ms: now,
2702        expires_at_epoch_ms: expires_at,
2703    })
2704}
2705
2706fn ensure_session_execution_lease_conn(
2707    conn: &Connection,
2708    session_id: &str,
2709    fence: &SessionExecutionLeaseFence,
2710    now: u64,
2711) -> Result<(), StoreError> {
2712    if fence.session_id != session_id {
2713        return Err(StoreError::SessionExecutionLeaseExpired {
2714            session_id: session_id.to_string(),
2715        });
2716    }
2717    let current = load_session_execution_lease_row_conn(conn, session_id)?;
2718    let Some(current) = current else {
2719        return Err(StoreError::SessionExecutionLeaseExpired {
2720            session_id: session_id.to_string(),
2721        });
2722    };
2723    if current
2724        .owner
2725        .as_ref()
2726        .is_some_and(|owner| owner.same_incarnation(&fence.owner))
2727        && current.lease_token.as_deref() == Some(fence.lease_token.as_str())
2728        && current.fencing_token == fence.fencing_token
2729        && current.expires_at_ms > now
2730    {
2731        Ok(())
2732    } else {
2733        Err(StoreError::SessionExecutionLeaseExpired {
2734            session_id: session_id.to_string(),
2735        })
2736    }
2737}
2738
2739fn release_session_execution_lease_conn(
2740    conn: &Connection,
2741    completion: &SessionExecutionLeaseCompletion,
2742) -> Result<(), StoreError> {
2743    conn.execute(
2744        "UPDATE session_execution_leases
2745         SET lease_owner_id = NULL,
2746             lease_owner_incarnation_id = NULL,
2747             lease_owner_liveness_json = NULL,
2748             lease_token = NULL,
2749             lease_claimed_at_ms = 0,
2750             lease_expires_at_ms = 0
2751         WHERE session_id = ?1
2752           AND lease_owner_id = ?2
2753           AND lease_owner_incarnation_id = ?3
2754           AND lease_token = ?4
2755           AND lease_fencing_token = ?5",
2756        params![
2757            completion.session_id,
2758            completion.owner.owner_id,
2759            completion.owner.incarnation_id,
2760            completion.lease_token,
2761            completion.fencing_token as i64
2762        ],
2763    )
2764    .map_err(sqlite_error)?;
2765    Ok(())
2766}