Skip to main content

chio_store_sqlite/receipt_store/
underwriting_credit.rs

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}