1use super::*;
2
3impl SqliteReceiptStore {
4 pub fn record_underwriting_decision(
5 &mut self,
6 decision: &SignedUnderwritingDecision,
7 ) -> Result<(), ReceiptStoreError> {
8 if !decision
9 .verify_signature()
10 .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
11 {
12 return Err(ReceiptStoreError::Conflict(
13 "underwriting decision signature verification failed".to_string(),
14 ));
15 }
16
17 let decision_owned = decision.clone();
18 self.writer_handle().run_write(move |connection| {
19 let decision = &decision_owned;
20 let artifact = &decision.body;
21 let tx = connection.transaction()?;
22 let existing = tx
23 .query_row(
24 "SELECT decision_id FROM underwriting_decisions WHERE decision_id = ?1",
25 params![artifact.decision_id],
26 |row| row.get::<_, String>(0),
27 )
28 .optional()?;
29 if existing.is_some() {
30 return Err(ReceiptStoreError::Conflict(format!(
31 "underwriting decision `{}` already exists",
32 artifact.decision_id
33 )));
34 }
35
36 if let Some(supersedes_decision_id) = artifact.supersedes_decision_id.as_deref() {
37 let state = tx
38 .query_row(
39 "SELECT lifecycle_state, superseded_by_decision_id
40 FROM underwriting_decisions
41 WHERE decision_id = ?1",
42 params![supersedes_decision_id],
43 |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?)),
44 )
45 .optional()?
46 .ok_or_else(|| {
47 ReceiptStoreError::NotFound(format!(
48 "superseded underwriting decision `{supersedes_decision_id}` not found"
49 ))
50 })?;
51 if state.0
52 != underwriting_lifecycle_state_label(
53 UnderwritingDecisionLifecycleState::Active,
54 )
55 || state.1.is_some()
56 {
57 return Err(ReceiptStoreError::Conflict(format!(
58 "underwriting decision `{supersedes_decision_id}` is not active"
59 )));
60 }
61 }
62
63 let premium_units = artifact
64 .premium
65 .quoted_amount
66 .as_ref()
67 .map(|amount| amount.units as i64);
68 tx.execute(
69 "INSERT INTO underwriting_decisions (
70 decision_id, issued_at, capability_id, subject_key, tool_server, tool_name,
71 outcome, lifecycle_state, review_state, risk_class, supersedes_decision_id,
72 superseded_by_decision_id, premium_units, raw_json, signer_key, signature
73 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, ?12, ?13, ?14, ?15)",
74 params![
75 artifact.decision_id,
76 artifact.issued_at as i64,
77 artifact.evaluation.input.filters.capability_id.as_deref(),
78 artifact.evaluation.input.filters.agent_subject.as_deref(),
79 artifact.evaluation.input.filters.tool_server.as_deref(),
80 artifact.evaluation.input.filters.tool_name.as_deref(),
81 underwriting_decision_outcome_label(artifact.evaluation.outcome),
82 underwriting_lifecycle_state_label(artifact.lifecycle_state),
83 underwriting_review_state_label(artifact.review_state),
84 underwriting_risk_class_label(artifact.evaluation.risk_class),
85 artifact.supersedes_decision_id.as_deref(),
86 premium_units,
87 serde_json::to_string(decision)?,
88 decision.signer_key.to_hex(),
89 decision.signature.to_hex(),
90 ],
91 )?;
92
93 if let Some(supersedes_decision_id) = artifact.supersedes_decision_id.as_deref() {
94 tx.execute(
95 "UPDATE underwriting_decisions
96 SET lifecycle_state = ?1, superseded_by_decision_id = ?2
97 WHERE decision_id = ?3",
98 params![
99 underwriting_lifecycle_state_label(
100 UnderwritingDecisionLifecycleState::Superseded,
101 ),
102 artifact.decision_id,
103 supersedes_decision_id,
104 ],
105 )?;
106 }
107
108 tx.commit()?;
109 Ok(())
110 })
111 }
112
113 pub fn create_underwriting_appeal(
114 &mut self,
115 request: &UnderwritingAppealCreateRequest,
116 ) -> Result<UnderwritingAppealRecord, ReceiptStoreError> {
117 let request_owned = request.clone();
118 self.writer_handle().run_write(move |connection| {
119 let request = &request_owned;
120 let tx = connection.transaction()?;
121 let exists = tx
122 .query_row(
123 "SELECT decision_id FROM underwriting_decisions WHERE decision_id = ?1",
124 params![request.decision_id],
125 |row| row.get::<_, String>(0),
126 )
127 .optional()?;
128 if exists.is_none() {
129 return Err(ReceiptStoreError::NotFound(format!(
130 "underwriting decision `{}` not found",
131 request.decision_id
132 )));
133 }
134 let open_appeal = tx
135 .query_row(
136 "SELECT appeal_id FROM underwriting_appeals
137 WHERE decision_id = ?1 AND status = ?2",
138 params![
139 request.decision_id,
140 underwriting_appeal_status_label(UnderwritingAppealStatus::Open)
141 ],
142 |row| row.get::<_, String>(0),
143 )
144 .optional()?;
145 if let Some(appeal_id) = open_appeal {
146 return Err(ReceiptStoreError::Conflict(format!(
147 "underwriting decision `{}` already has open appeal `{appeal_id}`",
148 request.decision_id
149 )));
150 }
151
152 let created_at = unix_now();
153 let appeal_id = format!(
154 "uwa-{}",
155 chio_core::sha256_hex(
156 &canonical_json_bytes(&(
157 &request.decision_id,
158 &request.requested_by,
159 &request.reason,
160 &request.note,
161 created_at,
162 ))
163 .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
164 )
165 );
166 let record = UnderwritingAppealRecord {
167 schema: chio_kernel::UNDERWRITING_APPEAL_SCHEMA.to_string(),
168 appeal_id: appeal_id.clone(),
169 decision_id: request.decision_id.clone(),
170 requested_by: request.requested_by.clone(),
171 reason: request.reason.clone(),
172 status: UnderwritingAppealStatus::Open,
173 created_at,
174 updated_at: created_at,
175 note: request.note.clone(),
176 resolved_by: None,
177 replacement_decision_id: None,
178 };
179 tx.execute(
180 "INSERT INTO underwriting_appeals (
181 appeal_id, decision_id, requested_by, reason, status, note,
182 created_at, updated_at, resolved_by, replacement_decision_id
183 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, NULL, NULL)",
184 params![
185 record.appeal_id,
186 record.decision_id,
187 record.requested_by,
188 record.reason,
189 underwriting_appeal_status_label(record.status),
190 record.note.as_deref(),
191 record.created_at as i64,
192 record.updated_at as i64,
193 ],
194 )?;
195 tx.commit()?;
196 Ok(record)
197 })
198 }
199
200 pub fn resolve_underwriting_appeal(
201 &mut self,
202 request: &UnderwritingAppealResolveRequest,
203 ) -> Result<UnderwritingAppealRecord, ReceiptStoreError> {
204 let request_owned = request.clone();
205 self.writer_handle().run_write(move |connection| {
206 let request = &request_owned;
207 let tx = connection.transaction()?;
208 let mut record = query_underwriting_appeal(&tx, &request.appeal_id)?.ok_or_else(|| {
209 ReceiptStoreError::NotFound(format!(
210 "underwriting appeal `{}` not found",
211 request.appeal_id
212 ))
213 })?;
214 if record.status != UnderwritingAppealStatus::Open {
215 return Err(ReceiptStoreError::Conflict(format!(
216 "underwriting appeal `{}` is already resolved",
217 request.appeal_id
218 )));
219 }
220
221 if let Some(replacement_decision_id) = request.replacement_decision_id.as_deref() {
222 if request.resolution != UnderwritingAppealResolution::Accepted {
223 return Err(ReceiptStoreError::Conflict(
224 "replacement underwriting decision may only be linked when an appeal is accepted"
225 .to_string(),
226 ));
227 }
228 let replacement = tx
229 .query_row(
230 "SELECT supersedes_decision_id FROM underwriting_decisions WHERE decision_id = ?1",
231 params![replacement_decision_id],
232 |row| row.get::<_, Option<String>>(0),
233 )
234 .optional()?
235 .ok_or_else(|| {
236 ReceiptStoreError::NotFound(format!(
237 "replacement underwriting decision `{replacement_decision_id}` not found"
238 ))
239 })?;
240 if replacement.as_deref() != Some(record.decision_id.as_str()) {
241 return Err(ReceiptStoreError::Conflict(format!(
242 "replacement underwriting decision `{replacement_decision_id}` does not supersede `{}`",
243 record.decision_id
244 )));
245 }
246 }
247
248 record.status = match request.resolution {
249 UnderwritingAppealResolution::Accepted => UnderwritingAppealStatus::Accepted,
250 UnderwritingAppealResolution::Rejected => UnderwritingAppealStatus::Rejected,
251 };
252 record.updated_at = unix_now();
253 record.note = request.note.clone().or(record.note);
254 record.resolved_by = Some(request.resolved_by.clone());
255 record.replacement_decision_id = request.replacement_decision_id.clone();
256
257 tx.execute(
258 "UPDATE underwriting_appeals
259 SET status = ?1, note = ?2, updated_at = ?3, resolved_by = ?4,
260 replacement_decision_id = ?5
261 WHERE appeal_id = ?6",
262 params![
263 underwriting_appeal_status_label(record.status),
264 record.note.as_deref(),
265 record.updated_at as i64,
266 record.resolved_by.as_deref(),
267 record.replacement_decision_id.as_deref(),
268 record.appeal_id,
269 ],
270 )?;
271 tx.commit()?;
272 Ok(record)
273 })
274 }
275
276 pub fn query_underwriting_decisions(
277 &self,
278 query: &UnderwritingDecisionQuery,
279 ) -> Result<UnderwritingDecisionListReport, ReceiptStoreError> {
280 let normalized = query.normalized();
281 let appeals = self.load_underwriting_appeals_by_decision()?;
282 let connection = self.connection()?;
283 let mut statement = connection.prepare(
284 "SELECT raw_json, lifecycle_state
285 FROM underwriting_decisions
286 ORDER BY issued_at DESC, decision_id DESC",
287 )?;
288 let rows = statement.query_map([], |row| {
289 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
290 })?;
291
292 let mut matching_decisions = 0_u64;
293 let mut active_decisions = 0_u64;
294 let mut superseded_decisions = 0_u64;
295 let mut open_appeals = 0_u64;
296 let mut accepted_appeals = 0_u64;
297 let mut rejected_appeals = 0_u64;
298 let mut total_quoted_premium_units = 0_u64;
299 let mut total_quoted_premium_currency = None;
300 let mut quoted_premium_totals_by_currency = BTreeMap::<String, u64>::new();
301 let mut decisions = Vec::new();
302
303 for row in rows {
304 let (raw_json, lifecycle_state_raw) = row?;
305 let decision: SignedUnderwritingDecision = serde_json::from_str(&raw_json)?;
306 let lifecycle_state =
307 parse_underwriting_lifecycle_state(&lifecycle_state_raw).map_err(|error| {
308 ReceiptStoreError::Conflict(format!(
309 "invalid underwriting decision lifecycle state `{lifecycle_state_raw}`: {error}"
310 ))
311 })?;
312 let decision_appeals = appeals
313 .get(decision.body.decision_id.as_str())
314 .cloned()
315 .unwrap_or_default();
316 let latest_appeal = decision_appeals.iter().max_by(|left, right| {
317 left.updated_at
318 .cmp(&right.updated_at)
319 .then(left.appeal_id.cmp(&right.appeal_id))
320 });
321 if !underwriting_decision_matches_query(
322 &decision,
323 lifecycle_state,
324 latest_appeal.map(|appeal| appeal.status),
325 &normalized,
326 ) {
327 continue;
328 }
329
330 matching_decisions += 1;
331 match lifecycle_state {
332 UnderwritingDecisionLifecycleState::Active => active_decisions += 1,
333 UnderwritingDecisionLifecycleState::Superseded => superseded_decisions += 1,
334 }
335 for appeal in &decision_appeals {
336 match appeal.status {
337 UnderwritingAppealStatus::Open => open_appeals += 1,
338 UnderwritingAppealStatus::Accepted => accepted_appeals += 1,
339 UnderwritingAppealStatus::Rejected => rejected_appeals += 1,
340 }
341 }
342 if let Some(quoted_amount) = decision.body.premium.quoted_amount.as_ref() {
343 let total = quoted_premium_totals_by_currency
344 .entry(quoted_amount.currency.clone())
345 .or_insert(0);
346 *total = total.saturating_add(quoted_amount.units);
347 }
348
349 if decisions.len() < normalized.limit_or_default() {
350 let open_appeal_count = decision_appeals
351 .iter()
352 .filter(|appeal| appeal.status == UnderwritingAppealStatus::Open)
353 .count() as u64;
354 decisions.push(UnderwritingDecisionRow {
355 decision,
356 lifecycle_state,
357 open_appeal_count,
358 latest_appeal_id: latest_appeal.map(|appeal| appeal.appeal_id.clone()),
359 latest_appeal_status: latest_appeal.map(|appeal| appeal.status),
360 });
361 }
362 }
363
364 if quoted_premium_totals_by_currency.len() == 1 {
365 if let Some((currency, units)) = quoted_premium_totals_by_currency
366 .iter()
367 .next()
368 .map(|(currency, units)| (currency.clone(), *units))
369 {
370 total_quoted_premium_units = units;
371 total_quoted_premium_currency = Some(currency);
372 }
373 }
374
375 Ok(UnderwritingDecisionListReport {
376 generated_at: unix_now(),
377 filters: normalized,
378 summary: UnderwritingDecisionSummary {
379 matching_decisions,
380 returned_decisions: decisions.len() as u64,
381 active_decisions,
382 superseded_decisions,
383 open_appeals,
384 accepted_appeals,
385 rejected_appeals,
386 total_quoted_premium_units,
387 total_quoted_premium_currency,
388 quoted_premium_totals_by_currency,
389 },
390 decisions,
391 })
392 }
393
394 pub fn record_credit_facility(
395 &mut self,
396 facility: &SignedCreditFacility,
397 ) -> Result<(), ReceiptStoreError> {
398 if !facility
399 .verify_signature()
400 .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
401 {
402 return Err(ReceiptStoreError::Conflict(
403 "credit facility signature verification failed".to_string(),
404 ));
405 }
406
407 let facility_owned = facility.clone();
408 self.writer_handle().run_write(move |connection| {
409 let facility = &facility_owned;
410 let artifact = &facility.body;
411 let tx = connection.transaction()?;
412 let existing = tx
413 .query_row(
414 "SELECT facility_id FROM credit_facilities WHERE facility_id = ?1",
415 params![artifact.facility_id],
416 |row| row.get::<_, String>(0),
417 )
418 .optional()?;
419 if existing.is_some() {
420 return Err(ReceiptStoreError::Conflict(format!(
421 "credit facility `{}` already exists",
422 artifact.facility_id
423 )));
424 }
425
426 if let Some(supersedes_facility_id) = artifact.supersedes_facility_id.as_deref() {
427 let state = tx
428 .query_row(
429 "SELECT lifecycle_state, superseded_by_facility_id, expires_at
430 FROM credit_facilities
431 WHERE facility_id = ?1",
432 params![supersedes_facility_id],
433 |row| {
434 Ok((
435 row.get::<_, String>(0)?,
436 row.get::<_, Option<String>>(1)?,
437 row.get::<_, i64>(2)?,
438 ))
439 },
440 )
441 .optional()?
442 .ok_or_else(|| {
443 ReceiptStoreError::NotFound(format!(
444 "superseded credit facility `{supersedes_facility_id}` not found"
445 ))
446 })?;
447 if state.0
448 != credit_facility_lifecycle_state_label(CreditFacilityLifecycleState::Active)
449 || state.1.is_some()
450 || state.2.max(0) as u64 <= unix_now()
451 {
452 return Err(ReceiptStoreError::Conflict(format!(
453 "credit facility `{supersedes_facility_id}` is not active"
454 )));
455 }
456 }
457
458 tx.execute(
459 "INSERT INTO credit_facilities (
460 facility_id, issued_at, expires_at, capability_id, subject_key, tool_server,
461 tool_name, disposition, lifecycle_state, supersedes_facility_id,
462 superseded_by_facility_id, raw_json, signer_key, signature
463 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL, ?11, ?12, ?13)",
464 params![
465 artifact.facility_id,
466 artifact.issued_at as i64,
467 artifact.expires_at as i64,
468 artifact.report.filters.capability_id.as_deref(),
469 artifact.report.filters.agent_subject.as_deref(),
470 artifact.report.filters.tool_server.as_deref(),
471 artifact.report.filters.tool_name.as_deref(),
472 credit_facility_disposition_label(artifact.report.disposition),
473 credit_facility_lifecycle_state_label(artifact.lifecycle_state),
474 artifact.supersedes_facility_id.as_deref(),
475 serde_json::to_string(facility)?,
476 facility.signer_key.to_hex(),
477 facility.signature.to_hex(),
478 ],
479 )?;
480
481 if let Some(supersedes_facility_id) = artifact.supersedes_facility_id.as_deref() {
482 tx.execute(
483 "UPDATE credit_facilities
484 SET lifecycle_state = ?1, superseded_by_facility_id = ?2
485 WHERE facility_id = ?3",
486 params![
487 credit_facility_lifecycle_state_label(
488 CreditFacilityLifecycleState::Superseded,
489 ),
490 artifact.facility_id,
491 supersedes_facility_id,
492 ],
493 )?;
494 }
495
496 tx.commit()?;
497 Ok(())
498 })
499 }
500
501 pub fn query_credit_facilities(
502 &self,
503 query: &CreditFacilityListQuery,
504 ) -> Result<CreditFacilityListReport, ReceiptStoreError> {
505 let normalized = query.normalized();
506 let now = unix_now();
507 let connection = self.connection()?;
508 let mut statement = connection.prepare(
509 "SELECT raw_json, lifecycle_state, superseded_by_facility_id
510 FROM credit_facilities
511 ORDER BY issued_at DESC, facility_id DESC",
512 )?;
513 let rows = statement.query_map([], |row| {
514 Ok((
515 row.get::<_, String>(0)?,
516 row.get::<_, String>(1)?,
517 row.get::<_, Option<String>>(2)?,
518 ))
519 })?;
520
521 let mut matching_facilities = 0_u64;
522 let mut active_facilities = 0_u64;
523 let mut superseded_facilities = 0_u64;
524 let mut denied_facilities = 0_u64;
525 let mut expired_facilities = 0_u64;
526 let mut granted_facilities = 0_u64;
527 let mut manual_review_facilities = 0_u64;
528 let mut facilities = Vec::new();
529
530 for row in rows {
531 let (raw_json, lifecycle_state_raw, superseded_by_facility_id) = row?;
532 let facility: SignedCreditFacility = serde_json::from_str(&raw_json)?;
533 let persisted_lifecycle = parse_credit_facility_lifecycle_state(&lifecycle_state_raw)
534 .map_err(|error| {
535 ReceiptStoreError::Conflict(format!(
536 "invalid credit facility lifecycle state `{lifecycle_state_raw}`: {error}"
537 ))
538 })?;
539 let lifecycle_state =
540 effective_credit_facility_lifecycle_state(&facility, persisted_lifecycle, now);
541 if !credit_facility_matches_query(&facility, lifecycle_state, &normalized) {
542 continue;
543 }
544
545 matching_facilities += 1;
546 match lifecycle_state {
547 CreditFacilityLifecycleState::Active => active_facilities += 1,
548 CreditFacilityLifecycleState::Superseded => superseded_facilities += 1,
549 CreditFacilityLifecycleState::Denied => denied_facilities += 1,
550 CreditFacilityLifecycleState::Expired => expired_facilities += 1,
551 }
552 match facility.body.report.disposition {
553 CreditFacilityDisposition::Grant => granted_facilities += 1,
554 CreditFacilityDisposition::ManualReview => manual_review_facilities += 1,
555 CreditFacilityDisposition::Deny => {}
556 }
557
558 if facilities.len() < normalized.limit_or_default() {
559 facilities.push(CreditFacilityRow {
560 facility,
561 lifecycle_state,
562 superseded_by_facility_id,
563 });
564 }
565 }
566
567 Ok(CreditFacilityListReport {
568 schema: CREDIT_FACILITY_LIST_REPORT_SCHEMA.to_string(),
569 generated_at: unix_now(),
570 query: normalized,
571 summary: CreditFacilityListSummary {
572 matching_facilities,
573 returned_facilities: facilities.len() as u64,
574 active_facilities,
575 superseded_facilities,
576 denied_facilities,
577 expired_facilities,
578 granted_facilities,
579 manual_review_facilities,
580 },
581 facilities,
582 })
583 }
584
585 pub fn record_credit_bond(&mut self, bond: &SignedCreditBond) -> Result<(), ReceiptStoreError> {
586 if !bond
587 .verify_signature()
588 .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
589 {
590 return Err(ReceiptStoreError::Conflict(
591 "credit bond signature verification failed".to_string(),
592 ));
593 }
594
595 let bond_owned = bond.clone();
596 self.writer_handle().run_write(move |connection| {
597 let bond = &bond_owned;
598 let artifact = &bond.body;
599 let tx = connection.transaction()?;
600 let existing = tx
601 .query_row(
602 "SELECT bond_id FROM credit_bonds WHERE bond_id = ?1",
603 params![artifact.bond_id],
604 |row| row.get::<_, String>(0),
605 )
606 .optional()?;
607 if existing.is_some() {
608 return Err(ReceiptStoreError::Conflict(format!(
609 "credit bond `{}` already exists",
610 artifact.bond_id
611 )));
612 }
613
614 if let Some(supersedes_bond_id) = artifact.supersedes_bond_id.as_deref() {
615 let state = tx
616 .query_row(
617 "SELECT lifecycle_state, superseded_by_bond_id, expires_at
618 FROM credit_bonds
619 WHERE bond_id = ?1",
620 params![supersedes_bond_id],
621 |row| {
622 Ok((
623 row.get::<_, String>(0)?,
624 row.get::<_, Option<String>>(1)?,
625 row.get::<_, i64>(2)?,
626 ))
627 },
628 )
629 .optional()?
630 .ok_or_else(|| {
631 ReceiptStoreError::NotFound(format!(
632 "superseded credit bond `{supersedes_bond_id}` not found"
633 ))
634 })?;
635 if state.0 != credit_bond_lifecycle_state_label(CreditBondLifecycleState::Active)
636 || state.1.is_some()
637 || state.2.max(0) as u64 <= unix_now()
638 {
639 return Err(ReceiptStoreError::Conflict(format!(
640 "credit bond `{supersedes_bond_id}` is not active"
641 )));
642 }
643 }
644
645 tx.execute(
646 "INSERT INTO credit_bonds (
647 bond_id, issued_at, expires_at, facility_id, capability_id, subject_key,
648 tool_server, tool_name, disposition, lifecycle_state, supersedes_bond_id,
649 superseded_by_bond_id, raw_json, signer_key, signature
650 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, ?12, ?13, ?14)",
651 params![
652 artifact.bond_id,
653 artifact.issued_at as i64,
654 artifact.expires_at as i64,
655 artifact.report.latest_facility_id.as_deref(),
656 artifact.report.filters.capability_id.as_deref(),
657 artifact.report.filters.agent_subject.as_deref(),
658 artifact.report.filters.tool_server.as_deref(),
659 artifact.report.filters.tool_name.as_deref(),
660 credit_bond_disposition_label(artifact.report.disposition),
661 credit_bond_lifecycle_state_label(artifact.lifecycle_state),
662 artifact.supersedes_bond_id.as_deref(),
663 serde_json::to_string(bond)?,
664 bond.signer_key.to_hex(),
665 bond.signature.to_hex(),
666 ],
667 )?;
668
669 if let Some(supersedes_bond_id) = artifact.supersedes_bond_id.as_deref() {
670 tx.execute(
671 "UPDATE credit_bonds
672 SET lifecycle_state = ?1, superseded_by_bond_id = ?2
673 WHERE bond_id = ?3",
674 params![
675 credit_bond_lifecycle_state_label(CreditBondLifecycleState::Superseded),
676 artifact.bond_id,
677 supersedes_bond_id,
678 ],
679 )?;
680 }
681
682 tx.commit()?;
683 Ok(())
684 })
685 }
686
687 pub fn query_credit_bonds(
688 &self,
689 query: &CreditBondListQuery,
690 ) -> Result<CreditBondListReport, ReceiptStoreError> {
691 let normalized = query.normalized();
692 let now = unix_now();
693 let connection = self.connection()?;
694 let mut statement = connection.prepare(
695 "SELECT raw_json, lifecycle_state, superseded_by_bond_id
696 FROM credit_bonds
697 ORDER BY issued_at DESC, bond_id DESC",
698 )?;
699 let rows = statement.query_map([], |row| {
700 Ok((
701 row.get::<_, String>(0)?,
702 row.get::<_, String>(1)?,
703 row.get::<_, Option<String>>(2)?,
704 ))
705 })?;
706
707 let mut matching_bonds = 0_u64;
708 let mut active_bonds = 0_u64;
709 let mut superseded_bonds = 0_u64;
710 let mut released_bonds = 0_u64;
711 let mut impaired_bonds = 0_u64;
712 let mut expired_bonds = 0_u64;
713 let mut locked_bonds = 0_u64;
714 let mut held_bonds = 0_u64;
715 let mut bonds = Vec::new();
716
717 for row in rows {
718 let (raw_json, lifecycle_state_raw, superseded_by_bond_id) = row?;
719 let bond: SignedCreditBond = serde_json::from_str(&raw_json)?;
720 let persisted_lifecycle = parse_credit_bond_lifecycle_state(&lifecycle_state_raw)
721 .map_err(|error| {
722 ReceiptStoreError::Conflict(format!(
723 "invalid credit bond lifecycle state `{lifecycle_state_raw}`: {error}"
724 ))
725 })?;
726 let lifecycle_state =
727 effective_credit_bond_lifecycle_state(&bond, persisted_lifecycle, now);
728 if !credit_bond_matches_query(&bond, lifecycle_state, &normalized) {
729 continue;
730 }
731
732 matching_bonds += 1;
733 match lifecycle_state {
734 CreditBondLifecycleState::Active => active_bonds += 1,
735 CreditBondLifecycleState::Superseded => superseded_bonds += 1,
736 CreditBondLifecycleState::Released => released_bonds += 1,
737 CreditBondLifecycleState::Impaired => impaired_bonds += 1,
738 CreditBondLifecycleState::Expired => expired_bonds += 1,
739 }
740 match bond.body.report.disposition {
741 CreditBondDisposition::Lock => locked_bonds += 1,
742 CreditBondDisposition::Hold => held_bonds += 1,
743 CreditBondDisposition::Release | CreditBondDisposition::Impair => {}
744 }
745
746 if bonds.len() < normalized.limit_or_default() {
747 bonds.push(CreditBondRow {
748 bond,
749 lifecycle_state,
750 superseded_by_bond_id,
751 });
752 }
753 }
754
755 Ok(CreditBondListReport {
756 schema: CREDIT_BOND_LIST_REPORT_SCHEMA.to_string(),
757 generated_at: unix_now(),
758 query: normalized,
759 summary: CreditBondListSummary {
760 matching_bonds,
761 returned_bonds: bonds.len() as u64,
762 active_bonds,
763 superseded_bonds,
764 released_bonds,
765 impaired_bonds,
766 expired_bonds,
767 locked_bonds,
768 held_bonds,
769 },
770 bonds,
771 })
772 }
773
774 pub fn record_credit_loss_lifecycle(
775 &mut self,
776 event: &SignedCreditLossLifecycle,
777 ) -> Result<(), ReceiptStoreError> {
778 if !event
779 .verify_signature()
780 .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
781 {
782 return Err(ReceiptStoreError::Conflict(
783 "credit loss lifecycle signature verification failed".to_string(),
784 ));
785 }
786
787 let event_owned = event.clone();
788 self.writer_handle().run_write(move |connection| {
789 let event = &event_owned;
790 let artifact = &event.body;
791 let tx = connection.transaction()?;
792 let existing = tx
793 .query_row(
794 "SELECT event_id FROM credit_loss_lifecycle WHERE event_id = ?1",
795 params![artifact.event_id],
796 |row| row.get::<_, String>(0),
797 )
798 .optional()?;
799 if existing.is_some() {
800 return Err(ReceiptStoreError::Conflict(format!(
801 "credit loss lifecycle `{}` already exists",
802 artifact.event_id
803 )));
804 }
805
806 let bond_exists = tx
807 .query_row(
808 "SELECT bond_id FROM credit_bonds WHERE bond_id = ?1",
809 params![artifact.bond_id],
810 |row| row.get::<_, String>(0),
811 )
812 .optional()?;
813 if bond_exists.is_none() {
814 return Err(ReceiptStoreError::NotFound(format!(
815 "credit bond `{}` not found",
816 artifact.bond_id
817 )));
818 }
819
820 tx.execute(
821 "INSERT INTO credit_loss_lifecycle (
822 event_id, issued_at, bond_id, facility_id, capability_id, subject_key,
823 tool_server, tool_name, event_kind, projected_bond_lifecycle_state,
824 raw_json, signer_key, signature
825 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
826 params![
827 artifact.event_id,
828 artifact.issued_at as i64,
829 artifact.bond_id,
830 artifact.report.summary.facility_id.as_deref(),
831 artifact.report.summary.capability_id.as_deref(),
832 artifact.report.summary.agent_subject.as_deref(),
833 artifact.report.summary.tool_server.as_deref(),
834 artifact.report.summary.tool_name.as_deref(),
835 credit_loss_lifecycle_event_kind_label(artifact.event_kind),
836 credit_bond_lifecycle_state_label(artifact.projected_bond_lifecycle_state),
837 serde_json::to_string(event)?,
838 event.signer_key.to_hex(),
839 event.signature.to_hex(),
840 ],
841 )?;
842
843 tx.execute(
844 "UPDATE credit_bonds
845 SET lifecycle_state = ?1
846 WHERE bond_id = ?2",
847 params![
848 credit_bond_lifecycle_state_label(artifact.projected_bond_lifecycle_state),
849 artifact.bond_id,
850 ],
851 )?;
852
853 tx.commit()?;
854 Ok(())
855 })
856 }
857
858 pub fn query_credit_loss_lifecycle(
859 &self,
860 query: &CreditLossLifecycleListQuery,
861 ) -> Result<CreditLossLifecycleListReport, ReceiptStoreError> {
862 let normalized = query.normalized();
863 let connection = self.connection()?;
864 let mut statement = connection.prepare(
865 "SELECT raw_json
866 FROM credit_loss_lifecycle
867 ORDER BY issued_at DESC, event_id DESC",
868 )?;
869 let rows = statement.query_map([], |row| row.get::<_, String>(0))?;
870
871 let mut matching_events = 0_u64;
872 let mut delinquency_events = 0_u64;
873 let mut recovery_events = 0_u64;
874 let mut reserve_release_events = 0_u64;
875 let mut reserve_slash_events = 0_u64;
876 let mut write_off_events = 0_u64;
877 let mut events = Vec::new();
878
879 for row in rows {
880 let raw_json = row?;
881 let event: SignedCreditLossLifecycle = serde_json::from_str(&raw_json)?;
882 let body = &event.body;
883 let summary = &body.report.summary;
884 if normalized
885 .event_id
886 .as_deref()
887 .is_some_and(|value| value != body.event_id)
888 {
889 continue;
890 }
891 if normalized
892 .bond_id
893 .as_deref()
894 .is_some_and(|value| value != body.bond_id)
895 {
896 continue;
897 }
898 if normalized
899 .facility_id
900 .as_deref()
901 .is_some_and(|value| summary.facility_id.as_deref() != Some(value))
902 {
903 continue;
904 }
905 if normalized
906 .capability_id
907 .as_deref()
908 .is_some_and(|value| summary.capability_id.as_deref() != Some(value))
909 {
910 continue;
911 }
912 if normalized
913 .agent_subject
914 .as_deref()
915 .is_some_and(|value| summary.agent_subject.as_deref() != Some(value))
916 {
917 continue;
918 }
919 if normalized
920 .tool_server
921 .as_deref()
922 .is_some_and(|value| summary.tool_server.as_deref() != Some(value))
923 {
924 continue;
925 }
926 if normalized
927 .tool_name
928 .as_deref()
929 .is_some_and(|value| summary.tool_name.as_deref() != Some(value))
930 {
931 continue;
932 }
933 if normalized
934 .event_kind
935 .is_some_and(|value| value != body.event_kind)
936 {
937 continue;
938 }
939
940 matching_events = matching_events.saturating_add(1);
941 match body.event_kind {
942 CreditLossLifecycleEventKind::Delinquency => {
943 delinquency_events = delinquency_events.saturating_add(1);
944 }
945 CreditLossLifecycleEventKind::Recovery => {
946 recovery_events = recovery_events.saturating_add(1);
947 }
948 CreditLossLifecycleEventKind::ReserveRelease => {
949 reserve_release_events = reserve_release_events.saturating_add(1);
950 }
951 CreditLossLifecycleEventKind::ReserveSlash => {
952 reserve_slash_events = reserve_slash_events.saturating_add(1);
953 }
954 CreditLossLifecycleEventKind::WriteOff => {
955 write_off_events = write_off_events.saturating_add(1);
956 }
957 }
958
959 if events.len() < normalized.limit_or_default() {
960 events.push(CreditLossLifecycleRow { event });
961 }
962 }
963
964 Ok(CreditLossLifecycleListReport {
965 schema: CREDIT_LOSS_LIFECYCLE_LIST_REPORT_SCHEMA.to_string(),
966 generated_at: unix_now(),
967 query: normalized,
968 summary: CreditLossLifecycleListSummary {
969 matching_events,
970 returned_events: events.len() as u64,
971 delinquency_events,
972 recovery_events,
973 reserve_release_events,
974 reserve_slash_events,
975 write_off_events,
976 },
977 events,
978 })
979 }
980}