1use 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
20pub 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
39pub struct SqliteIouEnvelopeStore {
43 pool: Pool<SqliteConnectionManager>,
44 writer: Option<crate::receipt_store::WriterHandle>,
48}
49
50impl SqliteIouEnvelopeStore {
51 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 pub fn open_alongside(
73 store: &crate::SqliteReceiptStore,
74 ) -> Result<Self, IouEnvelopeStoreError> {
75 let writer = store.writer_handle();
76 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 insert_envelope_on_connection(
208 connection,
209 &receipt_id,
210 &iou_id,
211 receipt_ts,
212 tenant_id.as_deref(),
213 amount,
214 ¤cy,
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 let manager = SqliteConnectionManager::file(path);
344 let pool = Pool::builder().max_size(2).build(manager).unwrap();
345 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 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 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 let failed_before = receipt_store
485 .flush_receipt_writes()
486 .unwrap()
487 .writer
488 .failed_total;
489
490 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}