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