Skip to main content

lash_sqlite_store/
persistence.rs

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