Skip to main content

chio_store_sqlite/
iou_store.rs

1//! SQLite-backed persistence for IOU envelopes.
2//!
3//! The `iou_envelope` table is keyed by `receipt_id` so a finalized
4//! receipt maps to exactly one row. Re-processing the same finalized
5//! receipt is idempotent: a byte-identical envelope returns `Ok(false)`;
6//! a different envelope returns [`IouEnvelopeStoreError::Conflict`].
7//!
8//! The migration is `CREATE TABLE IF NOT EXISTS` plus
9//! `CREATE INDEX IF NOT EXISTS`, so it can run repeatedly against a
10//! receipt-store database that already holds other tables.
11
12use std::sync::Arc;
13
14use chio_core::canonical::canonical_json_bytes;
15use chio_credit::{IouEnvelope, IouEnvelopeStore, IouEnvelopeStoreError};
16use r2d2::Pool;
17use r2d2_sqlite::SqliteConnectionManager;
18use rusqlite::{params, OptionalExtension};
19
20/// SQL migration applied by [`SqliteIouEnvelopeStore::open_with_pool`]
21/// to create the `iou_envelope` table.
22pub const IOU_ENVELOPE_MIGRATION: &str = r#"
23CREATE TABLE IF NOT EXISTS iou_envelope (
24    receipt_id TEXT PRIMARY KEY,
25    iou_id TEXT NOT NULL,
26    receipt_timestamp INTEGER NOT NULL,
27    tenant_id TEXT,
28    amount_units INTEGER NOT NULL,
29    currency TEXT NOT NULL,
30    issuer_key TEXT NOT NULL,
31    canonical_json TEXT NOT NULL
32);
33CREATE INDEX IF NOT EXISTS idx_iou_envelope_receipt_timestamp
34    ON iou_envelope(receipt_timestamp);
35CREATE INDEX IF NOT EXISTS idx_iou_envelope_tenant
36    ON iou_envelope(tenant_id);
37"#;
38
39/// SQLite-backed [`IouEnvelopeStore`] implementation. Wraps an
40/// existing connection pool from [`crate::SqliteReceiptStore`] so
41/// IOU writes share the same SQLite database and journal mode.
42pub struct SqliteIouEnvelopeStore {
43    pool: Pool<SqliteConnectionManager>,
44    /// Present when opened alongside a receipt store: all writes are
45    /// serialized through the receipt store's single writer connection.
46    /// `None` only for the standalone `open_with_pool` path.
47    writer: Option<crate::receipt_store::WriterHandle>,
48}
49
50impl SqliteIouEnvelopeStore {
51    /// Open a store backed by the same pool as a sibling receipt
52    /// store. Runs the additive migration if the table is absent.
53    pub fn open_with_pool(
54        pool: Pool<SqliteConnectionManager>,
55    ) -> Result<Self, IouEnvelopeStoreError> {
56        let connection = pool
57            .get()
58            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
59        connection
60            .execute_batch(IOU_ENVELOPE_MIGRATION)
61            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
62        Ok(Self { pool, writer: None })
63    }
64
65    /// Construct the store sharing the connection pool of an
66    /// existing [`crate::SqliteReceiptStore`]. The receipt store has
67    /// already configured WAL / synchronous=FULL on every connection
68    /// out of the pool, so no additional connection setup is needed.
69    /// Writes are routed through the receipt store's writer handle so
70    /// they serialize with receipt commits on the single writer
71    /// connection; reads keep using the reader pool.
72    pub fn open_alongside(
73        store: &crate::SqliteReceiptStore,
74    ) -> Result<Self, IouEnvelopeStoreError> {
75        let writer = store.writer_handle();
76        // Run the additive migration on the writer connection so the reader
77        // pool never executes DDL.
78        writer
79            .run_write(|connection| {
80                connection
81                    .execute_batch(IOU_ENVELOPE_MIGRATION)
82                    .map_err(chio_kernel::ReceiptStoreError::from)
83            })
84            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
85        Ok(Self {
86            pool: store.pool.clone(),
87            writer: Some(writer),
88        })
89    }
90}
91
92fn encode_envelope(envelope: &IouEnvelope) -> Result<Arc<[u8]>, IouEnvelopeStoreError> {
93    let canonical = canonical_json_bytes(envelope)
94        .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
95    Ok(Arc::from(canonical.into_boxed_slice()))
96}
97
98fn decode_envelope(canonical: &str) -> Result<IouEnvelope, IouEnvelopeStoreError> {
99    serde_json::from_str(canonical).map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))
100}
101
102#[allow(clippy::too_many_arguments)]
103fn insert_envelope_on_connection(
104    connection: &rusqlite::Connection,
105    receipt_id: &str,
106    iou_id: &str,
107    receipt_ts: i64,
108    tenant_id: Option<&str>,
109    amount: i64,
110    currency: &str,
111    issuer_key_str: &str,
112    canonical_str: &str,
113) -> Result<bool, IouEnvelopeStoreError> {
114    let inserted = connection
115        .execute(
116            r#"
117            INSERT INTO iou_envelope (
118                receipt_id,
119                iou_id,
120                receipt_timestamp,
121                tenant_id,
122                amount_units,
123                currency,
124                issuer_key,
125                canonical_json
126            ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
127            ON CONFLICT(receipt_id) DO NOTHING
128            "#,
129            params![
130                receipt_id,
131                iou_id,
132                receipt_ts,
133                tenant_id,
134                amount,
135                currency,
136                issuer_key_str,
137                canonical_str,
138            ],
139        )
140        .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
141    if inserted == 1 {
142        return Ok(true);
143    }
144
145    let existing = connection
146        .query_row(
147            "SELECT canonical_json FROM iou_envelope WHERE receipt_id = ?1",
148            params![receipt_id],
149            |row| row.get::<_, String>(0),
150        )
151        .optional()
152        .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
153    match existing {
154        Some(existing_canonical) if existing_canonical == canonical_str => Ok(false),
155        Some(_) => Err(IouEnvelopeStoreError::Conflict(format!(
156            "iou_envelope row for receipt_id={receipt_id} already exists with different bytes"
157        ))),
158        None => Err(IouEnvelopeStoreError::Backend(format!(
159            "iou_envelope conflict for receipt_id={receipt_id} but no row was readable"
160        ))),
161    }
162}
163
164impl IouEnvelopeStore for SqliteIouEnvelopeStore {
165    fn insert(&self, envelope: &IouEnvelope) -> Result<bool, IouEnvelopeStoreError> {
166        let canonical_bytes = encode_envelope(envelope)?;
167        let canonical_str = std::str::from_utf8(&canonical_bytes)
168            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
169        let issuer_key_str = serde_json::to_string(&envelope.body.issuer_key)
170            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
171        let amount: i64 =
172            envelope
173                .body
174                .amount_units
175                .try_into()
176                .map_err(|err: std::num::TryFromIntError| {
177                    IouEnvelopeStoreError::Backend(err.to_string())
178                })?;
179        let receipt_ts: i64 = envelope.body.receipt_timestamp.try_into().map_err(
180            |err: std::num::TryFromIntError| IouEnvelopeStoreError::Backend(err.to_string()),
181        )?;
182
183        match &self.writer {
184            Some(writer) => {
185                let receipt_id = envelope.body.receipt_id.clone();
186                let iou_id = envelope.body.iou_id.clone();
187                let tenant_id = envelope.body.tenant_id.clone();
188                let currency = envelope.body.currency.clone();
189                let issuer_key = issuer_key_str.clone();
190                let canonical = canonical_str.to_string();
191                writer
192                    .run_write(move |connection| {
193                        // Propagate the IOU insert result to the writer.
194                        // Wrapping the inner result in `Ok(...)` would make the
195                        // writer actor record EVERY insert as committed -
196                        // refreshing `committed_total`/`last_commit_unix_ms` -
197                        // even when the inner insert returned a conflict or
198                        // backend error, so the receipt writer health telemetry
199                        // would misreport failed IOU persistence. Surface the
200                        // failure to the actor as a
201                        // `ReceiptStoreError` (fail-closed) so it is counted as
202                        // `failed_total`. The reverse mapping below restores the
203                        // original `IouEnvelopeStoreError` variant for the caller,
204                        // so a mismatched re-insert still returns `Conflict`
205                        // exactly as the standalone path and the trait contract
206                        // require.
207                        insert_envelope_on_connection(
208                            connection,
209                            &receipt_id,
210                            &iou_id,
211                            receipt_ts,
212                            tenant_id.as_deref(),
213                            amount,
214                            &currency,
215                            &issuer_key,
216                            &canonical,
217                        )
218                        .map_err(|err| match err {
219                            IouEnvelopeStoreError::Conflict(message) => {
220                                chio_kernel::ReceiptStoreError::Conflict(message)
221                            }
222                            IouEnvelopeStoreError::Backend(message) => {
223                                chio_kernel::ReceiptStoreError::Canonical(message)
224                            }
225                        })
226                    })
227                    .map_err(|err| match err {
228                        chio_kernel::ReceiptStoreError::Conflict(message) => {
229                            IouEnvelopeStoreError::Conflict(message)
230                        }
231                        chio_kernel::ReceiptStoreError::Canonical(message) => {
232                            IouEnvelopeStoreError::Backend(message)
233                        }
234                        other => IouEnvelopeStoreError::Backend(other.to_string()),
235                    })
236            }
237            None => {
238                let connection = self
239                    .pool
240                    .get()
241                    .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
242                insert_envelope_on_connection(
243                    &connection,
244                    envelope.body.receipt_id.as_str(),
245                    envelope.body.iou_id.as_str(),
246                    receipt_ts,
247                    envelope.body.tenant_id.as_deref(),
248                    amount,
249                    envelope.body.currency.as_str(),
250                    issuer_key_str.as_str(),
251                    canonical_str,
252                )
253            }
254        }
255    }
256
257    fn get_by_receipt_id(
258        &self,
259        receipt_id: &str,
260    ) -> Result<Option<IouEnvelope>, IouEnvelopeStoreError> {
261        let connection = self
262            .pool
263            .get()
264            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
265        let row = connection
266            .query_row(
267                "SELECT canonical_json FROM iou_envelope WHERE receipt_id = ?1",
268                params![receipt_id],
269                |row| row.get::<_, String>(0),
270            )
271            .optional()
272            .map_err(|err| IouEnvelopeStoreError::Backend(err.to_string()))?;
273        match row {
274            Some(canonical) => Ok(Some(decode_envelope(&canonical)?)),
275            None => Ok(None),
276        }
277    }
278}
279
280#[cfg(test)]
281#[allow(clippy::unwrap_used, clippy::expect_used)]
282mod tests {
283    use super::*;
284    use chio_core::crypto::{sha256_hex, Ed25519Backend, Keypair};
285    use chio_core::receipt::{
286        body::ChioReceipt, body::ChioReceiptBody, decision::Decision, decision::ToolCallAction,
287        economics::FinancialReceiptMetadata, economics::SettlementStatus, kinds::TrustLevel,
288        metadata::GuardEvidence,
289    };
290    use chio_credit::{CreditEvaluatorHook, LocalCreditAccount};
291    use tempfile::tempdir;
292
293    fn make_priced_receipt(kp: &Keypair, receipt_id: &str, cost: u64) -> ChioReceipt {
294        let financial = FinancialReceiptMetadata {
295            grant_index: 0,
296            cost_charged: cost,
297            currency: "USD".to_string(),
298            budget_remaining: 1000 - cost,
299            budget_total: 1000,
300            delegation_depth: 1,
301            root_budget_holder: "tenant-a".to_string(),
302            payment_reference: None,
303            settlement_status: SettlementStatus::Pending,
304            cost_breakdown: None,
305            oracle_evidence: None,
306            attempted_cost: None,
307        };
308        let body = ChioReceiptBody {
309            id: receipt_id.to_string(),
310            timestamp: 1_710_000_000,
311            capability_id: "cap-001".to_string(),
312            tool_server: "srv".to_string(),
313            tool_name: "tool".to_string(),
314            action: ToolCallAction::from_parameters(serde_json::json!({})).unwrap(),
315            decision: Some(Decision::Allow),
316            receipt_kind: Default::default(),
317            boundary_class: Default::default(),
318            observation_outcome: None,
319            tool_origin: Default::default(),
320            redaction_mode: Default::default(),
321            actor_chain: Vec::new(),
322            content_hash: sha256_hex(b"{}"),
323            policy_hash: "policy".to_string(),
324            evidence: vec![GuardEvidence {
325                guard_name: "G".to_string(),
326                verdict: true,
327                details: None,
328            }],
329            metadata: Some(serde_json::json!({"financial": financial})),
330            trust_level: TrustLevel::default(),
331            tenant_id: Some("tenant-a".to_string()),
332            kernel_key: kp.public_key(),
333            bbs_projection_version: None,
334        };
335        ChioReceipt::sign(body, kp).unwrap()
336    }
337
338    fn open_store() -> SqliteIouEnvelopeStore {
339        let dir = tempdir().unwrap();
340        let path = dir.path().join("iou.sqlite");
341        // Construct a pool directly so the test does not require a
342        // sibling receipt store.
343        let manager = SqliteConnectionManager::file(path);
344        let pool = Pool::builder().max_size(2).build(manager).unwrap();
345        // Leak the tempdir so the file outlives the test.
346        std::mem::forget(dir);
347        SqliteIouEnvelopeStore::open_with_pool(pool).unwrap()
348    }
349
350    #[test]
351    fn insert_then_get_round_trip() {
352        let kp = Keypair::generate();
353        let account = LocalCreditAccount::new_with_trusted_kernel_keys(
354            Ed25519Backend::new(kp.clone()),
355            [kp.public_key()],
356        );
357        let receipt = make_priced_receipt(&kp, "rcpt-store-1", 250);
358        let envelope = account.evaluate(&receipt).unwrap().unwrap();
359        let store = open_store();
360        assert!(store.insert(&envelope).unwrap());
361        let fetched = store
362            .get_by_receipt_id(&receipt.id)
363            .unwrap()
364            .expect("envelope was inserted");
365        assert_eq!(fetched, envelope);
366    }
367
368    #[test]
369    fn duplicate_insert_is_idempotent() {
370        let kp = Keypair::generate();
371        let account = LocalCreditAccount::new_with_trusted_kernel_keys(
372            Ed25519Backend::new(kp.clone()),
373            [kp.public_key()],
374        );
375        let receipt = make_priced_receipt(&kp, "rcpt-store-2", 100);
376        let envelope = account.evaluate(&receipt).unwrap().unwrap();
377        let store = open_store();
378        assert!(store.insert(&envelope).unwrap());
379        assert!(!store.insert(&envelope).unwrap());
380    }
381
382    #[test]
383    fn conflicting_envelope_for_same_receipt_id_errors() {
384        let kp_a = Keypair::generate();
385        let kp_b = Keypair::generate();
386        let receipt_a = make_priced_receipt(&kp_a, "rcpt-store-3", 100);
387        let env_a = LocalCreditAccount::new_with_trusted_kernel_keys(
388            Ed25519Backend::new(kp_a.clone()),
389            [kp_a.public_key()],
390        )
391        .evaluate(&receipt_a)
392        .unwrap()
393        .unwrap();
394        let env_b = LocalCreditAccount::new_with_trusted_kernel_keys(
395            Ed25519Backend::new(kp_b),
396            [kp_a.public_key()],
397        )
398        .evaluate(&receipt_a)
399        .unwrap()
400        .unwrap();
401        assert_eq!(env_a.body.receipt_id, env_b.body.receipt_id);
402        assert_ne!(env_a.body.issuer_key, env_b.body.issuer_key);
403        let store = open_store();
404        assert!(store.insert(&env_a).unwrap());
405        match store.insert(&env_b) {
406            Err(IouEnvelopeStoreError::Conflict(_)) => {}
407            other => panic!("expected Conflict, got {other:?}"),
408        }
409    }
410
411    #[test]
412    fn get_missing_returns_none() {
413        let store = open_store();
414        assert!(store.get_by_receipt_id("nope").unwrap().is_none());
415    }
416
417    #[test]
418    fn open_alongside_routes_writes_through_the_receipt_writer() {
419        let dir = tempdir().unwrap();
420        let path = dir.path().join("iou-alongside.sqlite3");
421        let receipt_store = crate::SqliteReceiptStore::open(&path).unwrap();
422        let store = SqliteIouEnvelopeStore::open_alongside(&receipt_store).unwrap();
423        assert!(
424            store.writer.is_some(),
425            "open_alongside must carry the receipt writer handle"
426        );
427
428        let kp = Keypair::generate();
429        let account = LocalCreditAccount::new_with_trusted_kernel_keys(
430            Ed25519Backend::new(kp.clone()),
431            [kp.public_key()],
432        );
433        let receipt = make_priced_receipt(&kp, "rcpt-alongside-1", 42);
434        let envelope = account.evaluate(&receipt).unwrap().unwrap();
435        assert!(store.insert(&envelope).unwrap());
436        assert!(!store.insert(&envelope).unwrap());
437        let fetched = store
438            .get_by_receipt_id(&receipt.id)
439            .unwrap()
440            .expect("envelope was inserted");
441        assert_eq!(fetched, envelope);
442        std::mem::forget(dir);
443    }
444
445    #[test]
446    fn failed_writer_routed_insert_is_recorded_as_a_writer_failure() {
447        // When an IOU store is opened alongside a
448        // receipt store, a failed insert routed through the shared writer must
449        // surface as a writer FAILURE, not be swallowed as a committed write, so
450        // the receipt writer health telemetry stays accurate. The caller still
451        // receives the original Conflict variant.
452        let dir = tempdir().unwrap();
453        let path = dir.path().join("iou-writer-failure.sqlite3");
454        let receipt_store = crate::SqliteReceiptStore::open(&path).unwrap();
455        let store = SqliteIouEnvelopeStore::open_alongside(&receipt_store).unwrap();
456        assert!(store.writer.is_some());
457
458        // Two envelopes share a receipt_id but carry different bytes (different
459        // kernel signer), so the second insert conflicts.
460        let kp_a = Keypair::generate();
461        let kp_b = Keypair::generate();
462        let receipt = make_priced_receipt(&kp_a, "rcpt-writer-fail-1", 100);
463        let env_a = LocalCreditAccount::new_with_trusted_kernel_keys(
464            Ed25519Backend::new(kp_a.clone()),
465            [kp_a.public_key()],
466        )
467        .evaluate(&receipt)
468        .unwrap()
469        .unwrap();
470        let env_b = LocalCreditAccount::new_with_trusted_kernel_keys(
471            Ed25519Backend::new(kp_b),
472            [kp_a.public_key()],
473        )
474        .evaluate(&receipt)
475        .unwrap()
476        .unwrap();
477        assert_eq!(env_a.body.receipt_id, env_b.body.receipt_id);
478        assert_ne!(env_a.body.issuer_key, env_b.body.issuer_key);
479
480        assert!(store.insert(&env_a).unwrap());
481
482        // `flush_receipt_writes` is a writer barrier, so its snapshot reflects the
483        // fully-processed job outcome (no race with `record_write_job_outcome`).
484        let failed_before = receipt_store
485            .flush_receipt_writes()
486            .unwrap()
487            .writer
488            .failed_total;
489
490        // The conflicting insert must surface as a Conflict to the caller AND be
491        // recorded as a writer failure.
492        match store.insert(&env_b) {
493            Err(IouEnvelopeStoreError::Conflict(_)) => {}
494            other => panic!("expected Conflict, got {other:?}"),
495        }
496
497        let failed_after = receipt_store
498            .flush_receipt_writes()
499            .unwrap()
500            .writer
501            .failed_total;
502        assert_eq!(
503            failed_after,
504            failed_before + 1,
505            "a failed IOU insert must increment the receipt writer failed_total"
506        );
507        std::mem::forget(dir);
508    }
509}