Skip to main content

chio_store_sqlite/
tool_outcome_store.rs

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/// Result of one raw-invocation retention pass.
36#[derive(Debug, Clone, Copy, PartialEq, Eq)]
37pub struct ToolOutcomeCompactionSummary {
38    /// Raw invocation payloads cleared to NULL during this pass.
39    pub compacted: u64,
40    /// Payloads past the retention cutoff left in place because at least one
41    /// owning admission operation has not yet reached a terminal state.
42    pub retained_live: u64,
43}
44
45/// A stored raw invocation blob, or a marker that its payload was compacted away
46/// under retention while the digest and size were preserved.
47enum 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    /// Clears the raw invocation payload of every content-addressed blob whose
124    /// owning operations are all terminal and that was recorded at or before
125    /// `retention_cutoff_unix_ms`, preserving the digest and size that
126    /// verification depends on. The caller supplies the cutoff (for example from
127    /// a retention policy); this store never reads kernel configuration on its
128    /// own. Blobs whose owning operation is still live are left untouched, which
129    /// the schema triggers independently enforce.
130    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                // The payload was compacted under retention. Blobs are
214                // content-addressed, so an equal digest is an equal payload.
215                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(&current, 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, &current_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            &current_outcome,
451            &current_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    // Version 2 permits a compacted payload to be restored when a later owner
545    // supplies the same digest- and size-verified canonical bytes. Recreate the
546    // trigger transactionally before applying the canonical schema definition.
547    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        // A compacted blob keeps its digest and size but holds no payload, so
586        // there are no bytes to re-hash; the triggers freeze the retained
587        // columns against tampering.
588        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        // The payload was compacted under retention; its bytes are gone, but the
642        // digest and size are retained and re-checked in
643        // verify_tool_outcome_invariants, so there is nothing to re-derive here.
644        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;