1use std::sync::{Arc, Mutex, MutexGuard};
2
3use chio_core::canonical::canonical_json_bytes;
4use chio_core::sha256_hex;
5use chio_kernel::admission_operation::{
6 AdmissionOperationId, AdmissionOperationState, AdmissionOperationStoreError,
7 AdmissionOperationV1, AdmissionRecoveryLease, StoreMutationFence,
8};
9use chio_kernel::tool_outcome::{
10 CanonicalInvocationBlobV1, CanonicalResolvedOutputBlobV1,
11 PersistedPostReturnEvaluationRecordV1, PersistedToolOutcomeRecordV1,
12 PostReturnEvaluationRecordV1, QualifiedToolOutcomeStore, RawInvocationOutcomeV1,
13 ToolOutcomeInsertResultV1, ToolOutcomeRecordV1, ToolOutcomeStore, ToolOutcomeStoreError,
14};
15use rusqlite::{params, Connection, OptionalExtension, Transaction, TransactionBehavior};
16use serde::Serialize;
17
18use crate::admission_operation_store::{
19 advance_tool_outcome_tx, append_participant_update_tx, load_operation_for_participant_tx,
20 verify_active_owner, verify_trusted_time,
21};
22use crate::serving_owner::SqliteServingOwner;
23
24const TOOL_OUTCOME_SCHEMA_KEY: &str = "tool_outcome";
25pub(crate) const TOOL_OUTCOME_SUPPORTED_SCHEMA_VERSION: i32 = 2;
26const TOOL_OUTCOME_SCHEMA_ANCHORS: &[&str] = &[
27 "tool_outcomes",
28 "admission_operations",
29 "chio_serving_owner",
30];
31const TOOL_OUTCOME_SCHEMA: &str = include_str!("tool_outcome_store.sql");
32const MAX_OUTCOME_RECORD_BYTES: usize = 1024 * 1024;
33const MAX_EVALUATION_RECORD_BYTES: usize = 64 * 1024 * 1024;
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
37pub struct ToolOutcomeCompactionSummary {
38 pub compacted: u64,
40 pub retained_live: u64,
43}
44
45enum StoredInvocationBlob {
48 Present(CanonicalInvocationBlobV1),
49 Compacted,
50}
51
52#[derive(Clone)]
53pub struct SqliteToolOutcomeStore {
54 connection: Arc<Mutex<Connection>>,
55 serving_owner: Arc<SqliteServingOwner>,
56}
57
58impl SqliteToolOutcomeStore {
59 pub(crate) fn open_alongside(
60 connection: Arc<Mutex<Connection>>,
61 serving_owner: Arc<SqliteServingOwner>,
62 ) -> Self {
63 Self {
64 connection,
65 serving_owner,
66 }
67 }
68
69 fn connection(&self) -> Result<MutexGuard<'_, Connection>, ToolOutcomeStoreError> {
70 self.connection.lock().map_err(|_| {
71 ToolOutcomeStoreError::Unavailable("sqlite tool outcome lock poisoned".to_owned())
72 })
73 }
74
75 fn begin_read<'a>(
76 &self,
77 connection: &'a mut Connection,
78 ) -> Result<Transaction<'a>, ToolOutcomeStoreError> {
79 let transaction = connection
80 .transaction_with_behavior(TransactionBehavior::Deferred)
81 .map_err(sqlite_error)?;
82 verify_active_owner(&transaction, &self.serving_owner, None).map_err(admission_error)?;
83 self.serving_owner
84 .verify_authority_anchor(&transaction)
85 .map_err(|error| ToolOutcomeStoreError::Unavailable(error.to_string()))?;
86 Ok(transaction)
87 }
88
89 fn begin_write<'a>(
90 &self,
91 connection: &'a mut Connection,
92 fence: &StoreMutationFence,
93 trusted_now_unix_ms: u64,
94 ) -> Result<Transaction<'a>, ToolOutcomeStoreError> {
95 let transaction = connection
96 .transaction_with_behavior(TransactionBehavior::Immediate)
97 .map_err(sqlite_error)?;
98 verify_active_owner(&transaction, &self.serving_owner, Some(fence))
99 .map_err(admission_error)?;
100 self.serving_owner
101 .verify_authority_anchor(&transaction)
102 .map_err(|error| ToolOutcomeStoreError::Unavailable(error.to_string()))?;
103 verify_trusted_time(&transaction, trusted_now_unix_ms).map_err(admission_error)?;
104 Ok(transaction)
105 }
106
107 fn commit_write(&self, transaction: Transaction<'_>) -> Result<(), ToolOutcomeStoreError> {
108 transaction.commit().map_err(|error| {
109 ToolOutcomeStoreError::Unavailable(
110 self.serving_owner
111 .outcome_unknown(format!("sqlite tool outcome commit is unknown: {error}"))
112 .to_string(),
113 )
114 })
115 }
116
117 fn sync_after_write(&self, connection: &Connection) -> Result<(), ToolOutcomeStoreError> {
118 self.serving_owner
119 .sync_authority_anchor(connection)
120 .map_err(|error| ToolOutcomeStoreError::Unavailable(error.to_string()))
121 }
122
123 pub fn compact_retained_invocation_blobs(
131 &self,
132 retention_cutoff_unix_ms: u64,
133 active_fence: &StoreMutationFence,
134 trusted_now_unix_ms: u64,
135 ) -> Result<ToolOutcomeCompactionSummary, ToolOutcomeStoreError> {
136 let cutoff = sqlite_u64(retention_cutoff_unix_ms, "retention_cutoff_unix_ms")?;
137 let mut connection = self.connection()?;
138 let transaction = self.begin_write(&mut connection, active_fence, trusted_now_unix_ms)?;
139 let compacted = transaction
140 .execute(
141 r#"
142 UPDATE tool_outcome_blobs
143 SET canonical_bytes = NULL
144 WHERE canonical_bytes IS NOT NULL
145 AND recorded_at_unix_ms <= ?1
146 AND EXISTS (
147 SELECT 1 FROM tool_outcomes o
148 WHERE o.raw_output_digest = tool_outcome_blobs.digest
149 )
150 AND NOT EXISTS (
151 SELECT 1 FROM tool_outcomes o
152 JOIN admission_operations a ON a.operation_id = o.operation_id
153 WHERE o.raw_output_digest = tool_outcome_blobs.digest
154 AND a.terminal = 0
155 )
156 "#,
157 params![cutoff],
158 )
159 .map_err(sqlite_error)?;
160 let retained_live: i64 = transaction
161 .query_row(
162 r#"
163 SELECT COUNT(*) FROM tool_outcome_blobs b
164 WHERE b.canonical_bytes IS NOT NULL
165 AND b.recorded_at_unix_ms <= ?1
166 AND EXISTS (
167 SELECT 1 FROM tool_outcomes o
168 WHERE o.raw_output_digest = b.digest
169 )
170 AND EXISTS (
171 SELECT 1 FROM tool_outcomes o
172 JOIN admission_operations a ON a.operation_id = o.operation_id
173 WHERE o.raw_output_digest = b.digest
174 AND a.terminal = 0
175 )
176 "#,
177 params![cutoff],
178 |row| row.get(0),
179 )
180 .map_err(sqlite_error)?;
181 self.commit_write(transaction)?;
182 self.sync_after_write(&connection)?;
183 Ok(ToolOutcomeCompactionSummary {
184 compacted: u64::try_from(compacted).unwrap_or(0),
185 retained_live: u64::try_from(retained_live).unwrap_or(0),
186 })
187 }
188}
189
190impl ToolOutcomeStore for SqliteToolOutcomeStore {
191 fn record_tool_returned(
192 &self,
193 operation: &AdmissionOperationV1,
194 recovery_lease: &AdmissionRecoveryLease,
195 blob: &CanonicalInvocationBlobV1,
196 record: &ToolOutcomeRecordV1,
197 active_fence: &StoreMutationFence,
198 trusted_now_unix_ms: u64,
199 ) -> Result<ToolOutcomeInsertResultV1, ToolOutcomeStoreError> {
200 let mut connection = self.connection()?;
201 let transaction = self.begin_write(&mut connection, active_fence, trusted_now_unix_ms)?;
202 let stored_operation =
203 load_operation_for_participant_tx(&transaction, operation.binding().operation_id())
204 .map_err(admission_error)?
205 .ok_or(ToolOutcomeStoreError::NotFound)?;
206 if let Some(existing) = load_outcome_tx(&transaction, operation.binding().operation_id())? {
207 let stored_blob = load_blob_state_tx(&transaction, existing.raw_output_digest())?
208 .ok_or_else(|| invariant("tool outcome lost its canonical blob"))?;
209 let blob_matches = match &stored_blob {
210 StoredInvocationBlob::Present(existing_blob) => {
211 existing_blob.bytes() == blob.bytes()
212 }
213 StoredInvocationBlob::Compacted => {
216 existing.raw_output_digest() == blob.blob_ref().digest()
217 }
218 };
219 if !existing.same_immutable_outcome(record)
220 || !blob_matches
221 || stored_operation.tool_outcome_id() != Some(existing.outcome_id())
222 || !matches!(
223 stored_operation.state(),
224 AdmissionOperationState::Finalizing | AdmissionOperationState::Completed
225 )
226 {
227 return Err(ToolOutcomeStoreError::Conflict);
228 }
229 if let StoredInvocationBlob::Present(existing_blob) = &stored_blob {
230 existing
231 .validate_canonical_blob(&stored_operation, existing_blob)
232 .map_err(|error| invariant(error.to_string()))?;
233 }
234 transaction.commit().map_err(sqlite_error)?;
235 return Ok(ToolOutcomeInsertResultV1::ExactReplay {
236 outcome: existing,
237 operation: stored_operation,
238 });
239 }
240 if stored_operation != *operation {
241 return Err(ToolOutcomeStoreError::CasConflict);
242 }
243 record
244 .validate_for_store_insert(operation, blob, active_fence, trusted_now_unix_ms)
245 .map_err(|error| invariant(error.to_string()))?;
246 let outcome_json = encode_outcome(record)?;
247 let participant_digest = returned_participant_digest(
248 record,
249 record.raw_output_digest().as_str(),
250 &outcome_json,
251 )?;
252 insert_blob_tx(&transaction, blob, active_fence, trusted_now_unix_ms)?;
253 insert_outcome_tx(
254 &transaction,
255 record,
256 &outcome_json,
257 &participant_digest,
258 active_fence,
259 trusted_now_unix_ms,
260 )?;
261 let finalizing = advance_tool_outcome_tx(
262 &transaction,
263 &self.serving_owner,
264 operation,
265 recovery_lease,
266 record.outcome_id().clone(),
267 &participant_digest,
268 trusted_now_unix_ms,
269 )
270 .map_err(admission_error)?;
271 self.commit_write(transaction)?;
272 self.sync_after_write(&connection)?;
273 Ok(ToolOutcomeInsertResultV1::Inserted {
274 outcome: record.clone(),
275 operation: finalizing,
276 })
277 }
278
279 fn lookup_by_operation(
280 &self,
281 operation_id: &AdmissionOperationId,
282 ) -> Result<Option<ToolOutcomeRecordV1>, ToolOutcomeStoreError> {
283 let mut connection = self.connection()?;
284 let transaction = self.begin_read(&mut connection)?;
285 let outcome = load_outcome_tx(&transaction, operation_id)?;
286 transaction.commit().map_err(sqlite_error)?;
287 Ok(outcome)
288 }
289
290 fn load_raw_invocation_by_operation(
291 &self,
292 operation_id: &AdmissionOperationId,
293 ) -> Result<Option<RawInvocationOutcomeV1>, ToolOutcomeStoreError> {
294 let mut connection = self.connection()?;
295 let transaction = self.begin_read(&mut connection)?;
296 let raw = match load_outcome_tx(&transaction, operation_id)? {
297 Some(outcome) => load_blob_tx(&transaction, outcome.raw_output_digest())?
298 .map(|blob| RawInvocationOutcomeV1::from_canonical_bytes(blob.bytes()))
299 .transpose()
300 .map_err(|error| invariant(error.to_string()))?,
301 None => None,
302 };
303 transaction.commit().map_err(sqlite_error)?;
304 Ok(raw)
305 }
306
307 fn lookup_post_return_evaluation(
308 &self,
309 operation_id: &AdmissionOperationId,
310 ) -> Result<Option<PostReturnEvaluationRecordV1>, ToolOutcomeStoreError> {
311 let mut connection = self.connection()?;
312 let transaction = self.begin_read(&mut connection)?;
313 let evaluation = load_evaluation_tx(&transaction, operation_id)?;
314 transaction.commit().map_err(sqlite_error)?;
315 Ok(evaluation)
316 }
317
318 fn begin_post_return_evaluation(
319 &self,
320 recovery_lease: &AdmissionRecoveryLease,
321 record: &PostReturnEvaluationRecordV1,
322 active_fence: &StoreMutationFence,
323 trusted_now_unix_ms: u64,
324 ) -> Result<PostReturnEvaluationRecordV1, ToolOutcomeStoreError> {
325 let mut connection = self.connection()?;
326 let transaction = self.begin_write(&mut connection, active_fence, trusted_now_unix_ms)?;
327 let operation = load_operation_for_participant_tx(&transaction, record.operation_id())
328 .map_err(admission_error)?
329 .ok_or(ToolOutcomeStoreError::NotFound)?;
330 let outcome = load_outcome_tx(&transaction, record.operation_id())?
331 .ok_or(ToolOutcomeStoreError::NotFound)?;
332 require_finalizing_operation(&operation, &outcome)?;
333 record
334 .validate_against(&operation, &outcome)
335 .and_then(|_| record.validate_for_store_mutation(trusted_now_unix_ms))
336 .map_err(|error| invariant(error.to_string()))?;
337 if let Some(existing) = load_evaluation_tx(&transaction, record.operation_id())? {
338 if existing != *record {
339 return Err(ToolOutcomeStoreError::Conflict);
340 }
341 transaction.commit().map_err(sqlite_error)?;
342 return Ok(existing);
343 }
344 let evaluation_json = encode_evaluation(record)?;
345 let participant_digest = evaluation_participant_digest(record, &evaluation_json)?;
346 insert_evaluation_tx(
347 &transaction,
348 record,
349 outcome.outcome_id().as_str(),
350 &evaluation_json,
351 &participant_digest,
352 active_fence,
353 trusted_now_unix_ms,
354 )?;
355 append_participant_update_tx(
356 &transaction,
357 &self.serving_owner,
358 &operation,
359 recovery_lease,
360 &participant_digest,
361 trusted_now_unix_ms,
362 )
363 .map_err(admission_error)?;
364 self.commit_write(transaction)?;
365 self.sync_after_write(&connection)?;
366 Ok(record.clone())
367 }
368
369 fn stage_post_return_evaluation(
370 &self,
371 operation_id: &AdmissionOperationId,
372 expected_version: u64,
373 recovery_lease: &AdmissionRecoveryLease,
374 next: &PostReturnEvaluationRecordV1,
375 active_fence: &StoreMutationFence,
376 trusted_now_unix_ms: u64,
377 ) -> Result<PostReturnEvaluationRecordV1, ToolOutcomeStoreError> {
378 let mut connection = self.connection()?;
379 let transaction = self.begin_write(&mut connection, active_fence, trusted_now_unix_ms)?;
380 let operation = load_operation_for_participant_tx(&transaction, operation_id)
381 .map_err(admission_error)?
382 .ok_or(ToolOutcomeStoreError::NotFound)?;
383 let outcome =
384 load_outcome_tx(&transaction, operation_id)?.ok_or(ToolOutcomeStoreError::NotFound)?;
385 require_finalizing_operation(&operation, &outcome)?;
386 let current = load_evaluation_tx(&transaction, operation_id)?
387 .ok_or(ToolOutcomeStoreError::NotFound)?;
388 if current.version() != expected_version {
389 return Err(ToolOutcomeStoreError::CasConflict);
390 }
391 chio_kernel::tool_outcome::validate_evaluation_store_successor(¤t, next)
392 .and_then(|_| next.validate_against(&operation, &outcome))
393 .and_then(|_| next.validate_for_store_mutation(trusted_now_unix_ms))
394 .map_err(|error| invariant(error.to_string()))?;
395 let evaluation_json = encode_evaluation(next)?;
396 let participant_digest = evaluation_participant_digest(next, &evaluation_json)?;
397 update_evaluation_tx(
398 &transaction,
399 operation_id,
400 expected_version,
401 next,
402 &evaluation_json,
403 &participant_digest,
404 active_fence,
405 trusted_now_unix_ms,
406 )?;
407 append_participant_update_tx(
408 &transaction,
409 &self.serving_owner,
410 &operation,
411 recovery_lease,
412 &participant_digest,
413 trusted_now_unix_ms,
414 )
415 .map_err(admission_error)?;
416 self.commit_write(transaction)?;
417 self.sync_after_write(&connection)?;
418 Ok(next.clone())
419 }
420
421 fn finalize_post_return(
422 &self,
423 operation_id: &AdmissionOperationId,
424 expected_evaluation_version: u64,
425 recovery_lease: &AdmissionRecoveryLease,
426 terminal_evaluation: &PostReturnEvaluationRecordV1,
427 expected_outcome_version: u64,
428 terminal_outcome: &ToolOutcomeRecordV1,
429 resolved_output: Option<&CanonicalResolvedOutputBlobV1>,
430 active_fence: &StoreMutationFence,
431 trusted_now_unix_ms: u64,
432 ) -> Result<(PostReturnEvaluationRecordV1, ToolOutcomeRecordV1), ToolOutcomeStoreError> {
433 let mut connection = self.connection()?;
434 let transaction = self.begin_write(&mut connection, active_fence, trusted_now_unix_ms)?;
435 let operation = load_operation_for_participant_tx(&transaction, operation_id)
436 .map_err(admission_error)?
437 .ok_or(ToolOutcomeStoreError::NotFound)?;
438 let current_outcome =
439 load_outcome_tx(&transaction, operation_id)?.ok_or(ToolOutcomeStoreError::NotFound)?;
440 require_finalizing_operation(&operation, ¤t_outcome)?;
441 let current_evaluation = load_evaluation_tx(&transaction, operation_id)?
442 .ok_or(ToolOutcomeStoreError::NotFound)?;
443 if current_evaluation.version() != expected_evaluation_version
444 || current_outcome.version() != expected_outcome_version
445 {
446 return Err(ToolOutcomeStoreError::CasConflict);
447 }
448 chio_kernel::tool_outcome::validate_terminal_store_pair(
449 &operation,
450 ¤t_outcome,
451 ¤t_evaluation,
452 terminal_evaluation,
453 terminal_outcome,
454 resolved_output,
455 )
456 .and_then(|_| terminal_evaluation.validate_for_store_mutation(trusted_now_unix_ms))
457 .map_err(|error| invariant(error.to_string()))?;
458 let outcome_json = encode_outcome(terminal_outcome)?;
459 let evaluation_json = encode_evaluation(terminal_evaluation)?;
460 let participant_digest = finalization_participant_digest(
461 terminal_outcome,
462 terminal_evaluation,
463 &outcome_json,
464 &evaluation_json,
465 )?;
466 if let Some(blob) = resolved_output {
467 insert_blob_bytes_tx(
468 &transaction,
469 blob.blob_ref().digest().as_str(),
470 blob.bytes(),
471 active_fence,
472 trusted_now_unix_ms,
473 )?;
474 }
475 update_outcome_tx(
476 &transaction,
477 operation_id,
478 expected_outcome_version,
479 terminal_outcome,
480 &outcome_json,
481 &participant_digest,
482 active_fence,
483 trusted_now_unix_ms,
484 )?;
485 update_evaluation_tx(
486 &transaction,
487 operation_id,
488 expected_evaluation_version,
489 terminal_evaluation,
490 &evaluation_json,
491 &participant_digest,
492 active_fence,
493 trusted_now_unix_ms,
494 )?;
495 append_participant_update_tx(
496 &transaction,
497 &self.serving_owner,
498 &operation,
499 recovery_lease,
500 &participant_digest,
501 trusted_now_unix_ms,
502 )
503 .map_err(admission_error)?;
504 self.commit_write(transaction)?;
505 self.sync_after_write(&connection)?;
506 Ok((terminal_evaluation.clone(), terminal_outcome.clone()))
507 }
508
509 fn load_resolved_output_by_operation(
510 &self,
511 operation_id: &AdmissionOperationId,
512 ) -> Result<Option<CanonicalResolvedOutputBlobV1>, ToolOutcomeStoreError> {
513 let mut connection = self.connection()?;
514 let transaction = self.begin_read(&mut connection)?;
515 let outcome = load_outcome_tx(&transaction, operation_id)?;
516 let resolved = outcome
517 .as_ref()
518 .map(|outcome| load_resolved_blob_connection(&transaction, outcome))
519 .transpose()?
520 .flatten();
521 transaction.commit().map_err(sqlite_error)?;
522 Ok(resolved)
523 }
524}
525
526impl QualifiedToolOutcomeStore for SqliteToolOutcomeStore {}
527
528pub(crate) fn initialize_tool_outcome_schema(
529 connection: &mut Connection,
530) -> Result<(), ToolOutcomeStoreError> {
531 let on_disk = crate::check_schema_version(
532 connection,
533 TOOL_OUTCOME_SCHEMA_KEY,
534 TOOL_OUTCOME_SUPPORTED_SCHEMA_VERSION,
535 TOOL_OUTCOME_SCHEMA_ANCHORS,
536 )
537 .map_err(|error| invariant(error.to_string()))?;
538 if on_disk == TOOL_OUTCOME_SUPPORTED_SCHEMA_VERSION {
539 return verify_tool_outcome_invariants(connection);
540 }
541 let transaction = connection
542 .transaction_with_behavior(TransactionBehavior::Immediate)
543 .map_err(sqlite_error)?;
544 transaction
548 .execute_batch("DROP TRIGGER IF EXISTS tool_outcome_blobs_immutable;")
549 .map_err(sqlite_error)?;
550 transaction
551 .execute_batch(TOOL_OUTCOME_SCHEMA)
552 .map_err(sqlite_error)?;
553 crate::stamp_schema_version(
554 &transaction,
555 TOOL_OUTCOME_SCHEMA_KEY,
556 TOOL_OUTCOME_SUPPORTED_SCHEMA_VERSION,
557 )
558 .map_err(|error| invariant(error.to_string()))?;
559 verify_tool_outcome_invariants(&transaction)?;
560 transaction.commit().map_err(sqlite_error)
561}
562
563pub(crate) fn verify_tool_outcome_invariants(
564 connection: &Connection,
565) -> Result<(), ToolOutcomeStoreError> {
566 let expected = Connection::open_in_memory().map_err(sqlite_error)?;
567 expected
568 .execute_batch(TOOL_OUTCOME_SCHEMA)
569 .map_err(sqlite_error)?;
570 if tool_outcome_schema_catalog(connection)? != tool_outcome_schema_catalog(&expected)? {
571 return Err(invariant(
572 "tool outcome schema differs from the canonical definition",
573 ));
574 }
575 let mut blob_statement = connection
576 .prepare(
577 "SELECT digest, blob_size_bytes, canonical_bytes FROM tool_outcome_blobs ORDER BY digest",
578 )
579 .map_err(sqlite_error)?;
580 let mut blob_rows = blob_statement.query([]).map_err(sqlite_error)?;
581 while let Some(row) = blob_rows.next().map_err(sqlite_error)? {
582 let digest: String = row.get(0).map_err(sqlite_error)?;
583 let size: i64 = row.get(1).map_err(sqlite_error)?;
584 let bytes: Option<Vec<u8>> = row.get(2).map_err(sqlite_error)?;
585 if let Some(bytes) = bytes {
589 if usize::try_from(size).ok() != Some(bytes.len()) || sha256_hex(&bytes) != digest {
590 return Err(invariant("tool outcome blob digest is invalid"));
591 }
592 }
593 }
594 drop(blob_rows);
595 drop(blob_statement);
596 let mut statement = connection
597 .prepare("SELECT operation_id FROM tool_outcomes ORDER BY operation_id")
598 .map_err(sqlite_error)?;
599 let operation_ids = statement
600 .query_map([], |row| row.get::<_, String>(0))
601 .map_err(sqlite_error)?
602 .collect::<Result<Vec<_>, _>>()
603 .map_err(sqlite_error)?;
604 drop(statement);
605 for operation_id in operation_ids {
606 verify_outcome_projection(connection, &operation_id)?;
607 }
608 Ok(())
609}
610
611fn verify_outcome_projection(
612 connection: &Connection,
613 operation_id: &str,
614) -> Result<(), ToolOutcomeStoreError> {
615 let outcome = load_outcome_connection(connection, operation_id)?
616 .ok_or_else(|| invariant("tool outcome projection disappeared"))?;
617 let operation_json: Vec<u8> = connection
618 .query_row(
619 "SELECT operation_json FROM admission_operations WHERE operation_id = ?1",
620 [operation_id],
621 |row| row.get(0),
622 )
623 .map_err(sqlite_error)?;
624 let persisted = serde_json::from_slice(&operation_json)
625 .map_err(|error| invariant(format!("admission operation decode failed: {error}")))?;
626 let operation = AdmissionOperationV1::from_persisted(persisted)
627 .map_err(|error| invariant(error.to_string()))?;
628 if operation.tool_outcome_id() != Some(outcome.outcome_id()) {
629 return Err(invariant(
630 "tool outcome is not attached to its admission operation",
631 ));
632 }
633 outcome
634 .validate_against(&operation)
635 .map_err(|error| invariant(error.to_string()))?;
636 match load_blob_state_connection(connection, outcome.raw_output_digest())? {
637 None => return Err(invariant("tool outcome canonical blob is absent")),
638 Some(StoredInvocationBlob::Present(blob)) => outcome
639 .validate_canonical_blob(&operation, &blob)
640 .map_err(|error| invariant(error.to_string()))?,
641 Some(StoredInvocationBlob::Compacted) => {}
645 }
646 let returned_digest = returned_participant_digest(
647 &outcome,
648 outcome.raw_output_digest().as_str(),
649 &encode_outcome(&outcome)?,
650 )?;
651 let stored_outcome_digest: String = connection
652 .query_row(
653 "SELECT participant_digest FROM tool_outcomes WHERE operation_id = ?1",
654 [operation_id],
655 |row| row.get(0),
656 )
657 .map_err(sqlite_error)?;
658 let evaluation = load_evaluation_connection(connection, operation_id)?;
659 let (expected_outcome_digest, expected_evaluation_digest, expected_latest_digest) =
660 if let Some(evaluation) = &evaluation {
661 evaluation
662 .validate_against(&operation, &outcome)
663 .map_err(|error| invariant(error.to_string()))?;
664 let outcome_json = encode_outcome(&outcome)?;
665 let evaluation_json = encode_evaluation(evaluation)?;
666 if outcome.version() > 1 {
667 let digest = finalization_participant_digest(
668 &outcome,
669 evaluation,
670 &outcome_json,
671 &evaluation_json,
672 )?;
673 (digest.clone(), Some(digest.clone()), digest)
674 } else {
675 let evaluation_digest =
676 evaluation_participant_digest(evaluation, &evaluation_json)?;
677 (
678 returned_digest.clone(),
679 Some(evaluation_digest.clone()),
680 evaluation_digest,
681 )
682 }
683 } else {
684 if outcome.version() != 1 {
685 return Err(invariant(
686 "terminal tool outcome has no post-return evaluation",
687 ));
688 }
689 (returned_digest.clone(), None, returned_digest)
690 };
691 if stored_outcome_digest != expected_outcome_digest {
692 return Err(invariant(
693 "tool outcome row has an invalid participant commitment",
694 ));
695 }
696 let stored_evaluation_digest: Option<String> = connection
697 .query_row(
698 "SELECT participant_digest FROM post_return_evaluations WHERE operation_id = ?1",
699 [operation_id],
700 |row| row.get(0),
701 )
702 .optional()
703 .map_err(sqlite_error)?;
704 if stored_evaluation_digest != expected_evaluation_digest {
705 return Err(invariant(
706 "post-return evaluation row has an invalid participant commitment",
707 ));
708 }
709 let latest: Option<String> = connection
710 .query_row(
711 r#"
712 SELECT participant_digest FROM admission_operation_commits
713 WHERE operation_id = ?1 AND participant_digest IS NOT NULL
714 ORDER BY commit_sequence DESC LIMIT 1
715 "#,
716 [operation_id],
717 |row| row.get(0),
718 )
719 .optional()
720 .map_err(sqlite_error)?;
721 if latest.as_deref() != Some(expected_latest_digest.as_str()) {
722 return Err(invariant(
723 "tool outcome projection is not bound to the admission commit chain",
724 ));
725 }
726 let resolved = load_resolved_blob_connection(connection, &outcome)?;
727 if resolved.is_some() != outcome.resolved_output_ref().is_some() {
728 return Err(invariant(
729 "tool outcome resolved-output projection is incomplete",
730 ));
731 }
732 Ok(())
733}
734
735fn require_finalizing_operation(
736 operation: &AdmissionOperationV1,
737 outcome: &ToolOutcomeRecordV1,
738) -> Result<(), ToolOutcomeStoreError> {
739 if operation.state() != AdmissionOperationState::Finalizing
740 || operation.tool_outcome_id() != Some(outcome.outcome_id())
741 {
742 return Err(invariant(
743 "post-return evaluation requires the attached finalizing operation",
744 ));
745 }
746 Ok(())
747}
748
749fn insert_blob_tx(
750 transaction: &Transaction<'_>,
751 blob: &CanonicalInvocationBlobV1,
752 fence: &StoreMutationFence,
753 recorded_at_unix_ms: u64,
754) -> Result<(), ToolOutcomeStoreError> {
755 insert_blob_bytes_tx(
756 transaction,
757 blob.blob_ref().digest().as_str(),
758 blob.bytes(),
759 fence,
760 recorded_at_unix_ms,
761 )
762}
763
764fn insert_blob_bytes_tx(
765 transaction: &Transaction<'_>,
766 digest: &str,
767 bytes: &[u8],
768 fence: &StoreMutationFence,
769 recorded_at_unix_ms: u64,
770) -> Result<(), ToolOutcomeStoreError> {
771 if sha256_hex(bytes) != digest {
772 return Err(invariant(
773 "content-addressed blob digest does not match its bytes",
774 ));
775 }
776 let size = i64::try_from(bytes.len())
777 .map_err(|_| invariant("tool outcome blob size overflowed SQLite"))?;
778 transaction
779 .execute(
780 r#"
781 INSERT INTO tool_outcome_blobs (
782 digest, blob_size_bytes, canonical_bytes, recorded_at_unix_ms,
783 store_uuid, store_lease_id, store_owner_epoch
784 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
785 ON CONFLICT(digest) DO UPDATE SET
786 canonical_bytes = excluded.canonical_bytes
787 WHERE tool_outcome_blobs.canonical_bytes IS NULL
788 AND tool_outcome_blobs.blob_size_bytes = excluded.blob_size_bytes
789 "#,
790 params![
791 digest,
792 size,
793 bytes,
794 sqlite_u64(recorded_at_unix_ms, "recorded_at_unix_ms")?,
795 &fence.store_uuid,
796 &fence.lease_id,
797 sqlite_u64(fence.owner_epoch, "store_owner_epoch")?,
798 ],
799 )
800 .map_err(sqlite_error)?;
801 let (stored_size, stored): (i64, Option<Vec<u8>>) = transaction
802 .query_row(
803 "SELECT blob_size_bytes, canonical_bytes
804 FROM tool_outcome_blobs WHERE digest = ?1",
805 [digest],
806 |row| Ok((row.get(0)?, row.get(1)?)),
807 )
808 .map_err(sqlite_error)?;
809 if stored_size != size {
810 return Err(invariant(
811 "content-addressed blob size does not match its digest",
812 ));
813 }
814 match stored {
815 Some(stored) if stored == bytes => Ok(()),
816 Some(_) => Err(invariant("content-addressed blob digest collision")),
817 None => Err(invariant(
818 "content-addressed blob remained compacted after verified rehydration",
819 )),
820 }
821}
822
823fn load_resolved_blob_connection(
824 connection: &Connection,
825 outcome: &ToolOutcomeRecordV1,
826) -> Result<Option<CanonicalResolvedOutputBlobV1>, ToolOutcomeStoreError> {
827 let Some((expected, expected_size)) = outcome.resolved_output_ref() else {
828 return Ok(None);
829 };
830 let bytes = connection
831 .query_row(
832 "SELECT canonical_bytes FROM tool_outcome_blobs WHERE digest = ?1",
833 [expected.digest().as_str()],
834 |row| row.get::<_, Option<Vec<u8>>>(0),
835 )
836 .optional()
837 .map_err(sqlite_error)?
838 .flatten()
839 .ok_or_else(|| invariant("resolved tool output blob is absent"))?;
840 let blob = CanonicalResolvedOutputBlobV1::from_signing_preimage(bytes)
841 .map_err(|error| invariant(error.to_string()))?;
842 if blob.blob_ref() != expected || u64::try_from(blob.bytes().len()).ok() != Some(expected_size)
843 {
844 return Err(invariant(
845 "resolved tool output blob does not match its terminal record",
846 ));
847 }
848 Ok(Some(blob))
849}
850
851fn insert_outcome_tx(
852 transaction: &Transaction<'_>,
853 record: &ToolOutcomeRecordV1,
854 encoded: &[u8],
855 participant_digest: &str,
856 fence: &StoreMutationFence,
857 trusted_now_unix_ms: u64,
858) -> Result<(), ToolOutcomeStoreError> {
859 let persisted = record.to_persisted();
860 let inserted = transaction
861 .execute(
862 r#"
863 INSERT INTO tool_outcomes (
864 operation_id, outcome_id, request_id, raw_output_digest,
865 outcome_version, lifecycle_digest, participant_digest, outcome_json,
866 recorded_at_unix_ms, updated_at_unix_ms,
867 store_uuid, store_lease_id, store_owner_epoch
868 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)
869 "#,
870 params![
871 record.operation_id().as_str(),
872 record.outcome_id().as_str(),
873 persisted.request_id.as_str(),
874 record.raw_output_digest().as_str(),
875 sqlite_u64(record.version(), "outcome_version")?,
876 record.lifecycle_digest().as_str(),
877 participant_digest,
878 encoded,
879 sqlite_u64(record.recorded_at_unix_ms(), "recorded_at_unix_ms")?,
880 sqlite_u64(trusted_now_unix_ms, "trusted_now_unix_ms")?,
881 &fence.store_uuid,
882 &fence.lease_id,
883 sqlite_u64(fence.owner_epoch, "store_owner_epoch")?,
884 ],
885 )
886 .map_err(sqlite_error)?;
887 if inserted != 1 {
888 return Err(invariant("tool outcome insert did not affect one row"));
889 }
890 Ok(())
891}
892
893#[allow(clippy::too_many_arguments)]
894fn update_outcome_tx(
895 transaction: &Transaction<'_>,
896 operation_id: &AdmissionOperationId,
897 expected_version: u64,
898 next: &ToolOutcomeRecordV1,
899 encoded: &[u8],
900 participant_digest: &str,
901 fence: &StoreMutationFence,
902 trusted_now_unix_ms: u64,
903) -> Result<(), ToolOutcomeStoreError> {
904 let changed = transaction
905 .execute(
906 r#"
907 UPDATE tool_outcomes
908 SET outcome_version = ?1, lifecycle_digest = ?2,
909 participant_digest = ?3, outcome_json = ?4,
910 updated_at_unix_ms = ?5, store_uuid = ?6,
911 store_lease_id = ?7, store_owner_epoch = ?8
912 WHERE operation_id = ?9 AND outcome_version = ?10
913 "#,
914 params![
915 sqlite_u64(next.version(), "outcome_version")?,
916 next.lifecycle_digest().as_str(),
917 participant_digest,
918 encoded,
919 sqlite_u64(trusted_now_unix_ms, "trusted_now_unix_ms")?,
920 &fence.store_uuid,
921 &fence.lease_id,
922 sqlite_u64(fence.owner_epoch, "store_owner_epoch")?,
923 operation_id.as_str(),
924 sqlite_u64(expected_version, "expected_outcome_version")?,
925 ],
926 )
927 .map_err(sqlite_error)?;
928 if changed != 1 {
929 return Err(ToolOutcomeStoreError::CasConflict);
930 }
931 Ok(())
932}
933
934fn insert_evaluation_tx(
935 transaction: &Transaction<'_>,
936 record: &PostReturnEvaluationRecordV1,
937 outcome_id: &str,
938 encoded: &[u8],
939 participant_digest: &str,
940 fence: &StoreMutationFence,
941 trusted_now_unix_ms: u64,
942) -> Result<(), ToolOutcomeStoreError> {
943 let persisted = record.to_persisted();
944 let inserted = transaction
945 .execute(
946 r#"
947 INSERT INTO post_return_evaluations (
948 operation_id, evaluation_id, outcome_id, evaluation_version,
949 lifecycle_digest, participant_digest, evaluation_json,
950 created_at_unix_ms, updated_at_unix_ms,
951 store_uuid, store_lease_id, store_owner_epoch
952 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)
953 "#,
954 params![
955 record.operation_id().as_str(),
956 record.evaluation_id().as_str(),
957 outcome_id,
958 sqlite_u64(record.version(), "evaluation_version")?,
959 persisted.lifecycle_digest.as_str(),
960 participant_digest,
961 encoded,
962 sqlite_u64(persisted.trusted_time_unix_ms, "evaluation_trusted_time")?,
963 sqlite_u64(trusted_now_unix_ms, "trusted_now_unix_ms")?,
964 &fence.store_uuid,
965 &fence.lease_id,
966 sqlite_u64(fence.owner_epoch, "store_owner_epoch")?,
967 ],
968 )
969 .map_err(sqlite_error)?;
970 if inserted != 1 {
971 return Err(invariant(
972 "post-return evaluation insert did not affect one row",
973 ));
974 }
975 Ok(())
976}
977
978#[allow(clippy::too_many_arguments)]
979fn update_evaluation_tx(
980 transaction: &Transaction<'_>,
981 operation_id: &AdmissionOperationId,
982 expected_version: u64,
983 next: &PostReturnEvaluationRecordV1,
984 encoded: &[u8],
985 participant_digest: &str,
986 fence: &StoreMutationFence,
987 trusted_now_unix_ms: u64,
988) -> Result<(), ToolOutcomeStoreError> {
989 let persisted = next.to_persisted();
990 let changed = transaction
991 .execute(
992 r#"
993 UPDATE post_return_evaluations
994 SET evaluation_version = ?1, lifecycle_digest = ?2,
995 participant_digest = ?3, evaluation_json = ?4,
996 updated_at_unix_ms = ?5, store_uuid = ?6,
997 store_lease_id = ?7, store_owner_epoch = ?8
998 WHERE operation_id = ?9 AND evaluation_version = ?10
999 "#,
1000 params![
1001 sqlite_u64(next.version(), "evaluation_version")?,
1002 persisted.lifecycle_digest.as_str(),
1003 participant_digest,
1004 encoded,
1005 sqlite_u64(trusted_now_unix_ms, "trusted_now_unix_ms")?,
1006 &fence.store_uuid,
1007 &fence.lease_id,
1008 sqlite_u64(fence.owner_epoch, "store_owner_epoch")?,
1009 operation_id.as_str(),
1010 sqlite_u64(expected_version, "expected_evaluation_version")?,
1011 ],
1012 )
1013 .map_err(sqlite_error)?;
1014 if changed != 1 {
1015 return Err(ToolOutcomeStoreError::CasConflict);
1016 }
1017 Ok(())
1018}
1019
1020fn load_outcome_tx(
1021 transaction: &Transaction<'_>,
1022 operation_id: &AdmissionOperationId,
1023) -> Result<Option<ToolOutcomeRecordV1>, ToolOutcomeStoreError> {
1024 load_outcome_connection(transaction, operation_id.as_str())
1025}
1026
1027fn load_outcome_connection(
1028 connection: &Connection,
1029 operation_id: &str,
1030) -> Result<Option<ToolOutcomeRecordV1>, ToolOutcomeStoreError> {
1031 let row = connection
1032 .query_row(
1033 r#"
1034 SELECT outcome_id, request_id, raw_output_digest, outcome_version,
1035 lifecycle_digest, outcome_json, recorded_at_unix_ms
1036 FROM tool_outcomes WHERE operation_id = ?1
1037 "#,
1038 [operation_id],
1039 |row| {
1040 Ok((
1041 row.get::<_, String>(0)?,
1042 row.get::<_, String>(1)?,
1043 row.get::<_, String>(2)?,
1044 row.get::<_, i64>(3)?,
1045 row.get::<_, String>(4)?,
1046 row.get::<_, Vec<u8>>(5)?,
1047 row.get::<_, i64>(6)?,
1048 ))
1049 },
1050 )
1051 .optional()
1052 .map_err(sqlite_error)?;
1053 let Some((outcome_id, request_id, raw_digest, version, lifecycle, encoded, recorded_at)) = row
1054 else {
1055 return Ok(None);
1056 };
1057 let persisted: PersistedToolOutcomeRecordV1 = serde_json::from_slice(&encoded)
1058 .map_err(|error| invariant(format!("tool outcome decode failed: {error}")))?;
1059 let record = ToolOutcomeRecordV1::from_persisted(persisted)
1060 .map_err(|error| invariant(error.to_string()))?;
1061 if record.operation_id().as_str() != operation_id
1062 || record.outcome_id().as_str() != outcome_id
1063 || record.to_persisted().request_id.as_str() != request_id
1064 || record.raw_output_digest().as_str() != raw_digest
1065 || sqlite_u64(record.version(), "outcome_version")? != version
1066 || record.lifecycle_digest().as_str() != lifecycle
1067 || sqlite_u64(record.recorded_at_unix_ms(), "recorded_at_unix_ms")? != recorded_at
1068 || encode_outcome(&record)? != encoded
1069 {
1070 return Err(invariant(
1071 "tool outcome columns do not match canonical record",
1072 ));
1073 }
1074 Ok(Some(record))
1075}
1076
1077fn load_evaluation_tx(
1078 transaction: &Transaction<'_>,
1079 operation_id: &AdmissionOperationId,
1080) -> Result<Option<PostReturnEvaluationRecordV1>, ToolOutcomeStoreError> {
1081 load_evaluation_connection(transaction, operation_id.as_str())
1082}
1083
1084fn load_evaluation_connection(
1085 connection: &Connection,
1086 operation_id: &str,
1087) -> Result<Option<PostReturnEvaluationRecordV1>, ToolOutcomeStoreError> {
1088 let row = connection
1089 .query_row(
1090 r#"
1091 SELECT evaluation_id, outcome_id, evaluation_version,
1092 lifecycle_digest, evaluation_json
1093 FROM post_return_evaluations WHERE operation_id = ?1
1094 "#,
1095 [operation_id],
1096 |row| {
1097 Ok((
1098 row.get::<_, String>(0)?,
1099 row.get::<_, String>(1)?,
1100 row.get::<_, i64>(2)?,
1101 row.get::<_, String>(3)?,
1102 row.get::<_, Vec<u8>>(4)?,
1103 ))
1104 },
1105 )
1106 .optional()
1107 .map_err(sqlite_error)?;
1108 let Some((evaluation_id, outcome_id, version, lifecycle, encoded)) = row else {
1109 return Ok(None);
1110 };
1111 let persisted: PersistedPostReturnEvaluationRecordV1 = serde_json::from_slice(&encoded)
1112 .map_err(|error| invariant(format!("post-return evaluation decode failed: {error}")))?;
1113 let record = PostReturnEvaluationRecordV1::from_persisted(persisted)
1114 .map_err(|error| invariant(error.to_string()))?;
1115 let canonical = record.to_persisted();
1116 if record.operation_id().as_str() != operation_id
1117 || record.evaluation_id().as_str() != evaluation_id
1118 || canonical.tool_outcome_id.as_str() != outcome_id
1119 || sqlite_u64(record.version(), "evaluation_version")? != version
1120 || canonical.lifecycle_digest.as_str() != lifecycle
1121 || encode_evaluation(&record)? != encoded
1122 {
1123 return Err(invariant(
1124 "post-return evaluation columns do not match canonical record",
1125 ));
1126 }
1127 Ok(Some(record))
1128}
1129
1130fn load_blob_tx(
1131 transaction: &Transaction<'_>,
1132 digest: &chio_kernel::admission_operation::AdmissionDigest,
1133) -> Result<Option<CanonicalInvocationBlobV1>, ToolOutcomeStoreError> {
1134 load_blob_connection(transaction, digest)
1135}
1136
1137fn load_blob_connection(
1138 connection: &Connection,
1139 digest: &chio_kernel::admission_operation::AdmissionDigest,
1140) -> Result<Option<CanonicalInvocationBlobV1>, ToolOutcomeStoreError> {
1141 match load_blob_state_connection(connection, digest)? {
1142 None => Ok(None),
1143 Some(StoredInvocationBlob::Present(blob)) => Ok(Some(blob)),
1144 Some(StoredInvocationBlob::Compacted) => Err(compacted_blob_error(digest.as_str())),
1145 }
1146}
1147
1148fn load_blob_state_tx(
1149 transaction: &Transaction<'_>,
1150 digest: &chio_kernel::admission_operation::AdmissionDigest,
1151) -> Result<Option<StoredInvocationBlob>, ToolOutcomeStoreError> {
1152 load_blob_state_connection(transaction, digest)
1153}
1154
1155fn load_blob_state_connection(
1156 connection: &Connection,
1157 digest: &chio_kernel::admission_operation::AdmissionDigest,
1158) -> Result<Option<StoredInvocationBlob>, ToolOutcomeStoreError> {
1159 let stored: Option<Option<Vec<u8>>> = connection
1160 .query_row(
1161 "SELECT canonical_bytes FROM tool_outcome_blobs WHERE digest = ?1",
1162 [digest.as_str()],
1163 |row| row.get::<_, Option<Vec<u8>>>(0),
1164 )
1165 .optional()
1166 .map_err(sqlite_error)?;
1167 match stored {
1168 None => Ok(None),
1169 Some(None) => Ok(Some(StoredInvocationBlob::Compacted)),
1170 Some(Some(bytes)) => {
1171 let blob = RawInvocationOutcomeV1::from_canonical_bytes(&bytes)
1172 .and_then(|raw| raw.canonical_blob())
1173 .map_err(|error| invariant(error.to_string()))?;
1174 Ok(Some(StoredInvocationBlob::Present(blob)))
1175 }
1176 }
1177}
1178
1179fn compacted_blob_error(digest: &str) -> ToolOutcomeStoreError {
1180 invariant(format!(
1181 "tool outcome raw invocation blob `{digest}` was compacted under retention and is no longer available"
1182 ))
1183}
1184
1185fn encode_outcome(record: &ToolOutcomeRecordV1) -> Result<Vec<u8>, ToolOutcomeStoreError> {
1186 encode_bounded(
1187 "tool outcome",
1188 &record.to_persisted(),
1189 MAX_OUTCOME_RECORD_BYTES,
1190 )
1191}
1192
1193fn encode_evaluation(
1194 record: &PostReturnEvaluationRecordV1,
1195) -> Result<Vec<u8>, ToolOutcomeStoreError> {
1196 encode_bounded(
1197 "post-return evaluation",
1198 &record.to_persisted(),
1199 MAX_EVALUATION_RECORD_BYTES,
1200 )
1201}
1202
1203fn encode_bounded(
1204 label: &str,
1205 value: &impl Serialize,
1206 maximum: usize,
1207) -> Result<Vec<u8>, ToolOutcomeStoreError> {
1208 let encoded = canonical_json_bytes(value)
1209 .map_err(|error| invariant(format!("{label} encoding failed: {error}")))?;
1210 if encoded.is_empty() || encoded.len() > maximum {
1211 return Err(invariant(format!("{label} exceeds its storage bound")));
1212 }
1213 Ok(encoded)
1214}
1215
1216#[derive(Serialize)]
1217struct ParticipantCommitment<'a> {
1218 schema: &'static str,
1219 mutation: &'static str,
1220 operation_id: &'a str,
1221 outcome_id: &'a str,
1222 outcome_record_digest: Option<String>,
1223 raw_output_digest: Option<&'a str>,
1224 evaluation_id: Option<&'a str>,
1225 evaluation_record_digest: Option<String>,
1226}
1227
1228fn returned_participant_digest(
1229 record: &ToolOutcomeRecordV1,
1230 raw_output_digest: &str,
1231 outcome_json: &[u8],
1232) -> Result<String, ToolOutcomeStoreError> {
1233 participant_digest(&ParticipantCommitment {
1234 schema: "chio.tool-outcome-participant-commitment.v1",
1235 mutation: "record_tool_returned",
1236 operation_id: record.operation_id().as_str(),
1237 outcome_id: record.outcome_id().as_str(),
1238 outcome_record_digest: Some(sha256_hex(outcome_json)),
1239 raw_output_digest: Some(raw_output_digest),
1240 evaluation_id: None,
1241 evaluation_record_digest: None,
1242 })
1243}
1244
1245fn evaluation_participant_digest(
1246 record: &PostReturnEvaluationRecordV1,
1247 evaluation_json: &[u8],
1248) -> Result<String, ToolOutcomeStoreError> {
1249 let persisted = record.to_persisted();
1250 participant_digest(&ParticipantCommitment {
1251 schema: "chio.tool-outcome-participant-commitment.v1",
1252 mutation: "stage_post_return_evaluation",
1253 operation_id: record.operation_id().as_str(),
1254 outcome_id: persisted.tool_outcome_id.as_str(),
1255 outcome_record_digest: None,
1256 raw_output_digest: Some(persisted.raw_output_digest.as_str()),
1257 evaluation_id: Some(record.evaluation_id().as_str()),
1258 evaluation_record_digest: Some(sha256_hex(evaluation_json)),
1259 })
1260}
1261
1262fn finalization_participant_digest(
1263 outcome: &ToolOutcomeRecordV1,
1264 evaluation: &PostReturnEvaluationRecordV1,
1265 outcome_json: &[u8],
1266 evaluation_json: &[u8],
1267) -> Result<String, ToolOutcomeStoreError> {
1268 participant_digest(&ParticipantCommitment {
1269 schema: "chio.tool-outcome-participant-commitment.v1",
1270 mutation: "finalize_post_return",
1271 operation_id: outcome.operation_id().as_str(),
1272 outcome_id: outcome.outcome_id().as_str(),
1273 outcome_record_digest: Some(sha256_hex(outcome_json)),
1274 raw_output_digest: Some(outcome.raw_output_digest().as_str()),
1275 evaluation_id: Some(evaluation.evaluation_id().as_str()),
1276 evaluation_record_digest: Some(sha256_hex(evaluation_json)),
1277 })
1278}
1279
1280fn participant_digest(value: &impl Serialize) -> Result<String, ToolOutcomeStoreError> {
1281 canonical_json_bytes(value)
1282 .map(|bytes| sha256_hex(&bytes))
1283 .map_err(|error| invariant(format!("participant commitment encoding failed: {error}")))
1284}
1285
1286type SchemaCatalogEntry = (String, String, String, Option<String>);
1287
1288fn tool_outcome_schema_catalog(
1289 connection: &Connection,
1290) -> Result<Vec<SchemaCatalogEntry>, ToolOutcomeStoreError> {
1291 let mut statement = connection
1292 .prepare(
1293 r#"
1294 SELECT type, name, tbl_name, sql
1295 FROM sqlite_schema
1296 WHERE name GLOB 'tool_outcome*'
1297 OR tbl_name GLOB 'tool_outcome*'
1298 OR name GLOB 'post_return_evaluation*'
1299 OR tbl_name GLOB 'post_return_evaluation*'
1300 ORDER BY type, name, tbl_name
1301 "#,
1302 )
1303 .map_err(sqlite_error)?;
1304 let entries = statement
1305 .query_map([], |row| {
1306 Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
1307 })
1308 .map_err(sqlite_error)?
1309 .collect::<Result<Vec<_>, _>>()
1310 .map_err(sqlite_error)?;
1311 Ok(entries)
1312}
1313
1314fn sqlite_u64(value: u64, field: &str) -> Result<i64, ToolOutcomeStoreError> {
1315 i64::try_from(value).map_err(|_| invariant(format!("{field} exceeds SQLite INTEGER")))
1316}
1317
1318fn admission_error(error: AdmissionOperationStoreError) -> ToolOutcomeStoreError {
1319 match error {
1320 AdmissionOperationStoreError::Fenced => ToolOutcomeStoreError::Fenced,
1321 AdmissionOperationStoreError::NotFound => ToolOutcomeStoreError::NotFound,
1322 AdmissionOperationStoreError::Unavailable(detail)
1323 | AdmissionOperationStoreError::OutcomeUnknown(detail) => {
1324 ToolOutcomeStoreError::Unavailable(detail)
1325 }
1326 AdmissionOperationStoreError::Invariant(detail) => ToolOutcomeStoreError::Invariant(detail),
1327 AdmissionOperationStoreError::Operation(error) => invariant(error.to_string()),
1328 }
1329}
1330
1331fn sqlite_error(error: rusqlite::Error) -> ToolOutcomeStoreError {
1332 ToolOutcomeStoreError::Unavailable(error.to_string())
1333}
1334
1335fn invariant(detail: impl Into<String>) -> ToolOutcomeStoreError {
1336 ToolOutcomeStoreError::Invariant(detail.into())
1337}
1338
1339#[cfg(test)]
1340#[path = "tool_outcome_store_tests.rs"]
1341#[allow(clippy::expect_used, clippy::unwrap_used)]
1342mod tests;