Skip to main content

chio_kernel/kernel/
admission_coordinator.rs

1use std::fmt;
2use std::sync::{Arc, Mutex};
3
4use chio_log_redact::redacted;
5use serde::Serialize;
6use tracing::warn;
7
8#[path = "admission_coordinator/terminal.rs"]
9mod terminal;
10pub(crate) use terminal::DurableToolReturnInput;
11
12use super::*;
13use crate::admission_operation::{
14    verified_outcome_unknown_after_dispatch_projection,
15    verified_released_pre_dispatch_compensation_projection, AdmissionAttachment,
16    AdmissionBeginResult, AdmissionCompensationStatus, AdmissionCompletedProjection,
17    AdmissionDigest, AdmissionDispatchState, AdmissionIdentifier, AdmissionMutationGuard,
18    AdmissionMutationSequencer, AdmissionOperationBindingInputV1, AdmissionOperationBindingV1,
19    AdmissionOperationCommand, AdmissionOperationKind, AdmissionOperationState,
20    AdmissionOperationV1, AdmissionParticipantRequirements, AdmissionProjectionContext,
21    AdmissionReceiptMetadataV1, AdmissionReceiptSchema, AdmissionRequestBindingV1,
22    AdmissionTerminalProjection, AdmissionTerminalReplay, AuthenticatedRequestNamespace,
23    ObservationAttemptZero, PaymentTerminalEvidence, ProviderAttemptBindingV1,
24    QualifiedAdmissionOperationStoreExt, QualifiedChannelTerminalAuthority, SideEffectClass,
25    StoreMutationFence, VerifiedAdmissionReceipt, ADMISSION_RECEIPT_METADATA_KEY,
26    LOCAL_SYSTEM_TENANT_ID,
27};
28use crate::budget_store::{
29    BudgetAdmissionBinding, BudgetCaptureInvocationRequest, BudgetEventAuthority,
30    BudgetGuaranteeLevel, BudgetInvocationQuota, BudgetQuotaKey, BudgetQuotaProfile,
31    BudgetReconcileHoldDecision, BudgetReconcileHoldRequest,
32};
33use crate::receipt_store::QualifiedAdmissionProjectionStore;
34use crate::supplemental_quota::{
35    canonical_revocation_set_for_verified_claim, supplemental_authorization_artifact_digest,
36    verify_supplemental_quota, CanonicalRevocationSet, KernelVerifiedSupplementalQuotaClaim,
37    SupplementalQuotaError, SupplementalQuotaVerificationContext,
38    BROKER_CAPABILITY_EXECUTION_PROFILE,
39};
40use crate::tool_outcome::{
41    EvaluationModeV1, EvaluationPhaseV1, FrozenEvaluationStepV1, InvocationOutputV1,
42    InvocationStreamLimitsV1, PostReturnEvaluationRecordV1, PostReturnEvaluationStateV1,
43    PostReturnNormalizedRequestContextV1, QualifiedDurableOutcomeAuthority,
44    QualifiedToolOutcomeStore, RawInvocationOutcomeV1, ResolvedToolOutcomeV1,
45    SettlementDispositionV1, ToolOutcomeError, ToolOutcomeRecordV1, ToolOutcomeStoreError,
46    ToolOutcomeTerminalEvidenceV1, ToolOutcomeTransitionV1, VerifiedContractualZeroCharge,
47};
48
49const RECOVERY_LEASE_DURATION_MS: u64 = 60_000;
50const I_JSON_MAX_SAFE_INTEGER: u64 = (1_u64 << 53) - 1;
51
52#[derive(Clone)]
53pub(crate) struct DurableAdmissionRuntime {
54    store: Arc<dyn QualifiedAdmissionProjectionStore>,
55    outcome_store: Arc<dyn QualifiedToolOutcomeStore>,
56    channel_terminal_authority: Option<Arc<dyn QualifiedChannelTerminalAuthority>>,
57    fence: StoreMutationFence,
58    claimant_id: AdmissionIdentifier,
59    mutation_sequencer: AdmissionMutationSequencer,
60    startup_reconciled: Arc<Mutex<bool>>,
61}
62
63impl fmt::Debug for DurableAdmissionRuntime {
64    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
65        formatter
66            .debug_struct("DurableAdmissionRuntime")
67            .field("fence", &self.fence)
68            .field("claimant_id", &self.claimant_id)
69            .field(
70                "channel_terminal_authority",
71                &self.channel_terminal_authority.is_some(),
72            )
73            .finish_non_exhaustive()
74    }
75}
76
77impl DurableAdmissionRuntime {
78    pub(crate) fn new(
79        store: Arc<dyn QualifiedAdmissionProjectionStore>,
80        outcome_store: Arc<dyn QualifiedToolOutcomeStore>,
81        fence: StoreMutationFence,
82        kernel_id: &str,
83    ) -> Result<Self, crate::admission_operation::AdmissionOperationError> {
84        AdmissionIdentifier::try_new("store_uuid", fence.store_uuid.clone())?;
85        AdmissionIdentifier::try_new("store_lease_id", fence.lease_id.clone())?;
86        if fence.owner_epoch == 0 || fence.owner_epoch > I_JSON_MAX_SAFE_INTEGER {
87            return Err(crate::admission_operation::AdmissionOperationError::InvalidStoreFence);
88        }
89        let claimant_id =
90            AdmissionIdentifier::try_new("admission_claimant_id", format!("kernel:{kernel_id}"))?;
91        Ok(Self {
92            store,
93            outcome_store,
94            channel_terminal_authority: None,
95            mutation_sequencer: AdmissionMutationSequencer::for_fence(&fence)?,
96            startup_reconciled: Arc::new(Mutex::new(false)),
97            fence,
98            claimant_id,
99        })
100    }
101
102    fn lock_mutations(&self) -> Result<AdmissionMutationGuard<'_>, KernelError> {
103        self.mutation_sequencer
104            .lock()
105            .map_err(|error| KernelError::DurableAdmission(error.to_string()))
106    }
107
108    pub(super) fn set_channel_terminal_authority(
109        &mut self,
110        authority: Arc<dyn QualifiedChannelTerminalAuthority>,
111    ) {
112        self.channel_terminal_authority = Some(authority);
113    }
114
115    fn refresh_trusted_time(&self, requested_unix_ms: u64) -> u64 {
116        requested_unix_ms.max(current_unix_timestamp_ms()).max(1)
117    }
118
119    fn authority(&self) -> BudgetEventAuthority {
120        BudgetEventAuthority {
121            authority_id: self.fence.store_uuid.clone(),
122            lease_id: self.fence.lease_id.clone(),
123            lease_epoch: self.fence.owner_epoch,
124        }
125    }
126
127    fn qualified_terminal_records(
128        &self,
129        operation: &AdmissionOperationV1,
130    ) -> Result<(ToolOutcomeRecordV1, PostReturnEvaluationRecordV1), ToolOutcomeError> {
131        let unavailable =
132            || ToolOutcomeError::ReleaseAuthorityUnavailable("durable terminal outcome store");
133        let outcome = self
134            .outcome_store
135            .lookup_by_operation(operation.binding().operation_id())
136            .map_err(|_| unavailable())?
137            .ok_or_else(unavailable)?;
138        let evaluation = self
139            .outcome_store
140            .lookup_post_return_evaluation(operation.binding().operation_id())
141            .map_err(|_| unavailable())?
142            .ok_or_else(unavailable)?;
143        Ok((outcome, evaluation))
144    }
145}
146
147impl QualifiedDurableOutcomeAuthority for DurableAdmissionRuntime {
148    fn verify_terminal_outcome(
149        &self,
150        operation: &AdmissionOperationV1,
151        context: &AdmissionProjectionContext,
152    ) -> Result<ToolOutcomeTerminalEvidenceV1, ToolOutcomeError> {
153        let (outcome, evaluation) = self.qualified_terminal_records(operation)?;
154        ToolOutcomeTerminalEvidenceV1::from_records(operation, context, &outcome, &evaluation)
155    }
156
157    fn verify_contractual_zero_charge(
158        &self,
159        operation: &AdmissionOperationV1,
160        context: &AdmissionProjectionContext,
161    ) -> Result<VerifiedContractualZeroCharge, ToolOutcomeError> {
162        let (outcome, evaluation) = self.qualified_terminal_records(operation)?;
163        VerifiedContractualZeroCharge::from_records(operation, context, &outcome, &evaluation)
164    }
165}
166
167pub(crate) struct DurableToolAdmission {
168    pub(super) operation: AdmissionOperationV1,
169    aggregate_quota: Option<BudgetInvocationQuota>,
170    supplemental_quota: Option<KernelVerifiedSupplementalQuotaClaim>,
171}
172
173impl DurableToolAdmission {
174    pub(crate) fn operation(&self) -> &AdmissionOperationV1 {
175        &self.operation
176    }
177
178    pub(crate) fn operation_id(&self) -> &str {
179        self.operation.binding().operation_id().as_str()
180    }
181
182    pub(crate) fn budget_hold_id(&self, grant_index: usize) -> String {
183        format!("admission-budget:{}:{grant_index}", self.operation_id())
184    }
185
186    pub(crate) fn budget_authorize_event_id(&self, grant_index: usize) -> String {
187        format!("{}:authorize", self.budget_hold_id(grant_index))
188    }
189
190    pub(crate) fn permits_grant(&self, grant_index: usize) -> bool {
191        self.operation
192            .budget_hold_id()
193            .is_none_or(|hold_id| hold_id.as_str() == self.budget_hold_id(grant_index))
194    }
195
196    pub(crate) fn permits_matching_grant(&self, matching: &MatchingGrant<'_>) -> bool {
197        self.permits_grant(matching.index)
198            && if self.requires_payment() {
199                matching
200                    .grant
201                    .max_cost_per_invocation
202                    .as_ref()
203                    .is_some_and(|amount| amount.units != 0)
204            } else {
205                matching.grant.max_cost_per_invocation.is_none()
206                    && matching.grant.max_total_cost.is_none()
207            }
208    }
209
210    pub(crate) fn can_resume_captured_hold(&self) -> bool {
211        self.operation.state() == AdmissionOperationState::CapturePending
212    }
213
214    pub(crate) fn requires_payment(&self) -> bool {
215        self.operation.binding().participant_requirements().payment
216    }
217
218    pub(crate) fn state(&self) -> AdmissionOperationState {
219        self.operation.state()
220    }
221
222    pub(crate) fn supplemental_quota(&self) -> Option<&KernelVerifiedSupplementalQuotaClaim> {
223        self.supplemental_quota.as_ref()
224    }
225
226    pub(crate) fn aggregate_quota(&self) -> Option<&BudgetInvocationQuota> {
227        self.aggregate_quota.as_ref()
228    }
229}
230
231#[derive(Serialize)]
232struct ImmutableToolAdmissionRequest<'a> {
233    schema: &'static str,
234    server_id: &'a str,
235    tool_name: &'a str,
236    agent_id: &'a str,
237    arguments: &'a serde_json::Value,
238    governed_intent: &'a Option<chio_core::capability::governance::GovernedTransactionIntent>,
239    model_metadata: &'a Option<chio_core::capability::scope::ModelMetadata>,
240    federated_origin_kernel_id: &'a Option<String>,
241    matching_grants: Vec<ImmutableMatchingGrant<'a>>,
242    post_return_steps: &'a [FrozenEvaluationStepV1],
243}
244
245#[derive(Serialize)]
246struct ImmutableActiveResponseAdmissionRequest<'a> {
247    schema: &'static str,
248    governed_intent: &'a chio_core::capability::governance::GovernedTransactionIntent,
249    federated_origin_kernel_id: &'a Option<String>,
250    governed_intent_hash: &'a str,
251}
252
253#[derive(Serialize)]
254struct ImmutableMatchingGrant<'a> {
255    index: usize,
256    grant: &'a ToolGrant,
257}
258
259struct DurablePostReturnPlan {
260    hook_identities: Vec<crate::post_invocation::PostInvocationHookIdentity>,
261    frozen_steps: Vec<FrozenEvaluationStepV1>,
262}
263
264fn immutable_tool_admission_request_hash(
265    request: &ToolCallRequest,
266    matching_grants: &[MatchingGrant<'_>],
267    post_return_plan: &DurablePostReturnPlan,
268) -> Result<AdmissionDigest, KernelError> {
269    let immutable_request = ImmutableToolAdmissionRequest {
270        schema: "chio.tool-admission-request.v1",
271        server_id: &request.server_id,
272        tool_name: &request.tool_name,
273        agent_id: &request.agent_id,
274        arguments: &request.arguments,
275        governed_intent: &request.governed_intent,
276        model_metadata: &request.model_metadata,
277        federated_origin_kernel_id: &request.federated_origin_kernel_id,
278        matching_grants: matching_grants
279            .iter()
280            .map(|matching| ImmutableMatchingGrant {
281                index: matching.index,
282                grant: matching.grant,
283            })
284            .collect(),
285        post_return_steps: &post_return_plan.frozen_steps,
286    };
287    admission_digest("immutable_request_hash", &immutable_request)
288}
289
290impl ChioKernel {
291    pub(crate) fn load_durable_admission_receipt(
292        &self,
293        receipt_id: &str,
294    ) -> Result<Option<chio_core::receipt::body::ChioReceipt>, KernelError> {
295        let Some(runtime) = self.durable_admission_runtime.as_ref() else {
296            return Ok(None);
297        };
298        runtime
299            .store
300            .load_chio_receipt(receipt_id)
301            .map_err(|error| KernelError::DurableAdmission(error.to_string()))
302    }
303
304    pub fn reconcile_durable_admission_receipt_projections(&self) -> Result<usize, KernelError> {
305        const PAGE_LIMIT: usize = 256;
306
307        if self.receipt_store.is_none() {
308            return Ok(0);
309        }
310        let Some(runtime) = self.durable_admission_runtime.as_ref() else {
311            return Ok(0);
312        };
313        let mut after_receipt_id = None;
314        let mut reconciled = 0_usize;
315        loop {
316            let page = runtime
317                .store
318                .list_admission_receipts_after(after_receipt_id.as_deref(), PAGE_LIMIT)
319                .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
320            if page.is_empty() {
321                return Ok(reconciled);
322            }
323            if page.len() > PAGE_LIMIT {
324                return Err(KernelError::DurableAdmission(
325                    "admission receipt store exceeded the requested page limit".to_owned(),
326                ));
327            }
328            let mut previous = after_receipt_id.as_deref();
329            for receipt in &page {
330                if previous.is_some_and(|cursor| receipt.id.as_str() <= cursor) {
331                    return Err(KernelError::DurableAdmission(
332                        "admission receipt store returned a non-advancing page".to_owned(),
333                    ));
334                }
335                self.materialize_durable_admission_receipt(receipt)?;
336                previous = Some(receipt.id.as_str());
337            }
338            reconciled = reconciled.checked_add(page.len()).ok_or_else(|| {
339                KernelError::DurableAdmission("admission receipt count overflow".to_owned())
340            })?;
341            after_receipt_id = page.last().map(|receipt| receipt.id.clone());
342        }
343    }
344
345    pub fn reconcile_durable_admission_startup(&self) -> Result<usize, KernelError> {
346        let Some(runtime) = self.durable_admission_runtime.as_ref() else {
347            return Ok(0);
348        };
349        let mut reconciled = runtime.startup_reconciled.lock().map_err(|_| {
350            KernelError::DurableAdmission("startup reconciliation lock is poisoned".to_owned())
351        })?;
352        if *reconciled {
353            return Ok(0);
354        }
355        let operation_count = self.reconcile_recoverable_admissions()?;
356        let receipt_count = self.reconcile_durable_admission_receipt_projections()?;
357        let total = operation_count.checked_add(receipt_count).ok_or_else(|| {
358            KernelError::DurableAdmission("startup reconciliation count overflow".to_owned())
359        })?;
360        *reconciled = true;
361        Ok(total)
362    }
363
364    pub fn reconcile_recoverable_admissions(&self) -> Result<usize, KernelError> {
365        const PAGE_LIMIT: usize = 256;
366
367        let Some(runtime) = self.durable_admission_runtime.as_ref() else {
368            return Ok(0);
369        };
370        let trusted_now_unix_ms = runtime.refresh_trusted_time(current_unix_timestamp_ms());
371        let mut reconciled = 0_usize;
372        // An operation that cannot be reconciled is recorded and skipped rather
373        // than abandoning the sweep, so one wedged operation cannot hold up every
374        // other recoverable operation. The first failure is still returned once
375        // the sweep finishes, so callers keep failing closed on it.
376        let mut deferred_failure: Option<KernelError> = None;
377        loop {
378            let recoverable = runtime
379                .store
380                .list_recoverable(trusted_now_unix_ms, PAGE_LIMIT)
381                .map_err(durable_store_error)?;
382            if recoverable.len() > PAGE_LIMIT {
383                return Err(KernelError::DurableAdmission(
384                    "admission recovery store exceeded the requested page limit".to_owned(),
385                ));
386            }
387            if recoverable.is_empty() {
388                break;
389            }
390            let reconciled_before_page = reconciled;
391            for operation in recoverable {
392                match operation.state() {
393                    AdmissionOperationState::DispatchCommitted => {
394                        if let Err(error) = self.terminalize_dispatch_committed_admission(
395                            &operation,
396                            trusted_now_unix_ms,
397                        ) {
398                            warn!(
399                                operation_id = %operation.binding().operation_id().as_str(),
400                                reason = %redacted!(&error),
401                                audit_fault = "admission_recovery_terminalization_unresolved",
402                                "failed to terminalize a dispatch-committed admission"
403                            );
404                            deferred_failure.get_or_insert(error);
405                            continue;
406                        }
407                        reconciled = reconciled.checked_add(1).ok_or_else(|| {
408                            KernelError::DurableAdmission(
409                                "admission recovery count overflow".to_owned(),
410                            )
411                        })?;
412                    }
413                    AdmissionOperationState::Prepared
414                    | AdmissionOperationState::BrokerAttemptRegistered
415                    | AdmissionOperationState::BudgetAuthorized
416                    | AdmissionOperationState::ApprovalReserved
417                    | AdmissionOperationState::ReadyToDispatch
418                    | AdmissionOperationState::CapturePending => {
419                        // One operation that cannot be compensated must not abandon
420                        // the rest of the page: it stays recoverable for a later
421                        // sweep, and the remaining operations still reconcile.
422                        if let Err(error) = self.compensate_durable_admission_before_dispatch(
423                            &operation,
424                            serde_json::json!({
425                                "authority": "startup-recovery",
426                                "cause": "no-authoritative-budget-participant"
427                            }),
428                            trusted_now_unix_ms,
429                        ) {
430                            warn!(
431                                operation_id = %operation.binding().operation_id().as_str(),
432                                reason = %redacted!(&error),
433                                audit_fault = "admission_recovery_compensation_unresolved",
434                                "failed to compensate a recoverable admission"
435                            );
436                            deferred_failure.get_or_insert(error);
437                            continue;
438                        }
439                        reconciled = reconciled.checked_add(1).ok_or_else(|| {
440                            KernelError::DurableAdmission(
441                                "admission recovery count overflow".to_owned(),
442                            )
443                        })?;
444                    }
445                    AdmissionOperationState::ApprovalRequired => {
446                        deferred_failure.get_or_insert_with(|| {
447                            KernelError::DurableAdmission(
448                                "admission recovery store returned a quiescent approval-required operation"
449                                    .to_owned(),
450                            )
451                        });
452                    }
453                    AdmissionOperationState::Finalizing => {
454                        let mut admission = DurableToolAdmission {
455                            operation,
456                            aggregate_quota: None,
457                            supplemental_quota: None,
458                        };
459                        let tool_return = self.load_durable_tool_return(&admission)?;
460                        let Some(request) =
461                            tool_return.recovery_request().map_err(tool_outcome_error)?
462                        else {
463                            self.claim_admission_recovery(
464                                &admission.operation,
465                                trusted_now_unix_ms,
466                            )?;
467                            continue;
468                        };
469                        if let Err(error) = self.finalize_durable_tool_return(
470                            &mut admission,
471                            &request,
472                            &tool_return,
473                        ) {
474                            warn!(
475                                operation_id = %admission.operation.binding().operation_id().as_str(),
476                                reason = %redacted!(&error),
477                                audit_fault = "admission_recovery_finalization_unresolved",
478                                "failed to finalize a recoverable admission"
479                            );
480                            deferred_failure.get_or_insert(error);
481                            continue;
482                        }
483                        reconciled = reconciled.checked_add(1).ok_or_else(|| {
484                            KernelError::DurableAdmission(
485                                "admission recovery count overflow".to_owned(),
486                            )
487                        })?;
488                    }
489                    _ => {
490                        self.claim_admission_recovery(&operation, trusted_now_unix_ms)?;
491                    }
492                }
493            }
494            if reconciled == reconciled_before_page {
495                break;
496            }
497        }
498        if let Some(error) = deferred_failure {
499            return Err(error);
500        }
501        Ok(reconciled)
502    }
503
504    /// Terminalize a dispatch-committed admission whose outcome is unknown.
505    ///
506    /// Refuses when a durable tool outcome already exists, so this is a no-op on
507    /// an operation whose return did land. Used both by startup recovery and by
508    /// the post-dispatch drop path, where the evaluation future was cancelled
509    /// after the dispatch commit and would otherwise strand the operation until
510    /// the next process restart.
511    pub(crate) fn terminalize_dispatch_committed_admission(
512        &self,
513        operation: &AdmissionOperationV1,
514        trusted_now_unix_ms: u64,
515    ) -> Result<(), KernelError> {
516        let runtime = self.durable_runtime()?;
517        let _mutation_guard = runtime.lock_mutations()?;
518        if runtime
519            .outcome_store
520            .lookup_by_operation(operation.binding().operation_id())
521            .map_err(durable_outcome_store_error)?
522            .is_some()
523        {
524            return Err(KernelError::DurableAdmission(
525                "dispatch-committed admission already has a durable tool outcome".to_owned(),
526            ));
527        }
528        let lease = self.claim_admission_recovery(operation, trusted_now_unix_ms)?;
529        let context = AdmissionProjectionContext {
530            operation_id: operation.binding().operation_id().clone(),
531            request_id: operation.binding().request_id().clone(),
532            expected_operation_version: operation.version(),
533            trusted_time_unix_ms: trusted_now_unix_ms,
534            coordinator_lease_id: lease.coordinator_lease_id().clone(),
535            coordinator_lease_epoch: lease.coordinator_lease_epoch(),
536            store_fence: runtime.fence.clone(),
537        };
538        let projection = verified_outcome_unknown_after_dispatch_projection(operation, context)?;
539        let terminal = runtime
540            .store
541            .commit_admission_projection(&projection)
542            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
543        if terminal.operation_id != *operation.binding().operation_id()
544            || terminal.state != AdmissionOperationState::OutcomeUnknownAfterDispatch
545        {
546            return Err(KernelError::DurableAdmission(
547                "admission recovery committed a different terminal operation".to_owned(),
548            ));
549        }
550        Ok(())
551    }
552
553    pub(crate) fn begin_durable_tool_admission(
554        &self,
555        request: &ToolCallRequest,
556        matching_grants: &[MatchingGrant<'_>],
557        trusted_now_unix_ms: u64,
558    ) -> Result<Option<DurableToolAdmission>, KernelError> {
559        let aggregate_quota =
560            self.verify_aggregate_quota_for_admission(request, trusted_now_unix_ms / 1_000)?;
561        let cumulative_matching_grant_count = matching_grants
562            .iter()
563            .filter(|matching| {
564                matching.grant.constraints.iter().any(|constraint| {
565                    matches!(
566                        constraint,
567                        Constraint::RequireCumulativeApprovalAbove { .. }
568                    )
569                })
570            })
571            .count();
572        // Only a grant that can serve this request may force the structured path. An
573        // unrelated cumulative grant elsewhere in the capability must not withdraw an
574        // otherwise exempt call.
575        let requires_structured_admission = aggregate_quota.is_some()
576            || request.supplemental_authorization.is_some()
577            || cumulative_matching_grant_count != 0;
578        if request.supplemental_authorization.is_some()
579            && self.supplemental_quota_verifier.is_none()
580        {
581            return Err(KernelError::DurableAdmission(
582                SupplementalQuotaError::MissingVerifier.to_string(),
583            ));
584        }
585        let effect_class = if matching_grants.iter().all(|matching| {
586            matching.grant.max_cost_per_invocation.is_some()
587                || matching.grant.max_total_cost.is_some()
588        }) {
589            SideEffectClass::Monetary
590        } else if self
591            .tool_servers
592            .get(&request.server_id)
593            .is_some_and(|server| server.tool_is_read_only(&request.tool_name))
594        {
595            SideEffectClass::ReadOnly
596        } else {
597            SideEffectClass::SideEffecting
598        };
599        if !self.durable_admission_mode.covers(effect_class) {
600            if requires_structured_admission {
601                return Err(KernelError::DurableAdmission(
602                    "aggregate, cumulative, and supplemental authorization requires durable admission coverage"
603                        .to_string(),
604                ));
605            }
606            return Ok(None);
607        }
608        self.durable_stream_limits()?;
609        let Some(runtime) = self.durable_admission_runtime.as_ref() else {
610            if self.config.allow_ephemeral_receipt_log && !requires_structured_admission {
611                return Ok(None);
612            }
613            return Err(KernelError::DurableAdmission(
614                "no qualified admission operation store is configured".to_string(),
615            ));
616        };
617        let payment_required = effect_class == SideEffectClass::Monetary;
618        if payment_required {
619            let adapter = self.payment_adapter.as_ref().ok_or_else(|| {
620                KernelError::DurableAdmission(
621                    "durable monetary admission requires a qualified payment adapter".to_owned(),
622                )
623            })?;
624            if adapter.rail_id().is_empty()
625                || adapter.rail_id() == "unspecified"
626                || adapter.rail_mode().is_none()
627            {
628                return Err(KernelError::DurableAdmission(
629                    "durable monetary admission requires a recoverable payment rail identity"
630                        .to_owned(),
631                ));
632            }
633        }
634        if self.execution_nonce_config.is_some() {
635            return Err(KernelError::DurableAdmission(
636                "durable execution nonces require an atomic admission participant".to_owned(),
637            ));
638        }
639        let projection_capabilities = runtime.store.admission_projection_capabilities();
640        let observer_required = self.settlement_observer.is_some();
641        if !projection_capabilities.operation_terminal
642            || !projection_capabilities.tool_outcome
643            || (payment_required && !projection_capabilities.payment_terminal)
644            || (observer_required && !projection_capabilities.observation_attempt_zero)
645        {
646            return Err(KernelError::DurableAdmission(
647                "admission store lacks atomic terminal tool-outcome projection support".to_owned(),
648            ));
649        }
650        let post_return_plan = self.durable_post_return_plan()?;
651
652        let supplemental_authorization_artifact_digest = request
653            .supplemental_authorization
654            .as_ref()
655            .map(|authorization| {
656                supplemental_authorization_artifact_digest(
657                    authorization.signed_extension.as_bytes(),
658                )
659            });
660        let immutable_request_hash =
661            immutable_tool_admission_request_hash(request, matching_grants, &post_return_plan)?;
662        let action =
663            ToolCallAction::from_parameters(request.arguments.clone()).map_err(|error| {
664                KernelError::DurableAdmission(format!(
665                    "tool action parameters are invalid: {error}"
666                ))
667            })?;
668        let action_parameter_hash =
669            AdmissionDigest::try_new("action_parameter_hash", action.parameter_hash.clone())?;
670        let authorization_capability_hash =
671            admission_digest("authorization_capability_hash", &request.capability)?;
672        let policy_hash = AdmissionDigest::try_new("policy_hash", self.config.policy_hash.clone())
673            .map_err(|_| {
674                KernelError::DurableAdmission(
675                    "durable admission requires a canonical SHA-256 policy hash".to_owned(),
676                )
677            })?;
678        if cumulative_matching_grant_count != 0
679            && cumulative_matching_grant_count != matching_grants.len()
680        {
681            return Err(KernelError::DurableAdmission(
682                "matching grants disagree on cumulative approval requirements".to_owned(),
683            ));
684        }
685        let matching_grant_requires_cumulative_approval = cumulative_matching_grant_count != 0;
686        let requirements = AdmissionParticipantRequirements {
687            broker_attempt: true,
688            budget_capture: true,
689            approval: matching_grant_requires_cumulative_approval
690                || request.approval_token.is_some()
691                || !request.approval_tokens.is_empty(),
692            payment: payment_required,
693            observation_attempt_zero: observer_required,
694            ..AdmissionParticipantRequirements::NONE
695        };
696        let coordinator_authority_id = AdmissionIdentifier::try_new(
697            "coordinator_authority_id",
698            runtime.fence.store_uuid.clone(),
699        )?;
700        let namespace = match self.receipt_tenant_id_for_request(Some(&request.request_id)) {
701            Some(tenant_id) => AuthenticatedRequestNamespace::from_authentication_context(
702                coordinator_authority_id,
703                tenant_id,
704            )?,
705            None => AuthenticatedRequestNamespace::for_local_system(coordinator_authority_id)?,
706        };
707        let binding = AdmissionOperationBindingV1::new(AdmissionOperationBindingInputV1 {
708            kind: AdmissionOperationKind::ToolDispatch,
709            namespace,
710            request_id: AdmissionIdentifier::try_new("request_id", request.request_id.clone())?,
711            capability_id: AdmissionIdentifier::try_new(
712                "capability_id",
713                request.capability.id.clone(),
714            )?,
715            authorization_capability_hash,
716            request_binding: AdmissionRequestBindingV1::new_with_action_parameter_hash(
717                immutable_request_hash,
718                action_parameter_hash,
719                requirements,
720            )?,
721            policy_hash,
722            effect_class,
723        })?;
724        let supplemental_quota = self.verify_supplemental_quota_for_admission(
725            request,
726            &binding,
727            trusted_now_unix_ms / 1000,
728        )?;
729        let prepared = AdmissionOperationV1::prepare(binding, runtime.fence.owner_epoch)?;
730        let _mutation_guard = runtime.lock_mutations()?;
731        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
732        let operation = match runtime
733            .store
734            .begin(&prepared, &runtime.fence, trusted_now_unix_ms)
735            .map_err(durable_store_error)?
736        {
737            AdmissionBeginResult::Created(operation) => operation,
738            AdmissionBeginResult::ExactReplay { operation, .. }
739                if matches!(
740                    operation.state(),
741                    AdmissionOperationState::Prepared
742                        | AdmissionOperationState::BrokerAttemptRegistered
743                        | AdmissionOperationState::ApprovalRequired
744                        | AdmissionOperationState::BudgetAuthorized
745                        | AdmissionOperationState::ApprovalReserved
746                        | AdmissionOperationState::ReadyToDispatch
747                        | AdmissionOperationState::CapturePending
748                        | AdmissionOperationState::Finalizing
749                        | AdmissionOperationState::Completed
750                ) =>
751            {
752                operation
753            }
754            AdmissionBeginResult::ExactReplay { operation, .. } => {
755                return Err(KernelError::DurableAdmission(format!(
756                    "request replay is retained in state {:?}",
757                    operation.state()
758                )));
759            }
760            AdmissionBeginResult::Conflict {
761                existing_operation_id,
762            } => {
763                return Err(KernelError::DurableAdmission(format!(
764                    "request id conflicts with retained operation {}",
765                    existing_operation_id.as_str()
766                )));
767            }
768        };
769        let expected_operation_id = operation.binding().operation_id().as_str();
770        let expected_attempt_id = format!("attempt:{expected_operation_id}");
771        let expected_transport_id = format!("kernel-tool-server:{}", request.server_id);
772        let operation = match operation.state() {
773            AdmissionOperationState::Prepared => {
774                let expected_attempt = ProviderAttemptBindingV1 {
775                    operation_id: expected_operation_id.to_owned(),
776                    attempt_id: expected_attempt_id,
777                    transport_id: expected_transport_id,
778                    transport_key_epoch: runtime.fence.owner_epoch,
779                };
780                expected_attempt.validate().map_err(|error| {
781                    KernelError::DurableAdmission(format!(
782                        "provider attempt binding is invalid: {error}"
783                    ))
784                })?;
785                let mut attachments = Vec::with_capacity(2);
786                if let Some(digest) = supplemental_authorization_artifact_digest.as_ref() {
787                    attachments.push(AdmissionAttachment::SupplementalAuthorizationDigest(
788                        AdmissionDigest::try_new(
789                            "supplemental_authorization_digest",
790                            digest.clone(),
791                        )?,
792                    ));
793                }
794                attachments.push(AdmissionAttachment::BrokerAttempt(expected_attempt));
795                self.apply_admission_command(
796                    operation,
797                    attachments,
798                    AdmissionOperationState::BrokerAttemptRegistered,
799                    trusted_now_unix_ms,
800                )?
801            }
802            _ if operation.provider_attempt().is_some_and(|attempt| {
803                attempt.operation_id == expected_operation_id
804                    && attempt.attempt_id == expected_attempt_id
805                    && attempt.transport_id == expected_transport_id
806                    && attempt.transport_key_epoch <= operation.coordinator_lease_epoch()
807            }) =>
808            {
809                operation
810            }
811            _ => {
812                return Err(KernelError::DurableAdmission(
813                    "retained provider attempt does not match this dispatch".to_string(),
814                ));
815            }
816        };
817        if operation
818            .supplemental_authorization_digest()
819            .map(AdmissionDigest::as_str)
820            != supplemental_authorization_artifact_digest.as_deref()
821        {
822            return Err(KernelError::DurableAdmission(
823                "retained supplemental authorization digest does not match request".to_string(),
824            ));
825        }
826        Ok(Some(DurableToolAdmission {
827            operation,
828            aggregate_quota,
829            supplemental_quota,
830        }))
831    }
832
833    pub(crate) fn begin_durable_active_response_admission(
834        &self,
835        request: &crate::governed_active_response::GovernedActiveResponseRequest,
836        governed_intent_hash: &str,
837        trusted_now_unix_ms: u64,
838    ) -> Result<(DurableToolAdmission, bool), KernelError> {
839        let runtime = self.durable_runtime()?;
840        if !runtime
841            .store
842            .admission_projection_capabilities()
843            .operation_terminal
844        {
845            return Err(KernelError::DurableAdmission(
846                "active-response admission store lacks atomic terminal projection support"
847                    .to_owned(),
848            ));
849        }
850        let immutable_request_hash = admission_digest(
851            "immutable_request_hash",
852            &ImmutableActiveResponseAdmissionRequest {
853                schema: "chio.governed-active-response-admission.v1",
854                governed_intent: &request.governed_intent,
855                federated_origin_kernel_id: &request.federated_origin_kernel_id,
856                governed_intent_hash,
857            },
858        )?;
859        let authorization_capability_hash = admission_digest(
860            "authorization_capability_hash",
861            &request.operator_capability,
862        )?;
863        let policy_hash = AdmissionDigest::try_new("policy_hash", self.config.policy_hash.clone())
864            .map_err(|_| {
865                KernelError::DurableAdmission(
866                    "durable admission requires a canonical SHA-256 policy hash".to_owned(),
867                )
868            })?;
869        let requirements = AdmissionParticipantRequirements {
870            approval: true,
871            ..AdmissionParticipantRequirements::NONE
872        };
873        let coordinator_authority_id = AdmissionIdentifier::try_new(
874            "coordinator_authority_id",
875            runtime.fence.store_uuid.clone(),
876        )?;
877        let namespace = match self.receipt_tenant_id_for_request(Some(&request.request_id)) {
878            Some(tenant_id) => AuthenticatedRequestNamespace::from_authentication_context(
879                coordinator_authority_id,
880                tenant_id,
881            )?,
882            None => AuthenticatedRequestNamespace::for_local_system(coordinator_authority_id)?,
883        };
884        let binding = AdmissionOperationBindingV1::new(AdmissionOperationBindingInputV1 {
885            kind: AdmissionOperationKind::GovernedActiveResponse,
886            namespace,
887            request_id: AdmissionIdentifier::try_new("request_id", request.request_id.clone())?,
888            capability_id: AdmissionIdentifier::try_new(
889                "capability_id",
890                request.operator_capability.id.clone(),
891            )?,
892            authorization_capability_hash,
893            request_binding: AdmissionRequestBindingV1::new(immutable_request_hash, requirements)?,
894            policy_hash,
895            effect_class: SideEffectClass::SideEffecting,
896        })?;
897        let prepared = AdmissionOperationV1::prepare(binding, runtime.fence.owner_epoch)?;
898        let _mutation_guard = runtime.lock_mutations()?;
899        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
900        let (operation, created_by_this_attempt) = match runtime
901            .store
902            .begin(&prepared, &runtime.fence, trusted_now_unix_ms)
903            .map_err(durable_store_error)?
904        {
905            AdmissionBeginResult::Created(operation) => (operation, true),
906            AdmissionBeginResult::ExactReplay { operation, .. }
907                if matches!(
908                    operation.state(),
909                    AdmissionOperationState::Prepared
910                        | AdmissionOperationState::ApprovalReserved
911                        | AdmissionOperationState::ReadyToDispatch
912                        | AdmissionOperationState::DispatchCommitted
913                ) =>
914            {
915                (operation, false)
916            }
917            AdmissionBeginResult::ExactReplay { operation, .. } => {
918                return Err(KernelError::DurableAdmission(format!(
919                    "active-response request replay is retained in state {:?}",
920                    operation.state()
921                )));
922            }
923            AdmissionBeginResult::Conflict {
924                existing_operation_id,
925            } => {
926                return Err(KernelError::DurableAdmission(format!(
927                    "active-response request id conflicts with retained operation {}",
928                    existing_operation_id.as_str()
929                )));
930            }
931        };
932        Ok((
933            DurableToolAdmission {
934                operation,
935                aggregate_quota: None,
936                supplemental_quota: None,
937            },
938            created_by_this_attempt,
939        ))
940    }
941
942    fn durable_post_return_plan(&self) -> Result<DurablePostReturnPlan, KernelError> {
943        let hook_identities = self
944            .post_invocation_pipeline
945            .durable_identities()
946            .map_err(KernelError::DurableAdmission)?;
947        let mut frozen_steps = Vec::with_capacity(hook_identities.len() + 1);
948        frozen_steps.push(FrozenEvaluationStepV1 {
949            phase: EvaluationPhaseV1::OutputGuard,
950            position: 0,
951            component_id: AdmissionIdentifier::try_new(
952                "component_id",
953                "kernel-output-materialization",
954            )?,
955            component_version: AdmissionIdentifier::try_new("component_version", "v1")?,
956            implementation_digest: AdmissionDigest::try_new(
957                "implementation_digest",
958                sha256_hex(b"chio.kernel-output-materialization.v1"),
959            )?,
960            mode: EvaluationModeV1::Pure,
961        });
962        for (index, identity) in hook_identities.iter().enumerate() {
963            let position = u32::try_from(index + 1).map_err(|_| {
964                KernelError::DurableAdmission(
965                    "post-invocation pipeline has too many durable steps".to_owned(),
966                )
967            })?;
968            frozen_steps.push(FrozenEvaluationStepV1 {
969                phase: EvaluationPhaseV1::OutputGuard,
970                position,
971                component_id: AdmissionIdentifier::try_new(
972                    "component_id",
973                    identity.component_id(),
974                )?,
975                component_version: AdmissionIdentifier::try_new(
976                    "component_version",
977                    identity.component_version(),
978                )?,
979                implementation_digest: AdmissionDigest::try_new(
980                    "implementation_digest",
981                    identity.implementation_digest(),
982                )?,
983                mode: EvaluationModeV1::Pure,
984            });
985        }
986        Ok(DurablePostReturnPlan {
987            hook_identities,
988            frozen_steps,
989        })
990    }
991
992    fn verify_supplemental_quota_for_admission(
993        &self,
994        request: &ToolCallRequest,
995        binding: &AdmissionOperationBindingV1,
996        now: u64,
997    ) -> Result<Option<KernelVerifiedSupplementalQuotaClaim>, KernelError> {
998        let Some(authorization) = request.supplemental_authorization.as_ref() else {
999            return Ok(None);
1000        };
1001        let runtime = self.supplemental_quota_verifier.as_ref().ok_or_else(|| {
1002            KernelError::DurableAdmission(SupplementalQuotaError::MissingVerifier.to_string())
1003        })?;
1004
1005        #[derive(Serialize)]
1006        struct ToolDestination<'a> {
1007            server_id: &'a str,
1008            tool_name: &'a str,
1009        }
1010
1011        let normalized_destination = String::from_utf8(
1012            canonical_json_bytes(&ToolDestination {
1013                server_id: &request.server_id,
1014                tool_name: &request.tool_name,
1015            })
1016            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?,
1017        )
1018        .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1019        let capability_digest = sha256_hex(
1020            &canonical_json_bytes(&request.capability)
1021                .map_err(|error| KernelError::DurableAdmission(error.to_string()))?,
1022        );
1023        let arguments_hash = sha256_hex(
1024            &canonical_json_bytes(&request.arguments)
1025                .map_err(|error| KernelError::DurableAdmission(error.to_string()))?,
1026        );
1027        let mut negotiated_features = self
1028            .capability_negotiation_for_remote(request.federated_origin_kernel_id.as_deref(), now)
1029            .map_err(KernelError::DurableAdmission)?;
1030        if request.federated_origin_kernel_id.is_none() {
1031            negotiated_features
1032                .features
1033                .insert(BROKER_CAPABILITY_EXECUTION_PROFILE.to_string(), true);
1034        }
1035        let context = SupplementalQuotaVerificationContext {
1036            capability_id: request.capability.id.clone(),
1037            capability_digest,
1038            request_namespace_digest: binding.request_namespace_digest().as_str().to_string(),
1039            operation_id: binding.operation_id().as_str().to_string(),
1040            subject: request.capability.subject.clone(),
1041            request_id: request.request_id.clone(),
1042            normalized_destination,
1043            arguments_hash,
1044            negotiated_profile: BROKER_CAPABILITY_EXECUTION_PROFILE.to_string(),
1045            negotiated_features,
1046            verifier_binding: runtime.binding().clone(),
1047        };
1048        verify_supplemental_quota(
1049            Some(runtime.verifier()),
1050            authorization.signed_extension.as_bytes(),
1051            &context,
1052            now,
1053        )
1054        .map(Some)
1055        .map_err(|error| KernelError::DurableAdmission(error.to_string()))
1056    }
1057
1058    fn verify_aggregate_quota_for_admission(
1059        &self,
1060        request: &ToolCallRequest,
1061        now: u64,
1062    ) -> Result<Option<BudgetInvocationQuota>, KernelError> {
1063        use chio_core::capability::aggregate_invocation::{
1064            verify_aggregate_invocation_budget, AggregateInvocationScope,
1065        };
1066
1067        if request.capability.aggregate_invocation_budget.is_none() {
1068            return Ok(None);
1069        }
1070        let peer = self
1071            .capability_negotiation_for_remote(request.federated_origin_kernel_id.as_deref(), now)
1072            .map_err(KernelError::DurableAdmission)?;
1073        let direct_root = self
1074            .negotiated_capability_root(&request.capability, &peer)
1075            .map_err(KernelError::DurableAdmission)?;
1076        let verified = verify_aggregate_invocation_budget(
1077            &request.capability,
1078            &self.trusted_issuer_keys(),
1079            direct_root.as_ref(),
1080        )
1081        .map_err(|error| KernelError::DurableAdmission(error.to_string()))?
1082        .ok_or_else(|| {
1083            KernelError::DurableAdmission(
1084                "aggregate invocation budget did not produce a verified quota".to_string(),
1085            )
1086        })?;
1087        let profile = match verified.scope {
1088            AggregateInvocationScope::Capability => {
1089                BudgetQuotaProfile::AggregateCapabilityInvocation
1090            }
1091            AggregateInvocationScope::DelegationFamily => {
1092                BudgetQuotaProfile::AggregateFamilyInvocation
1093            }
1094        };
1095        Ok(Some(BudgetInvocationQuota {
1096            key: BudgetQuotaKey {
1097                profile,
1098                owner_id: verified.owner_id,
1099                grant_index: None,
1100            },
1101            max_invocations: verified.max_invocations,
1102        }))
1103    }
1104
1105    pub(crate) fn durable_budget_binding(
1106        &self,
1107        admission: &DurableToolAdmission,
1108        capability: &CapabilityToken,
1109    ) -> Result<(BudgetAdmissionBinding, BudgetEventAuthority), KernelError> {
1110        let runtime = self.durable_runtime()?;
1111        let ancestor_capability_ids = capability
1112            .delegation_chain
1113            .iter()
1114            .map(|link| link.capability_id.clone())
1115            .collect::<Vec<_>>();
1116        let revocation_set = match admission.supplemental_quota() {
1117            Some(claim) => canonical_revocation_set_for_verified_claim(
1118                &capability.id,
1119                &ancestor_capability_ids,
1120                claim,
1121            ),
1122            None => {
1123                let mut revocation_ids = Vec::with_capacity(ancestor_capability_ids.len() + 1);
1124                revocation_ids.push(capability.id.clone());
1125                revocation_ids.extend(ancestor_capability_ids);
1126                CanonicalRevocationSet::canonicalize(revocation_ids)
1127            }
1128        }
1129        .map_err(|error| {
1130            KernelError::DurableAdmission(format!("capability revocation set is invalid: {error}"))
1131        })?;
1132        let supplemental = admission.supplemental_quota();
1133        let last_observed_revocation = if supplemental.is_some() {
1134            let observation = self
1135                .revocation_store
1136                .observe_revocation(&capability.id)
1137                .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1138            if observation.revoked {
1139                return Err(KernelError::DurableAdmission(
1140                    "supplemental authorization leaf capability is revoked".to_string(),
1141                ));
1142            }
1143            let commit = observation.commit.ok_or_else(|| {
1144                KernelError::DurableAdmission(
1145                    "supplemental authorization requires atomic revocation observation".to_string(),
1146                )
1147            })?;
1148            if !matches!(
1149                commit.guarantee_level,
1150                BudgetGuaranteeLevel::SingleNodeAtomic | BudgetGuaranteeLevel::HaLinearizable
1151            ) {
1152                return Err(KernelError::DurableAdmission(
1153                    "supplemental authorization revocation observation is not atomic".to_string(),
1154                ));
1155            }
1156            Some(commit)
1157        } else {
1158            None
1159        };
1160        let authorization_artifact_digests = supplemental
1161            .map(|claim| vec![claim.authorization_artifact_digest().to_string()])
1162            .unwrap_or_default();
1163        Ok((
1164            BudgetAdmissionBinding {
1165                operation_id: admission.operation_id().to_string(),
1166                revocation_set,
1167                authorization_artifact_digests,
1168                last_observed_revocation,
1169                supplemental_verifier_id: supplemental
1170                    .map(|claim| claim.verifier_binding().verifier_identity.clone()),
1171                supplemental_verifier_config_digest: supplemental
1172                    .map(|claim| claim.verifier_binding().configuration_digest.clone()),
1173                supplemental_authorization_artifact_digest: supplemental
1174                    .map(|claim| claim.authorization_artifact_digest().to_string()),
1175                supplemental_authorization_expires_at: supplemental
1176                    .map(KernelVerifiedSupplementalQuotaClaim::expires_at),
1177            },
1178            runtime.authority(),
1179        ))
1180    }
1181
1182    pub(crate) fn authorize_durable_budget_hold(
1183        &self,
1184        admission: &mut DurableToolAdmission,
1185        request: crate::budget_store::BudgetAuthorizeHoldRequest,
1186        payment_journal: Option<crate::payment::PaymentJournalRecord>,
1187        trusted_now_unix_ms: u64,
1188    ) -> Result<crate::budget_store::BudgetAuthorizeHoldDecision, KernelError> {
1189        let runtime = self.durable_runtime()?;
1190        let _mutation_guard = runtime.lock_mutations()?;
1191        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1192        let expected = admission.operation.clone();
1193        let hold_id = request.hold_id.clone().ok_or_else(|| {
1194            KernelError::DurableAdmission(
1195                "combined durable authorization omitted its budget hold identity".to_owned(),
1196            )
1197        })?;
1198        if request.authority.as_ref() != Some(&runtime.authority()) {
1199            return Err(KernelError::DurableAdmission(
1200                "combined durable authorization authority does not match the admission fence"
1201                    .to_owned(),
1202            ));
1203        }
1204        let recovery_lease =
1205            self.claim_admission_recovery(&admission.operation, trusted_now_unix_ms)?;
1206        let authorization = runtime
1207            .store
1208            .authorize_budget_and_commit_admission(
1209                &admission.operation,
1210                &recovery_lease,
1211                request,
1212                payment_journal,
1213                None,
1214                &runtime.fence,
1215                trusted_now_unix_ms,
1216            )
1217            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1218        if authorization.operation.binding() != expected.binding() {
1219            return Err(KernelError::DurableAdmission(
1220                "combined durable authorization changed the immutable operation binding".to_owned(),
1221            ));
1222        }
1223        match &authorization.decision {
1224            crate::budget_store::BudgetAuthorizeHoldDecision::Authorized(_) => {
1225                if !matches!(
1226                    authorization.operation.state(),
1227                    AdmissionOperationState::BudgetAuthorized
1228                        | AdmissionOperationState::ApprovalReserved
1229                        | AdmissionOperationState::ReadyToDispatch
1230                        | AdmissionOperationState::CapturePending
1231                        | AdmissionOperationState::DispatchCommitted
1232                        | AdmissionOperationState::Finalizing
1233                        | AdmissionOperationState::Completed
1234                ) || authorization
1235                    .operation
1236                    .budget_hold_id()
1237                    .is_none_or(|bound| bound.as_str() != hold_id)
1238                {
1239                    return Err(KernelError::DurableAdmission(
1240                        "combined durable authorization returned an unbound operation".to_owned(),
1241                    ));
1242                }
1243            }
1244            crate::budget_store::BudgetAuthorizeHoldDecision::Denied(_)
1245            | crate::budget_store::BudgetAuthorizeHoldDecision::ApprovalRequired(_)
1246                if authorization.operation == expected => {}
1247            crate::budget_store::BudgetAuthorizeHoldDecision::AlreadyCaptured(_)
1248                if authorization.operation.state() == AdmissionOperationState::CapturePending
1249                    && authorization
1250                        .operation
1251                        .budget_hold_id()
1252                        .is_some_and(|bound| bound.as_str() == hold_id) => {}
1253            _ => {
1254                return Err(KernelError::DurableAdmission(format!(
1255                    "combined durable authorization returned incompatible operation state {:?}",
1256                    authorization.operation.state()
1257                )));
1258            }
1259        }
1260        admission.operation = authorization.operation;
1261        Ok(authorization.decision)
1262    }
1263
1264    pub(crate) fn reserve_durable_approval_set(
1265        &self,
1266        admission: &mut DurableToolAdmission,
1267        verified: &VerifiedApprovalReservation,
1268        trusted_now_unix_ms: u64,
1269    ) -> Result<(), KernelError> {
1270        let proposal_hash = AdmissionDigest::try_new(
1271            "threshold_proposal_hash",
1272            verified.threshold_proposal_hash.clone(),
1273        )?;
1274        let approval_set_hash =
1275            AdmissionDigest::try_new("approval_set_hash", verified.approval_set_hash.clone())?;
1276        if matches!(
1277            admission.operation.state(),
1278            AdmissionOperationState::ApprovalReserved
1279                | AdmissionOperationState::ReadyToDispatch
1280                | AdmissionOperationState::CapturePending
1281                | AdmissionOperationState::DispatchCommitted
1282                | AdmissionOperationState::Finalizing
1283                | AdmissionOperationState::Completed
1284        ) {
1285            let proposal_matches =
1286                admission.operation.threshold_proposal_hash() == Some(&proposal_hash);
1287            let set_matches = admission.operation.approval_set_hash() == Some(&approval_set_hash);
1288            if proposal_matches && set_matches {
1289                return Ok(());
1290            }
1291            return Err(KernelError::DurableAdmission(
1292                "retained approval reservation does not match the verified approval set".to_owned(),
1293            ));
1294        }
1295        let required_source = match admission.operation.binding().kind() {
1296            AdmissionOperationKind::ToolDispatch => AdmissionOperationState::BudgetAuthorized,
1297            AdmissionOperationKind::GovernedActiveResponse => AdmissionOperationState::Prepared,
1298            AdmissionOperationKind::GovernedEconomicMutation => {
1299                return Err(KernelError::DurableAdmission(
1300                    "economic mutation admission does not accept threshold approval sets"
1301                        .to_owned(),
1302                ));
1303            }
1304        };
1305        if admission.operation.state() != required_source {
1306            return Err(KernelError::DurableAdmission(format!(
1307                "approval reservation requires {required_source:?}, found {:?}",
1308                admission.operation.state()
1309            )));
1310        }
1311        let runtime = self.durable_runtime()?;
1312        let _mutation_guard = runtime.lock_mutations()?;
1313        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1314        let attachments = vec![
1315            AdmissionAttachment::ThresholdProposalHash(proposal_hash),
1316            AdmissionAttachment::ApprovalSetHash(approval_set_hash),
1317        ];
1318        admission.operation = if let Some(replay) = verified.threshold_replay.as_ref() {
1319            let expires_at_unix_ms = trusted_now_unix_ms
1320                .checked_add(RECOVERY_LEASE_DURATION_MS)
1321                .ok_or_else(|| {
1322                    KernelError::DurableAdmission("recovery lease expiration overflowed".to_owned())
1323                })?;
1324            let lease = runtime
1325                .store
1326                .claim_recovery(
1327                    admission.operation.binding().operation_id(),
1328                    admission.operation.version(),
1329                    &runtime.claimant_id,
1330                    trusted_now_unix_ms,
1331                    expires_at_unix_ms,
1332                    &runtime.fence,
1333                )
1334                .map_err(durable_store_error)?;
1335            let command = AdmissionOperationCommand::new(
1336                admission.operation.binding().operation_id().clone(),
1337                admission.operation.version(),
1338                lease,
1339                attachments,
1340                Some(AdmissionOperationState::ApprovalReserved),
1341                None,
1342                None,
1343            )?;
1344            runtime
1345                .store
1346                .reserve_threshold_approval_and_commit_admission(
1347                    &command,
1348                    replay,
1349                    trusted_now_unix_ms,
1350                )
1351                .map(|result| result.into_operation())
1352                .map_err(durable_store_error)?
1353        } else {
1354            self.apply_admission_command(
1355                admission.operation.clone(),
1356                attachments,
1357                AdmissionOperationState::ApprovalReserved,
1358                trusted_now_unix_ms,
1359            )?
1360        };
1361        Ok(())
1362    }
1363
1364    pub(crate) fn load_durable_payment_journal(
1365        &self,
1366        admission: &DurableToolAdmission,
1367    ) -> Result<crate::payment::PaymentJournalRecord, KernelError> {
1368        let runtime = self.durable_runtime()?;
1369        let _mutation_guard = runtime.lock_mutations()?;
1370        let journal = runtime
1371            .store
1372            .load_payment_journal(admission.operation_id(), &runtime.fence)
1373            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?
1374            .ok_or_else(|| {
1375                KernelError::DurableAdmission(
1376                    "durable payment participant is absent for the admission operation".to_owned(),
1377                )
1378            })?;
1379        journal
1380            .validate()
1381            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1382        if journal.operation_id != admission.operation_id() {
1383            return Err(KernelError::DurableAdmission(
1384                "durable payment participant changed operation identity".to_owned(),
1385            ));
1386        }
1387        Ok(journal)
1388    }
1389
1390    pub(crate) fn advance_durable_payment_journal(
1391        &self,
1392        admission: &DurableToolAdmission,
1393        expected: &crate::payment::PaymentJournalRecord,
1394        transition: &crate::payment::PaymentJournalTransition,
1395        trusted_now_unix_ms: u64,
1396    ) -> Result<crate::payment::PaymentJournalRecord, KernelError> {
1397        let runtime = self.durable_runtime()?;
1398        let _mutation_guard = runtime.lock_mutations()?;
1399        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1400        let recovery_lease =
1401            self.claim_admission_recovery(&admission.operation, trusted_now_unix_ms)?;
1402        let journal = runtime
1403            .store
1404            .advance_payment_journal(crate::receipt_store::AdmissionPaymentJournalAdvance {
1405                operation: &admission.operation,
1406                recovery_lease: &recovery_lease,
1407                expected,
1408                transition,
1409                release_evidence: None,
1410                active_fence: &runtime.fence,
1411                trusted_now_unix_ms,
1412            })
1413            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1414        journal
1415            .validate()
1416            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1417        Ok(journal)
1418    }
1419
1420    pub(crate) fn mark_durable_capture_pending(
1421        &self,
1422        admission: &mut DurableToolAdmission,
1423        trusted_now_unix_ms: u64,
1424    ) -> Result<(), KernelError> {
1425        let runtime = self.durable_runtime()?;
1426        let _mutation_guard = runtime.lock_mutations()?;
1427        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1428        if matches!(
1429            admission.operation.state(),
1430            AdmissionOperationState::BudgetAuthorized | AdmissionOperationState::ApprovalReserved
1431        ) {
1432            admission.operation = self.apply_admission_command(
1433                admission.operation.clone(),
1434                Vec::new(),
1435                AdmissionOperationState::ReadyToDispatch,
1436                trusted_now_unix_ms,
1437            )?;
1438        }
1439        if admission.operation.state() == AdmissionOperationState::ReadyToDispatch {
1440            admission.operation = self.apply_admission_command(
1441                admission.operation.clone(),
1442                Vec::new(),
1443                AdmissionOperationState::CapturePending,
1444                trusted_now_unix_ms,
1445            )?;
1446        }
1447        if admission.operation.state() != AdmissionOperationState::CapturePending {
1448            return Err(KernelError::DurableAdmission(format!(
1449                "capture cannot start from state {:?}",
1450                admission.operation.state()
1451            )));
1452        }
1453        Ok(())
1454    }
1455
1456    pub(crate) fn commit_durable_dispatch(
1457        &self,
1458        admission: &mut DurableToolAdmission,
1459        trusted_now_unix_ms: u64,
1460    ) -> Result<(), KernelError> {
1461        let runtime = self.durable_runtime()?;
1462        let _mutation_guard = runtime.lock_mutations()?;
1463        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1464        if admission.operation.state() == AdmissionOperationState::DispatchCommitted {
1465            return Ok(());
1466        }
1467        if admission.operation.binding().kind() == AdmissionOperationKind::GovernedActiveResponse
1468            && admission.operation.state() == AdmissionOperationState::ApprovalReserved
1469        {
1470            admission.operation = self.apply_admission_command(
1471                admission.operation.clone(),
1472                Vec::new(),
1473                AdmissionOperationState::ReadyToDispatch,
1474                trusted_now_unix_ms,
1475            )?;
1476        }
1477        let required_source = match admission.operation.binding().kind() {
1478            AdmissionOperationKind::ToolDispatch => AdmissionOperationState::CapturePending,
1479            AdmissionOperationKind::GovernedActiveResponse => {
1480                AdmissionOperationState::ReadyToDispatch
1481            }
1482            AdmissionOperationKind::GovernedEconomicMutation => {
1483                return Err(KernelError::DurableAdmission(
1484                    "economic mutation admission does not use dispatch commitment".to_owned(),
1485                ));
1486            }
1487        };
1488        if admission.operation.state() != required_source {
1489            return Err(KernelError::DurableAdmission(format!(
1490                "dispatch cannot commit from state {:?}; expected {required_source:?}",
1491                admission.operation.state()
1492            )));
1493        }
1494        admission.operation = self.apply_admission_command(
1495            admission.operation.clone(),
1496            Vec::new(),
1497            AdmissionOperationState::DispatchCommitted,
1498            trusted_now_unix_ms,
1499        )?;
1500        Ok(())
1501    }
1502
1503    pub(crate) fn compensate_durable_admission_before_dispatch(
1504        &self,
1505        operation: &AdmissionOperationV1,
1506        verifier_policy: serde_json::Value,
1507        trusted_now_unix_ms: u64,
1508    ) -> Result<(), KernelError> {
1509        let runtime = self.durable_runtime()?;
1510        let _mutation_guard = runtime.lock_mutations()?;
1511        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1512        let current = runtime
1513            .store
1514            .load_by_operation_id(operation.binding().operation_id())
1515            .map_err(durable_store_error)?
1516            .ok_or_else(|| {
1517                KernelError::DurableAdmission(
1518                    "pre-dispatch admission disappeared during compensation".to_owned(),
1519                )
1520            })?;
1521        if &current != operation
1522            || current.state().is_terminal()
1523            || current.dispatch_commit().is_some()
1524        {
1525            return Err(KernelError::DurableAdmission(
1526                "pre-dispatch compensation operation changed".to_owned(),
1527            ));
1528        }
1529        let lease = self.claim_admission_recovery(&current, trusted_now_unix_ms)?;
1530        let context = AdmissionProjectionContext {
1531            operation_id: current.binding().operation_id().clone(),
1532            request_id: current.binding().request_id().clone(),
1533            expected_operation_version: current.version(),
1534            trusted_time_unix_ms: trusted_now_unix_ms,
1535            coordinator_lease_id: lease.coordinator_lease_id().clone(),
1536            coordinator_lease_epoch: lease.coordinator_lease_epoch(),
1537            store_fence: runtime.fence.clone(),
1538        };
1539        if current
1540            .attachment(crate::admission_operation::AdmissionAttachmentKind::PaymentParticipant)
1541            .is_some()
1542        {
1543            let mut journal = runtime
1544                .store
1545                .load_payment_journal(current.binding().operation_id().as_str(), &runtime.fence)
1546                .map_err(|error| KernelError::DurableAdmission(error.to_string()))?
1547                .ok_or_else(|| {
1548                    KernelError::DurableAdmission(
1549                        "pre-dispatch payment journal disappeared".to_owned(),
1550                    )
1551                })?;
1552            if journal.state == crate::payment::PaymentJournalState::HoldPlaced {
1553                journal = runtime
1554                    .store
1555                    .advance_payment_journal(crate::receipt_store::AdmissionPaymentJournalAdvance {
1556                        operation: &current,
1557                        recovery_lease: &lease,
1558                        expected: &journal,
1559                        transition:
1560                            &crate::payment::PaymentJournalTransition::CancelBeforeAuthorization,
1561                        release_evidence: None,
1562                        active_fence: &runtime.fence,
1563                        trusted_now_unix_ms,
1564                    })
1565                    .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1566            }
1567            if journal.state == crate::payment::PaymentJournalState::Authorized {
1568                // The rail hold was authorized before dispatch, so the tool never
1569                // ran and the authorization must be released rather than left held.
1570                // Drive the durable release the same way the live cleanup path does:
1571                // record the no-effect release authority, advance the journal to
1572                // Settling, release on the rail, then settle. The release proof is
1573                // built from the acquired-participant snapshot, which the terminal
1574                // compensation projection below also accepts.
1575                let adapter = self.payment_adapter.as_ref().ok_or_else(|| {
1576                    KernelError::DurableAdmission(
1577                        "authorized pre-dispatch hold has no configured payment adapter".to_owned(),
1578                    )
1579                })?;
1580                let authorization_id = journal.authorization_id.clone().ok_or_else(|| {
1581                    KernelError::DurableAdmission(
1582                        "authorized payment journal omitted its authorization".to_owned(),
1583                    )
1584                })?;
1585                let proof =
1586                    crate::tool_outcome::VerifiedPreDispatchNoEffect::from_qualified_released_operation_snapshot(
1587                        &current,
1588                        &context,
1589                        verifier_policy.clone(),
1590                    )
1591                    .map_err(tool_outcome_error)?;
1592                let evidence = crate::tool_outcome::MonetaryReleaseAuthority::NoEffect(
1593                    crate::tool_outcome::VerifiedNoEffectProof::BeforeDispatch(proof),
1594                )
1595                .evidence_bundle()
1596                .map_err(tool_outcome_error)?;
1597                let persisted = evidence.to_persisted();
1598                let authority = crate::payment::PaymentReleaseAuthorityBinding {
1599                    kind: crate::payment::PaymentReleaseAuthorityKind::PreDispatchNoEffect,
1600                    operation_id: persisted.operation_id.as_str().to_owned(),
1601                    operation_version: persisted.operation_version,
1602                    evidence_id: persisted.evidence_id.as_str().to_owned(),
1603                    evidence_digest: persisted.bundle_digest.as_str().to_owned(),
1604                };
1605                journal = runtime
1606                    .store
1607                    .advance_payment_journal(crate::receipt_store::AdmissionPaymentJournalAdvance {
1608                        operation: &current,
1609                        recovery_lease: &lease,
1610                        expected: &journal,
1611                        transition: &crate::payment::PaymentJournalTransition::BeginRelease {
1612                            authority,
1613                        },
1614                        release_evidence: Some(&evidence),
1615                        active_fence: &runtime.fence,
1616                        trusted_now_unix_ms,
1617                    })
1618                    .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1619                let result = adapter
1620                    .release(&authorization_id, current.binding().request_id().as_str())
1621                    .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1622                if result.settlement_status != crate::payment::RailSettlementStatus::Released {
1623                    return Err(KernelError::DurableAdmission(
1624                        "pre-dispatch rail release was not confirmed".to_owned(),
1625                    ));
1626                }
1627                journal = runtime
1628                    .store
1629                    .advance_payment_journal(crate::receipt_store::AdmissionPaymentJournalAdvance {
1630                        operation: &current,
1631                        recovery_lease: &lease,
1632                        expected: &journal,
1633                        transition:
1634                            &crate::payment::PaymentJournalTransition::SettlementCompleted {
1635                                transaction_id: result.transaction_id,
1636                            },
1637                        release_evidence: None,
1638                        active_fence: &runtime.fence,
1639                        trusted_now_unix_ms,
1640                    })
1641                    .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1642            }
1643            let released = journal.state == crate::payment::PaymentJournalState::Settled
1644                && journal.settle_action == Some(crate::payment::PaymentSettleAction::Release);
1645            let cancelled_before_authorization = journal.state
1646                == crate::payment::PaymentJournalState::Closed
1647                && journal.authorization_id.is_none();
1648            if !released && !cancelled_before_authorization {
1649                return Err(KernelError::DurableAdmission(
1650                    "pre-dispatch payment release is not durable".to_owned(),
1651                ));
1652            }
1653        }
1654        let projection = verified_released_pre_dispatch_compensation_projection(
1655            &current,
1656            context,
1657            verifier_policy,
1658        )?;
1659        let terminal = runtime
1660            .store
1661            .commit_admission_projection(&projection)
1662            .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1663        if terminal.operation_id != *current.binding().operation_id()
1664            || terminal.state != AdmissionOperationState::CompensatedBeforeDispatch
1665        {
1666            return Err(KernelError::DurableAdmission(
1667                "pre-dispatch compensation committed a different terminal operation".to_owned(),
1668            ));
1669        }
1670        Ok(())
1671    }
1672
1673    pub(crate) fn capture_and_commit_durable_dispatch(
1674        &self,
1675        admission: &mut DurableToolAdmission,
1676        capability: &CapabilityToken,
1677        budget_mutation: &mut PreExecutionBudgetMutation,
1678        trusted_now_unix_ms: u64,
1679    ) -> Result<(), KernelError> {
1680        let runtime = self.durable_runtime()?;
1681        let _mutation_guard = runtime.lock_mutations()?;
1682        let trusted_now_unix_ms = runtime.refresh_trusted_time(trusted_now_unix_ms);
1683        let charge = budget_mutation.durable_hold_result_mut().ok_or_else(|| {
1684            KernelError::DurableAdmission(
1685                "combined dispatch commit requires an authorized budget hold".to_owned(),
1686            )
1687        })?;
1688        let request = BudgetCaptureInvocationRequest {
1689            capability_id: capability.id.clone(),
1690            grant_index: charge.grant_index,
1691            hold_id: charge.budget_hold_id.clone(),
1692            event_id: charge.capture_invocation_event_id(),
1693            trusted_time: None,
1694            authority: charge.authorize_metadata.authority.clone(),
1695        };
1696        match admission.operation.state() {
1697            AdmissionOperationState::CapturePending => {
1698                let recovery_lease =
1699                    self.claim_admission_recovery(&admission.operation, trusted_now_unix_ms)?;
1700                let capture = runtime
1701                    .store
1702                    .capture_invocation_and_commit_dispatch(
1703                        &admission.operation,
1704                        &recovery_lease,
1705                        request,
1706                        &runtime.fence,
1707                        trusted_now_unix_ms,
1708                    )
1709                    .map_err(|error| KernelError::DurableAdmission(error.to_string()))?;
1710                admission.operation = capture.operation;
1711                let mutation = match capture.decision {
1712                    crate::budget_store::BudgetInvocationCaptureDecision::Captured(mutation)
1713                    | crate::budget_store::BudgetInvocationCaptureDecision::AlreadyCaptured(
1714                        mutation,
1715                    ) => mutation,
1716                };
1717                charge.invocation_capture = Some(Box::new(mutation));
1718            }
1719            AdmissionOperationState::DispatchCommitted => {}
1720            state => {
1721                return Err(KernelError::DurableAdmission(format!(
1722                    "combined dispatch capture cannot resume from state {state:?}"
1723                )));
1724            }
1725        }
1726        Ok(())
1727    }
1728
1729    fn claim_admission_recovery(
1730        &self,
1731        operation: &AdmissionOperationV1,
1732        trusted_now_unix_ms: u64,
1733    ) -> Result<crate::admission_operation::AdmissionRecoveryLease, KernelError> {
1734        let runtime = self.durable_runtime()?;
1735        let expires_at_unix_ms = trusted_now_unix_ms
1736            .checked_add(RECOVERY_LEASE_DURATION_MS)
1737            .ok_or_else(|| {
1738                KernelError::DurableAdmission("recovery lease expiration overflowed".to_owned())
1739            })?;
1740        runtime
1741            .store
1742            .claim_recovery(
1743                operation.binding().operation_id(),
1744                operation.version(),
1745                &runtime.claimant_id,
1746                trusted_now_unix_ms,
1747                expires_at_unix_ms,
1748                &runtime.fence,
1749            )
1750            .map_err(durable_store_error)
1751    }
1752
1753    pub(super) fn apply_admission_command(
1754        &self,
1755        operation: AdmissionOperationV1,
1756        attachments: Vec<AdmissionAttachment>,
1757        next_state: AdmissionOperationState,
1758        trusted_now_unix_ms: u64,
1759    ) -> Result<AdmissionOperationV1, KernelError> {
1760        let runtime = self.durable_runtime()?;
1761        let expires_at_unix_ms = trusted_now_unix_ms
1762            .checked_add(RECOVERY_LEASE_DURATION_MS)
1763            .ok_or_else(|| {
1764                KernelError::DurableAdmission("recovery lease expiration overflowed".to_string())
1765            })?;
1766        let lease = runtime
1767            .store
1768            .claim_recovery(
1769                operation.binding().operation_id(),
1770                operation.version(),
1771                &runtime.claimant_id,
1772                trusted_now_unix_ms,
1773                expires_at_unix_ms,
1774                &runtime.fence,
1775            )
1776            .map_err(durable_store_error)?;
1777        let command = AdmissionOperationCommand::new(
1778            operation.binding().operation_id().clone(),
1779            operation.version(),
1780            lease,
1781            attachments,
1782            Some(next_state),
1783            None,
1784            None,
1785        )?;
1786        runtime
1787            .store
1788            .compare_and_swap(&command, trusted_now_unix_ms)
1789            .map(|result| result.into_operation())
1790            .map_err(durable_store_error)
1791    }
1792
1793    fn durable_runtime(&self) -> Result<&DurableAdmissionRuntime, KernelError> {
1794        self.durable_admission_runtime.as_ref().ok_or_else(|| {
1795            KernelError::DurableAdmission(
1796                "qualified admission operation store is unavailable".to_string(),
1797            )
1798        })
1799    }
1800
1801    fn durable_stream_limits(&self) -> Result<InvocationStreamLimitsV1, KernelError> {
1802        InvocationStreamLimitsV1::new(
1803            self.config.max_stream_total_bytes,
1804            self.config.memory_budget.max_stream_chunks,
1805            self.config.max_stream_duration_secs,
1806        )
1807        .map_err(|error| KernelError::DurableAdmission(error.to_string()))
1808    }
1809}
1810
1811fn admission_digest(
1812    field: &'static str,
1813    value: &impl Serialize,
1814) -> Result<AdmissionDigest, KernelError> {
1815    let canonical = canonical_json_bytes(value)
1816        .map_err(|error| KernelError::DurableAdmission(format!("{field}: {error}")))?;
1817    AdmissionDigest::try_new(field, sha256_hex(&canonical)).map_err(KernelError::from)
1818}
1819
1820fn invocation_output_to_server_output(output: &InvocationOutputV1) -> ToolServerOutput {
1821    let stream = |chunks: &[serde_json::Value]| ToolCallStream {
1822        chunks: chunks
1823            .iter()
1824            .cloned()
1825            .map(|data| ToolCallChunk { data })
1826            .collect(),
1827    };
1828    match output {
1829        InvocationOutputV1::Value { value } => ToolServerOutput::Value(value.clone()),
1830        InvocationOutputV1::CompleteStream { chunks } => {
1831            ToolServerOutput::Stream(ToolServerStreamResult::Complete(stream(chunks)))
1832        }
1833        InvocationOutputV1::IncompleteStream { chunks, reason } => {
1834            ToolServerOutput::Stream(ToolServerStreamResult::Incomplete {
1835                stream: stream(chunks),
1836                reason: reason.clone(),
1837            })
1838        }
1839    }
1840}
1841
1842fn durable_store_error(
1843    error: crate::admission_operation::AdmissionOperationStoreError,
1844) -> KernelError {
1845    KernelError::DurableAdmission(error.to_string())
1846}
1847
1848fn durable_outcome_store_error(error: ToolOutcomeStoreError) -> KernelError {
1849    KernelError::DurableAdmission(error.to_string())
1850}
1851
1852fn tool_outcome_error(error: crate::tool_outcome::ToolOutcomeError) -> KernelError {
1853    KernelError::DurableAdmission(error.to_string())
1854}