1use 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
20pub 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
94pub struct SqliteSettlementOutcomeStore {
96 pool: Pool<SqliteConnectionManager>,
97 writer: Option<crate::receipt_store::WriterHandle>,
98 binding: SettlementStoreBinding,
99}
100
101impl SqliteSettlementOutcomeStore {
102 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 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
286pub(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 ¤t,
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}