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 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}