1use chio_core::canonical::canonical_json_bytes;
19use chio_settle::{DeadLetterRecord, SETTLE_DEAD_LETTER_SCHEMA};
20use r2d2::Pool;
21use r2d2_sqlite::SqliteConnectionManager;
22use rusqlite::{params, OptionalExtension};
23use serde::Deserialize;
24use thiserror::Error;
25
26pub const SETTLE_DEAD_LETTERS_MIGRATION: &str = r#"
29CREATE TABLE IF NOT EXISTS settle_dead_letters (
30 receipt_id TEXT PRIMARY KEY,
31 finalized_at INTEGER NOT NULL,
32 attempts INTEGER NOT NULL,
33 reason TEXT NOT NULL,
34 pipeline_error TEXT,
35 canonical_json TEXT NOT NULL,
36 recorded_at INTEGER NOT NULL DEFAULT (strftime('%s','now'))
37);
38CREATE INDEX IF NOT EXISTS idx_settle_dead_letters_finalized_at
39 ON settle_dead_letters(finalized_at);
40"#;
41
42#[derive(Debug, Error)]
44pub enum DeadLetterStoreError {
45 #[error("dead letter backend error: {0}")]
47 Backend(String),
48 #[error("dead letter conflict: {0}")]
51 Conflict(String),
52 #[error("invalid dead letter record: {0}")]
54 InvalidRecord(String),
55}
56
57fn dead_letter_row_error(error: rusqlite::Error) -> DeadLetterStoreError {
58 match error {
59 rusqlite::Error::FromSqlConversionFailure(..)
60 | rusqlite::Error::IntegralValueOutOfRange(..)
61 | rusqlite::Error::InvalidColumnType(..)
62 | rusqlite::Error::Utf8Error(..) => DeadLetterStoreError::InvalidRecord(
63 "dead-letter row contains an invalid SQLite value".to_string(),
64 ),
65 other => DeadLetterStoreError::Backend(other.to_string()),
66 }
67}
68
69struct StoredDeadLetterRow {
70 receipt_id: String,
71 finalized_at: i64,
72 attempts: i64,
73 reason: String,
74 pipeline_error: Option<String>,
75 canonical: String,
76}
77
78#[derive(Deserialize)]
79struct DeadLetterSchemaProbe {
80 schema: String,
81}
82
83pub struct SqliteDeadLetterStore {
87 pool: Pool<SqliteConnectionManager>,
88 writer: Option<crate::receipt_store::WriterHandle>,
89}
90
91fn encode_dead_letter(record: &DeadLetterRecord) -> Result<(i64, Vec<u8>), DeadLetterStoreError> {
92 if !record.has_supported_schema() {
93 return Err(DeadLetterStoreError::InvalidRecord(
94 "unsupported programmatic settlement dead-letter schema".to_string(),
95 ));
96 }
97 if record.receipt_id.is_empty() {
98 return Err(DeadLetterStoreError::InvalidRecord(
99 "receipt_id must not be empty".to_string(),
100 ));
101 }
102 if record.attempts == 0 {
103 return Err(DeadLetterStoreError::InvalidRecord(
104 "attempts must be at least one".to_string(),
105 ));
106 }
107
108 let finalized_at =
109 record
110 .finalized_at
111 .try_into()
112 .map_err(|err: std::num::TryFromIntError| {
113 DeadLetterStoreError::InvalidRecord(err.to_string())
114 })?;
115 let canonical = canonical_json_bytes(record)
116 .map_err(|err| DeadLetterStoreError::InvalidRecord(err.to_string()))?;
117 Ok((finalized_at, canonical))
118}
119
120fn read_stored_dead_letter_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<StoredDeadLetterRow> {
121 Ok(StoredDeadLetterRow {
122 receipt_id: row.get(0)?,
123 finalized_at: row.get(1)?,
124 attempts: row.get(2)?,
125 reason: row.get(3)?,
126 pipeline_error: row.get(4)?,
127 canonical: row.get(5)?,
128 })
129}
130
131fn decode_stored_dead_letter(
132 row: StoredDeadLetterRow,
133) -> Result<DeadLetterRecord, DeadLetterStoreError> {
134 let finalized_at: u64 =
135 row.finalized_at
136 .try_into()
137 .map_err(|_: std::num::TryFromIntError| {
138 DeadLetterStoreError::InvalidRecord(
139 "persisted finalized_at is outside the u64 range".to_string(),
140 )
141 })?;
142 let attempts: u32 = row
143 .attempts
144 .try_into()
145 .map_err(|_: std::num::TryFromIntError| {
146 DeadLetterStoreError::InvalidRecord(
147 "persisted attempts is outside the u32 range".to_string(),
148 )
149 })?;
150 let schema: DeadLetterSchemaProbe = serde_json::from_str(&row.canonical)
151 .map_err(|err| DeadLetterStoreError::InvalidRecord(err.to_string()))?;
152 if schema.schema != SETTLE_DEAD_LETTER_SCHEMA {
153 return Err(DeadLetterStoreError::InvalidRecord(
154 "unsupported persisted settlement dead-letter schema".to_string(),
155 ));
156 }
157 let record: DeadLetterRecord = serde_json::from_str(&row.canonical)
158 .map_err(|err| DeadLetterStoreError::InvalidRecord(err.to_string()))?;
159 let (_, canonical) = encode_dead_letter(&record)?;
160 let columns_match = record.receipt_id == row.receipt_id
161 && record.finalized_at == finalized_at
162 && record.attempts == attempts
163 && record.reason.code().as_str() == row.reason
164 && row.pipeline_error.is_none()
165 && canonical.as_slice() == row.canonical.as_bytes();
166 if columns_match {
167 Ok(record)
168 } else {
169 Err(DeadLetterStoreError::InvalidRecord(
170 "SQLite columns do not match canonical dead-letter bytes".to_string(),
171 ))
172 }
173}
174
175fn select_dead_letter_row(
176 connection: &rusqlite::Connection,
177 receipt_id: &str,
178) -> Result<Option<StoredDeadLetterRow>, DeadLetterStoreError> {
179 connection
180 .query_row(
181 "SELECT receipt_id, finalized_at, attempts, reason, pipeline_error, canonical_json \
182 FROM settle_dead_letters WHERE receipt_id = ?1",
183 params![receipt_id],
184 read_stored_dead_letter_row,
185 )
186 .optional()
187 .map_err(dead_letter_row_error)
188}
189
190pub(crate) fn insert_dead_letter_on_connection(
195 connection: &rusqlite::Connection,
196 record: &DeadLetterRecord,
197) -> Result<bool, DeadLetterStoreError> {
198 let (finalized_at, canonical) = encode_dead_letter(record)?;
199 let canonical = std::str::from_utf8(&canonical)
200 .map_err(|err| DeadLetterStoreError::InvalidRecord(err.to_string()))?;
201 let inserted = connection
202 .execute(
203 "INSERT INTO settle_dead_letters \
204 (receipt_id, finalized_at, attempts, reason, pipeline_error, canonical_json) \
205 VALUES (?1, ?2, ?3, ?4, ?5, ?6) \
206 ON CONFLICT(receipt_id) DO NOTHING",
207 params![
208 record.receipt_id.as_str(),
209 finalized_at,
210 i64::from(record.attempts),
211 record.reason.code().as_str(),
212 Option::<&str>::None,
213 canonical,
214 ],
215 )
216 .map_err(|err| {
217 if err.sqlite_error_code() == Some(rusqlite::ErrorCode::ConstraintViolation) {
218 DeadLetterStoreError::Conflict(
219 "dead letter conflicts with active settlement work".to_string(),
220 )
221 } else {
222 DeadLetterStoreError::Backend(err.to_string())
223 }
224 })?;
225 if inserted == 1 {
226 return Ok(true);
227 }
228
229 let existing = select_dead_letter_row(connection, record.receipt_id.as_str())?;
230 match existing {
231 Some(existing) if existing.canonical.as_bytes() == canonical.as_bytes() => {
232 decode_stored_dead_letter(existing)?;
233 Ok(false)
234 }
235 Some(_) => Err(DeadLetterStoreError::Conflict(format!(
236 "settle_dead_letters row for receipt_id={} already exists with different bytes",
237 record.receipt_id
238 ))),
239 None => Err(DeadLetterStoreError::Backend(format!(
240 "settle_dead_letters conflict for receipt_id={} but no row was readable",
241 record.receipt_id
242 ))),
243 }
244}
245
246pub(crate) fn read_dead_letter_on_connection(
248 connection: &rusqlite::Connection,
249 receipt_id: &str,
250) -> Result<Option<DeadLetterRecord>, DeadLetterStoreError> {
251 select_dead_letter_row(connection, receipt_id)?
252 .map(decode_stored_dead_letter)
253 .transpose()
254}
255
256fn insert_dead_letter_transaction(
257 connection: &mut rusqlite::Connection,
258 record: &DeadLetterRecord,
259) -> Result<bool, DeadLetterStoreError> {
260 let tx = connection
261 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
262 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
263 let inserted = insert_dead_letter_on_connection(&tx, record)?;
264 tx.commit()
265 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
266 Ok(inserted)
267}
268
269fn clear_dead_letter_on_connection(
270 connection: &rusqlite::Connection,
271 receipt_id: &str,
272) -> Result<bool, DeadLetterStoreError> {
273 connection
274 .execute(
275 "DELETE FROM settle_dead_letters WHERE receipt_id = ?1",
276 params![receipt_id],
277 )
278 .map(|affected| affected > 0)
279 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))
280}
281
282const WRITER_BACKEND_TAG: &str = "chio-store-sqlite/dead-letter/backend:";
283const WRITER_CONFLICT_TAG: &str = "chio-store-sqlite/dead-letter/conflict:";
284const WRITER_INVALID_RECORD_TAG: &str = "chio-store-sqlite/dead-letter/invalid-record:";
285
286fn dead_letter_error_into_receipt_store(
287 error: DeadLetterStoreError,
288) -> chio_kernel::ReceiptStoreError {
289 match error {
290 DeadLetterStoreError::Conflict(message) => {
291 chio_kernel::ReceiptStoreError::Conflict(format!("{WRITER_CONFLICT_TAG}{message}"))
292 }
293 DeadLetterStoreError::InvalidRecord(message) => chio_kernel::ReceiptStoreError::Canonical(
294 format!("{WRITER_INVALID_RECORD_TAG}{message}"),
295 ),
296 DeadLetterStoreError::Backend(message) => {
297 chio_kernel::ReceiptStoreError::Pool(format!("{WRITER_BACKEND_TAG}{message}"))
298 }
299 }
300}
301
302fn dead_letter_error_from_receipt_store(
303 error: chio_kernel::ReceiptStoreError,
304) -> DeadLetterStoreError {
305 match error {
306 chio_kernel::ReceiptStoreError::Conflict(message) => {
307 if let Some(message) = message.strip_prefix(WRITER_CONFLICT_TAG) {
308 DeadLetterStoreError::Conflict(message.to_string())
309 } else {
310 DeadLetterStoreError::Backend(format!("receipt writer conflict: {message}"))
311 }
312 }
313 chio_kernel::ReceiptStoreError::Canonical(message) => {
314 if let Some(message) = message.strip_prefix(WRITER_INVALID_RECORD_TAG) {
315 DeadLetterStoreError::InvalidRecord(message.to_string())
316 } else {
317 DeadLetterStoreError::Backend(format!("receipt writer canonical error: {message}"))
318 }
319 }
320 chio_kernel::ReceiptStoreError::Pool(message) => {
321 if let Some(message) = message.strip_prefix(WRITER_BACKEND_TAG) {
322 DeadLetterStoreError::Backend(message.to_string())
323 } else {
324 DeadLetterStoreError::Backend(format!("receipt writer pool error: {message}"))
325 }
326 }
327 other => DeadLetterStoreError::Backend(other.to_string()),
328 }
329}
330
331impl SqliteDeadLetterStore {
332 pub fn open_with_pool(
335 pool: Pool<SqliteConnectionManager>,
336 ) -> Result<Self, DeadLetterStoreError> {
337 let connection = pool
338 .get()
339 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
340 connection
341 .execute_batch(SETTLE_DEAD_LETTERS_MIGRATION)
342 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
343 Ok(Self { pool, writer: None })
344 }
345
346 pub fn open_alongside(store: &crate::SqliteReceiptStore) -> Result<Self, DeadLetterStoreError> {
349 let writer = store.writer_handle();
350 writer
351 .run_write(|connection| {
352 connection
353 .execute_batch(SETTLE_DEAD_LETTERS_MIGRATION)
354 .map_err(chio_kernel::ReceiptStoreError::from)
355 })
356 .map_err(dead_letter_error_from_receipt_store)?;
357 Ok(Self {
358 pool: store.pool.clone(),
359 writer: Some(writer),
360 })
361 }
362
363 pub fn insert(&self, record: &DeadLetterRecord) -> Result<bool, DeadLetterStoreError> {
367 match &self.writer {
368 Some(writer) => {
369 let record = record.clone();
370 writer
371 .run_write(move |connection| {
372 insert_dead_letter_transaction(connection, &record)
373 .map_err(dead_letter_error_into_receipt_store)
374 })
375 .map_err(dead_letter_error_from_receipt_store)
376 }
377 None => {
378 let mut connection = self
379 .pool
380 .get()
381 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
382 insert_dead_letter_transaction(&mut connection, record)
383 }
384 }
385 }
386
387 pub fn get(&self, receipt_id: &str) -> Result<Option<DeadLetterRecord>, DeadLetterStoreError> {
389 let connection = self
390 .pool
391 .get()
392 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
393 read_dead_letter_on_connection(&connection, receipt_id)
394 }
395
396 pub fn list(&self) -> Result<Vec<DeadLetterRecord>, DeadLetterStoreError> {
399 let connection = self
400 .pool
401 .get()
402 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
403 let mut stmt = connection
404 .prepare(
405 "SELECT receipt_id, finalized_at, attempts, reason, pipeline_error, canonical_json \
406 FROM settle_dead_letters \
407 ORDER BY finalized_at ASC, receipt_id ASC",
408 )
409 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
410 let rows = stmt
411 .query_map([], read_stored_dead_letter_row)
412 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
413 let mut out = Vec::new();
414 for row in rows {
415 let row = row.map_err(dead_letter_row_error)?;
416 out.push(decode_stored_dead_letter(row)?);
417 }
418 Ok(out)
419 }
420
421 pub fn clear(&self, receipt_id: &str) -> Result<bool, DeadLetterStoreError> {
423 match &self.writer {
424 Some(writer) => {
425 let receipt_id = receipt_id.to_string();
426 writer
427 .run_write(move |connection| {
428 clear_dead_letter_on_connection(connection, &receipt_id)
429 .map_err(dead_letter_error_into_receipt_store)
430 })
431 .map_err(dead_letter_error_from_receipt_store)
432 }
433 None => {
434 let connection = self
435 .pool
436 .get()
437 .map_err(|err| DeadLetterStoreError::Backend(err.to_string()))?;
438 clear_dead_letter_on_connection(&connection, receipt_id)
439 }
440 }
441 }
442}
443
444#[cfg(test)]
445mod tests {
446 use super::*;
447
448 use chio_core::canonical::canonical_json_bytes;
449 use chio_settle::{DeadLetterRecord, SettlementFailureCode, SettlementFailureReason};
450 use r2d2::Pool;
451 use r2d2_sqlite::SqliteConnectionManager;
452 use rusqlite::params;
453 use tempfile::tempdir;
454
455 use chio_test_support::prelude::*;
456
457 const UNBOUNDED_V1_FIXTURE: &str = r#"{"schema":"chio.settle.dead-letter.v1","receipt_id":"old-1","finalized_at":17,"attempts":2,"reason":"rpc unavailable","pipeline_error":"settlement pipeline error: rpc unavailable"}"#;
458
459 fn pool() -> Pool<SqliteConnectionManager> {
460 let manager = SqliteConnectionManager::memory();
461 Pool::builder()
462 .max_size(2)
463 .build(manager)
464 .test_expect("test pool builds")
465 }
466
467 fn sample_record(receipt_id: &str, attempts: u32) -> DeadLetterRecord {
468 DeadLetterRecord::new(
469 receipt_id,
470 100,
471 attempts,
472 SettlementFailureReason::from_detail(SettlementFailureCode::Rpc, "connection refused"),
473 )
474 }
475
476 #[test]
477 fn migration_is_idempotent() {
478 let pool = pool();
479 SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("first open");
480 SqliteDeadLetterStore::open_with_pool(pool).test_expect("second open");
481 }
482
483 #[test]
484 fn insert_persists_and_get_round_trips() {
485 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
486 let record = sample_record("rcpt-1", 4);
487 assert!(store.insert(&record).test_expect("insert ok"));
488 let loaded = store
489 .get("rcpt-1")
490 .test_expect("get ok")
491 .test_expect("row present");
492 assert_eq!(loaded, record);
493 }
494
495 #[test]
496 fn insert_writes_rfc_8785_bytes_and_bounded_columns() {
497 let pool = pool();
498 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
499 let record = sample_record("rcpt-canonical", 4);
500 store.insert(&record).test_expect("insert succeeds");
501
502 let connection = pool.get().test_expect("connection opens");
503 let (reason, pipeline_error, canonical): (String, Option<String>, String) = connection
504 .query_row(
505 "SELECT reason, pipeline_error, canonical_json FROM settle_dead_letters WHERE receipt_id = ?1",
506 params![record.receipt_id.as_str()],
507 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
508 )
509 .test_expect("row loads");
510 let expected = canonical_json_bytes(&record).test_expect("record canonicalizes");
511
512 assert_eq!(reason, "rpc");
513 assert_eq!(pipeline_error, None);
514 assert_eq!(canonical.as_bytes(), expected);
515 assert!(canonical.starts_with("{\"attempts\":"));
516 }
517
518 #[test]
519 fn insert_rejects_an_unsupported_schema() {
520 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
521 let mut record = sample_record("rcpt-schema", 1);
522 record.schema = "chio.settle.dead-letter.v99".to_string();
523
524 assert!(matches!(
525 store.insert(&record),
526 Err(DeadLetterStoreError::InvalidRecord(message))
527 if message.contains("schema")
528 ));
529 assert!(store
530 .get("rcpt-schema")
531 .test_expect("get succeeds")
532 .is_none());
533 }
534
535 #[test]
536 fn insert_is_idempotent_on_byte_identical_replays() {
537 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
538 let record = sample_record("rcpt-2", 3);
539 assert!(store
540 .insert(&record)
541 .test_expect("first insert returns true"));
542 assert!(!store
543 .insert(&record)
544 .test_expect("second insert returns false"));
545 }
546
547 #[test]
548 fn byte_identical_replay_rejects_mismatched_sql_columns() {
549 let pool = pool();
550 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
551 let record = sample_record("rcpt-tampered-columns", 3);
552 store.insert(&record).test_expect("first insert succeeds");
553 pool.get()
554 .test_expect("connection opens")
555 .execute(
556 "UPDATE settle_dead_letters SET reason = ?1, pipeline_error = ?2 \
557 WHERE receipt_id = ?3",
558 params!["backend", "raw pipeline detail", record.receipt_id.as_str()],
559 )
560 .test_expect("row columns change");
561
562 assert!(matches!(
563 store.insert(&record),
564 Err(DeadLetterStoreError::InvalidRecord(_))
565 ));
566 assert!(matches!(
567 store.get(record.receipt_id.as_str()),
568 Err(DeadLetterStoreError::InvalidRecord(_))
569 ));
570 assert!(matches!(
571 store.list(),
572 Err(DeadLetterStoreError::InvalidRecord(_))
573 ));
574 }
575
576 #[test]
577 fn malformed_persisted_numeric_type_is_invalid_record() {
578 let pool = pool();
579 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
580 let record = sample_record("malformed-number", 1);
581 store.insert(&record).test_expect("record inserts");
582 pool.get()
583 .test_expect("connection opens")
584 .execute(
585 "UPDATE settle_dead_letters SET finalized_at = 'not-a-number' \
586 WHERE receipt_id = ?1",
587 [record.receipt_id.as_str()],
588 )
589 .test_expect("malformed numeric value persists");
590
591 assert!(matches!(
592 store.get(&record.receipt_id),
593 Err(DeadLetterStoreError::InvalidRecord(_))
594 ));
595 }
596
597 #[test]
598 fn insert_with_different_bytes_returns_conflict() {
599 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
600 let record = sample_record("rcpt-3", 2);
601 store.insert(&record).test_expect("first insert");
602 let mut conflicting = record.clone();
603 conflicting.reason =
604 SettlementFailureReason::from_detail(SettlementFailureCode::Rpc, "different reason");
605 let err = store
606 .insert(&conflicting)
607 .test_expect_err("byte-different second insert errors");
608 match err {
609 DeadLetterStoreError::Conflict(message) => {
610 assert!(message.contains("rcpt-3"));
611 }
612 other => panic!("expected Conflict, got {other:?}"),
613 }
614 }
615
616 #[test]
617 fn unbounded_v1_body_fails_closed() {
618 let pool = pool();
619 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
620 let connection = pool.get().test_expect("connection opens");
621 connection
622 .execute(
623 "INSERT INTO settle_dead_letters \
624 (receipt_id, finalized_at, attempts, reason, pipeline_error, canonical_json) \
625 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
626 params![
627 "old-1",
628 17_i64,
629 2_i64,
630 "rpc unavailable",
631 "settlement pipeline error: rpc unavailable",
632 UNBOUNDED_V1_FIXTURE,
633 ],
634 )
635 .test_expect("old row inserts");
636 drop(connection);
637
638 assert!(matches!(
639 store.get("old-1"),
640 Err(DeadLetterStoreError::InvalidRecord(_))
641 ));
642 }
643
644 #[test]
645 fn unknown_persisted_schema_fails_closed() {
646 let pool = pool();
647 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
648 let connection = pool.get().test_expect("connection opens");
649 connection
650 .execute(
651 "INSERT INTO settle_dead_letters \
652 (receipt_id, finalized_at, attempts, reason, canonical_json) \
653 VALUES (?1, ?2, ?3, ?4, ?5)",
654 params![
655 "unknown-1",
656 17_i64,
657 1_i64,
658 "rpc",
659 r#"{"schema":"chio.settle.dead-letter.v99","receipt_id":"unknown-1","finalized_at":17,"attempts":1,"reason":"rpc"}"#,
660 ],
661 )
662 .test_expect("unknown row inserts");
663 drop(connection);
664
665 assert!(matches!(
666 store.get("unknown-1"),
667 Err(DeadLetterStoreError::InvalidRecord(message))
668 if message.contains("unsupported")
669 ));
670 }
671
672 #[test]
673 fn noncanonical_body_fails_closed_on_read() {
674 let pool = pool();
675 let store = SqliteDeadLetterStore::open_with_pool(pool.clone()).test_expect("store opens");
676 let record = sample_record("noncanonical", 1);
677 let noncanonical = serde_json::to_string(&record).test_expect("record serializes");
678 assert_ne!(
679 noncanonical.as_bytes(),
680 canonical_json_bytes(&record)
681 .test_expect("record canonicalizes")
682 .as_slice()
683 );
684 pool.get()
685 .test_expect("connection opens")
686 .execute(
687 "INSERT INTO settle_dead_letters \
688 (receipt_id, finalized_at, attempts, reason, pipeline_error, canonical_json) \
689 VALUES (?1, ?2, ?3, ?4, NULL, ?5)",
690 params![
691 record.receipt_id.as_str(),
692 100_i64,
693 1_i64,
694 record.reason.code().as_str(),
695 noncanonical,
696 ],
697 )
698 .test_expect("noncanonical row inserts");
699
700 assert!(matches!(
701 store.get(record.receipt_id.as_str()),
702 Err(DeadLetterStoreError::InvalidRecord(_))
703 ));
704 }
705
706 #[test]
707 fn list_orders_by_finalization_time_then_receipt_id() {
708 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
709 let mut a = sample_record("rcpt-b", 1);
710 a.finalized_at = 5;
711 let mut b = sample_record("rcpt-a", 1);
712 b.finalized_at = 5;
713 let mut c = sample_record("rcpt-c", 1);
714 c.finalized_at = 10;
715 store.insert(&a).test_expect("insert a");
716 store.insert(&b).test_expect("insert b");
717 store.insert(&c).test_expect("insert c");
718 let listed = store.list().test_expect("list ok");
719 assert_eq!(
720 listed
721 .iter()
722 .map(|record| record.receipt_id.as_str())
723 .collect::<Vec<_>>(),
724 vec!["rcpt-a", "rcpt-b", "rcpt-c"]
725 );
726 }
727
728 #[test]
729 fn clear_removes_existing_row() {
730 let store = SqliteDeadLetterStore::open_with_pool(pool()).test_expect("store opens");
731 let record = sample_record("rcpt-4", 1);
732 store.insert(&record).test_expect("insert");
733 assert!(store.clear("rcpt-4").test_expect("clear ok"));
734 assert!(store.get("rcpt-4").test_expect("get ok").is_none());
735 assert!(!store.clear("rcpt-4").test_expect("idempotent clear ok"));
736 }
737
738 #[test]
739 fn open_alongside_routes_failures_through_receipt_writer_health() {
740 let dir = tempdir().test_expect("temporary directory creates");
741 let path = dir.path().join("dead-letter-writer.sqlite3");
742 let receipt_store =
743 crate::SqliteReceiptStore::open(&path).test_expect("receipt store opens");
744 let store = SqliteDeadLetterStore::open_alongside(&receipt_store)
745 .test_expect("dead-letter store opens");
746 assert!(store.writer.is_some());
747
748 let first = sample_record("rcpt-writer", 1);
749 let mut conflicting = first.clone();
750 conflicting.attempts = 2;
751 store.insert(&first).test_expect("first insert succeeds");
752 let failed_before = receipt_store
753 .flush_receipt_writes()
754 .test_expect("writer flushes")
755 .writer
756 .failed_total;
757 assert!(matches!(
758 store.insert(&conflicting),
759 Err(DeadLetterStoreError::Conflict(_))
760 ));
761 let failed_after = receipt_store
762 .flush_receipt_writes()
763 .test_expect("writer flushes after conflict")
764 .writer
765 .failed_total;
766
767 assert_eq!(failed_after, failed_before + 1);
768 }
769
770 #[test]
771 fn untagged_receipt_writer_errors_remain_backend_errors() {
772 for error in [
773 chio_kernel::ReceiptStoreError::Canonical("writer job panicked".to_string()),
774 chio_kernel::ReceiptStoreError::Conflict("head resync conflict".to_string()),
775 ] {
776 assert!(matches!(
777 dead_letter_error_from_receipt_store(error),
778 DeadLetterStoreError::Backend(_)
779 ));
780 }
781 }
782}