1use std::{fmt, sync::Arc};
2
3use serde::Serialize;
4
5use crate::budget_store::{
6 BudgetEventAuthority, BudgetGuaranteeLevel, BudgetInvocationQuotaUsage, BudgetQuotaProfile,
7 MAX_INVOCATION_QUOTAS_PER_ADMISSION,
8};
9use crate::supplemental_quota::CanonicalRevocationSet;
10
11use super::*;
12
13pub const ADMISSION_CAPTURE_GLOBAL_GENESIS_DIGEST: &str =
14 "0000000000000000000000000000000000000000000000000000000000000000";
15pub const ADMISSION_CAPTURE_MUTATION_KIND: &str = "combined_admission_capture";
16const ADMISSION_CAPTURE_RESPONSE_DOMAIN: &[u8] = b"chio.admission-capture-response.v1\0";
17#[cfg(test)]
18const ADMISSION_CAPTURE_QUALIFICATION_DOMAIN: &[u8] = b"chio.admission-capture-qualification.v1\0";
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub enum AdmissionCaptureDenialReason {
22 Revoked,
23 AuthorizationExpired,
24 InvocationQuotaExhausted,
25 MonetaryLimitExceeded,
26 ApprovalNotAuthorized,
27 ParticipantBindingMismatch,
28}
29
30impl AdmissionCaptureDenialReason {
31 const fn as_str(self) -> &'static str {
32 match self {
33 Self::Revoked => "revoked",
34 Self::AuthorizationExpired => "authorization_expired",
35 Self::InvocationQuotaExhausted => "invocation_quota_exhausted",
36 Self::MonetaryLimitExceeded => "monetary_limit_exceeded",
37 Self::ApprovalNotAuthorized => "approval_not_authorized",
38 Self::ParticipantBindingMismatch => "participant_binding_mismatch",
39 }
40 }
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
44pub struct AdmissionCaptureRequestV1 {
45 pub operation_id: AdmissionOperationId,
46 pub expected_operation_version: u64,
47 pub coordinator_lease_epoch: u64,
48 pub capability_id: AdmissionIdentifier,
49 pub grant_index: u32,
50 pub hold_id: AdmissionIdentifier,
51 pub event_id: AdmissionIdentifier,
52 pub revocation_set: CanonicalRevocationSet,
53 pub authorization_artifact_digests: Vec<AdmissionDigest>,
54 pub authorization_expires_at_unix_ms: u64,
55 pub expected_invocation_quota_usages: Vec<BudgetInvocationQuotaUsage>,
56 pub authority: BudgetEventAuthority,
57 pub previous_global_commit_sequence: u64,
58 pub previous_global_commit_digest: AdmissionDigest,
59 pub store_fence: StoreMutationFence,
60}
61
62impl AdmissionCaptureRequestV1 {
63 pub fn validate(&self) -> Result<(), AdmissionOperationError> {
64 validate_positive_ijson(
65 "expected_operation_version",
66 self.expected_operation_version,
67 )?;
68 validate_positive_ijson("coordinator_lease_epoch", self.coordinator_lease_epoch)?;
69 validate_positive_ijson(
70 "authorization_expires_at_unix_ms",
71 self.authorization_expires_at_unix_ms,
72 )?;
73 validate_store_fence(&self.store_fence)?;
74 validate_artifact_digests(&self.authorization_artifact_digests)?;
75 validate_authority_binding(&self.authority, &self.store_fence)?;
76 validate_global_predecessor(
77 self.previous_global_commit_sequence,
78 &self.previous_global_commit_digest,
79 )?;
80 validate_quota_usages(
81 &self.expected_invocation_quota_usages,
82 &self.capability_id,
83 self.grant_index,
84 )
85 }
86
87 pub fn validate_against(
88 &self,
89 operation: &AdmissionOperationV1,
90 ) -> Result<(), AdmissionOperationError> {
91 self.validate()?;
92 operation.validate()?;
93 if operation.state != AdmissionOperationState::CapturePending
94 || operation.version != self.expected_operation_version
95 || operation.coordinator_lease_epoch != self.coordinator_lease_epoch
96 {
97 return Err(AdmissionOperationError::CapturePreconditionMismatch);
98 }
99 if operation.binding.operation_id != self.operation_id
100 || operation.binding.capability_id != self.capability_id
101 || operation
102 .attachments
103 .existing_matches(&AdmissionAttachment::BudgetHoldId(self.hold_id.clone()))
104 != Some(true)
105 {
106 return Err(AdmissionOperationError::CaptureBindingMismatch);
107 }
108 Ok(())
109 }
110}
111
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub enum AdmissionCaptureDisposition {
114 Captured,
115 Denied(AdmissionCaptureDenialReason),
116}
117
118#[derive(Debug, Clone, PartialEq, Eq)]
119pub struct AdmissionCombinedCaptureCommit {
120 pub authority: BudgetEventAuthority,
121 pub guarantee_level: BudgetGuaranteeLevel,
122 pub previous_global_commit_sequence: u64,
123 pub previous_global_commit_digest: AdmissionDigest,
124 pub global_commit_sequence: u64,
125 pub global_commit_digest: AdmissionDigest,
126 pub store_fence: StoreMutationFence,
127}
128
129impl AdmissionCombinedCaptureCommit {
130 pub fn validate(&self) -> Result<(), AdmissionOperationError> {
131 validate_authority_binding(&self.authority, &self.store_fence)?;
132 validate_global_predecessor(
133 self.previous_global_commit_sequence,
134 &self.previous_global_commit_digest,
135 )?;
136 validate_positive_ijson("global_commit_sequence", self.global_commit_sequence)?;
137 if self.guarantee_level != BudgetGuaranteeLevel::SingleNodeAtomic
138 || self.global_commit_sequence
139 != self
140 .previous_global_commit_sequence
141 .checked_add(1)
142 .ok_or(AdmissionOperationError::CaptureCommitMismatch)?
143 || self.global_commit_digest == self.previous_global_commit_digest
144 {
145 return Err(AdmissionOperationError::CaptureCommitMismatch);
146 }
147 Ok(())
148 }
149}
150
151#[derive(Debug, Clone, PartialEq, Eq)]
152pub struct AdmissionCaptureRecordV1 {
153 pub operation_id: AdmissionOperationId,
154 pub operation_version: u64,
155 pub coordinator_lease_epoch: u64,
156 pub capability_id: AdmissionIdentifier,
157 pub grant_index: u32,
158 pub hold_id: AdmissionIdentifier,
159 pub event_id: AdmissionIdentifier,
160 pub revocation_set: CanonicalRevocationSet,
161 pub authorization_artifact_digests: Vec<AdmissionDigest>,
162 pub invocation_quota_usages: Vec<BudgetInvocationQuotaUsage>,
163 pub authorization_expires_at_unix_ms: u64,
164 pub authority_time_unix_ms: u64,
165 pub combined_commit: AdmissionCombinedCaptureCommit,
166 pub disposition: AdmissionCaptureDisposition,
167}
168
169impl AdmissionCaptureRecordV1 {
170 pub fn validate(&self) -> Result<(), AdmissionOperationError> {
173 validate_positive_ijson("operation_version", self.operation_version)?;
174 validate_positive_ijson("coordinator_lease_epoch", self.coordinator_lease_epoch)?;
175 validate_positive_ijson(
176 "authorization_expires_at_unix_ms",
177 self.authorization_expires_at_unix_ms,
178 )?;
179 validate_positive_ijson("authority_time_unix_ms", self.authority_time_unix_ms)?;
180 validate_artifact_digests(&self.authorization_artifact_digests)?;
181 validate_quota_usages(
182 &self.invocation_quota_usages,
183 &self.capability_id,
184 self.grant_index,
185 )?;
186 self.combined_commit.validate()?;
187 let expired = self.authority_time_unix_ms >= self.authorization_expires_at_unix_ms;
188 let expiry_disposition_matches = match self.disposition {
189 AdmissionCaptureDisposition::Captured => !expired,
190 AdmissionCaptureDisposition::Denied(
191 AdmissionCaptureDenialReason::AuthorizationExpired,
192 ) => expired,
193 AdmissionCaptureDisposition::Denied(_) => !expired,
194 };
195 if !expiry_disposition_matches {
196 return Err(AdmissionOperationError::CaptureExpiryMismatch);
197 }
198 Ok(())
199 }
200
201 fn matches_request(&self, request: &AdmissionCaptureRequestV1) -> bool {
202 self.operation_id == request.operation_id
203 && self.operation_version == request.expected_operation_version
204 && self.coordinator_lease_epoch == request.coordinator_lease_epoch
205 && self.capability_id == request.capability_id
206 && self.grant_index == request.grant_index
207 && self.hold_id == request.hold_id
208 && self.event_id == request.event_id
209 && self.revocation_set == request.revocation_set
210 && self.authorization_artifact_digests == request.authorization_artifact_digests
211 && self.authorization_expires_at_unix_ms == request.authorization_expires_at_unix_ms
212 && self.combined_commit.authority == request.authority
213 && self.combined_commit.previous_global_commit_sequence
214 == request.previous_global_commit_sequence
215 && self.combined_commit.previous_global_commit_digest
216 == request.previous_global_commit_digest
217 && self.combined_commit.store_fence == request.store_fence
218 }
219
220 fn quota_effect_matches_request(&self, request: &AdmissionCaptureRequestV1) -> bool {
221 match self.disposition {
222 AdmissionCaptureDisposition::Captured => {
223 request
224 .expected_invocation_quota_usages
225 .iter()
226 .zip(&self.invocation_quota_usages)
227 .all(|(before, after)| {
228 before.quota == after.quota
229 && before.reserved_invocations.checked_sub(1)
230 == Some(after.reserved_invocations)
231 && before.captured_invocations.checked_add(1)
232 == Some(after.captured_invocations)
233 })
234 && self.invocation_quota_usages.len()
235 == request.expected_invocation_quota_usages.len()
236 }
237 AdmissionCaptureDisposition::Denied(_) => {
238 self.invocation_quota_usages == request.expected_invocation_quota_usages
239 }
240 }
241 }
242}
243
244#[derive(Debug, Clone, Copy, PartialEq, Eq)]
245enum AdmissionCaptureResponseKind {
246 Captured,
247 AlreadyCaptured,
248 Denied,
249}
250
251impl AdmissionCaptureResponseKind {
252 const fn as_str(self) -> &'static str {
253 match self {
254 Self::Captured => "captured",
255 Self::AlreadyCaptured => "already_captured",
256 Self::Denied => "denied",
257 }
258 }
259}
260
261#[derive(Debug, Clone, PartialEq, Eq)]
262pub struct AdmissionCaptureAuthorityResponseV1 {
263 record: AdmissionCaptureRecordV1,
264 response_digest: AdmissionDigest,
265}
266
267impl AdmissionCaptureAuthorityResponseV1 {
268 fn record(&self) -> &AdmissionCaptureRecordV1 {
269 &self.record
270 }
271
272 fn response_digest(&self) -> &AdmissionDigest {
273 &self.response_digest
274 }
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
278pub enum AdmissionCaptureDecision {
279 Captured(AdmissionCaptureAuthorityResponseV1),
280 AlreadyCaptured(AdmissionCaptureAuthorityResponseV1),
281 Denied(AdmissionCaptureAuthorityResponseV1),
282}
283
284impl AdmissionCaptureDecision {
285 pub fn untrusted_captured(
287 record: AdmissionCaptureRecordV1,
288 ) -> Result<Self, AdmissionOperationError> {
289 Self::build(AdmissionCaptureResponseKind::Captured, record)
290 }
291
292 pub fn untrusted_already_captured(
294 record: AdmissionCaptureRecordV1,
295 ) -> Result<Self, AdmissionOperationError> {
296 Self::build(AdmissionCaptureResponseKind::AlreadyCaptured, record)
297 }
298
299 pub fn untrusted_denied(
301 record: AdmissionCaptureRecordV1,
302 ) -> Result<Self, AdmissionOperationError> {
303 Self::build(AdmissionCaptureResponseKind::Denied, record)
304 }
305
306 fn build(
307 kind: AdmissionCaptureResponseKind,
308 record: AdmissionCaptureRecordV1,
309 ) -> Result<Self, AdmissionOperationError> {
310 record.validate()?;
311 validate_disposition(kind, &record.disposition)?;
312 let response = AdmissionCaptureAuthorityResponseV1 {
313 response_digest: capture_response_digest(kind, &record)?,
314 record,
315 };
316 Ok(match kind {
317 AdmissionCaptureResponseKind::Captured => Self::Captured(response),
318 AdmissionCaptureResponseKind::AlreadyCaptured => Self::AlreadyCaptured(response),
319 AdmissionCaptureResponseKind::Denied => Self::Denied(response),
320 })
321 }
322
323 fn validate_for(
324 &self,
325 request: &AdmissionCaptureRequestV1,
326 ) -> Result<(), AdmissionOperationError> {
327 request.validate()?;
328 self.record().validate()?;
329 validate_disposition(self.kind(), &self.record().disposition)?;
330 if capture_response_digest(self.kind(), self.record())? != *self.response_digest() {
331 return Err(AdmissionOperationError::CaptureResponseDigestMismatch);
332 }
333 if !self.record().matches_request(request) {
334 return Err(AdmissionOperationError::CaptureBindingMismatch);
335 }
336 if !self.record().quota_effect_matches_request(request) {
337 return Err(AdmissionOperationError::CaptureQuotaMismatch);
338 }
339 Ok(())
340 }
341
342 fn kind(&self) -> AdmissionCaptureResponseKind {
343 match self {
344 Self::Captured(_) => AdmissionCaptureResponseKind::Captured,
345 Self::AlreadyCaptured(_) => AdmissionCaptureResponseKind::AlreadyCaptured,
346 Self::Denied(_) => AdmissionCaptureResponseKind::Denied,
347 }
348 }
349
350 fn response(&self) -> &AdmissionCaptureAuthorityResponseV1 {
351 match self {
352 Self::Captured(response) | Self::AlreadyCaptured(response) | Self::Denied(response) => {
353 response
354 }
355 }
356 }
357
358 fn record(&self) -> &AdmissionCaptureRecordV1 {
359 self.response().record()
360 }
361
362 fn response_digest(&self) -> &AdmissionDigest {
363 self.response().response_digest()
364 }
365}
366
367#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
368pub enum AdmissionCaptureError {
369 #[error("combined admission capture is unavailable: {0}")]
370 Unavailable(String),
371 #[error("combined admission capture was fenced")]
372 Fenced,
373 #[error("combined admission capture durable outcome is unknown: {0}")]
374 OutcomeUnknown(String),
375 #[error("combined admission capture invariant failed: {0}")]
376 Invariant(String),
377 #[error(transparent)]
378 Operation(#[from] AdmissionOperationError),
379}
380
381pub trait AdmissionCaptureAuthority: Send + Sync {
384 fn capture(
385 &self,
386 request: &AdmissionCaptureRequestV1,
387 ) -> Result<AdmissionCaptureDecision, AdmissionCaptureError>;
388
389 fn lookup_by_operation(
390 &self,
391 operation_id: &AdmissionOperationId,
392 ) -> Result<Option<AdmissionCaptureRecordV1>, AdmissionCaptureError>;
393}
394
395trait QualifiedAdmissionCaptureAdapter: AdmissionCaptureAuthority {
396 fn trusted_time_unix_ms(&self) -> Result<u64, AdmissionCaptureError>;
397
398 fn verifies_global_commit_link(
399 &self,
400 binding: &CaptureCommitVerification<'_>,
401 ) -> Result<bool, AdmissionCaptureError>;
402}
403
404struct CaptureCommitVerification<'a> {
405 mutation_kind: &'static str,
406 operation_id: &'a AdmissionOperationId,
407 event_id: &'a AdmissionIdentifier,
408 response_digest: &'a AdmissionDigest,
409 previous_sequence: u64,
410 previous_digest: &'a AdmissionDigest,
411 sequence: u64,
412 digest: &'a AdmissionDigest,
413}
414
415impl CaptureCommitVerification<'_> {
416 fn is_well_formed(&self) -> bool {
417 self.mutation_kind == ADMISSION_CAPTURE_MUTATION_KIND
418 && !self.operation_id.as_str().is_empty()
419 && !self.event_id.as_str().is_empty()
420 && !self.response_digest.as_str().is_empty()
421 && self.previous_sequence <= I_JSON_MAX_SAFE_INTEGER
422 && self.previous_sequence.checked_add(1) == Some(self.sequence)
423 && self.previous_digest != self.digest
424 }
425}
426
427pub struct QualifiedAdmissionCaptureVerifier {
428 authority: Arc<dyn QualifiedAdmissionCaptureAdapter>,
429 authority_binding: BudgetEventAuthority,
430 store_fence: StoreMutationFence,
431 qualification_digest: AdmissionDigest,
432}
433
434impl fmt::Debug for QualifiedAdmissionCaptureVerifier {
435 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
436 formatter
437 .debug_struct("QualifiedAdmissionCaptureVerifier")
438 .field("authority_binding", &self.authority_binding)
439 .field("store_fence", &self.store_fence)
440 .field("qualification_digest", &self.qualification_digest)
441 .finish_non_exhaustive()
442 }
443}
444
445impl QualifiedAdmissionCaptureVerifier {
446 #[must_use]
447 pub fn qualification_digest(&self) -> &AdmissionDigest {
448 &self.qualification_digest
449 }
450}
451
452#[derive(Debug, Clone, PartialEq, Eq)]
453pub struct QualifiedAdmissionCaptureDecision {
454 decision: AdmissionCaptureDecision,
455 qualification_digest: AdmissionDigest,
456 verified_at_unix_ms: u64,
457 replayed: bool,
458}
459
460impl QualifiedAdmissionCaptureDecision {
461 #[must_use]
462 pub fn record(&self) -> &AdmissionCaptureRecordV1 {
463 self.decision.record()
464 }
465
466 #[must_use]
467 pub fn response_digest(&self) -> &AdmissionDigest {
468 self.decision.response_digest()
469 }
470
471 #[must_use]
472 pub fn qualification_digest(&self) -> &AdmissionDigest {
473 &self.qualification_digest
474 }
475
476 #[must_use]
477 pub fn disposition(&self) -> &AdmissionCaptureDisposition {
478 &self.decision.record().disposition
479 }
480
481 #[must_use]
482 pub fn was_replay(&self) -> bool {
483 self.replayed
484 }
485
486 #[must_use]
487 pub fn verified_at_unix_ms(&self) -> u64 {
488 self.verified_at_unix_ms
489 }
490}
491
492pub fn capture_with_qualified_verifier(
493 operation: &AdmissionOperationV1,
494 request: &AdmissionCaptureRequestV1,
495 verifier: &QualifiedAdmissionCaptureVerifier,
496) -> Result<QualifiedAdmissionCaptureDecision, AdmissionCaptureError> {
497 validate_qualified_request(operation, request, verifier)?;
498 let decision = verifier.authority.capture(request)?;
499 let replayed = matches!(decision, AdmissionCaptureDecision::AlreadyCaptured(_));
500 qualify_capture_decision(request, verifier, decision, replayed)
501}
502
503pub fn lookup_capture_with_qualified_verifier(
504 operation: &AdmissionOperationV1,
505 request: &AdmissionCaptureRequestV1,
506 verifier: &QualifiedAdmissionCaptureVerifier,
507) -> Result<Option<QualifiedAdmissionCaptureDecision>, AdmissionCaptureError> {
508 validate_qualified_request(operation, request, verifier)?;
509 let Some(record) = verifier
510 .authority
511 .lookup_by_operation(&request.operation_id)?
512 else {
513 return Ok(None);
514 };
515 let decision = match record.disposition {
516 AdmissionCaptureDisposition::Captured => {
517 AdmissionCaptureDecision::untrusted_already_captured(record)?
518 }
519 AdmissionCaptureDisposition::Denied(_) => {
520 AdmissionCaptureDecision::untrusted_denied(record)?
521 }
522 };
523 qualify_capture_decision(request, verifier, decision, true).map(Some)
524}
525
526fn validate_qualified_request(
527 operation: &AdmissionOperationV1,
528 request: &AdmissionCaptureRequestV1,
529 verifier: &QualifiedAdmissionCaptureVerifier,
530) -> Result<(), AdmissionCaptureError> {
531 request.validate_against(operation)?;
532 if request.authority != verifier.authority_binding
533 || request.store_fence != verifier.store_fence
534 {
535 return Err(AdmissionOperationError::CaptureAuthorityMismatch.into());
536 }
537 Ok(())
538}
539
540fn qualify_capture_decision(
541 request: &AdmissionCaptureRequestV1,
542 verifier: &QualifiedAdmissionCaptureVerifier,
543 decision: AdmissionCaptureDecision,
544 replayed: bool,
545) -> Result<QualifiedAdmissionCaptureDecision, AdmissionCaptureError> {
546 decision.validate_for(request)?;
547 let record = decision.record();
548 let verified_at_unix_ms = verifier.authority.trusted_time_unix_ms()?;
549 validate_positive_ijson("capture_verified_at_unix_ms", verified_at_unix_ms)?;
550 if record.authority_time_unix_ms > verified_at_unix_ms {
551 return Err(AdmissionOperationError::CaptureAuthorityTimeMismatch.into());
552 }
553 let commit = &record.combined_commit;
554 let commit_verification = CaptureCommitVerification {
555 mutation_kind: ADMISSION_CAPTURE_MUTATION_KIND,
556 operation_id: &record.operation_id,
557 event_id: &record.event_id,
558 response_digest: decision.response_digest(),
559 previous_sequence: commit.previous_global_commit_sequence,
560 previous_digest: &commit.previous_global_commit_digest,
561 sequence: commit.global_commit_sequence,
562 digest: &commit.global_commit_digest,
563 };
564 if !commit_verification.is_well_formed()
565 || !verifier
566 .authority
567 .verifies_global_commit_link(&commit_verification)?
568 {
569 return Err(AdmissionOperationError::CaptureCommitMismatch.into());
570 }
571 Ok(QualifiedAdmissionCaptureDecision {
572 decision,
573 qualification_digest: verifier.qualification_digest.clone(),
574 verified_at_unix_ms,
575 replayed,
576 })
577}
578
579fn validate_disposition(
580 kind: AdmissionCaptureResponseKind,
581 disposition: &AdmissionCaptureDisposition,
582) -> Result<(), AdmissionOperationError> {
583 let valid = match kind {
584 AdmissionCaptureResponseKind::Captured | AdmissionCaptureResponseKind::AlreadyCaptured => {
585 *disposition == AdmissionCaptureDisposition::Captured
586 }
587 AdmissionCaptureResponseKind::Denied => {
588 matches!(disposition, AdmissionCaptureDisposition::Denied(_))
589 }
590 };
591 if valid {
592 Ok(())
593 } else {
594 Err(AdmissionOperationError::CaptureDispositionMismatch)
595 }
596}
597
598fn validate_authority_binding(
599 authority: &BudgetEventAuthority,
600 store_fence: &StoreMutationFence,
601) -> Result<(), AdmissionOperationError> {
602 authority
603 .validate()
604 .map_err(|_| AdmissionOperationError::CaptureAuthorityMismatch)?;
605 validate_positive_ijson("capture_authority_epoch", authority.lease_epoch)?;
606 validate_store_fence(store_fence)?;
607 if authority.authority_id != store_fence.store_uuid
608 || authority.lease_id != store_fence.lease_id
609 || authority.lease_epoch != store_fence.owner_epoch
610 {
611 return Err(AdmissionOperationError::CaptureAuthorityMismatch);
612 }
613 Ok(())
614}
615
616fn validate_global_predecessor(
617 sequence: u64,
618 digest: &AdmissionDigest,
619) -> Result<(), AdmissionOperationError> {
620 if sequence > I_JSON_MAX_SAFE_INTEGER
621 || (sequence == 0) != (digest.as_str() == ADMISSION_CAPTURE_GLOBAL_GENESIS_DIGEST)
622 {
623 return Err(AdmissionOperationError::CaptureCommitMismatch);
624 }
625 Ok(())
626}
627
628fn validate_quota_usages(
629 usages: &[BudgetInvocationQuotaUsage],
630 capability_id: &AdmissionIdentifier,
631 grant_index: u32,
632) -> Result<(), AdmissionOperationError> {
633 if usages.len() > MAX_INVOCATION_QUOTAS_PER_ADMISSION
634 || usages
635 .windows(2)
636 .any(|pair| pair[0].quota.key >= pair[1].quota.key)
637 {
638 return Err(AdmissionOperationError::CaptureQuotaMismatch);
639 }
640 for usage in usages {
641 usage
642 .validate()
643 .map_err(|_| AdmissionOperationError::CaptureQuotaMismatch)?;
644 AdmissionIdentifier::try_new("capture_quota_owner_id", usage.quota.key.owner_id.clone())
645 .map_err(|_| AdmissionOperationError::CaptureQuotaMismatch)?;
646 if usage.quota.key.profile == BudgetQuotaProfile::GrantInvocation
647 && (usage.quota.key.owner_id != capability_id.as_str()
648 || usage.quota.key.grant_index != Some(grant_index))
649 {
650 return Err(AdmissionOperationError::CaptureQuotaMismatch);
651 }
652 }
653 Ok(())
654}
655
656#[derive(Serialize)]
657struct CaptureQuotaUsageBody<'a> {
658 profile: &'static str,
659 owner_id: &'a str,
660 grant_index: Option<u32>,
661 max_invocations: u32,
662 reserved_invocations: u32,
663 captured_invocations: u32,
664}
665
666#[derive(Serialize)]
667struct CaptureResponseBody<'a> {
668 response_kind: &'static str,
669 operation_id: &'a AdmissionOperationId,
670 operation_version: u64,
671 coordinator_lease_epoch: u64,
672 capability_id: &'a AdmissionIdentifier,
673 grant_index: u32,
674 hold_id: &'a AdmissionIdentifier,
675 event_id: &'a AdmissionIdentifier,
676 revocation_ids: &'a [String],
677 revocation_digest: &'a str,
678 authorization_artifact_digests: &'a [AdmissionDigest],
679 invocation_quota_usages: Vec<CaptureQuotaUsageBody<'a>>,
680 authorization_expires_at_unix_ms: u64,
681 authority_time_unix_ms: u64,
682 authority_id: &'a str,
683 authority_lease_id: &'a str,
684 authority_epoch: u64,
685 guarantee_level: &'static str,
686 previous_global_commit_sequence: u64,
687 previous_global_commit_digest: &'a AdmissionDigest,
688 global_commit_sequence: u64,
689 global_commit_digest: &'a AdmissionDigest,
690 store_fence: &'a StoreMutationFence,
691 disposition: &'static str,
692 denial_reason: Option<&'static str>,
693}
694
695fn capture_response_digest(
696 kind: AdmissionCaptureResponseKind,
697 record: &AdmissionCaptureRecordV1,
698) -> Result<AdmissionDigest, AdmissionOperationError> {
699 let invocation_quota_usages = record
700 .invocation_quota_usages
701 .iter()
702 .map(|usage| CaptureQuotaUsageBody {
703 profile: usage.quota.key.profile.as_str(),
704 owner_id: &usage.quota.key.owner_id,
705 grant_index: usage.quota.key.grant_index,
706 max_invocations: usage.quota.max_invocations,
707 reserved_invocations: usage.reserved_invocations,
708 captured_invocations: usage.captured_invocations,
709 })
710 .collect();
711 let (disposition, denial_reason) = match record.disposition {
712 AdmissionCaptureDisposition::Captured => ("captured", None),
713 AdmissionCaptureDisposition::Denied(reason) => ("denied", Some(reason.as_str())),
714 };
715 let commit = &record.combined_commit;
716 let canonical = canonical_json_bytes(&CaptureResponseBody {
717 response_kind: kind.as_str(),
718 operation_id: &record.operation_id,
719 operation_version: record.operation_version,
720 coordinator_lease_epoch: record.coordinator_lease_epoch,
721 capability_id: &record.capability_id,
722 grant_index: record.grant_index,
723 hold_id: &record.hold_id,
724 event_id: &record.event_id,
725 revocation_ids: record.revocation_set.ids(),
726 revocation_digest: record.revocation_set.digest(),
727 authorization_artifact_digests: &record.authorization_artifact_digests,
728 invocation_quota_usages,
729 authorization_expires_at_unix_ms: record.authorization_expires_at_unix_ms,
730 authority_time_unix_ms: record.authority_time_unix_ms,
731 authority_id: &commit.authority.authority_id,
732 authority_lease_id: &commit.authority.lease_id,
733 authority_epoch: commit.authority.lease_epoch,
734 guarantee_level: commit.guarantee_level.as_str(),
735 previous_global_commit_sequence: commit.previous_global_commit_sequence,
736 previous_global_commit_digest: &commit.previous_global_commit_digest,
737 global_commit_sequence: commit.global_commit_sequence,
738 global_commit_digest: &commit.global_commit_digest,
739 store_fence: &commit.store_fence,
740 disposition,
741 denial_reason,
742 })
743 .map_err(|error| AdmissionOperationError::CanonicalJson(error.to_string()))?;
744 let mut bytes = Vec::with_capacity(ADMISSION_CAPTURE_RESPONSE_DOMAIN.len() + canonical.len());
745 bytes.extend_from_slice(ADMISSION_CAPTURE_RESPONSE_DOMAIN);
746 bytes.extend_from_slice(&canonical);
747 AdmissionDigest::try_new("capture_response_digest", sha256_hex(&bytes))
748}
749
750#[cfg(test)]
751#[derive(Clone)]
752pub(crate) struct TestCaptureQualification {
753 pub(crate) authority: BudgetEventAuthority,
754 pub(crate) store_fence: StoreMutationFence,
755 pub(crate) verified_at_unix_ms: u64,
756 pub(crate) operation_id: AdmissionOperationId,
757 pub(crate) event_id: AdmissionIdentifier,
758 pub(crate) response_digest: AdmissionDigest,
759 pub(crate) previous_global_commit_sequence: u64,
760 pub(crate) previous_global_commit_digest: AdmissionDigest,
761 pub(crate) global_commit_sequence: u64,
762 pub(crate) global_commit_digest: AdmissionDigest,
763}
764
765#[cfg(test)]
766struct TestQualifiedAdmissionCaptureAdapter {
767 authority: Arc<dyn AdmissionCaptureAuthority>,
768 qualification: TestCaptureQualification,
769}
770
771#[cfg(test)]
772impl AdmissionCaptureAuthority for TestQualifiedAdmissionCaptureAdapter {
773 fn capture(
774 &self,
775 request: &AdmissionCaptureRequestV1,
776 ) -> Result<AdmissionCaptureDecision, AdmissionCaptureError> {
777 self.authority.capture(request)
778 }
779
780 fn lookup_by_operation(
781 &self,
782 operation_id: &AdmissionOperationId,
783 ) -> Result<Option<AdmissionCaptureRecordV1>, AdmissionCaptureError> {
784 self.authority.lookup_by_operation(operation_id)
785 }
786}
787
788#[cfg(test)]
789impl QualifiedAdmissionCaptureAdapter for TestQualifiedAdmissionCaptureAdapter {
790 fn trusted_time_unix_ms(&self) -> Result<u64, AdmissionCaptureError> {
791 Ok(self.qualification.verified_at_unix_ms)
792 }
793
794 fn verifies_global_commit_link(
795 &self,
796 binding: &CaptureCommitVerification<'_>,
797 ) -> Result<bool, AdmissionCaptureError> {
798 Ok(binding.mutation_kind == ADMISSION_CAPTURE_MUTATION_KIND
799 && binding.operation_id == &self.qualification.operation_id
800 && binding.event_id == &self.qualification.event_id
801 && binding.response_digest == &self.qualification.response_digest
802 && binding.previous_sequence == self.qualification.previous_global_commit_sequence
803 && binding.previous_digest == &self.qualification.previous_global_commit_digest
804 && binding.sequence == self.qualification.global_commit_sequence
805 && binding.digest == &self.qualification.global_commit_digest)
806 }
807}
808
809#[cfg(test)]
810pub(crate) fn qualify_capture_authority_for_test(
811 authority: Arc<dyn AdmissionCaptureAuthority>,
812 qualification: TestCaptureQualification,
813) -> Result<QualifiedAdmissionCaptureVerifier, AdmissionOperationError> {
814 validate_authority_binding(&qualification.authority, &qualification.store_fence)?;
815 validate_positive_ijson(
816 "capture_verified_at_unix_ms",
817 qualification.verified_at_unix_ms,
818 )?;
819 validate_global_predecessor(
820 qualification.previous_global_commit_sequence,
821 &qualification.previous_global_commit_digest,
822 )?;
823 if qualification.global_commit_sequence
824 != qualification
825 .previous_global_commit_sequence
826 .checked_add(1)
827 .ok_or(AdmissionOperationError::CaptureCommitMismatch)?
828 || qualification.global_commit_digest == qualification.previous_global_commit_digest
829 {
830 return Err(AdmissionOperationError::CaptureCommitMismatch);
831 }
832 #[derive(Serialize)]
833 struct QualificationBody<'a> {
834 authority_id: &'a str,
835 authority_lease_id: &'a str,
836 authority_epoch: u64,
837 store_fence: &'a StoreMutationFence,
838 verifier_id: &'static str,
839 rollback_anchor_id: &'static str,
840 trusted_clock_id: &'static str,
841 }
842 let canonical = canonical_json_bytes(&QualificationBody {
843 authority_id: &qualification.authority.authority_id,
844 authority_lease_id: &qualification.authority.lease_id,
845 authority_epoch: qualification.authority.lease_epoch,
846 store_fence: &qualification.store_fence,
847 verifier_id: "test-combined-capture-verifier",
848 rollback_anchor_id: "test-global-authority-anchor",
849 trusted_clock_id: "test-authority-clock",
850 })
851 .map_err(|error| AdmissionOperationError::CanonicalJson(error.to_string()))?;
852 let mut bytes =
853 Vec::with_capacity(ADMISSION_CAPTURE_QUALIFICATION_DOMAIN.len() + canonical.len());
854 bytes.extend_from_slice(ADMISSION_CAPTURE_QUALIFICATION_DOMAIN);
855 bytes.extend_from_slice(&canonical);
856 Ok(QualifiedAdmissionCaptureVerifier {
857 authority: Arc::new(TestQualifiedAdmissionCaptureAdapter {
858 authority,
859 qualification: qualification.clone(),
860 }),
861 authority_binding: qualification.authority,
862 store_fence: qualification.store_fence,
863 qualification_digest: AdmissionDigest::try_new(
864 "capture_qualification_digest",
865 sha256_hex(&bytes),
866 )?,
867 })
868}
869
870#[cfg(test)]
871#[path = "capture_tests.rs"]
872mod tests;