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