Skip to main content

chio_store_sqlite/
settle_attempts.rs

1//! SQLite-backed leased settlement work.
2
3use chio_core::hashing::sha256;
4use chio_kernel::ReceiptStoreError;
5use chio_settle::{
6    classify_attempt, validate_settlement_claim, DeadLetterRecord, RetryDecision, RetryPolicy,
7    SettlementAttemptClaim, SettlementFailureCode, SettlementOutcomeStore, SettlementRoute,
8    SettlementRouteError, SettlementRoutingInput, SettlementStoreBinding,
9};
10use r2d2::Pool;
11use r2d2_sqlite::SqliteConnectionManager;
12use rusqlite::{params, OptionalExtension, TransactionBehavior};
13use uuid::Uuid;
14
15use crate::dead_letters::{
16    insert_dead_letter_on_connection, read_dead_letter_on_connection, DeadLetterStoreError,
17    SETTLE_DEAD_LETTERS_MIGRATION,
18};
19
20/// Additive schema for pending and retryable settlement work.
21pub const SETTLE_ATTEMPTS_MIGRATION: &str = r#"
22CREATE TABLE IF NOT EXISTS settle_attempts (
23    receipt_id           TEXT PRIMARY KEY,
24    finalized_at         INTEGER NOT NULL CHECK (finalized_at >= 0),
25    work_kind            TEXT NOT NULL CHECK (
26                            work_kind IN ('pending_observation', 'retry_scheduled')
27                         ),
28    attempts             INTEGER NOT NULL CHECK (attempts BETWEEN 0 AND 4294967295),
29    next_visible_at_ms   INTEGER NOT NULL CHECK (next_visible_at_ms >= 0),
30    row_version          INTEGER NOT NULL CHECK (row_version >= 0),
31    lease_owner          TEXT CHECK (
32                            lease_owner IS NULL OR
33                            (length(lease_owner) BETWEEN 1 AND 128)
34                         ),
35    lease_token          TEXT CHECK (
36                            lease_token IS NULL OR
37                            (length(lease_token) BETWEEN 1 AND 128)
38                         ),
39    lease_until_ms       INTEGER CHECK (lease_until_ms IS NULL OR lease_until_ms >= 0),
40    reason_code          TEXT,
41    reason_detail_sha256 BLOB CHECK (
42                            reason_detail_sha256 IS NULL OR
43                            length(reason_detail_sha256) = 32
44                         ),
45    updated_at_ms        INTEGER NOT NULL CHECK (updated_at_ms >= 0),
46    CHECK ((lease_owner IS NULL AND lease_token IS NULL AND lease_until_ms IS NULL) OR
47           (lease_owner IS NOT NULL AND lease_token IS NOT NULL AND lease_until_ms IS NOT NULL)),
48    CHECK ((work_kind = 'pending_observation' AND attempts = 0 AND
49            reason_code IS NULL AND reason_detail_sha256 IS NULL) OR
50           (work_kind = 'retry_scheduled' AND attempts > 0 AND
51            reason_code IS NOT NULL AND reason_detail_sha256 IS NOT NULL))
52);
53
54CREATE INDEX IF NOT EXISTS idx_settle_attempts_visible
55    ON settle_attempts(next_visible_at_ms, lease_until_ms, receipt_id);
56
57CREATE TRIGGER IF NOT EXISTS trg_settle_attempts_reject_terminal_insert
58BEFORE INSERT ON settle_attempts
59WHEN EXISTS (
60    SELECT 1 FROM settle_dead_letters WHERE receipt_id = NEW.receipt_id
61)
62BEGIN
63    SELECT RAISE(ABORT, 'settlement receipt already dead-lettered');
64END;
65
66CREATE TRIGGER IF NOT EXISTS trg_settle_attempts_reject_terminal_update
67BEFORE UPDATE OF receipt_id ON settle_attempts
68WHEN EXISTS (
69    SELECT 1 FROM settle_dead_letters WHERE receipt_id = NEW.receipt_id
70)
71BEGIN
72    SELECT RAISE(ABORT, 'settlement receipt already dead-lettered');
73END;
74
75CREATE TRIGGER IF NOT EXISTS trg_settle_dead_letters_reject_attempt_insert
76BEFORE INSERT ON settle_dead_letters
77WHEN EXISTS (
78    SELECT 1 FROM settle_attempts WHERE receipt_id = NEW.receipt_id
79)
80BEGIN
81    SELECT RAISE(ABORT, 'settlement receipt still has active work');
82END;
83
84CREATE TRIGGER IF NOT EXISTS trg_settle_dead_letters_reject_attempt_update
85BEFORE UPDATE OF receipt_id ON settle_dead_letters
86WHEN EXISTS (
87    SELECT 1 FROM settle_attempts WHERE receipt_id = NEW.receipt_id
88)
89BEGIN
90    SELECT RAISE(ABORT, 'settlement receipt still has active work');
91END;
92"#;
93
94/// SQLite implementation of the leased settlement outcome store.
95pub struct SqliteSettlementOutcomeStore {
96    pool: Pool<SqliteConnectionManager>,
97    writer: Option<crate::receipt_store::WriterHandle>,
98    binding: SettlementStoreBinding,
99}
100
101impl SqliteSettlementOutcomeStore {
102    /// Open a standalone store for tests and operator tooling.
103    pub fn open_with_pool(
104        pool: Pool<SqliteConnectionManager>,
105    ) -> Result<Self, SettlementRouteError> {
106        let connection = pool.get().map_err(backend_error)?;
107        connection
108            .execute_batch(SETTLE_DEAD_LETTERS_MIGRATION)
109            .map_err(backend_error)?;
110        connection
111            .execute_batch(SETTLE_ATTEMPTS_MIGRATION)
112            .map_err(backend_error)?;
113        Ok(Self {
114            pool,
115            writer: None,
116            binding: new_store_binding(),
117        })
118    }
119
120    /// Open a store sharing the receipt store's single writer and binding.
121    pub fn open_alongside(store: &crate::SqliteReceiptStore) -> Result<Self, SettlementRouteError> {
122        let writer = store.writer_handle();
123        let Some(binding) = writer.settlement_store_binding() else {
124            return Err(invalid_record("receipt store lacks settlement projection"));
125        };
126        writer
127            .run_write(|connection| {
128                connection.execute_batch(SETTLE_DEAD_LETTERS_MIGRATION)?;
129                connection.execute_batch(SETTLE_ATTEMPTS_MIGRATION)?;
130                Ok(())
131            })
132            .map_err(receipt_error_to_route)?;
133        Ok(Self {
134            pool: store.pool.clone(),
135            writer: Some(writer),
136            binding,
137        })
138    }
139
140    fn run_write<T, F>(&self, job: F) -> Result<T, SettlementRouteError>
141    where
142        T: Send + 'static,
143        F: FnOnce(&mut rusqlite::Connection) -> Result<T, SettlementRouteError> + Send + 'static,
144    {
145        match &self.writer {
146            Some(writer) => writer
147                .run_write(move |connection| job(connection).map_err(route_error_to_receipt))
148                .map_err(receipt_error_to_route),
149            None => {
150                let mut connection = self.pool.get().map_err(backend_error)?;
151                job(&mut connection)
152            }
153        }
154    }
155}
156
157fn new_store_binding() -> SettlementStoreBinding {
158    SettlementStoreBinding::from_digest(*sha256(Uuid::now_v7().as_bytes()).as_bytes())
159}
160
161fn backend_error(error: impl std::fmt::Display) -> SettlementRouteError {
162    SettlementRouteError::Backend {
163        detail: error.to_string(),
164    }
165}
166
167fn conflict(detail: impl Into<String>) -> SettlementRouteError {
168    SettlementRouteError::Conflict {
169        detail: detail.into(),
170    }
171}
172
173fn invalid_record(detail: impl Into<String>) -> SettlementRouteError {
174    SettlementRouteError::InvalidRecord {
175        detail: detail.into(),
176    }
177}
178
179const WRITER_ROUTE_BACKEND_TAG: &str = "chio-store-sqlite/settlement-route/backend:";
180const WRITER_ROUTE_CONFLICT_TAG: &str = "chio-store-sqlite/settlement-route/conflict:";
181const WRITER_ROUTE_INVALID_TAG: &str = "chio-store-sqlite/settlement-route/invalid-record:";
182
183fn route_error_to_receipt(error: SettlementRouteError) -> ReceiptStoreError {
184    match error {
185        SettlementRouteError::Backend { detail } => {
186            ReceiptStoreError::Pool(format!("{WRITER_ROUTE_BACKEND_TAG}{detail}"))
187        }
188        SettlementRouteError::Conflict { detail } => {
189            ReceiptStoreError::Conflict(format!("{WRITER_ROUTE_CONFLICT_TAG}{detail}"))
190        }
191        SettlementRouteError::InvalidRecord { detail } => {
192            ReceiptStoreError::InvalidOutcome(format!("{WRITER_ROUTE_INVALID_TAG}{detail}"))
193        }
194    }
195}
196
197fn receipt_error_to_route(error: ReceiptStoreError) -> SettlementRouteError {
198    match error {
199        ReceiptStoreError::Pool(detail) => match detail.strip_prefix(WRITER_ROUTE_BACKEND_TAG) {
200            Some(detail) => backend_error(detail),
201            None => backend_error(format!("receipt writer pool error: {detail}")),
202        },
203        ReceiptStoreError::Conflict(detail) => {
204            match detail.strip_prefix(WRITER_ROUTE_CONFLICT_TAG) {
205                Some(detail) => conflict(detail),
206                None => backend_error(format!("receipt writer conflict: {detail}")),
207            }
208        }
209        ReceiptStoreError::InvalidOutcome(detail) => {
210            match detail.strip_prefix(WRITER_ROUTE_INVALID_TAG) {
211                Some(detail) => invalid_record(detail),
212                None => backend_error(format!("receipt writer invalid outcome: {detail}")),
213            }
214        }
215        other => backend_error(other),
216    }
217}
218
219fn attempt_projection_read_error(error: rusqlite::Error) -> ReceiptStoreError {
220    match error {
221        rusqlite::Error::FromSqlConversionFailure(..)
222        | rusqlite::Error::IntegralValueOutOfRange(..)
223        | rusqlite::Error::InvalidColumnType(..)
224        | rusqlite::Error::Utf8Error(..) => ReceiptStoreError::Conflict(
225            "settlement attempt zero contains an invalid SQLite value".to_string(),
226        ),
227        other => ReceiptStoreError::Sqlite(other),
228    }
229}
230
231fn dead_letter_error_to_route(error: DeadLetterStoreError) -> SettlementRouteError {
232    match error {
233        DeadLetterStoreError::Backend(detail) => backend_error(detail),
234        DeadLetterStoreError::Conflict(detail) => conflict(detail),
235        DeadLetterStoreError::InvalidRecord(detail) => invalid_record(detail),
236    }
237}
238
239fn sqlite_i64(value: u64, field: &'static str) -> Result<i64, SettlementRouteError> {
240    value
241        .try_into()
242        .map_err(|_| invalid_record(format!("{field} exceeds SQLite integer range")))
243}
244
245fn stored_u64(value: i64, field: &'static str) -> Result<u64, SettlementRouteError> {
246    value
247        .try_into()
248        .map_err(|_| invalid_record(format!("{field} is negative")))
249}
250
251fn stored_u32(value: i64, field: &'static str) -> Result<u32, SettlementRouteError> {
252    value
253        .try_into()
254        .map_err(|_| invalid_record(format!("{field} is outside the u32 range")))
255}
256
257fn attempt_row_error(error: rusqlite::Error) -> SettlementRouteError {
258    match error {
259        rusqlite::Error::FromSqlConversionFailure(..)
260        | rusqlite::Error::IntegralValueOutOfRange(..)
261        | rusqlite::Error::InvalidColumnType(..)
262        | rusqlite::Error::Utf8Error(..) => {
263            invalid_record("settlement attempt row contains an invalid SQLite value")
264        }
265        other => backend_error(other),
266    }
267}
268
269fn validate_claim_input(
270    worker_id: &str,
271    now_ms: u64,
272    lease_ms: u64,
273    limit: usize,
274) -> Result<(i64, i64), SettlementRouteError> {
275    validate_settlement_claim(worker_id, lease_ms, limit)
276        .map_err(|error| invalid_record(error.to_string()))?;
277    let lease_until_ms = now_ms
278        .checked_add(lease_ms)
279        .ok_or_else(|| invalid_record("settlement lease deadline overflows u64"))?;
280    Ok((
281        sqlite_i64(now_ms, "claim time")?,
282        sqlite_i64(lease_until_ms, "lease deadline")?,
283    ))
284}
285
286/// Insert attempt-zero work inside the caller's receipt transaction.
287pub(crate) fn insert_attempt_zero_tx(
288    tx: &rusqlite::Transaction<'_>,
289    receipt_id: &str,
290    finalized_at: u64,
291    next_visible_at_ms: u64,
292) -> Result<(), ReceiptStoreError> {
293    let finalized_at: i64 = finalized_at
294        .try_into()
295        .map_err(|_: std::num::TryFromIntError| {
296            ReceiptStoreError::Canonical(
297                "settlement finalization time exceeds SQLite integer range".to_string(),
298            )
299        })?;
300    let next_visible_at_ms: i64 =
301        next_visible_at_ms
302            .try_into()
303            .map_err(|_: std::num::TryFromIntError| {
304                ReceiptStoreError::Canonical(
305                    "settlement visibility time exceeds SQLite integer range".to_string(),
306                )
307            })?;
308    let inserted = tx
309        .execute(
310            "INSERT INTO settle_attempts (\
311            receipt_id, finalized_at, work_kind, attempts, next_visible_at_ms, row_version, \
312            lease_owner, lease_token, lease_until_ms, reason_code, reason_detail_sha256, \
313            updated_at_ms\
314         ) VALUES (?1, ?2, 'pending_observation', 0, ?3, 0, NULL, NULL, NULL, NULL, NULL, ?3)",
315            params![receipt_id, finalized_at, next_visible_at_ms],
316        )
317        .map_err(|error| {
318            if error.sqlite_error_code() == Some(rusqlite::ErrorCode::ConstraintViolation) {
319                ReceiptStoreError::Conflict(
320                    "settlement attempt zero conflicts with durable state".to_string(),
321                )
322            } else {
323                ReceiptStoreError::Sqlite(error)
324            }
325        })?;
326    if inserted != 1 {
327        return Err(ReceiptStoreError::Conflict(
328            "settlement attempt zero did not insert exactly one row".to_string(),
329        ));
330    }
331    let stored = tx
332        .query_row(
333            "SELECT finalized_at, work_kind, attempts, next_visible_at_ms, row_version, \
334                    lease_owner, lease_token, lease_until_ms, reason_code, \
335                    reason_detail_sha256, updated_at_ms \
336             FROM settle_attempts WHERE receipt_id = ?1",
337            [receipt_id],
338            |row| {
339                Ok((
340                    row.get::<_, i64>(0)?,
341                    row.get::<_, String>(1)?,
342                    row.get::<_, i64>(2)?,
343                    row.get::<_, i64>(3)?,
344                    row.get::<_, i64>(4)?,
345                    row.get::<_, Option<String>>(5)?,
346                    row.get::<_, Option<String>>(6)?,
347                    row.get::<_, Option<i64>>(7)?,
348                    row.get::<_, Option<String>>(8)?,
349                    row.get::<_, Option<Vec<u8>>>(9)?,
350                    row.get::<_, i64>(10)?,
351                ))
352            },
353        )
354        .optional()
355        .map_err(attempt_projection_read_error)?;
356    let expected = (
357        finalized_at,
358        "pending_observation".to_string(),
359        0,
360        next_visible_at_ms,
361        0,
362        None,
363        None,
364        None,
365        None,
366        None,
367        next_visible_at_ms,
368    );
369    if stored != Some(expected) {
370        return Err(ReceiptStoreError::Conflict(
371            "settlement attempt zero does not match its inserted projection".to_string(),
372        ));
373    }
374    Ok(())
375}
376
377struct AttemptState {
378    finalized_at: u64,
379    attempts: u32,
380    row_version: u64,
381    next_visible_at_ms: u64,
382    lease_until_ms: Option<u64>,
383}
384
385struct ClaimedAttemptState {
386    finalized_at: u64,
387    attempts: u32,
388    row_version: u64,
389    lease_owner: String,
390    lease_token: String,
391    lease_until_ms: u64,
392}
393
394fn validate_work_shape(
395    work_kind: &str,
396    attempts: u32,
397    reason_code: Option<String>,
398    reason_detail_sha256: Option<Vec<u8>>,
399) -> Result<(), SettlementRouteError> {
400    match (work_kind, attempts, reason_code, reason_detail_sha256) {
401        ("pending_observation", 0, None, None) => Ok(()),
402        ("retry_scheduled", attempt, Some(code), Some(digest)) if attempt > 0 => {
403            SettlementFailureCode::try_from(code.as_str())
404                .map_err(|error| invalid_record(error.to_string()))?;
405            let _: [u8; 32] = digest
406                .try_into()
407                .map_err(|_| invalid_record("settlement reason digest is not 32 bytes"))?;
408            Ok(())
409        }
410        _ => Err(invalid_record(
411            "settlement attempt work shape is inconsistent",
412        )),
413    }
414}
415
416fn validate_lease_shape(
417    lease_owner: Option<&str>,
418    lease_token: Option<&str>,
419    lease_until_ms: Option<i64>,
420) -> Result<Option<u64>, SettlementRouteError> {
421    match (lease_owner, lease_token, lease_until_ms) {
422        (None, None, None) => Ok(None),
423        (Some(owner), Some(token), Some(deadline))
424            if !owner.is_empty()
425                && owner.len() <= chio_settle::MAX_SETTLEMENT_WORKER_ID_BYTES
426                && !token.is_empty()
427                && token.len() <= 128 =>
428        {
429            stored_u64(deadline, "lease_until_ms").map(Some)
430        }
431        _ => Err(invalid_record(
432            "settlement attempt lease fields are inconsistent",
433        )),
434    }
435}
436
437fn read_attempt_state(
438    tx: &rusqlite::Transaction<'_>,
439    receipt_id: &str,
440) -> Result<Option<AttemptState>, SettlementRouteError> {
441    let overlaps_terminal = tx
442        .query_row(
443            "SELECT EXISTS(\
444                 SELECT 1 FROM settle_attempts AS attempt \
445                 INNER JOIN settle_dead_letters AS terminal USING (receipt_id) \
446                 WHERE attempt.receipt_id = ?1\
447             )",
448            [receipt_id],
449            |row| row.get::<_, bool>(0),
450        )
451        .map_err(attempt_row_error)?;
452    if overlaps_terminal {
453        return Err(invalid_record(
454            "settlement attempt overlaps terminal dead-letter state",
455        ));
456    }
457
458    let stored = tx
459        .query_row(
460            "SELECT finalized_at, attempts, row_version, next_visible_at_ms, lease_owner, \
461                    lease_token, lease_until_ms, work_kind, reason_code, \
462                    reason_detail_sha256, updated_at_ms \
463             FROM settle_attempts WHERE receipt_id = ?1",
464            [receipt_id],
465            |row| {
466                Ok((
467                    row.get::<_, i64>(0)?,
468                    row.get::<_, i64>(1)?,
469                    row.get::<_, i64>(2)?,
470                    row.get::<_, i64>(3)?,
471                    row.get::<_, Option<String>>(4)?,
472                    row.get::<_, Option<String>>(5)?,
473                    row.get::<_, Option<i64>>(6)?,
474                    row.get::<_, String>(7)?,
475                    row.get::<_, Option<String>>(8)?,
476                    row.get::<_, Option<Vec<u8>>>(9)?,
477                    row.get::<_, i64>(10)?,
478                ))
479            },
480        )
481        .optional()
482        .map_err(attempt_row_error)?;
483    stored
484        .map(
485            |(
486                finalized_at,
487                attempts,
488                row_version,
489                next_visible_at_ms,
490                lease_owner,
491                lease_token,
492                lease_until_ms,
493                work_kind,
494                reason_code,
495                reason_detail_sha256,
496                updated_at_ms,
497            )| {
498                let attempts = stored_u32(attempts, "attempts")?;
499                validate_work_shape(&work_kind, attempts, reason_code, reason_detail_sha256)?;
500                let lease_until_ms = validate_lease_shape(
501                    lease_owner.as_deref(),
502                    lease_token.as_deref(),
503                    lease_until_ms,
504                )?;
505                stored_u64(updated_at_ms, "updated_at_ms")?;
506                Ok(AttemptState {
507                    finalized_at: stored_u64(finalized_at, "finalized_at")?,
508                    attempts,
509                    row_version: stored_u64(row_version, "row_version")?,
510                    next_visible_at_ms: stored_u64(next_visible_at_ms, "next_visible_at_ms")?,
511                    lease_until_ms,
512                })
513            },
514        )
515        .transpose()
516}
517
518fn claim_one_tx(
519    tx: &rusqlite::Transaction<'_>,
520    receipt_id: &str,
521    worker_id: &str,
522    now_ms: u64,
523    now_i64: i64,
524    lease_until_ms: u64,
525    lease_until_i64: i64,
526) -> Result<Option<SettlementAttemptClaim>, SettlementRouteError> {
527    let Some(state) = read_attempt_state(tx, receipt_id)? else {
528        return Ok(None);
529    };
530    if state.next_visible_at_ms > now_ms
531        || matches!(state.lease_until_ms, Some(deadline) if deadline > now_ms)
532    {
533        return Ok(None);
534    }
535    let row_version = state
536        .row_version
537        .checked_add(1)
538        .ok_or_else(|| invalid_record("settlement row version overflows u64"))?;
539    let row_version_i64 = sqlite_i64(row_version, "row_version")?;
540    let lease_token = Uuid::now_v7().to_string();
541    let affected = tx
542        .execute(
543            "UPDATE settle_attempts SET row_version = ?1, lease_owner = ?2, lease_token = ?3, \
544                lease_until_ms = ?4, updated_at_ms = ?5 \
545             WHERE receipt_id = ?6 AND row_version = ?7 AND next_visible_at_ms <= ?5 \
546               AND (lease_until_ms IS NULL OR lease_until_ms <= ?5)",
547            params![
548                row_version_i64,
549                worker_id,
550                lease_token.as_str(),
551                lease_until_i64,
552                now_i64,
553                receipt_id,
554                sqlite_i64(state.row_version, "row_version")?,
555            ],
556        )
557        .map_err(backend_error)?;
558    if affected != 1 {
559        return Err(conflict("settlement claim changed before lease commit"));
560    }
561    Ok(Some(SettlementAttemptClaim {
562        receipt_id: receipt_id.to_string(),
563        finalized_at: state.finalized_at,
564        attempts: state.attempts,
565        row_version,
566        lease_owner: worker_id.to_string(),
567        lease_token,
568        lease_until_ms,
569    }))
570}
571
572fn read_claimed_attempt(
573    tx: &rusqlite::Transaction<'_>,
574    receipt_id: &str,
575) -> Result<Option<ClaimedAttemptState>, SettlementRouteError> {
576    let stored = tx
577        .query_row(
578            "SELECT finalized_at, attempts, row_version, lease_owner, lease_token, \
579                    lease_until_ms, work_kind, reason_code, reason_detail_sha256, \
580                    next_visible_at_ms, updated_at_ms \
581             FROM settle_attempts WHERE receipt_id = ?1",
582            [receipt_id],
583            |row| {
584                Ok((
585                    row.get::<_, i64>(0)?,
586                    row.get::<_, i64>(1)?,
587                    row.get::<_, i64>(2)?,
588                    row.get::<_, Option<String>>(3)?,
589                    row.get::<_, Option<String>>(4)?,
590                    row.get::<_, Option<i64>>(5)?,
591                    row.get::<_, String>(6)?,
592                    row.get::<_, Option<String>>(7)?,
593                    row.get::<_, Option<Vec<u8>>>(8)?,
594                    row.get::<_, i64>(9)?,
595                    row.get::<_, i64>(10)?,
596                ))
597            },
598        )
599        .optional()
600        .map_err(attempt_row_error)?;
601    let Some((
602        finalized_at,
603        attempts,
604        row_version,
605        lease_owner,
606        lease_token,
607        lease_until_ms,
608        work_kind,
609        reason_code,
610        reason_detail_sha256,
611        next_visible_at_ms,
612        updated_at_ms,
613    )) = stored
614    else {
615        return Ok(None);
616    };
617    let finalized_at = stored_u64(finalized_at, "finalized_at")?;
618    let attempts = stored_u32(attempts, "attempts")?;
619    let row_version = stored_u64(row_version, "row_version")?;
620    stored_u64(next_visible_at_ms, "next_visible_at_ms")?;
621    stored_u64(updated_at_ms, "updated_at_ms")?;
622    let Some(validated_lease_until_ms) = validate_lease_shape(
623        lease_owner.as_deref(),
624        lease_token.as_deref(),
625        lease_until_ms,
626    )?
627    else {
628        return Err(invalid_record(
629            "claimed settlement row has an incomplete lease",
630        ));
631    };
632    let (Some(lease_owner), Some(lease_token)) = (lease_owner, lease_token) else {
633        return Err(invalid_record(
634            "claimed settlement row has an incomplete lease",
635        ));
636    };
637    validate_work_shape(&work_kind, attempts, reason_code, reason_detail_sha256)?;
638    Ok(Some(ClaimedAttemptState {
639        finalized_at,
640        attempts,
641        row_version,
642        lease_owner,
643        lease_token,
644        lease_until_ms: validated_lease_until_ms,
645    }))
646}
647
648fn verify_live_claim(
649    tx: &rusqlite::Transaction<'_>,
650    claim: &SettlementAttemptClaim,
651    observed_at_ms: u64,
652) -> Result<(), SettlementRouteError> {
653    let Some(stored) = read_claimed_attempt(tx, &claim.receipt_id)? else {
654        return Err(conflict("settlement attempt no longer exists"));
655    };
656    if stored.finalized_at != claim.finalized_at
657        || stored.attempts != claim.attempts
658        || stored.row_version != claim.row_version
659        || stored.lease_owner != claim.lease_owner
660        || stored.lease_token != claim.lease_token
661        || stored.lease_until_ms != claim.lease_until_ms
662    {
663        return Err(conflict(
664            "settlement claim does not match durable lease state",
665        ));
666    }
667    if stored.lease_until_ms <= observed_at_ms {
668        return Err(conflict("settlement claim lease has expired"));
669    }
670    Ok(())
671}
672
673fn expected_dead_letter(
674    claim: &SettlementAttemptClaim,
675    decision: &RetryDecision,
676) -> Result<Option<(DeadLetterRecord, u32)>, SettlementRouteError> {
677    let RetryDecision::DeadLetter { reason } = decision else {
678        return Ok(None);
679    };
680    let attempts = claim
681        .attempts
682        .checked_add(1)
683        .ok_or_else(|| invalid_record("settlement attempt count overflows u32"))?;
684    Ok(Some((
685        DeadLetterRecord::new(
686            claim.receipt_id.clone(),
687            claim.finalized_at,
688            attempts,
689            reason.clone(),
690        ),
691        attempts,
692    )))
693}
694
695fn delete_claimed_attempt(
696    tx: &rusqlite::Transaction<'_>,
697    claim: &SettlementAttemptClaim,
698    observed_at_i64: i64,
699) -> Result<(), SettlementRouteError> {
700    let affected = tx
701        .execute(
702            "DELETE FROM settle_attempts WHERE receipt_id = ?1 AND finalized_at = ?2 \
703             AND attempts = ?3 AND row_version = ?4 AND lease_owner = ?5 \
704             AND lease_token = ?6 AND lease_until_ms = ?7 AND lease_until_ms > ?8",
705            params![
706                claim.receipt_id.as_str(),
707                sqlite_i64(claim.finalized_at, "finalized_at")?,
708                i64::from(claim.attempts),
709                sqlite_i64(claim.row_version, "row_version")?,
710                claim.lease_owner.as_str(),
711                claim.lease_token.as_str(),
712                sqlite_i64(claim.lease_until_ms, "lease_until_ms")?,
713                observed_at_i64,
714            ],
715        )
716        .map_err(backend_error)?;
717    if affected != 1 {
718        return Err(conflict("settlement claim changed before delete"));
719    }
720    Ok(())
721}
722
723fn schedule_retry(
724    tx: &rusqlite::Transaction<'_>,
725    claim: &SettlementAttemptClaim,
726    observed_at_ms: u64,
727    observed_at_i64: i64,
728    attempt: u32,
729    backoff: std::time::Duration,
730    reason: &chio_settle::SettlementFailureReason,
731) -> Result<u64, SettlementRouteError> {
732    let backoff_ms: u64 = backoff
733        .as_millis()
734        .try_into()
735        .map_err(|_| invalid_record("settlement backoff exceeds u64"))?;
736    let next_visible_at_ms = observed_at_ms
737        .checked_add(backoff_ms)
738        .ok_or_else(|| invalid_record("settlement visibility deadline overflows u64"))?;
739    let next_visible_at_i64 = sqlite_i64(next_visible_at_ms, "next_visible_at_ms")?;
740    let affected = tx
741        .execute(
742            "UPDATE settle_attempts SET work_kind = 'retry_scheduled', attempts = ?1, \
743                next_visible_at_ms = ?2, lease_owner = NULL, lease_token = NULL, \
744                lease_until_ms = NULL, reason_code = ?3, reason_detail_sha256 = ?4, \
745                updated_at_ms = ?5 \
746             WHERE receipt_id = ?6 AND finalized_at = ?7 AND attempts = ?8 \
747               AND row_version = ?9 AND lease_owner = ?10 AND lease_token = ?11 \
748               AND lease_until_ms = ?12 AND lease_until_ms > ?5",
749            params![
750                i64::from(attempt),
751                next_visible_at_i64,
752                reason.code().as_str(),
753                reason.detail_sha256().as_slice(),
754                observed_at_i64,
755                claim.receipt_id.as_str(),
756                sqlite_i64(claim.finalized_at, "finalized_at")?,
757                i64::from(claim.attempts),
758                sqlite_i64(claim.row_version, "row_version")?,
759                claim.lease_owner.as_str(),
760                claim.lease_token.as_str(),
761                sqlite_i64(claim.lease_until_ms, "lease_until_ms")?,
762            ],
763        )
764        .map_err(backend_error)?;
765    if affected != 1 {
766        return Err(conflict("settlement claim changed before retry scheduling"));
767    }
768    Ok(next_visible_at_ms)
769}
770
771impl SettlementOutcomeStore for SqliteSettlementOutcomeStore {
772    fn settlement_store_binding(&self) -> SettlementStoreBinding {
773        self.binding
774    }
775
776    fn claim_receipt(
777        &self,
778        receipt_id: &str,
779        worker_id: &str,
780        now_ms: u64,
781        lease_ms: u64,
782    ) -> Result<Option<SettlementAttemptClaim>, SettlementRouteError> {
783        let (now_i64, lease_until_i64) = validate_claim_input(worker_id, now_ms, lease_ms, 1)?;
784        let lease_until_ms = now_ms
785            .checked_add(lease_ms)
786            .ok_or_else(|| invalid_record("settlement lease deadline overflows u64"))?;
787        let receipt_id = receipt_id.to_string();
788        let worker_id = worker_id.to_string();
789        self.run_write(move |connection| {
790            let tx = connection
791                .transaction_with_behavior(TransactionBehavior::Immediate)
792                .map_err(backend_error)?;
793            let claim = claim_one_tx(
794                &tx,
795                &receipt_id,
796                &worker_id,
797                now_ms,
798                now_i64,
799                lease_until_ms,
800                lease_until_i64,
801            )?;
802            tx.commit().map_err(backend_error)?;
803            Ok(claim)
804        })
805    }
806
807    fn claim_due(
808        &self,
809        worker_id: &str,
810        now_ms: u64,
811        lease_ms: u64,
812        limit: usize,
813    ) -> Result<Vec<SettlementAttemptClaim>, SettlementRouteError> {
814        let (now_i64, lease_until_i64) = validate_claim_input(worker_id, now_ms, lease_ms, limit)?;
815        let lease_until_ms = now_ms
816            .checked_add(lease_ms)
817            .ok_or_else(|| invalid_record("settlement lease deadline overflows u64"))?;
818        let sql_limit: i64 = limit
819            .try_into()
820            .map_err(|_| invalid_record("settlement claim limit exceeds SQLite range"))?;
821        let worker_id = worker_id.to_string();
822        self.run_write(move |connection| {
823            let tx = connection
824                .transaction_with_behavior(TransactionBehavior::Immediate)
825                .map_err(backend_error)?;
826            let receipt_ids = {
827                let mut statement = tx
828                    .prepare(
829                        "SELECT receipt_id FROM settle_attempts \
830                         WHERE next_visible_at_ms <= ?1 \
831                           AND (lease_until_ms IS NULL OR lease_until_ms <= ?1) \
832                         ORDER BY next_visible_at_ms ASC, receipt_id ASC LIMIT ?2",
833                    )
834                    .map_err(backend_error)?;
835                let rows = statement
836                    .query_map(params![now_i64, sql_limit], |row| row.get::<_, String>(0))
837                    .map_err(backend_error)?;
838                rows.collect::<Result<Vec<_>, _>>()
839                    .map_err(attempt_row_error)?
840            };
841            let mut claims = Vec::with_capacity(receipt_ids.len());
842            for receipt_id in receipt_ids {
843                let claim = claim_one_tx(
844                    &tx,
845                    &receipt_id,
846                    &worker_id,
847                    now_ms,
848                    now_i64,
849                    lease_until_ms,
850                    lease_until_i64,
851                )?
852                .ok_or_else(|| conflict("due settlement row changed before batch claim"))?;
853                claims.push(claim);
854            }
855            tx.commit().map_err(backend_error)?;
856            Ok(claims)
857        })
858    }
859
860    fn record_claimed_outcome(
861        &self,
862        claim: &SettlementAttemptClaim,
863        outcome: &SettlementRoutingInput,
864        policy: RetryPolicy,
865        observed_at_ms: u64,
866    ) -> Result<SettlementRoute, SettlementRouteError> {
867        policy
868            .validate()
869            .map_err(|error| invalid_record(error.to_string()))?;
870        let observed_at_i64 = sqlite_i64(observed_at_ms, "observed_at_ms")?;
871        let decision = classify_attempt(&policy, claim.attempts, outcome);
872        let terminal = expected_dead_letter(claim, &decision)?;
873        let claim = claim.clone();
874        self.run_write(move |connection| {
875            let tx = connection
876                .transaction_with_behavior(TransactionBehavior::Immediate)
877                .map_err(backend_error)?;
878
879            if read_dead_letter_on_connection(&tx, &claim.receipt_id)
880                .map_err(dead_letter_error_to_route)?
881                .is_some()
882            {
883                let Some((record, attempts)) = terminal.as_ref() else {
884                    return Err(conflict("terminal settlement state blocks this outcome"));
885                };
886                let inserted = insert_dead_letter_on_connection(&tx, record)
887                    .map_err(dead_letter_error_to_route)?;
888                if inserted {
889                    return Err(backend_error(
890                        "terminal settlement row disappeared during replay",
891                    ));
892                }
893                tx.commit().map_err(backend_error)?;
894                return Ok(SettlementRoute::DeadLettered {
895                    attempts: *attempts,
896                });
897            }
898
899            verify_live_claim(&tx, &claim, observed_at_ms)?;
900            let route = match decision {
901                RetryDecision::Accepted | RetryDecision::Skip { .. } => {
902                    delete_claimed_attempt(&tx, &claim, observed_at_i64)?;
903                    SettlementRoute::NoAction
904                }
905                RetryDecision::Retry {
906                    attempt,
907                    backoff,
908                    reason,
909                } => {
910                    let next_visible_at_ms = schedule_retry(
911                        &tx,
912                        &claim,
913                        observed_at_ms,
914                        observed_at_i64,
915                        attempt,
916                        backoff,
917                        &reason,
918                    )?;
919                    SettlementRoute::RetryScheduled {
920                        attempt,
921                        next_visible_at_ms,
922                    }
923                }
924                RetryDecision::DeadLetter { .. } => {
925                    let Some((record, attempts)) = terminal.as_ref() else {
926                        return Err(invalid_record(
927                            "dead-letter decision lacks a terminal record",
928                        ));
929                    };
930                    delete_claimed_attempt(&tx, &claim, observed_at_i64)?;
931                    let inserted = insert_dead_letter_on_connection(&tx, record)
932                        .map_err(dead_letter_error_to_route)?;
933                    if !inserted {
934                        return Err(conflict(
935                            "terminal settlement row appeared during claimed transition",
936                        ));
937                    }
938                    SettlementRoute::DeadLettered {
939                        attempts: *attempts,
940                    }
941                }
942            };
943            tx.commit().map_err(backend_error)?;
944            Ok(route)
945        })
946    }
947}
948
949#[cfg(test)]
950mod tests {
951    use std::sync::{Arc, Barrier};
952    use std::thread;
953    use std::time::Duration;
954
955    use chio_kernel::ReceiptStore;
956    use chio_settle::{
957        RetryPolicy, SettlementFailureCode, SettlementFailureReason, SettlementOutcomeStore,
958        SettlementRoute, SettlementRouteError, SettlementRoutingInput, SettlementSkipReason,
959    };
960    use chio_test_support::prelude::*;
961    use r2d2::Pool;
962    use r2d2_sqlite::SqliteConnectionManager;
963
964    use super::*;
965    use crate::dead_letters::{DeadLetterStoreError, SqliteDeadLetterStore};
966
967    fn pool() -> Pool<SqliteConnectionManager> {
968        Pool::builder()
969            .max_size(1)
970            .build(SqliteConnectionManager::memory())
971            .test_expect("test pool builds")
972    }
973
974    fn file_pool(path: &std::path::Path) -> Pool<SqliteConnectionManager> {
975        let manager = SqliteConnectionManager::file(path).with_init(|connection| {
976            connection.busy_timeout(Duration::from_secs(5))?;
977            connection.pragma_update(None, "journal_mode", "WAL")?;
978            Ok(())
979        });
980        Pool::builder()
981            .max_size(1)
982            .build(manager)
983            .test_expect("file-backed test pool builds")
984    }
985
986    fn seed(pool: &Pool<SqliteConnectionManager>, receipt_id: &str, visible_at_ms: u64) {
987        let mut connection = pool.get().test_expect("test connection");
988        let tx = connection.transaction().test_expect("test transaction");
989        insert_attempt_zero_tx(&tx, receipt_id, 10, visible_at_ms)
990            .test_expect("attempt zero inserts");
991        tx.commit().test_expect("attempt zero commits");
992    }
993
994    fn rpc_failure(detail: &str) -> SettlementFailureReason {
995        SettlementFailureReason::from_detail(SettlementFailureCode::Rpc, detail)
996    }
997
998    fn dead_letter(receipt_id: &str) -> DeadLetterRecord {
999        DeadLetterRecord::new(
1000            receipt_id,
1001            10,
1002            1,
1003            SettlementFailureReason::from_detail(
1004                SettlementFailureCode::InvalidBinding,
1005                "binding mismatch",
1006            ),
1007        )
1008    }
1009
1010    fn seed_overlap(pool: &Pool<SqliteConnectionManager>, attempt_id: &str, terminal_id: &str) {
1011        seed(pool, attempt_id, 0);
1012        let dead_letters =
1013            SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("dead letters open");
1014        assert!(dead_letters
1015            .insert(&dead_letter(terminal_id))
1016            .test_expect("dead letter inserts"));
1017
1018        let connection = pool.get().test_expect("test connection");
1019        connection
1020            .execute_batch("DROP TRIGGER IF EXISTS trg_settle_attempts_reject_terminal_update;")
1021            .test_expect("update guard drops");
1022        connection
1023            .execute(
1024                "UPDATE settle_attempts SET receipt_id = ?1 WHERE receipt_id = ?2",
1025                params![terminal_id, attempt_id],
1026            )
1027            .test_expect("corrupt overlap seeds");
1028        connection
1029            .execute_batch(SETTLE_ATTEMPTS_MIGRATION)
1030            .test_expect("attempt migration reinstalls");
1031    }
1032
1033    fn assert_attempt_unleased(pool: &Pool<SqliteConnectionManager>, receipt_id: &str) {
1034        let state: (i64, Option<String>, Option<String>, Option<i64>) = pool
1035            .get()
1036            .test_expect("test connection")
1037            .query_row(
1038                "SELECT row_version, lease_owner, lease_token, lease_until_ms \
1039                 FROM settle_attempts WHERE receipt_id = ?1",
1040                [receipt_id],
1041                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
1042            )
1043            .test_expect("attempt state reads");
1044        assert_eq!(state, (0, None, None, None));
1045    }
1046
1047    fn claim(
1048        store: &SqliteSettlementOutcomeStore,
1049        receipt_id: &str,
1050        now_ms: u64,
1051        lease_ms: u64,
1052    ) -> chio_settle::SettlementAttemptClaim {
1053        store
1054            .claim_receipt(receipt_id, "worker", now_ms, lease_ms)
1055            .test_expect("claim succeeds")
1056            .test_expect("claim is present")
1057    }
1058
1059    fn attempt_count(pool: &Pool<SqliteConnectionManager>, receipt_id: &str) -> i64 {
1060        pool.get()
1061            .test_expect("test connection")
1062            .query_row(
1063                "SELECT COUNT(*) FROM settle_attempts WHERE receipt_id = ?1",
1064                [receipt_id],
1065                |row| row.get(0),
1066            )
1067            .test_expect("attempt count reads")
1068    }
1069
1070    fn dead_letter_count(pool: &Pool<SqliteConnectionManager>, receipt_id: &str) -> i64 {
1071        pool.get()
1072            .test_expect("test connection")
1073            .query_row(
1074                "SELECT COUNT(*) FROM settle_dead_letters WHERE receipt_id = ?1",
1075                [receipt_id],
1076                |row| row.get(0),
1077            )
1078            .test_expect("dead letter count reads")
1079    }
1080
1081    #[test]
1082    fn settlement_migration_is_idempotent() {
1083        let pool = pool();
1084        SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("first open");
1085        SqliteSettlementOutcomeStore::open_with_pool(pool).test_expect("second open");
1086    }
1087
1088    #[test]
1089    fn attempt_zero_requires_one_inserted_row() {
1090        let pool = pool();
1091        let mut connection = pool.get().test_expect("test connection");
1092        connection
1093            .execute_batch(SETTLE_DEAD_LETTERS_MIGRATION)
1094            .test_expect("dead-letter schema installs");
1095        connection
1096            .execute_batch(SETTLE_ATTEMPTS_MIGRATION)
1097            .test_expect("attempt schema installs");
1098        connection
1099            .execute_batch(
1100                "CREATE TRIGGER swallow_settlement_attempt \
1101                 BEFORE INSERT ON settle_attempts BEGIN SELECT RAISE(IGNORE); END;",
1102            )
1103            .test_expect("no-op settlement trigger installs");
1104        let tx = connection.transaction().test_expect("test transaction");
1105
1106        assert!(matches!(
1107            insert_attempt_zero_tx(&tx, "receipt-1", 1, 1),
1108            Err(ReceiptStoreError::Conflict(_))
1109        ));
1110    }
1111
1112    #[test]
1113    fn alongside_store_copies_the_receipt_writer_binding() {
1114        let directory = tempfile::tempdir().test_expect("temporary directory");
1115        let receipts = crate::SqliteReceiptStore::open(directory.path().join("receipts.db"))
1116            .test_expect("receipt store opens");
1117        let alongside =
1118            SqliteSettlementOutcomeStore::open_alongside(&receipts).test_expect("alongside opens");
1119        let standalone = SqliteSettlementOutcomeStore::open_with_pool(receipts.pool.clone())
1120            .test_expect("standalone opens");
1121        let receipt_binding = ReceiptStore::settlement_store_binding(&receipts)
1122            .test_expect("receipt binding is available");
1123
1124        assert_eq!(alongside.settlement_store_binding(), receipt_binding);
1125        assert_ne!(standalone.settlement_store_binding(), receipt_binding);
1126    }
1127
1128    #[test]
1129    fn writer_error_transport_preserves_only_tagged_route_classes() {
1130        let cases = [
1131            backend_error("backend"),
1132            conflict("conflict"),
1133            invalid_record("invalid"),
1134        ];
1135        for error in cases {
1136            let class = error.class();
1137            assert_eq!(
1138                receipt_error_to_route(route_error_to_receipt(error)).class(),
1139                class
1140            );
1141        }
1142        assert!(matches!(
1143            receipt_error_to_route(ReceiptStoreError::Conflict("writer conflict".to_string())),
1144            SettlementRouteError::Backend { .. }
1145        ));
1146    }
1147
1148    #[test]
1149    fn settlement_tables_reject_attempt_dead_letter_overlap() {
1150        let pool = pool();
1151        SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1152        seed(&pool, "receipt-1", 0);
1153        let dead_letters =
1154            SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("dead letters open");
1155        let record = dead_letter("receipt-1");
1156
1157        assert!(matches!(
1158            dead_letters.insert(&record),
1159            Err(DeadLetterStoreError::Conflict(_))
1160        ));
1161        pool.get()
1162            .test_expect("test connection")
1163            .execute(
1164                "DELETE FROM settle_attempts WHERE receipt_id = ?1",
1165                ["receipt-1"],
1166            )
1167            .test_expect("attempt deletes");
1168        assert!(dead_letters
1169            .insert(&record)
1170            .test_expect("dead letter inserts"));
1171
1172        let mut connection = pool.get().test_expect("test connection");
1173        let tx = connection.transaction().test_expect("test transaction");
1174        assert!(matches!(
1175            insert_attempt_zero_tx(&tx, "receipt-1", 10, 0),
1176            Err(ReceiptStoreError::Conflict(_))
1177        ));
1178    }
1179
1180    #[test]
1181    fn settlement_tables_reject_receipt_id_updates_that_create_overlap() {
1182        let pool = pool();
1183        SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1184        seed(&pool, "attempt-1", 0);
1185        seed(&pool, "attempt-2", 0);
1186        let dead_letters =
1187            SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("dead letters open");
1188        assert!(dead_letters
1189            .insert(&dead_letter("terminal-1"))
1190            .test_expect("first dead letter inserts"));
1191        assert!(dead_letters
1192            .insert(&dead_letter("terminal-2"))
1193            .test_expect("second dead letter inserts"));
1194        let connection = pool.get().test_expect("test connection");
1195
1196        let attempt_error = connection
1197            .execute(
1198                "UPDATE settle_attempts SET receipt_id = 'terminal-1' \
1199                 WHERE receipt_id = 'attempt-1'",
1200                [],
1201            )
1202            .test_expect_err("attempt overlap update rejects");
1203        assert_eq!(
1204            attempt_error.sqlite_error_code(),
1205            Some(rusqlite::ErrorCode::ConstraintViolation)
1206        );
1207
1208        let terminal_error = connection
1209            .execute(
1210                "UPDATE settle_dead_letters SET receipt_id = 'attempt-2' \
1211                 WHERE receipt_id = 'terminal-2'",
1212                [],
1213            )
1214            .test_expect_err("terminal overlap update rejects");
1215        assert_eq!(
1216            terminal_error.sqlite_error_code(),
1217            Some(rusqlite::ErrorCode::ConstraintViolation)
1218        );
1219    }
1220
1221    #[test]
1222    fn claim_receipt_rejects_preexisting_attempt_dead_letter_overlap() {
1223        let pool = pool();
1224        let store =
1225            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1226        seed_overlap(&pool, "attempt-1", "terminal-1");
1227
1228        assert!(matches!(
1229            store.claim_receipt("terminal-1", "worker", 0, 100),
1230            Err(SettlementRouteError::InvalidRecord { .. })
1231        ));
1232        assert_attempt_unleased(&pool, "terminal-1");
1233    }
1234
1235    #[test]
1236    fn claim_due_rejects_preexisting_attempt_dead_letter_overlap() {
1237        let pool = pool();
1238        let store =
1239            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1240        seed(&pool, "a-valid", 0);
1241        seed_overlap(&pool, "attempt-1", "z-terminal");
1242
1243        assert!(matches!(
1244            store.claim_due("worker", 0, 100, 2),
1245            Err(SettlementRouteError::InvalidRecord { .. })
1246        ));
1247        assert_attempt_unleased(&pool, "a-valid");
1248        assert_attempt_unleased(&pool, "z-terminal");
1249    }
1250
1251    #[test]
1252    fn claim_receipt_increments_version_and_uses_fresh_token() {
1253        let pool = pool();
1254        let store =
1255            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1256        seed(&pool, "receipt-1", 5);
1257
1258        let first = store
1259            .claim_receipt("receipt-1", "worker-1", 5, 100)
1260            .test_expect("first claim")
1261            .test_expect("claim present");
1262        assert_eq!(first.receipt_id, "receipt-1");
1263        assert_eq!(first.finalized_at, 10);
1264        assert_eq!(first.attempts, 0);
1265        assert_eq!(first.row_version, 1);
1266        assert_eq!(first.lease_owner, "worker-1");
1267        assert_eq!(first.lease_until_ms, 105);
1268
1269        assert!(store
1270            .claim_receipt("receipt-1", "worker-2", 104, 100)
1271            .test_expect("live lease check")
1272            .is_none());
1273        let second = store
1274            .claim_receipt("receipt-1", "worker-2", 105, 100)
1275            .test_expect("expired lease claim")
1276            .test_expect("expired row claimed");
1277        assert_eq!(second.row_version, 2);
1278        assert_eq!(second.lease_owner, "worker-2");
1279        assert_ne!(second.lease_token, first.lease_token);
1280    }
1281
1282    #[test]
1283    fn claim_due_orders_and_limits_due_rows() {
1284        let pool = pool();
1285        let store =
1286            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1287        seed(&pool, "receipt-b", 5);
1288        seed(&pool, "receipt-a", 5);
1289        seed(&pool, "receipt-c", 6);
1290
1291        let claims = store
1292            .claim_due("worker", 6, 100, 2)
1293            .test_expect("due claim");
1294        assert_eq!(
1295            claims
1296                .iter()
1297                .map(|claim| claim.receipt_id.as_str())
1298                .collect::<Vec<_>>(),
1299            vec!["receipt-a", "receipt-b"]
1300        );
1301    }
1302
1303    #[test]
1304    fn invalid_claim_bounds_do_not_mutate_work() {
1305        let pool = pool();
1306        let store =
1307            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1308        seed(&pool, "receipt-1", 0);
1309
1310        for result in [
1311            store.claim_receipt("receipt-1", "", 0, 1).map(|_| ()),
1312            store.claim_receipt("receipt-1", "worker", 0, 0).map(|_| ()),
1313            store.claim_due("worker", 0, 1, 0).map(|_| ()),
1314        ] {
1315            assert!(matches!(
1316                result,
1317                Err(SettlementRouteError::InvalidRecord { .. })
1318            ));
1319        }
1320
1321        let claim = store
1322            .claim_receipt("receipt-1", "worker", 0, 1)
1323            .test_expect("valid claim")
1324            .test_expect("row remains claimable");
1325        assert_eq!(claim.row_version, 1);
1326    }
1327
1328    #[test]
1329    fn row_version_overflow_is_invalid_without_mutation() {
1330        let pool = pool();
1331        let store =
1332            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1333        seed(&pool, "receipt-1", 0);
1334        pool.get()
1335            .test_expect("test connection")
1336            .execute(
1337                "UPDATE settle_attempts SET row_version = ?1 WHERE receipt_id = ?2",
1338                rusqlite::params![i64::MAX, "receipt-1"],
1339            )
1340            .test_expect("seed row version overflow");
1341
1342        assert!(matches!(
1343            store.claim_receipt("receipt-1", "worker", 0, 1),
1344            Err(SettlementRouteError::InvalidRecord { .. })
1345        ));
1346        let version: i64 = pool
1347            .get()
1348            .test_expect("test connection")
1349            .query_row(
1350                "SELECT row_version FROM settle_attempts WHERE receipt_id = ?1",
1351                ["receipt-1"],
1352                |row| row.get(0),
1353            )
1354            .test_expect("row version reads");
1355        assert_eq!(version, i64::MAX);
1356    }
1357
1358    #[test]
1359    fn unknown_persisted_reason_fails_before_claim_mutation() {
1360        let pool = pool();
1361        let store =
1362            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1363        seed(&pool, "receipt-1", 0);
1364        pool.get()
1365            .test_expect("test connection")
1366            .execute(
1367                "UPDATE settle_attempts SET work_kind = 'retry_scheduled', attempts = 1, \
1368                    reason_code = 'unknown_code', reason_detail_sha256 = zeroblob(32) \
1369                 WHERE receipt_id = ?1",
1370                ["receipt-1"],
1371            )
1372            .test_expect("corrupt retry reason");
1373
1374        assert!(matches!(
1375            store.claim_receipt("receipt-1", "worker", 0, 1),
1376            Err(SettlementRouteError::InvalidRecord { .. })
1377        ));
1378        let state = pool
1379            .get()
1380            .test_expect("test connection")
1381            .query_row(
1382                "SELECT row_version, lease_owner FROM settle_attempts WHERE receipt_id = ?1",
1383                ["receipt-1"],
1384                |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?)),
1385            )
1386            .test_expect("corrupt row remains");
1387        assert_eq!(state, (0, None));
1388    }
1389
1390    #[test]
1391    fn malformed_persisted_numeric_type_is_invalid_record() {
1392        let pool = pool();
1393        let store =
1394            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1395        seed(&pool, "receipt-1", 0);
1396        pool.get()
1397            .test_expect("test connection")
1398            .execute(
1399                "UPDATE settle_attempts SET work_kind = 'retry_scheduled', attempts = 1.5, \
1400                    reason_code = 'rpc', reason_detail_sha256 = zeroblob(32) \
1401                 WHERE receipt_id = ?1",
1402                ["receipt-1"],
1403            )
1404            .test_expect("malformed numeric value persists");
1405
1406        assert!(matches!(
1407            store.claim_receipt("receipt-1", "worker", 0, 1),
1408            Err(SettlementRouteError::InvalidRecord { .. })
1409        ));
1410    }
1411
1412    #[test]
1413    fn retries_preserve_reason_and_millisecond_backoff() {
1414        let pool = pool();
1415        let store =
1416            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1417        seed(&pool, "receipt-1", 0);
1418        let reason = rpc_failure("temporary outage");
1419        let first_claim = claim(&store, "receipt-1", 0, 1_000);
1420
1421        assert_eq!(
1422            store
1423                .record_claimed_outcome(
1424                    &first_claim,
1425                    &SettlementRoutingInput::Retryable {
1426                        reason: reason.clone(),
1427                    },
1428                    RetryPolicy::default(),
1429                    100,
1430                )
1431                .test_expect("first retry records"),
1432            SettlementRoute::RetryScheduled {
1433                attempt: 1,
1434                next_visible_at_ms: 350,
1435            }
1436        );
1437        let row = pool
1438            .get()
1439            .test_expect("test connection")
1440            .query_row(
1441                "SELECT attempts, next_visible_at_ms, reason_code, reason_detail_sha256, \
1442                        lease_owner, lease_token, lease_until_ms \
1443                 FROM settle_attempts WHERE receipt_id = ?1",
1444                ["receipt-1"],
1445                |row| {
1446                    Ok((
1447                        row.get::<_, i64>(0)?,
1448                        row.get::<_, i64>(1)?,
1449                        row.get::<_, String>(2)?,
1450                        row.get::<_, Vec<u8>>(3)?,
1451                        row.get::<_, Option<String>>(4)?,
1452                        row.get::<_, Option<String>>(5)?,
1453                        row.get::<_, Option<i64>>(6)?,
1454                    ))
1455                },
1456            )
1457            .test_expect("retry row reads");
1458        assert_eq!(row.0, 1);
1459        assert_eq!(row.1, 350);
1460        assert_eq!(row.2, "rpc");
1461        assert_eq!(row.3.as_slice(), reason.detail_sha256());
1462        assert_eq!((row.4, row.5, row.6), (None, None, None));
1463
1464        let second_claim = claim(&store, "receipt-1", 350, 1_000);
1465        assert_eq!(
1466            store
1467                .record_claimed_outcome(
1468                    &second_claim,
1469                    &SettlementRoutingInput::Retryable { reason },
1470                    RetryPolicy::default(),
1471                    400,
1472                )
1473                .test_expect("second retry records"),
1474            SettlementRoute::RetryScheduled {
1475                attempt: 2,
1476                next_visible_at_ms: 900,
1477            }
1478        );
1479    }
1480
1481    #[test]
1482    fn permanent_outcome_atomically_replaces_attempt_with_dead_letter() {
1483        let pool = pool();
1484        let store =
1485            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1486        seed(&pool, "receipt-1", 0);
1487        let claim = claim(&store, "receipt-1", 0, 100);
1488
1489        assert_eq!(
1490            store
1491                .record_claimed_outcome(
1492                    &claim,
1493                    &SettlementRoutingInput::Permanent {
1494                        reason: SettlementFailureReason::from_detail(
1495                            SettlementFailureCode::InvalidBinding,
1496                            "binding mismatch",
1497                        ),
1498                    },
1499                    RetryPolicy::default(),
1500                    1,
1501                )
1502                .test_expect("permanent outcome records"),
1503            SettlementRoute::DeadLettered { attempts: 1 }
1504        );
1505        assert_eq!(attempt_count(&pool, "receipt-1"), 0);
1506        assert_eq!(dead_letter_count(&pool, "receipt-1"), 1);
1507    }
1508
1509    #[test]
1510    fn accepted_and_skipped_remove_only_live_claimed_work() {
1511        for outcome in [
1512            SettlementRoutingInput::Accepted,
1513            SettlementRoutingInput::Skipped {
1514                reason: SettlementSkipReason::ZeroCharge,
1515            },
1516        ] {
1517            let pool = pool();
1518            let store = SqliteSettlementOutcomeStore::open_with_pool(pool.clone())
1519                .test_expect("store opens");
1520            seed(&pool, "receipt-1", 0);
1521            let claim = claim(&store, "receipt-1", 0, 100);
1522
1523            assert_eq!(
1524                store
1525                    .record_claimed_outcome(&claim, &outcome, RetryPolicy::default(), 1)
1526                    .test_expect("terminal cleanup records"),
1527                SettlementRoute::NoAction
1528            );
1529            assert_eq!(attempt_count(&pool, "receipt-1"), 0);
1530            assert_eq!(dead_letter_count(&pool, "receipt-1"), 0);
1531        }
1532    }
1533
1534    #[test]
1535    fn retry_exhaustion_is_terminal_and_exact_replay_is_idempotent() {
1536        let pool = pool();
1537        let store =
1538            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1539        seed(&pool, "receipt-1", 0);
1540        let reason = rpc_failure("still unavailable");
1541        let outcome = SettlementRoutingInput::Retryable {
1542            reason: reason.clone(),
1543        };
1544        let policy = RetryPolicy {
1545            max_retries: 0,
1546            ..RetryPolicy::default()
1547        };
1548        let claim = claim(&store, "receipt-1", 0, 100);
1549
1550        assert_eq!(
1551            store
1552                .record_claimed_outcome(&claim, &outcome, policy, 1)
1553                .test_expect("exhaustion records"),
1554            SettlementRoute::DeadLettered { attempts: 1 }
1555        );
1556        assert_eq!(
1557            store
1558                .record_claimed_outcome(&claim, &outcome, policy, 1)
1559                .test_expect("exact replay is idempotent"),
1560            SettlementRoute::DeadLettered { attempts: 1 }
1561        );
1562        assert!(matches!(
1563            store.record_claimed_outcome(
1564                &claim,
1565                &SettlementRoutingInput::Retryable {
1566                    reason: rpc_failure("different failure"),
1567                },
1568                policy,
1569                1,
1570            ),
1571            Err(SettlementRouteError::Conflict { .. })
1572        ));
1573        assert!(matches!(
1574            store.record_claimed_outcome(&claim, &SettlementRoutingInput::Accepted, policy, 1,),
1575            Err(SettlementRouteError::Conflict { .. })
1576        ));
1577        assert_eq!(attempt_count(&pool, "receipt-1"), 0);
1578        assert_eq!(dead_letter_count(&pool, "receipt-1"), 1);
1579    }
1580
1581    #[test]
1582    fn stale_or_expired_claim_cannot_commit() {
1583        let pool = pool();
1584        let store =
1585            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1586        seed(&pool, "receipt-1", 0);
1587        let stale = claim(&store, "receipt-1", 0, 10);
1588        let current = claim(&store, "receipt-1", 10, 10);
1589
1590        assert!(matches!(
1591            store.record_claimed_outcome(
1592                &stale,
1593                &SettlementRoutingInput::Accepted,
1594                RetryPolicy::default(),
1595                10,
1596            ),
1597            Err(SettlementRouteError::Conflict { .. })
1598        ));
1599        assert!(matches!(
1600            store.record_claimed_outcome(
1601                &current,
1602                &SettlementRoutingInput::Accepted,
1603                RetryPolicy::default(),
1604                current.lease_until_ms,
1605            ),
1606            Err(SettlementRouteError::Conflict { .. })
1607        ));
1608        assert_eq!(attempt_count(&pool, "receipt-1"), 1);
1609    }
1610
1611    #[test]
1612    fn invalid_policy_and_visibility_overflow_preserve_claimed_state() {
1613        let pool = pool();
1614        let store =
1615            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1616        seed(&pool, "receipt-1", 0);
1617        let now_ms = (i64::MAX as u64) - 600;
1618        let claim = claim(&store, "receipt-1", now_ms, 600);
1619        let invalid_policy = RetryPolicy {
1620            initial_backoff_ms: 0,
1621            ..RetryPolicy::default()
1622        };
1623        let outcome = SettlementRoutingInput::Retryable {
1624            reason: rpc_failure("temporary"),
1625        };
1626
1627        assert!(matches!(
1628            store.record_claimed_outcome(&claim, &outcome, invalid_policy, now_ms + 1),
1629            Err(SettlementRouteError::InvalidRecord { .. })
1630        ));
1631        assert!(matches!(
1632            store.record_claimed_outcome(
1633                &claim,
1634                &outcome,
1635                RetryPolicy {
1636                    initial_backoff_ms: 1_000,
1637                    backoff_cap_ms: 1_000,
1638                    ..RetryPolicy::default()
1639                },
1640                now_ms + 100,
1641            ),
1642            Err(SettlementRouteError::InvalidRecord { .. })
1643        ));
1644        assert_eq!(attempt_count(&pool, "receipt-1"), 1);
1645        assert_eq!(dead_letter_count(&pool, "receipt-1"), 0);
1646    }
1647
1648    #[test]
1649    fn attempt_count_overflow_is_invalid_without_mutation() {
1650        let pool = pool();
1651        let store =
1652            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1653        seed(&pool, "receipt-1", 0);
1654        pool.get()
1655            .test_expect("test connection")
1656            .execute(
1657                "UPDATE settle_attempts SET work_kind = 'retry_scheduled', attempts = ?1, \
1658                    reason_code = 'rpc', reason_detail_sha256 = zeroblob(32) \
1659                 WHERE receipt_id = ?2",
1660                rusqlite::params![i64::from(u32::MAX), "receipt-1"],
1661            )
1662            .test_expect("maximum attempt count persists");
1663        let claim = claim(&store, "receipt-1", 0, 100);
1664
1665        assert!(matches!(
1666            store.record_claimed_outcome(
1667                &claim,
1668                &SettlementRoutingInput::Permanent {
1669                    reason: SettlementFailureReason::from_detail(
1670                        SettlementFailureCode::InvalidBinding,
1671                        "binding mismatch",
1672                    ),
1673                },
1674                RetryPolicy::default(),
1675                1,
1676            ),
1677            Err(SettlementRouteError::InvalidRecord { .. })
1678        ));
1679        let durable = pool
1680            .get()
1681            .test_expect("test connection")
1682            .query_row(
1683                "SELECT attempts, row_version, lease_token FROM settle_attempts \
1684                 WHERE receipt_id = ?1",
1685                ["receipt-1"],
1686                |row| {
1687                    Ok((
1688                        row.get::<_, i64>(0)?,
1689                        row.get::<_, i64>(1)?,
1690                        row.get::<_, String>(2)?,
1691                    ))
1692                },
1693            )
1694            .test_expect("claimed row remains");
1695        assert_eq!(durable, (i64::from(u32::MAX), 1, claim.lease_token));
1696        assert_eq!(dead_letter_count(&pool, "receipt-1"), 0);
1697    }
1698
1699    #[test]
1700    fn dead_letter_insert_failure_rolls_back_attempt_delete() {
1701        let pool = pool();
1702        let store =
1703            SqliteSettlementOutcomeStore::open_with_pool(pool.clone()).test_expect("store opens");
1704        seed(&pool, "receipt-1", 0);
1705        let claim = claim(&store, "receipt-1", 0, 100);
1706        pool.get()
1707            .test_expect("test connection")
1708            .execute_batch(
1709                "CREATE TRIGGER reject_settlement_dead_letter \
1710                 BEFORE INSERT ON settle_dead_letters BEGIN \
1711                     SELECT RAISE(ABORT, 'forced dead-letter failure'); \
1712                 END;",
1713            )
1714            .test_expect("rejecting trigger installs");
1715
1716        assert!(matches!(
1717            store.record_claimed_outcome(
1718                &claim,
1719                &SettlementRoutingInput::Permanent {
1720                    reason: SettlementFailureReason::from_detail(
1721                        SettlementFailureCode::InvalidBinding,
1722                        "binding mismatch",
1723                    ),
1724                },
1725                RetryPolicy::default(),
1726                1,
1727            ),
1728            Err(SettlementRouteError::Conflict { .. })
1729        ));
1730        let durable = pool
1731            .get()
1732            .test_expect("test connection")
1733            .query_row(
1734                "SELECT row_version, lease_owner, lease_token, lease_until_ms \
1735                 FROM settle_attempts WHERE receipt_id = ?1",
1736                ["receipt-1"],
1737                |row| {
1738                    Ok((
1739                        row.get::<_, i64>(0)?,
1740                        row.get::<_, String>(1)?,
1741                        row.get::<_, String>(2)?,
1742                        row.get::<_, i64>(3)?,
1743                    ))
1744                },
1745            )
1746            .test_expect("claimed row survives rollback");
1747        assert_eq!(
1748            durable,
1749            (
1750                i64::try_from(claim.row_version).test_expect("row version fits"),
1751                claim.lease_owner,
1752                claim.lease_token,
1753                i64::try_from(claim.lease_until_ms).test_expect("lease deadline fits"),
1754            )
1755        );
1756        assert_eq!(dead_letter_count(&pool, "receipt-1"), 0);
1757    }
1758
1759    #[test]
1760    fn concurrent_terminal_transitions_commit_exactly_one_record() {
1761        let database = tempfile::NamedTempFile::new().test_expect("temporary database");
1762        let pool_one = file_pool(database.path());
1763        let pool_two = file_pool(database.path());
1764        let store_one = Arc::new(
1765            SqliteSettlementOutcomeStore::open_with_pool(pool_one.clone())
1766                .test_expect("first store opens"),
1767        );
1768        let store_two = Arc::new(
1769            SqliteSettlementOutcomeStore::open_with_pool(pool_two)
1770                .test_expect("second store opens"),
1771        );
1772        seed(&pool_one, "receipt-1", 0);
1773        let claim = claim(&store_one, "receipt-1", 0, 1_000);
1774        let barrier = Arc::new(Barrier::new(2));
1775
1776        let mut handles = Vec::new();
1777        for (store, detail) in [
1778            (Arc::clone(&store_one), "failure-a"),
1779            (Arc::clone(&store_two), "failure-b"),
1780        ] {
1781            let claim = claim.clone();
1782            let barrier = Arc::clone(&barrier);
1783            handles.push(thread::spawn(move || {
1784                barrier.wait();
1785                store.record_claimed_outcome(
1786                    &claim,
1787                    &SettlementRoutingInput::Permanent {
1788                        reason: SettlementFailureReason::from_detail(
1789                            SettlementFailureCode::InvalidBinding,
1790                            detail,
1791                        ),
1792                    },
1793                    RetryPolicy::default(),
1794                    1,
1795                )
1796            }));
1797        }
1798        let results = handles
1799            .into_iter()
1800            .map(|handle| handle.join().test_expect("transition thread joins"))
1801            .collect::<Vec<_>>();
1802
1803        assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
1804        assert_eq!(
1805            results
1806                .iter()
1807                .filter(|result| matches!(result, Err(SettlementRouteError::Conflict { .. })))
1808                .count(),
1809            1
1810        );
1811        assert_eq!(attempt_count(&pool_one, "receipt-1"), 0);
1812        assert_eq!(dead_letter_count(&pool_one, "receipt-1"), 1);
1813    }
1814}