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 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 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 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 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 ¤t != 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(¤t, 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: ¤t,
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 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 ¤t,
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: ¤t,
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: ¤t,
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 ¤t,
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}