Skip to main content

chio_store_sqlite/receipt_store/reports/
reconciliation.rs

1// Settlement and metered-billing reconciliation upserts and the metered-billing reconciliation report.
2
3use super::*;
4
5impl SqliteReceiptStore {
6    pub fn upsert_settlement_reconciliation(
7        &self,
8        receipt_id: &str,
9        reconciliation_state: SettlementReconciliationState,
10        note: Option<&str>,
11    ) -> Result<i64, ReceiptStoreError> {
12        let exists = self
13            .connection()?
14            .query_row(
15                "SELECT 1 FROM chio_tool_receipts WHERE receipt_id = ?1",
16                params![receipt_id],
17                |row| row.get::<_, i64>(0),
18            )
19            .optional()?;
20        if exists.is_none() {
21            return Err(ReceiptStoreError::NotFound(format!(
22                "receipt {receipt_id} does not exist"
23            )));
24        }
25
26        let updated_at = unix_timestamp_now_i64();
27        let receipt_id_owned = receipt_id.to_string();
28        let note_owned = note.map(ToString::to_string);
29        self.writer_handle().run_write(move |connection| {
30            connection.execute(
31                r#"
32                INSERT INTO settlement_reconciliations (
33                    receipt_id,
34                    reconciliation_state,
35                    note,
36                    updated_at
37                ) VALUES (?1, ?2, ?3, ?4)
38                ON CONFLICT(receipt_id) DO UPDATE SET
39                    reconciliation_state = excluded.reconciliation_state,
40                    note = excluded.note,
41                    updated_at = excluded.updated_at
42                "#,
43                params![
44                    receipt_id_owned,
45                    settlement_reconciliation_state_text(reconciliation_state),
46                    note_owned,
47                    updated_at
48                ],
49            )?;
50            Ok(())
51        })?;
52
53        Ok(updated_at)
54    }
55
56    pub fn upsert_metered_billing_reconciliation(
57        &self,
58        receipt_id: &str,
59        evidence: &MeteredBillingEvidenceRecord,
60        reconciliation_state: MeteredBillingReconciliationState,
61        note: Option<&str>,
62    ) -> Result<i64, ReceiptStoreError> {
63        let (seq, raw_json) = self
64            .connection()?
65            .query_row(
66                "SELECT seq, raw_json FROM chio_tool_receipts WHERE receipt_id = ?1",
67                params![receipt_id],
68                |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
69            )
70            .optional()?
71            .ok_or_else(|| {
72                ReceiptStoreError::NotFound(format!("receipt {receipt_id} does not exist"))
73            })?;
74        let receipt = decode_verified_chio_receipt(
75            &raw_json,
76            "persisted tool receipt",
77            Some(seq.max(0) as u64),
78        )?;
79        let governed = extract_governed_transaction_metadata(&receipt).ok_or_else(|| {
80            ReceiptStoreError::Conflict(format!(
81                "receipt {receipt_id} does not carry governed transaction metadata"
82            ))
83        })?;
84        if governed.metered_billing.is_none() {
85            return Err(ReceiptStoreError::Conflict(format!(
86                "receipt {receipt_id} does not carry metered billing context"
87            )));
88        }
89
90        let existing_receipt = self
91            .connection()?
92            .query_row(
93                r#"
94                SELECT receipt_id
95                FROM metered_billing_reconciliations
96                WHERE adapter_kind = ?1 AND evidence_id = ?2
97                "#,
98                params![
99                    &evidence.usage_evidence.evidence_kind,
100                    &evidence.usage_evidence.evidence_id
101                ],
102                |row| row.get::<_, String>(0),
103            )
104            .optional()?;
105        if let Some(existing_receipt) = existing_receipt {
106            if existing_receipt != receipt_id {
107                return Err(ReceiptStoreError::Conflict(format!(
108                    "metered billing evidence {}/{} is already attached to receipt {}",
109                    evidence.usage_evidence.evidence_kind,
110                    evidence.usage_evidence.evidence_id,
111                    existing_receipt
112                )));
113            }
114        }
115
116        let updated_at = unix_timestamp_now_i64();
117        let receipt_id_owned = receipt_id.to_string();
118        let evidence_owned = evidence.clone();
119        let note_owned = note.map(ToString::to_string);
120        self.writer_handle().run_write(move |connection| {
121            let evidence = &evidence_owned;
122            connection.execute(
123                r#"
124                INSERT INTO metered_billing_reconciliations (
125                    receipt_id,
126                    adapter_kind,
127                    evidence_id,
128                    observed_units,
129                    billed_cost_units,
130                    billed_cost_currency,
131                    evidence_sha256,
132                    recorded_at,
133                    reconciliation_state,
134                    note,
135                    updated_at
136                ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
137                ON CONFLICT(receipt_id) DO UPDATE SET
138                    adapter_kind = excluded.adapter_kind,
139                    evidence_id = excluded.evidence_id,
140                    observed_units = excluded.observed_units,
141                    billed_cost_units = excluded.billed_cost_units,
142                    billed_cost_currency = excluded.billed_cost_currency,
143                    evidence_sha256 = excluded.evidence_sha256,
144                    recorded_at = excluded.recorded_at,
145                    reconciliation_state = excluded.reconciliation_state,
146                    note = excluded.note,
147                    updated_at = excluded.updated_at
148                "#,
149                params![
150                    receipt_id_owned,
151                    &evidence.usage_evidence.evidence_kind,
152                    &evidence.usage_evidence.evidence_id,
153                    evidence.usage_evidence.observed_units as i64,
154                    evidence.billed_cost.units as i64,
155                    &evidence.billed_cost.currency,
156                    evidence.usage_evidence.evidence_sha256.as_deref(),
157                    evidence.recorded_at as i64,
158                    metered_billing_reconciliation_state_text(reconciliation_state),
159                    note_owned,
160                    updated_at
161                ],
162            )?;
163            Ok(())
164        })?;
165
166        Ok(updated_at)
167    }
168
169    pub fn query_metered_billing_reconciliation_report(
170        &self,
171        query: &OperatorReportQuery,
172    ) -> Result<MeteredBillingReconciliationReport, ReceiptStoreError> {
173        require_admin_receipt_read_context(
174            query.read_context.as_ref(),
175            "metered billing reconciliation report",
176        )?;
177        let capability_id = query.capability_id.as_deref();
178        let tool_server = query.tool_server.as_deref();
179        let tool_name = query.tool_name.as_deref();
180        let since = query.since.map(|value| value as i64);
181        let until = query.until.map(|value| value as i64);
182        let agent_subject = query.agent_subject.as_deref();
183        let row_limit = query.metered_limit_or_default();
184
185        let summary = self.query_metered_billing_summary(query)?;
186
187        let rows_sql = r#"
188            SELECT
189                r.seq,
190                r.raw_json,
191                COALESCE(r.subject_key, cl.subject_key),
192                mbr.adapter_kind,
193                mbr.evidence_id,
194                mbr.observed_units,
195                mbr.billed_cost_units,
196                mbr.billed_cost_currency,
197                mbr.evidence_sha256,
198                mbr.recorded_at,
199                COALESCE(mbr.reconciliation_state, 'open'),
200                mbr.note,
201                mbr.updated_at
202            FROM chio_tool_receipts r
203            LEFT JOIN capability_lineage cl ON r.capability_id = cl.capability_id
204            LEFT JOIN metered_billing_reconciliations mbr ON r.receipt_id = mbr.receipt_id
205            WHERE json_type(r.raw_json, '$.metadata.governed_transaction.metered_billing') = 'object'
206              AND (?1 IS NULL OR r.capability_id = ?1)
207              AND (?2 IS NULL OR r.tool_server = ?2)
208              AND (?3 IS NULL OR r.tool_name = ?3)
209              AND (?4 IS NULL OR r.timestamp >= ?4)
210              AND (?5 IS NULL OR r.timestamp <= ?5)
211              AND (?6 IS NULL OR COALESCE(r.subject_key, cl.subject_key) = ?6)
212            ORDER BY r.timestamp DESC, r.seq DESC
213            LIMIT ?7
214        "#;
215
216        let connection = self.connection()?;
217        let mut stmt = connection.prepare(rows_sql)?;
218        let rows = stmt.query_map(
219            params![
220                capability_id,
221                tool_server,
222                tool_name,
223                since,
224                until,
225                agent_subject,
226                row_limit as i64
227            ],
228            |row| {
229                Ok((
230                    row.get::<_, i64>(0)?,
231                    row.get::<_, String>(1)?,
232                    row.get::<_, Option<String>>(2)?,
233                    row.get::<_, Option<String>>(3)?,
234                    row.get::<_, Option<String>>(4)?,
235                    row.get::<_, Option<i64>>(5)?,
236                    row.get::<_, Option<i64>>(6)?,
237                    row.get::<_, Option<String>>(7)?,
238                    row.get::<_, Option<String>>(8)?,
239                    row.get::<_, Option<i64>>(9)?,
240                    row.get::<_, String>(10)?,
241                    row.get::<_, Option<String>>(11)?,
242                    row.get::<_, Option<i64>>(12)?,
243                ))
244            },
245        )?;
246
247        let mut receipts = Vec::new();
248        for row in rows {
249            let (
250                seq,
251                raw_json,
252                subject_key,
253                adapter_kind,
254                evidence_id,
255                observed_units,
256                billed_cost_units,
257                billed_cost_currency,
258                evidence_sha256,
259                recorded_at,
260                reconciliation_state_text,
261                note,
262                updated_at,
263            ) = row?;
264            let receipt = decode_verified_chio_receipt(
265                &raw_json,
266                "persisted tool receipt",
267                Some(seq.max(0) as u64),
268            )?;
269            let governed = extract_governed_transaction_metadata(&receipt).ok_or_else(|| {
270                ReceiptStoreError::Canonical(format!(
271                    "receipt {} is missing governed transaction metadata",
272                    receipt.id
273                ))
274            })?;
275            let metered = governed.metered_billing.ok_or_else(|| {
276                ReceiptStoreError::Canonical(format!(
277                    "receipt {} is missing metered billing metadata",
278                    receipt.id
279                ))
280            })?;
281            let financial = extract_financial_metadata(&receipt);
282            let evidence = metered_billing_evidence_record_from_columns(
283                adapter_kind,
284                evidence_id,
285                observed_units,
286                billed_cost_units,
287                billed_cost_currency,
288                evidence_sha256,
289                recorded_at,
290            );
291            let reconciliation_state =
292                parse_metered_billing_reconciliation_state(&reconciliation_state_text)?;
293            let analysis = analyze_metered_billing_reconciliation(
294                &metered,
295                financial.as_ref(),
296                evidence.as_ref(),
297                reconciliation_state,
298            );
299            let budget_authority = receipt.financial_budget_authority_metadata();
300
301            receipts.push(MeteredBillingReconciliationRow {
302                receipt_id: receipt.id,
303                timestamp: receipt.timestamp,
304                capability_id: receipt.capability_id,
305                subject_key,
306                tool_server: receipt.tool_server,
307                tool_name: receipt.tool_name,
308                settlement_mode: metered.settlement_mode,
309                provider: metered.quote.provider.clone(),
310                quote_id: metered.quote.quote_id.clone(),
311                billing_unit: metered.quote.billing_unit.clone(),
312                quoted_units: metered.quote.quoted_units,
313                quoted_cost: metered.quote.quoted_cost.clone(),
314                max_billed_units: metered.max_billed_units,
315                financial_cost_charged: financial.as_ref().map(|value| value.cost_charged),
316                financial_currency: financial.as_ref().map(|value| value.currency.clone()),
317                budget_authority,
318                evidence,
319                reconciliation_state,
320                action_required: analysis.action_required,
321                evidence_missing: analysis.evidence_missing,
322                exceeds_quoted_units: analysis.exceeds_quoted_units,
323                exceeds_max_billed_units: analysis.exceeds_max_billed_units,
324                exceeds_quoted_cost: analysis.exceeds_quoted_cost,
325                financial_mismatch: analysis.financial_mismatch,
326                note,
327                updated_at: updated_at.map(|value| value.max(0) as u64),
328            });
329        }
330
331        Ok(MeteredBillingReconciliationReport {
332            summary: MeteredBillingReconciliationSummary {
333                matching_receipts: summary.metered_receipts,
334                returned_receipts: receipts.len() as u64,
335                evidence_attached_receipts: summary.evidence_attached_receipts,
336                missing_evidence_receipts: summary.missing_evidence_receipts,
337                over_quoted_units_receipts: summary.over_quoted_units_receipts,
338                over_max_billed_units_receipts: summary.over_max_billed_units_receipts,
339                over_quoted_cost_receipts: summary.over_quoted_cost_receipts,
340                financial_mismatch_receipts: summary.financial_mismatch_receipts,
341                actionable_receipts: summary.actionable_receipts,
342                reconciled_receipts: summary.reconciled_receipts,
343                truncated: summary.metered_receipts > receipts.len() as u64,
344            },
345            receipts,
346        })
347    }
348}