1use ed25519_dalek::{Signature, Verifier, VerifyingKey};
12use serde::Serialize;
13use serde_json::Value;
14use std::collections::{BTreeMap, HashMap};
15use std::sync::Mutex;
16use std::time::{SystemTime, UNIX_EPOCH};
17
18use traverse_contracts::{
19 CanonicalProposal, CapabilityContract, DataClassification, EffectClass, EgressPolicy,
20 MappingSource, ProposalNode, is_automatic_eligible,
21};
22use traverse_registry::{ApplicationBundleManifest, CapabilityRegistry, LookupScope};
23
24use crate::{
25 PlacementTarget, Runtime, RuntimeContext, RuntimeIntent, RuntimeLookup, RuntimeLookupScope,
26 RuntimeRequest, RuntimeResultStatus,
27};
28
29#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
34pub struct ProposalCrossValidationError {
35 pub code: ProposalCrossValidationErrorCode,
36 pub message: String,
37 pub path: String,
38}
39
40#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
41#[serde(rename_all = "snake_case")]
42pub enum ProposalCrossValidationErrorCode {
43 UndeclaredCapability,
44 CapabilityNotFound,
45 ArtifactDigestMismatch,
46 IncompatibleMappingSchema,
47 UndeclaredDataClassification,
48 DataClassificationOverAccepted,
49 EgressDeniedForClassifiedMapping,
50}
51
52#[derive(Debug, Clone, PartialEq, Eq)]
53pub struct ProposalCrossValidationFailure {
54 pub errors: Vec<ProposalCrossValidationError>,
55}
56
57#[derive(Debug, Clone, PartialEq)]
59pub struct ResolvedProposalNode {
60 pub node_id: String,
61 pub contract: CapabilityContract,
62}
63
64#[allow(clippy::too_many_lines)]
84pub fn validate_proposal_against_host_state(
85 canonical: &CanonicalProposal,
86 manifest: &ApplicationBundleManifest,
87 registry: &CapabilityRegistry,
88) -> Result<Vec<ResolvedProposalNode>, ProposalCrossValidationFailure> {
89 let mut errors = Vec::new();
90 let mut resolved: BTreeMap<String, ResolvedProposalNode> = BTreeMap::new();
91
92 for (index, node) in canonical.proposal.nodes.iter().enumerate() {
93 let path = format!("$.nodes[{index}]");
94 if !manifest_declares_capability(manifest, node) {
95 errors.push(cross_error(
96 ProposalCrossValidationErrorCode::UndeclaredCapability,
97 &path,
98 &format!(
99 "capability '{}@{}' is not declared in the application manifest",
100 node.capability_id, node.capability_version
101 ),
102 ));
103 continue;
104 }
105
106 let Some(capability) = registry.find_exact(
107 LookupScope::PreferPrivate,
108 &node.capability_id,
109 &node.capability_version,
110 ) else {
111 errors.push(cross_error(
112 ProposalCrossValidationErrorCode::CapabilityNotFound,
113 &path,
114 &format!(
115 "capability '{}@{}' was not found in the registry",
116 node.capability_id, node.capability_version
117 ),
118 ));
119 continue;
120 };
121
122 let registry_digest = capability
123 .artifact
124 .digests
125 .binary_digest
126 .clone()
127 .unwrap_or_else(|| capability.artifact.digests.source_digest.clone());
128 if registry_digest != node.artifact_digest {
129 errors.push(cross_error(
130 ProposalCrossValidationErrorCode::ArtifactDigestMismatch,
131 &format!("{path}.artifact_digest"),
132 &format!(
133 "proposal pins digest '{}' but the registry resolves '{}@{}' to digest '{registry_digest}'",
134 node.artifact_digest, node.capability_id, node.capability_version
135 ),
136 ));
137 continue;
138 }
139
140 resolved.insert(
141 node.node_id.clone(),
142 ResolvedProposalNode {
143 node_id: node.node_id.clone(),
144 contract: capability.contract,
145 },
146 );
147 }
148
149 if !errors.is_empty() {
150 return Err(ProposalCrossValidationFailure { errors });
151 }
152
153 for (index, mapping) in canonical.proposal.mappings.iter().enumerate() {
154 let path = format!("$.mappings[{index}]");
155 let Some(target) = resolved.get(&mapping.target_node_id) else {
156 continue;
157 };
158
159 let source_output_schema = match &mapping.source {
160 MappingSource::InitialInput => None,
161 MappingSource::Node { node_id } => {
162 resolved.get(node_id).map(|n| &n.contract.outputs.schema)
163 }
164 };
165 if let Some(source_schema) = source_output_schema
166 && !mapping_schema_compatible(
167 source_schema,
168 &mapping.source_path,
169 &target.contract.inputs.schema,
170 &mapping.target_path,
171 )
172 {
173 errors.push(cross_error(
174 ProposalCrossValidationErrorCode::IncompatibleMappingSchema,
175 &path,
176 &format!(
177 "source path '{}' and target path '{}' declare incompatible JSON Schema types",
178 mapping.source_path, mapping.target_path
179 ),
180 ));
181 }
182
183 let produced_classification = match &mapping.source {
184 MappingSource::InitialInput => None,
185 MappingSource::Node { node_id } => resolved.get(node_id).and_then(|source| {
186 classification_at_path(
187 &source.contract.risk.data_flow.produced_data_classifications,
188 &mapping.source_path,
189 )
190 }),
191 };
192 let Some(produced_classification) = produced_classification else {
193 continue;
197 };
198 let accepted_classification = classification_at_path(
199 &target.contract.risk.data_flow.accepted_data_classifications,
200 &mapping.target_path,
201 );
202 let Some(accepted_classification) = accepted_classification else {
203 errors.push(cross_error(
204 ProposalCrossValidationErrorCode::UndeclaredDataClassification,
205 &path,
206 &format!(
207 "target path '{}' on node '{}' has no declared accepted_data_classifications entry; \
208 schema compatibility alone does not authorize disclosure (spec 109 FR-011)",
209 mapping.target_path, mapping.target_node_id
210 ),
211 ));
212 continue;
213 };
214 if produced_classification > accepted_classification {
215 errors.push(cross_error(
216 ProposalCrossValidationErrorCode::DataClassificationOverAccepted,
217 &path,
218 &format!(
219 "source path '{}' produces data classified above what target path '{}' on node '{}' accepts",
220 mapping.source_path, mapping.target_path, mapping.target_node_id
221 ),
222 ));
223 continue;
224 }
225 let is_external_or_irreversible_effect = matches!(
226 target.contract.risk.effect_class,
227 EffectClass::ExternalEffect | EffectClass::IrreversibleEffect
228 );
229 let egress_is_denied = target.contract.risk.data_flow.egress_policy == EgressPolicy::Denied;
230 if produced_classification > DataClassification::Public
231 && is_external_or_irreversible_effect
232 && egress_is_denied
233 {
234 errors.push(cross_error(
235 ProposalCrossValidationErrorCode::EgressDeniedForClassifiedMapping,
236 &path,
237 &format!(
238 "node '{}' has an external/irreversible effect with a denied egress policy and \
239 cannot receive classified data from '{}'",
240 mapping.target_node_id, mapping.source_path
241 ),
242 ));
243 }
244 }
245
246 if !errors.is_empty() {
247 return Err(ProposalCrossValidationFailure { errors });
248 }
249
250 Ok(canonical
251 .execution_order
252 .iter()
253 .filter_map(|node_id| resolved.get(node_id).cloned())
254 .collect())
255}
256
257fn manifest_declares_capability(manifest: &ApplicationBundleManifest, node: &ProposalNode) -> bool {
258 manifest.components.iter().any(|component| {
259 component.manifest.capability_id == node.capability_id
260 && component.manifest.capability_version == node.capability_version
261 })
262}
263
264fn classification_at_path(
265 classifications: &[traverse_contracts::FieldDataClassification],
266 path: &str,
267) -> Option<DataClassification> {
268 classifications
269 .iter()
270 .find(|entry| entry.field_path == path)
271 .map(|entry| entry.classification)
272}
273
274fn mapping_schema_compatible(
282 source_schema: &Value,
283 source_path: &str,
284 target_schema: &Value,
285 target_path: &str,
286) -> bool {
287 let source_fragment = resolve_schema_pointer(source_schema, source_path);
288 let target_fragment = resolve_schema_pointer(target_schema, target_path);
289 match (
290 source_fragment
291 .and_then(|f| f.get("type"))
292 .and_then(Value::as_str),
293 target_fragment
294 .and_then(|f| f.get("type"))
295 .and_then(Value::as_str),
296 ) {
297 (Some(source_type), Some(target_type)) => source_type == target_type,
298 _ => true,
299 }
300}
301
302fn resolve_schema_pointer<'a>(schema: &'a Value, pointer: &str) -> Option<&'a Value> {
303 let mut current = schema;
304 for segment in pointer.split('/').filter(|s| !s.is_empty()) {
305 current = current
306 .get("properties")
307 .and_then(|properties| properties.get(segment))
308 .or_else(|| current.get("items"))?;
309 }
310 Some(current)
311}
312
313fn cross_error(
314 code: ProposalCrossValidationErrorCode,
315 path: &str,
316 message: &str,
317) -> ProposalCrossValidationError {
318 ProposalCrossValidationError {
319 code,
320 message: message.to_string(),
321 path: path.to_string(),
322 }
323}
324
325#[derive(Debug, Clone, PartialEq, Eq)]
330pub enum AuthorizationDecision {
331 Automatic,
332 Approved(Box<ApprovalTokenClaims>),
333}
334
335#[must_use]
339pub fn proposal_is_automatic_eligible(resolved_nodes: &[ResolvedProposalNode]) -> bool {
340 resolved_nodes
341 .iter()
342 .all(|node| is_automatic_eligible(&node.contract.risk))
343}
344
345#[derive(Debug, Clone, PartialEq, Eq)]
348pub struct ApprovalTokenClaims {
349 pub token_id: String,
350 pub issuer: String,
351 pub key_id: String,
352 pub audience: String,
353 pub principal: String,
354 pub workspace_id: String,
355 pub proposal_digest: String,
356 pub snapshot_digest: String,
357 pub permitted_effects: Vec<EffectClass>,
358 pub permitted_connectors: Vec<String>,
359 pub max_use_count: u32,
360 pub expiry_unix: i64,
361}
362
363#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
364#[serde(rename_all = "snake_case")]
365pub enum ApprovalTokenErrorCode {
366 Malformed,
367 AlgorithmNotAllowed,
368 UnknownKeyId,
369 SignatureVerificationFailed,
370 IssuerMismatch,
371 AudienceMismatch,
372 Expired,
373 WorkspaceMismatch,
374 ProposalDigestMismatch,
375 SnapshotDigestMismatch,
376 UseCountExhausted,
377 Revoked,
378 StoreUnavailable,
379}
380
381#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
382pub struct ApprovalTokenError {
383 pub code: ApprovalTokenErrorCode,
384 pub message: String,
385}
386
387pub struct ApprovalTokenVerificationContext<'a> {
391 pub expected_issuer: &'a str,
392 pub expected_audience: &'a str,
393 pub expected_workspace_id: &'a str,
394 pub expected_proposal_digest: &'a str,
395 pub expected_snapshot_digest: &'a str,
396 pub verifying_keys_by_key_id: &'a HashMap<String, VerifyingKey>,
397}
398
399const APPROVAL_TOKEN_ALLOWED_ALG: &str = "EdDSA";
400
401#[allow(clippy::too_many_lines)]
411pub fn verify_approval_token(
412 token: &str,
413 context: &ApprovalTokenVerificationContext<'_>,
414) -> Result<ApprovalTokenClaims, ApprovalTokenError> {
415 let mut parts = token.split('.');
416 let (Some(header_b64), Some(payload_b64), Some(signature_b64), None) =
417 (parts.next(), parts.next(), parts.next(), parts.next())
418 else {
419 return Err(token_error(
420 ApprovalTokenErrorCode::Malformed,
421 "approval token must have exactly three dot-separated segments",
422 ));
423 };
424
425 let header_bytes = base64url_decode(header_b64)
426 .map_err(|msg| token_error(ApprovalTokenErrorCode::Malformed, &msg))?;
427 let header: Value = serde_json::from_slice(&header_bytes).map_err(|e| {
428 token_error(
429 ApprovalTokenErrorCode::Malformed,
430 &format!("invalid header: {e}"),
431 )
432 })?;
433 let alg = header
434 .get("alg")
435 .and_then(Value::as_str)
436 .unwrap_or_default();
437 if alg != APPROVAL_TOKEN_ALLOWED_ALG {
438 return Err(token_error(
439 ApprovalTokenErrorCode::AlgorithmNotAllowed,
440 &format!("alg '{alg}' is not allowed; only {APPROVAL_TOKEN_ALLOWED_ALG} is accepted"),
441 ));
442 }
443 let key_id = header
444 .get("kid")
445 .and_then(Value::as_str)
446 .ok_or_else(|| token_error(ApprovalTokenErrorCode::Malformed, "header missing 'kid'"))?
447 .to_string();
448 let verifying_key = context
449 .verifying_keys_by_key_id
450 .get(&key_id)
451 .ok_or_else(|| {
452 token_error(
453 ApprovalTokenErrorCode::UnknownKeyId,
454 &format!("no verification key configured for key id '{key_id}'"),
455 )
456 })?;
457
458 let signature_bytes = base64url_decode(signature_b64)
459 .map_err(|msg| token_error(ApprovalTokenErrorCode::SignatureVerificationFailed, &msg))?;
460 let signature_array = <[u8; 64]>::try_from(signature_bytes.as_slice()).map_err(|_| {
461 token_error(
462 ApprovalTokenErrorCode::SignatureVerificationFailed,
463 "signature must be 64 bytes",
464 )
465 })?;
466 let signature = Signature::from_bytes(&signature_array);
467 let signing_input = format!("{header_b64}.{payload_b64}");
468 verifying_key
469 .verify(signing_input.as_bytes(), &signature)
470 .map_err(|_| {
471 token_error(
472 ApprovalTokenErrorCode::SignatureVerificationFailed,
473 "signature verification failed",
474 )
475 })?;
476
477 let payload_bytes = base64url_decode(payload_b64)
478 .map_err(|msg| token_error(ApprovalTokenErrorCode::Malformed, &msg))?;
479 let payload: Value = serde_json::from_slice(&payload_bytes).map_err(|e| {
480 token_error(
481 ApprovalTokenErrorCode::Malformed,
482 &format!("invalid payload: {e}"),
483 )
484 })?;
485
486 let claims = parse_approval_token_claims(&payload, &key_id)?;
487
488 if claims.issuer != context.expected_issuer {
489 return Err(token_error(
490 ApprovalTokenErrorCode::IssuerMismatch,
491 "token issuer does not match the expected issuer",
492 ));
493 }
494 if claims.audience != context.expected_audience {
495 return Err(token_error(
496 ApprovalTokenErrorCode::AudienceMismatch,
497 "token audience does not match the expected audience",
498 ));
499 }
500 if claims.workspace_id != context.expected_workspace_id {
501 return Err(token_error(
502 ApprovalTokenErrorCode::WorkspaceMismatch,
503 "token workspace_id does not match the proposal's workspace_id",
504 ));
505 }
506 if claims.proposal_digest != context.expected_proposal_digest {
507 return Err(token_error(
508 ApprovalTokenErrorCode::ProposalDigestMismatch,
509 "token is not bound to this exact proposal digest",
510 ));
511 }
512 if claims.snapshot_digest != context.expected_snapshot_digest {
513 return Err(token_error(
514 ApprovalTokenErrorCode::SnapshotDigestMismatch,
515 "token is not bound to the current pinned snapshot digest",
516 ));
517 }
518 let now = unix_now();
519 if claims.expiry_unix <= now {
520 return Err(token_error(
521 ApprovalTokenErrorCode::Expired,
522 "token is expired",
523 ));
524 }
525
526 Ok(claims)
527}
528
529fn parse_approval_token_claims(
530 payload: &Value,
531 key_id: &str,
532) -> Result<ApprovalTokenClaims, ApprovalTokenError> {
533 let get_str = |field: &str| -> Result<String, ApprovalTokenError> {
534 payload
535 .get(field)
536 .and_then(Value::as_str)
537 .filter(|s| !s.trim().is_empty())
538 .map(ToString::to_string)
539 .ok_or_else(|| {
540 token_error(
541 ApprovalTokenErrorCode::Malformed,
542 &format!("payload missing required non-empty claim '{field}'"),
543 )
544 })
545 };
546 let permitted_effects = payload
547 .get("permitted_effects")
548 .and_then(Value::as_array)
549 .map(|values| {
550 values
551 .iter()
552 .filter_map(Value::as_str)
553 .filter_map(parse_effect_class)
554 .collect()
555 })
556 .unwrap_or_default();
557 let permitted_connectors = payload
558 .get("permitted_connectors")
559 .and_then(Value::as_array)
560 .map(|values| {
561 values
562 .iter()
563 .filter_map(Value::as_str)
564 .map(ToString::to_string)
565 .collect()
566 })
567 .unwrap_or_default();
568 let max_use_count = payload
569 .get("max_use_count")
570 .and_then(Value::as_u64)
571 .and_then(|v| u32::try_from(v).ok())
572 .ok_or_else(|| {
573 token_error(
574 ApprovalTokenErrorCode::Malformed,
575 "payload missing required u32 claim 'max_use_count'",
576 )
577 })?;
578 let expiry_unix = payload.get("exp").and_then(Value::as_i64).ok_or_else(|| {
579 token_error(
580 ApprovalTokenErrorCode::Malformed,
581 "payload missing required claim 'exp'",
582 )
583 })?;
584
585 Ok(ApprovalTokenClaims {
586 token_id: get_str("jti")?,
587 issuer: get_str("iss")?,
588 key_id: key_id.to_string(),
589 audience: get_str("aud")?,
590 principal: get_str("sub")?,
591 workspace_id: get_str("workspace_id")?,
592 proposal_digest: get_str("proposal_digest")?,
593 snapshot_digest: get_str("snapshot_digest")?,
594 permitted_effects,
595 permitted_connectors,
596 max_use_count,
597 expiry_unix,
598 })
599}
600
601fn parse_effect_class(value: &str) -> Option<EffectClass> {
602 match value {
603 "pure_read" => Some(EffectClass::PureRead),
604 "state_write" => Some(EffectClass::StateWrite),
605 "external_effect" => Some(EffectClass::ExternalEffect),
606 "irreversible_effect" => Some(EffectClass::IrreversibleEffect),
607 _ => None,
608 }
609}
610
611fn token_error(code: ApprovalTokenErrorCode, message: &str) -> ApprovalTokenError {
612 ApprovalTokenError {
613 code,
614 message: message.to_string(),
615 }
616}
617
618fn unix_now() -> i64 {
619 SystemTime::now()
620 .duration_since(UNIX_EPOCH)
621 .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
622}
623
624fn base64url_decode(input: &str) -> Result<Vec<u8>, String> {
625 const ALPHABET: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_";
626 let mut lookup = [255u8; 256];
627 for (index, byte) in ALPHABET.iter().enumerate() {
628 lookup[*byte as usize] = u8::try_from(index).unwrap_or(0);
629 }
630
631 let bytes = input.as_bytes();
632 let mut out = Vec::with_capacity(bytes.len() * 3 / 4 + 3);
633 let mut buffer: u32 = 0;
634 let mut bits: u32 = 0;
635 for &byte in bytes {
636 let value = lookup[byte as usize];
637 if value == 255 {
638 return Err("invalid base64url character".to_string());
639 }
640 buffer = (buffer << 6) | u32::from(value);
641 bits += 6;
642 if bits >= 8 {
643 bits -= 8;
644 out.push(u8::try_from((buffer >> bits) & 0xFF).unwrap_or(0));
645 }
646 }
647 Ok(out)
648}
649
650pub struct ApprovalTokenStore {
655 used: Mutex<HashMap<String, TokenUsageRecord>>,
656}
657
658#[derive(Debug, Clone, Copy, Default)]
659struct TokenUsageRecord {
660 use_count: u32,
661 revoked: bool,
662}
663
664impl Default for ApprovalTokenStore {
665 fn default() -> Self {
666 Self::new()
667 }
668}
669
670impl ApprovalTokenStore {
671 #[must_use]
672 pub fn new() -> Self {
673 Self {
674 used: Mutex::new(HashMap::new()),
675 }
676 }
677
678 pub fn check_and_record_use(
687 &self,
688 claims: &ApprovalTokenClaims,
689 ) -> Result<(), ApprovalTokenError> {
690 let mut used = self.used.lock().map_err(|_| {
691 token_error(
692 ApprovalTokenErrorCode::StoreUnavailable,
693 "approval token store is unavailable; failing closed",
694 )
695 })?;
696 let record = used.entry(claims.token_id.clone()).or_default();
697 if record.revoked {
698 return Err(token_error(
699 ApprovalTokenErrorCode::Revoked,
700 "token has been revoked",
701 ));
702 }
703 if record.use_count >= claims.max_use_count {
704 return Err(token_error(
705 ApprovalTokenErrorCode::UseCountExhausted,
706 "token has reached its maximum use count",
707 ));
708 }
709 record.use_count += 1;
710 Ok(())
711 }
712
713 pub fn revoke(&self, token_id: &str) {
718 if let Ok(mut used) = self.used.lock() {
719 used.entry(token_id.to_string()).or_default().revoked = true;
720 }
721 }
722}
723
724#[derive(Debug, Clone, Copy, PartialEq, Eq)]
729pub struct QuotaLimits {
730 pub max_concurrent_per_principal: u32,
731 pub max_concurrent_per_app: u32,
732 pub max_concurrent_per_workspace: u32,
733}
734
735pub const DEFAULT_MAX_CONCURRENT_PER_PRINCIPAL: u32 = 4;
736pub const DEFAULT_MAX_CONCURRENT_PER_APP: u32 = 16;
737pub const DEFAULT_MAX_CONCURRENT_PER_WORKSPACE: u32 = 32;
738
739impl Default for QuotaLimits {
740 fn default() -> Self {
741 Self {
742 max_concurrent_per_principal: DEFAULT_MAX_CONCURRENT_PER_PRINCIPAL,
743 max_concurrent_per_app: DEFAULT_MAX_CONCURRENT_PER_APP,
744 max_concurrent_per_workspace: DEFAULT_MAX_CONCURRENT_PER_WORKSPACE,
745 }
746 }
747}
748
749#[derive(Debug, Clone, PartialEq, Eq)]
750pub struct QuotaDenial {
751 pub scope: &'static str,
752 pub message: String,
753}
754
755#[derive(Debug)]
758pub struct QuotaReservation<'a> {
759 tracker: &'a QuotaTracker,
760 principal: String,
761 app_id: String,
762 workspace_id: String,
763}
764
765impl Drop for QuotaReservation<'_> {
766 fn drop(&mut self) {
767 self.tracker
768 .release(&self.principal, &self.app_id, &self.workspace_id);
769 }
770}
771
772#[derive(Debug, Default)]
775#[allow(clippy::struct_field_names)]
776pub struct QuotaTracker {
777 principal_slots: Mutex<HashMap<String, u32>>,
778 app_slots: Mutex<HashMap<String, u32>>,
779 workspace_slots: Mutex<HashMap<String, u32>>,
780}
781
782impl QuotaTracker {
783 #[must_use]
784 pub fn new() -> Self {
785 Self::default()
786 }
787
788 pub fn reserve(
798 &self,
799 principal: &str,
800 app_id: &str,
801 workspace_id: &str,
802 limits: &QuotaLimits,
803 ) -> Result<QuotaReservation<'_>, QuotaDenial> {
804 reserve_dimension(
805 &self.principal_slots,
806 principal,
807 limits.max_concurrent_per_principal,
808 "principal",
809 )?;
810 if let Err(denial) = reserve_dimension(
811 &self.app_slots,
812 app_id,
813 limits.max_concurrent_per_app,
814 "app",
815 ) {
816 release_dimension(&self.principal_slots, principal);
817 return Err(denial);
818 }
819 if let Err(denial) = reserve_dimension(
820 &self.workspace_slots,
821 workspace_id,
822 limits.max_concurrent_per_workspace,
823 "workspace",
824 ) {
825 release_dimension(&self.principal_slots, principal);
826 release_dimension(&self.app_slots, app_id);
827 return Err(denial);
828 }
829
830 Ok(QuotaReservation {
831 tracker: self,
832 principal: principal.to_string(),
833 app_id: app_id.to_string(),
834 workspace_id: workspace_id.to_string(),
835 })
836 }
837
838 fn release(&self, principal: &str, app_id: &str, workspace_id: &str) {
839 release_dimension(&self.principal_slots, principal);
840 release_dimension(&self.app_slots, app_id);
841 release_dimension(&self.workspace_slots, workspace_id);
842 }
843}
844
845fn reserve_dimension(
846 counts: &Mutex<HashMap<String, u32>>,
847 key: &str,
848 limit: u32,
849 scope: &'static str,
850) -> Result<(), QuotaDenial> {
851 let mut counts = counts.lock().map_err(|_| QuotaDenial {
852 scope: "store",
853 message: "quota tracker is unavailable; failing closed".to_string(),
854 })?;
855 let count = counts.entry(key.to_string()).or_insert(0);
856 if *count >= limit {
857 return Err(QuotaDenial {
858 scope,
859 message: format!("{scope} '{key}' is at its concurrency limit of {limit}"),
860 });
861 }
862 *count += 1;
863 Ok(())
864}
865
866fn release_dimension(counts: &Mutex<HashMap<String, u32>>, key: &str) {
867 if let Ok(mut counts) = counts.lock()
868 && let Some(count) = counts.get_mut(key)
869 {
870 *count = count.saturating_sub(1);
871 }
872}
873
874#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
879#[serde(rename_all = "snake_case")]
880pub enum ProposalNodeStatus {
881 Succeeded,
882 Failed,
883 SkippedAfterEarlierFailure,
884}
885
886#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
887pub struct ProposalNodeOutcome {
888 pub node_id: String,
889 pub capability_id: String,
890 pub capability_version: String,
891 pub artifact_digest: String,
892 pub status: ProposalNodeStatus,
893 pub error_code: Option<String>,
894}
895
896#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
897#[serde(rename_all = "snake_case")]
898pub enum ProposalTerminalState {
899 Succeeded,
900 Failed,
901 Cancelled,
902 Expired,
903 AuthorizationRevoked,
904}
905
906#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
907pub struct AuthorizationSummary {
908 pub automatic: bool,
909 pub approval_token_id: Option<String>,
910}
911
912#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
916pub struct ProposalTrace {
917 pub proposal_id: String,
918 pub proposal_digest: String,
919 pub snapshot_digest: String,
920 pub authorization: AuthorizationSummary,
921 pub node_outcomes: Vec<ProposalNodeOutcome>,
922 pub mapping_paths: Vec<(String, String)>,
923 pub terminal_state: ProposalTerminalState,
924}
925
926#[must_use]
931#[allow(clippy::too_many_lines)]
932pub fn execute_proposal<E: crate::LocalExecutor>(
933 runtime: &Runtime<E>,
934 canonical: &CanonicalProposal,
935 resolved_nodes: &[ResolvedProposalNode],
936 authorization: AuthorizationSummary,
937 proposal_digest: &str,
938 snapshot_digest: &str,
939) -> ProposalTrace {
940 let contracts_by_node: HashMap<&str, &CapabilityContract> = resolved_nodes
941 .iter()
942 .map(|n| (n.node_id.as_str(), &n.contract))
943 .collect();
944 let nodes_by_id: HashMap<&str, &ProposalNode> = canonical
945 .proposal
946 .nodes
947 .iter()
948 .map(|n| (n.node_id.as_str(), n))
949 .collect();
950
951 let mut outputs: HashMap<String, Value> = HashMap::new();
952 let mut outcomes = Vec::with_capacity(canonical.execution_order.len());
953 let mut failed = false;
954
955 for node_id in &canonical.execution_order {
956 let Some(node) = nodes_by_id.get(node_id.as_str()) else {
957 continue;
958 };
959 if failed {
960 outcomes.push(ProposalNodeOutcome {
961 node_id: node_id.clone(),
962 capability_id: node.capability_id.clone(),
963 capability_version: node.capability_version.clone(),
964 artifact_digest: node.artifact_digest.clone(),
965 status: ProposalNodeStatus::SkippedAfterEarlierFailure,
966 error_code: None,
967 });
968 continue;
969 }
970
971 let input = assemble_node_input(canonical, node_id, &outputs);
972 let request = build_node_execution_request(canonical, node, node_id, input);
973 let outcome = runtime.execute(request);
974 match outcome.result.status {
975 RuntimeResultStatus::Completed => {
976 if let Some(output) = outcome.result.output.clone() {
977 outputs.insert(node_id.clone(), output);
978 }
979 outcomes.push(ProposalNodeOutcome {
980 node_id: node_id.clone(),
981 capability_id: node.capability_id.clone(),
982 capability_version: node.capability_version.clone(),
983 artifact_digest: node.artifact_digest.clone(),
984 status: ProposalNodeStatus::Succeeded,
985 error_code: None,
986 });
987 }
988 RuntimeResultStatus::Error => {
989 failed = true;
990 outcomes.push(ProposalNodeOutcome {
991 node_id: node_id.clone(),
992 capability_id: node.capability_id.clone(),
993 capability_version: node.capability_version.clone(),
994 artifact_digest: node.artifact_digest.clone(),
995 status: ProposalNodeStatus::Failed,
996 error_code: outcome
997 .result
998 .error
999 .as_ref()
1000 .map(|error| format!("{:?}", error.code)),
1001 });
1002 }
1003 }
1004 let _ = contracts_by_node.get(node_id.as_str());
1005 }
1006
1007 let terminal_state = if failed {
1008 ProposalTerminalState::Failed
1009 } else {
1010 ProposalTerminalState::Succeeded
1011 };
1012
1013 ProposalTrace {
1014 proposal_id: canonical.proposal.proposal_id.clone(),
1015 proposal_digest: proposal_digest.to_string(),
1016 snapshot_digest: snapshot_digest.to_string(),
1017 authorization,
1018 node_outcomes: outcomes,
1019 mapping_paths: canonical
1020 .proposal
1021 .mappings
1022 .iter()
1023 .map(|m| (m.source_path.clone(), m.target_path.clone()))
1024 .collect(),
1025 terminal_state,
1026 }
1027}
1028
1029pub(crate) fn assemble_node_input(
1033 canonical: &CanonicalProposal,
1034 node_id: &str,
1035 outputs: &HashMap<String, Value>,
1036) -> Value {
1037 let mut input = Value::Object(serde_json::Map::new());
1038 for mapping in &canonical.proposal.mappings {
1039 if mapping.target_node_id != *node_id {
1040 continue;
1041 }
1042 let value = match &mapping.source {
1043 MappingSource::InitialInput => {
1044 pointer_get(&canonical.proposal.initial_input, &mapping.source_path)
1045 }
1046 MappingSource::Node { node_id: source_id } => outputs
1047 .get(source_id)
1048 .and_then(|output| pointer_get(output, &mapping.source_path)),
1049 };
1050 if let Some(value) = value {
1051 pointer_set(&mut input, &mapping.target_path, value.clone());
1052 }
1053 }
1054 input
1055}
1056
1057pub(crate) fn build_node_execution_request(
1060 canonical: &CanonicalProposal,
1061 node: &ProposalNode,
1062 node_id: &str,
1063 input: Value,
1064) -> RuntimeRequest {
1065 RuntimeRequest {
1066 kind: "runtime_request".to_string(),
1067 schema_version: "1.0.0".to_string(),
1068 request_id: format!("{}-{node_id}", canonical.proposal.proposal_id),
1069 intent: RuntimeIntent {
1070 capability_id: Some(node.capability_id.clone()),
1071 capability_version: Some(node.capability_version.clone()),
1072 version_range: None,
1073 intent_key: None,
1074 },
1075 input,
1076 lookup: RuntimeLookup {
1077 scope: RuntimeLookupScope::PreferPrivate,
1078 allow_ambiguity: false,
1079 },
1080 context: RuntimeContext {
1081 requested_target: PlacementTarget::Local,
1082 correlation_id: Some(canonical.proposal.proposal_id.clone()),
1083 caller: Some("workflow_proposal".to_string()),
1084 traceparent: None,
1085 tracestate: None,
1086 metadata: None,
1087 identity: None,
1088 },
1089 governing_spec: "006-runtime-request-execution".to_string(),
1090 }
1091}
1092
1093pub(crate) fn pointer_get<'a>(value: &'a Value, pointer: &str) -> Option<&'a Value> {
1094 value.pointer(pointer)
1095}
1096
1097pub(crate) fn pointer_set(target: &mut Value, pointer: &str, new_value: Value) {
1098 let segments: Vec<&str> = pointer.split('/').filter(|s| !s.is_empty()).collect();
1099 *target = set_at_segments(std::mem::take(target), &segments, new_value);
1100}
1101
1102fn set_at_segments(current: Value, segments: &[&str], new_value: Value) -> Value {
1108 let Some((head, rest)) = segments.split_first() else {
1109 return new_value;
1110 };
1111 let mut map = match current {
1112 Value::Object(map) => map,
1113 _ => serde_json::Map::new(),
1114 };
1115 let child = map.remove(*head).unwrap_or(Value::Null);
1116 map.insert((*head).to_string(), set_at_segments(child, rest, new_value));
1117 Value::Object(map)
1118}
1119
1120#[cfg(test)]
1121#[allow(clippy::expect_used)]
1122#[allow(clippy::panic)]
1123mod tests {
1124 use super::*;
1125 use crate::security::RuntimeSecurityConfig;
1126 use crate::{
1127 LocalExecutionFailure, LocalExecutionFailureCode, LocalExecutionOutput, LocalExecutor,
1128 Runtime,
1129 };
1130 use ed25519_dalek::{Signer, SigningKey};
1131 use serde_json::json;
1132 use traverse_contracts::{
1133 BinaryFormat as ContractBinaryFormat, CanonicalProposal, CapabilityContract,
1134 DataClassification, DataFlowPolicy, DeterminismClass, EffectClass, EgressPolicy,
1135 Entrypoint, EntrypointKind, Execution, ExecutionConstraints, ExecutionTarget,
1136 FieldDataClassification, FilesystemAccess, HostApiAccess, Lifecycle, ManifestReference,
1137 MappingSource, NetworkAccess, Owner, ProposalEdge, ProposalLimits, ProposalMapping,
1138 ProposalNode, Provenance, ProvenanceSource, ReliabilityMetadata, RiskMetadata,
1139 SchemaContainer, ServiceType, SideEffect, SideEffectKind, WorkflowProposal,
1140 canonicalize_proposal, proposal_digest,
1141 };
1142 use traverse_registry::{
1143 ApplicationBundleManifest, ApplicationComponent, ApplicationComponentRef,
1144 ApplicationEffectiveConfig, ArtifactDigests, BinaryFormat as RegistryBinaryFormat,
1145 BinaryReference, CapabilityArtifactRecord, CapabilityRegistration, CapabilityRegistry,
1146 ComponentExecutionMode, ComposabilityMetadata, CompositionKind, CompositionPattern,
1147 ImplementationKind, RegistryProvenance, RegistryScope, SourceKind, SourceReference,
1148 WasmComponentManifest,
1149 };
1150
1151 fn automatic_risk() -> RiskMetadata {
1152 RiskMetadata {
1153 effect_class: EffectClass::PureRead,
1154 determinism_class: DeterminismClass::Deterministic,
1155 data_flow: DataFlowPolicy::default(),
1156 reliability: ReliabilityMetadata {
1157 idempotency_required: false,
1158 retryable: true,
1159 compensation_available: false,
1160 },
1161 }
1162 }
1163
1164 fn contract(
1165 id: &str,
1166 version: &str,
1167 outputs_schema: Value,
1168 inputs_schema: Value,
1169 risk: RiskMetadata,
1170 ) -> CapabilityContract {
1171 let (namespace, name) = id.rsplit_once('.').unwrap_or(("test", id));
1172 CapabilityContract {
1173 kind: "capability_contract".to_string(),
1174 schema_version: "1.0.0".to_string(),
1175 id: id.to_string(),
1176 namespace: namespace.to_string(),
1177 name: name.to_string(),
1178 version: version.to_string(),
1179 lifecycle: Lifecycle::Active,
1180 owner: Owner {
1181 team: "traverse-core".to_string(),
1182 contact: "enrico.piovesan10@gmail.com".to_string(),
1183 },
1184 summary: "Test capability for proposal lifecycle validation.".to_string(),
1185 description: "Portable test capability used to validate proposal cross-checks."
1186 .to_string(),
1187 inputs: SchemaContainer {
1188 schema: inputs_schema,
1189 },
1190 outputs: SchemaContainer {
1191 schema: outputs_schema,
1192 },
1193 preconditions: Vec::new(),
1194 postconditions: Vec::new(),
1195 side_effects: vec![SideEffect {
1196 kind: SideEffectKind::MemoryOnly,
1197 description: "No durable side effect.".to_string(),
1198 }],
1199 emits: Vec::new(),
1200 consumes: Vec::new(),
1201 permissions: Vec::new(),
1202 execution: Execution {
1203 binary_format: ContractBinaryFormat::Wasm,
1204 entrypoint: Entrypoint {
1205 kind: EntrypointKind::WasiCommand,
1206 command: "run".to_string(),
1207 },
1208 preferred_targets: vec![ExecutionTarget::Local],
1209 constraints: ExecutionConstraints {
1210 host_api_access: HostApiAccess::None,
1211 network_access: NetworkAccess::Forbidden,
1212 filesystem_access: FilesystemAccess::None,
1213 },
1214 },
1215 policies: Vec::new(),
1216 dependencies: Vec::new(),
1217 provenance: Provenance {
1218 source: ProvenanceSource::Greenfield,
1219 author: "test".to_string(),
1220 created_at: "2026-08-23T00:00:00Z".to_string(),
1221 spec_ref: None,
1222 adr_refs: Vec::new(),
1223 exception_refs: Vec::new(),
1224 },
1225 evidence: Vec::new(),
1226 service_type: ServiceType::Stateless,
1227 permitted_targets: vec![ExecutionTarget::Local],
1228 event_trigger: None,
1229 connector_requirements: Vec::new(),
1230 state_schema: None,
1231 use_cases: Vec::new(),
1232 risk,
1233 }
1234 }
1235
1236 fn artifact(digest: &str) -> CapabilityArtifactRecord {
1237 CapabilityArtifactRecord {
1238 artifact_ref: format!("artifact:{digest}"),
1239 implementation_kind: ImplementationKind::Executable,
1240 source: SourceReference {
1241 kind: SourceKind::Git,
1242 location: "https://example.invalid/repo".to_string(),
1243 },
1244 binary: Some(BinaryReference {
1245 format: RegistryBinaryFormat::Wasm,
1246 location: format!("artifacts/{digest}/capability.wasm"),
1247 signature: None,
1248 }),
1249 workflow_ref: None,
1250 digests: ArtifactDigests {
1251 source_digest: format!("src-{digest}"),
1252 binary_digest: Some(digest.to_string()),
1253 },
1254 provenance: RegistryProvenance {
1255 source: "test".to_string(),
1256 author: "test".to_string(),
1257 created_at: "2026-08-23T00:00:00Z".to_string(),
1258 },
1259 }
1260 }
1261
1262 fn registry_with(
1263 entries: Vec<(CapabilityContract, CapabilityArtifactRecord)>,
1264 ) -> CapabilityRegistry {
1265 let mut registry = CapabilityRegistry::new();
1266 for (contract, artifact) in entries {
1267 let outcome = registry.register(CapabilityRegistration {
1268 scope: RegistryScope::Public,
1269 contract,
1270 contract_path: "registry/test/contract.json".to_string(),
1271 artifact,
1272 registered_at: "2026-08-23T00:00:00Z".to_string(),
1273 tags: Vec::new(),
1274 composability: ComposabilityMetadata {
1275 kind: CompositionKind::Atomic,
1276 patterns: vec![CompositionPattern::Sequential],
1277 provides: Vec::new(),
1278 requires: Vec::new(),
1279 },
1280 governing_spec: "005-capability-registry".to_string(),
1281 validator_version: "0.1.0".to_string(),
1282 });
1283 assert!(outcome.is_ok(), "registration must succeed: {outcome:?}");
1284 }
1285 registry
1286 }
1287
1288 fn manifest_declaring(components: &[(&str, &str)]) -> ApplicationBundleManifest {
1289 ApplicationBundleManifest {
1290 app_id: "test-app".to_string(),
1291 version: "1.0.0".to_string(),
1292 schema_version: "1.0.0".to_string(),
1293 workspace_defaults: json!({}),
1294 components: components
1295 .iter()
1296 .map(|(capability_id, capability_version)| ApplicationComponent {
1297 reference: ApplicationComponentRef {
1298 component_id: capability_id.to_string(),
1299 version: (*capability_version).to_string(),
1300 digest: "sha256:component-digest".to_string(),
1301 manifest_path: "component.manifest.json".to_string(),
1302 },
1303 manifest_path: "component.manifest.json".into(),
1304 manifest: WasmComponentManifest {
1305 component_id: capability_id.to_string(),
1306 version: (*capability_version).to_string(),
1307 schema_version: "1.0.0".to_string(),
1308 execution_mode: ComponentExecutionMode::Wasm,
1309 capability_id: (*capability_id).to_string(),
1310 capability_version: (*capability_version).to_string(),
1311 contract_path: None,
1312 registry_ref: None,
1313 wasm_binary_path: None,
1314 wasm_digest: None,
1315 platforms: vec!["local".to_string()],
1316 wrapper_path: None,
1317 runtime_constraints: json!({}),
1318 permitted_targets: vec![ExecutionTarget::Local],
1319 dependencies: Vec::new(),
1320 connector_requirements: Vec::new(),
1321 validation_evidence: Vec::new(),
1322 executable_pin: None,
1323 },
1324 contract_path: "contract.json".into(),
1325 contract: contract(
1326 capability_id,
1327 capability_version,
1328 json!({"type": "object"}),
1329 json!({"type": "object"}),
1330 automatic_risk(),
1331 ),
1332 wasm_binary_path: None,
1333 verified_wasm_digest: None,
1334 })
1335 .collect(),
1336 workflows: Vec::new(),
1337 connector_bindings: Vec::new(),
1338 model_dependencies: Vec::new(),
1339 config_schema: json!({}),
1340 default_config: json!({}),
1341 effective_config: ApplicationEffectiveConfig {
1342 values: json!({}),
1343 redacted_secret_keys: Vec::new(),
1344 },
1345 placement_policy: json!({}),
1346 public_surfaces: Vec::new(),
1347 state_machine: None,
1348 }
1349 }
1350
1351 fn linear_proposal_source() -> WorkflowProposal {
1352 WorkflowProposal {
1353 kind: "workflow_proposal".to_string(),
1354 schema_version: "1.0.0".to_string(),
1355 proposal_id: "proposal-001".to_string(),
1356 workspace_id: "workspace-001".to_string(),
1357 app_manifest: ManifestReference {
1358 app_id: "test-app".to_string(),
1359 app_version: "1.0.0".to_string(),
1360 manifest_digest: "sha256:manifest-digest".to_string(),
1361 },
1362 nodes: vec![
1363 ProposalNode {
1364 node_id: "a".to_string(),
1365 capability_id: "test.produce".to_string(),
1366 capability_version: "1.0.0".to_string(),
1367 artifact_digest: "digest-a".to_string(),
1368 },
1369 ProposalNode {
1370 node_id: "b".to_string(),
1371 capability_id: "test.consume".to_string(),
1372 capability_version: "1.0.0".to_string(),
1373 artifact_digest: "digest-b".to_string(),
1374 },
1375 ],
1376 edges: vec![ProposalEdge {
1377 from_node_id: "a".to_string(),
1378 to_node_id: "b".to_string(),
1379 }],
1380 mappings: vec![ProposalMapping {
1381 source: MappingSource::Node {
1382 node_id: "a".to_string(),
1383 },
1384 source_path: "/value".to_string(),
1385 target_node_id: "b".to_string(),
1386 target_path: "/value".to_string(),
1387 }],
1388 initial_input: json!({}),
1389 }
1390 }
1391
1392 fn linear_canonical() -> CanonicalProposal {
1393 canonicalize_proposal(linear_proposal_source(), &ProposalLimits::default())
1394 .expect("linear proposal must canonicalize")
1395 }
1396
1397 fn initial_input_proposal_source() -> WorkflowProposal {
1398 WorkflowProposal {
1399 kind: "workflow_proposal".to_string(),
1400 schema_version: "1.0.0".to_string(),
1401 proposal_id: "proposal-002".to_string(),
1402 workspace_id: "workspace-001".to_string(),
1403 app_manifest: ManifestReference {
1404 app_id: "test-app".to_string(),
1405 app_version: "1.0.0".to_string(),
1406 manifest_digest: "sha256:manifest-digest".to_string(),
1407 },
1408 nodes: vec![ProposalNode {
1409 node_id: "x".to_string(),
1410 capability_id: "test.consume".to_string(),
1411 capability_version: "1.0.0".to_string(),
1412 artifact_digest: "digest-x".to_string(),
1413 }],
1414 edges: Vec::new(),
1415 mappings: vec![ProposalMapping {
1416 source: MappingSource::InitialInput,
1417 source_path: "/value".to_string(),
1418 target_node_id: "x".to_string(),
1419 target_path: "/nested/value".to_string(),
1420 }],
1421 initial_input: json!({"value": "from-caller"}),
1422 }
1423 }
1424
1425 fn initial_input_canonical() -> CanonicalProposal {
1426 canonicalize_proposal(initial_input_proposal_source(), &ProposalLimits::default())
1427 .expect("initial-input proposal must canonicalize")
1428 }
1429
1430 #[test]
1433 fn validates_a_well_formed_proposal_against_manifest_and_registry() {
1434 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1435 let registry = registry_with(vec![
1436 (
1437 contract(
1438 "test.produce",
1439 "1.0.0",
1440 json!({"type": "object"}),
1441 json!({"type": "object"}),
1442 automatic_risk(),
1443 ),
1444 artifact("digest-a"),
1445 ),
1446 (
1447 contract(
1448 "test.consume",
1449 "1.0.0",
1450 json!({"type": "object"}),
1451 json!({"type": "object"}),
1452 automatic_risk(),
1453 ),
1454 artifact("digest-b"),
1455 ),
1456 ]);
1457
1458 let resolved =
1459 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1460 .expect("well-formed proposal must validate");
1461 assert_eq!(resolved.len(), 2);
1462 }
1463
1464 #[test]
1465 fn rejects_a_capability_not_declared_in_the_manifest() {
1466 let manifest = manifest_declaring(&[("test.produce", "1.0.0")]); let registry = registry_with(vec![
1468 (
1469 contract(
1470 "test.produce",
1471 "1.0.0",
1472 json!({"type": "object"}),
1473 json!({"type": "object"}),
1474 automatic_risk(),
1475 ),
1476 artifact("digest-a"),
1477 ),
1478 (
1479 contract(
1480 "test.consume",
1481 "1.0.0",
1482 json!({"type": "object"}),
1483 json!({"type": "object"}),
1484 automatic_risk(),
1485 ),
1486 artifact("digest-b"),
1487 ),
1488 ]);
1489
1490 let failure =
1491 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1492 .expect_err("undeclared capability must be rejected");
1493 assert!(
1494 failure
1495 .errors
1496 .iter()
1497 .any(|e| e.code == ProposalCrossValidationErrorCode::UndeclaredCapability)
1498 );
1499 }
1500
1501 #[test]
1502 fn rejects_a_capability_not_found_in_the_registry() {
1503 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1504 let registry = registry_with(vec![(
1505 contract(
1506 "test.produce",
1507 "1.0.0",
1508 json!({"type": "object"}),
1509 json!({"type": "object"}),
1510 automatic_risk(),
1511 ),
1512 artifact("digest-a"),
1513 )]); let failure =
1516 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1517 .expect_err("unregistered capability must be rejected");
1518 assert!(
1519 failure
1520 .errors
1521 .iter()
1522 .any(|e| e.code == ProposalCrossValidationErrorCode::CapabilityNotFound)
1523 );
1524 }
1525
1526 #[test]
1527 fn rejects_an_artifact_digest_mismatch() {
1528 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1529 let registry = registry_with(vec![
1530 (
1531 contract(
1532 "test.produce",
1533 "1.0.0",
1534 json!({"type": "object"}),
1535 json!({"type": "object"}),
1536 automatic_risk(),
1537 ),
1538 artifact("wrong-digest"),
1539 ),
1540 (
1541 contract(
1542 "test.consume",
1543 "1.0.0",
1544 json!({"type": "object"}),
1545 json!({"type": "object"}),
1546 automatic_risk(),
1547 ),
1548 artifact("digest-b"),
1549 ),
1550 ]);
1551
1552 let failure =
1553 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1554 .expect_err("digest mismatch must be rejected");
1555 assert!(
1556 failure
1557 .errors
1558 .iter()
1559 .any(|e| e.code == ProposalCrossValidationErrorCode::ArtifactDigestMismatch)
1560 );
1561 }
1562
1563 #[test]
1564 fn rejects_incompatible_mapping_schema_types() {
1565 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1566 let registry = registry_with(vec![
1567 (
1568 contract(
1569 "test.produce",
1570 "1.0.0",
1571 json!({"type": "object", "properties": {"value": {"type": "string"}}}),
1572 json!({"type": "object"}),
1573 automatic_risk(),
1574 ),
1575 artifact("digest-a"),
1576 ),
1577 (
1578 contract(
1579 "test.consume",
1580 "1.0.0",
1581 json!({"type": "object"}),
1582 json!({"type": "object", "properties": {"value": {"type": "integer"}}}),
1583 automatic_risk(),
1584 ),
1585 artifact("digest-b"),
1586 ),
1587 ]);
1588
1589 let failure =
1590 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1591 .expect_err("incompatible mapping schema types must be rejected");
1592 assert!(
1593 failure
1594 .errors
1595 .iter()
1596 .any(|e| e.code == ProposalCrossValidationErrorCode::IncompatibleMappingSchema)
1597 );
1598 }
1599
1600 fn classification_risk(
1601 produced: &[(&str, DataClassification)],
1602 accepted: &[(&str, DataClassification)],
1603 egress_policy: EgressPolicy,
1604 effect_class: EffectClass,
1605 ) -> RiskMetadata {
1606 RiskMetadata {
1607 effect_class,
1608 determinism_class: DeterminismClass::Deterministic,
1609 data_flow: DataFlowPolicy {
1610 accepted_data_classifications: accepted
1611 .iter()
1612 .map(|(path, classification)| FieldDataClassification {
1613 field_path: (*path).to_string(),
1614 classification: *classification,
1615 })
1616 .collect(),
1617 produced_data_classifications: produced
1618 .iter()
1619 .map(|(path, classification)| FieldDataClassification {
1620 field_path: (*path).to_string(),
1621 classification: *classification,
1622 })
1623 .collect(),
1624 egress_policy,
1625 },
1626 reliability: ReliabilityMetadata {
1627 idempotency_required: false,
1628 retryable: true,
1629 compensation_available: false,
1630 },
1631 }
1632 }
1633
1634 #[test]
1635 fn rejects_a_mapping_target_with_no_declared_accepted_classification() {
1636 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1637 let registry = registry_with(vec![
1638 (
1639 contract(
1640 "test.produce",
1641 "1.0.0",
1642 json!({"type": "object"}),
1643 json!({"type": "object"}),
1644 classification_risk(
1645 &[("/value", DataClassification::Public)],
1646 &[],
1647 EgressPolicy::Denied,
1648 EffectClass::PureRead,
1649 ),
1650 ),
1651 artifact("digest-a"),
1652 ),
1653 (
1654 contract(
1655 "test.consume",
1656 "1.0.0",
1657 json!({"type": "object"}),
1658 json!({"type": "object"}),
1659 automatic_risk(), ),
1661 artifact("digest-b"),
1662 ),
1663 ]);
1664
1665 let failure =
1666 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1667 .expect_err("undeclared target classification must be rejected");
1668 assert!(
1669 failure
1670 .errors
1671 .iter()
1672 .any(|e| e.code == ProposalCrossValidationErrorCode::UndeclaredDataClassification)
1673 );
1674 }
1675
1676 #[test]
1677 fn rejects_a_mapping_whose_produced_classification_exceeds_accepted() {
1678 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1679 let registry = registry_with(vec![
1680 (
1681 contract(
1682 "test.produce",
1683 "1.0.0",
1684 json!({"type": "object"}),
1685 json!({"type": "object"}),
1686 classification_risk(
1687 &[("/value", DataClassification::Confidential)],
1688 &[],
1689 EgressPolicy::Denied,
1690 EffectClass::PureRead,
1691 ),
1692 ),
1693 artifact("digest-a"),
1694 ),
1695 (
1696 contract(
1697 "test.consume",
1698 "1.0.0",
1699 json!({"type": "object"}),
1700 json!({"type": "object"}),
1701 classification_risk(
1702 &[],
1703 &[("/value", DataClassification::Public)],
1704 EgressPolicy::Denied,
1705 EffectClass::PureRead,
1706 ),
1707 ),
1708 artifact("digest-b"),
1709 ),
1710 ]);
1711
1712 let failure =
1713 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1714 .expect_err("over-classified mapping must be rejected");
1715 assert!(
1716 failure
1717 .errors
1718 .iter()
1719 .any(|e| e.code == ProposalCrossValidationErrorCode::DataClassificationOverAccepted)
1720 );
1721 }
1722
1723 #[test]
1724 fn rejects_classified_data_into_an_egress_denied_external_effect_node() {
1725 let manifest = manifest_declaring(&[("test.produce", "1.0.0"), ("test.consume", "1.0.0")]);
1726 let registry = registry_with(vec![
1727 (
1728 contract(
1729 "test.produce",
1730 "1.0.0",
1731 json!({"type": "object"}),
1732 json!({"type": "object"}),
1733 classification_risk(
1734 &[("/value", DataClassification::Internal)],
1735 &[],
1736 EgressPolicy::Denied,
1737 EffectClass::PureRead,
1738 ),
1739 ),
1740 artifact("digest-a"),
1741 ),
1742 (
1743 contract(
1744 "test.consume",
1745 "1.0.0",
1746 json!({"type": "object"}),
1747 json!({"type": "object"}),
1748 classification_risk(
1749 &[],
1750 &[("/value", DataClassification::Internal)],
1751 EgressPolicy::Denied,
1752 EffectClass::ExternalEffect,
1753 ),
1754 ),
1755 artifact("digest-b"),
1756 ),
1757 ]);
1758
1759 let failure =
1760 validate_proposal_against_host_state(&linear_canonical(), &manifest, ®istry)
1761 .expect_err(
1762 "classified data into an egress-denied external-effect node must be rejected",
1763 );
1764 assert!(
1765 failure
1766 .errors
1767 .iter()
1768 .any(|e| e.code
1769 == ProposalCrossValidationErrorCode::EgressDeniedForClassifiedMapping)
1770 );
1771 }
1772
1773 #[test]
1774 fn accepts_a_mapping_sourced_from_initial_input_with_no_classification_check() {
1775 let manifest = manifest_declaring(&[("test.consume", "1.0.0")]);
1776 let registry = registry_with(vec![(
1777 contract(
1778 "test.consume",
1779 "1.0.0",
1780 json!({"type": "object"}),
1781 json!({"type": "object"}),
1782 automatic_risk(),
1783 ),
1784 artifact("digest-x"),
1785 )]);
1786
1787 let resolved =
1788 validate_proposal_against_host_state(&initial_input_canonical(), &manifest, ®istry)
1789 .expect("an initial-input-sourced mapping needs no capability-to-capability classification check");
1790 assert_eq!(resolved.len(), 1);
1791 }
1792
1793 #[test]
1794 fn validate_proposal_against_host_state_skips_a_mapping_whose_target_node_is_not_resolved() {
1795 let manifest = manifest_declaring(&[("test.produce", "1.0.0")]);
1802 let registry = registry_with(vec![(
1803 contract(
1804 "test.produce",
1805 "1.0.0",
1806 json!({"type": "object"}),
1807 json!({"type": "object"}),
1808 automatic_risk(),
1809 ),
1810 artifact("digest-a"),
1811 )]);
1812
1813 let mut proposal = linear_proposal_source();
1814 proposal.nodes.truncate(1);
1815 proposal.edges.clear();
1816 proposal.mappings[0].target_node_id = "ghost".to_string();
1817 let canonical = CanonicalProposal {
1818 execution_order: vec!["a".to_string()],
1819 proposal,
1820 };
1821
1822 let resolved = validate_proposal_against_host_state(&canonical, &manifest, ®istry)
1823 .expect("a dangling mapping target must not itself fail cross-validation");
1824 assert_eq!(resolved.len(), 1);
1825 }
1826
1827 #[test]
1830 fn proposal_is_automatic_eligible_true_when_every_node_is_automatic_eligible() {
1831 let nodes = vec![
1832 ResolvedProposalNode {
1833 node_id: "a".to_string(),
1834 contract: contract(
1835 "test.produce",
1836 "1.0.0",
1837 json!({}),
1838 json!({}),
1839 automatic_risk(),
1840 ),
1841 },
1842 ResolvedProposalNode {
1843 node_id: "b".to_string(),
1844 contract: contract(
1845 "test.consume",
1846 "1.0.0",
1847 json!({}),
1848 json!({}),
1849 automatic_risk(),
1850 ),
1851 },
1852 ];
1853 assert!(proposal_is_automatic_eligible(&nodes));
1854 }
1855
1856 #[test]
1857 fn proposal_is_automatic_eligible_false_when_any_node_is_not() {
1858 let mut non_automatic = automatic_risk();
1859 non_automatic.effect_class = EffectClass::StateWrite;
1860 let nodes = vec![
1861 ResolvedProposalNode {
1862 node_id: "a".to_string(),
1863 contract: contract(
1864 "test.produce",
1865 "1.0.0",
1866 json!({}),
1867 json!({}),
1868 automatic_risk(),
1869 ),
1870 },
1871 ResolvedProposalNode {
1872 node_id: "b".to_string(),
1873 contract: contract("test.consume", "1.0.0", json!({}), json!({}), non_automatic),
1874 },
1875 ];
1876 assert!(!proposal_is_automatic_eligible(&nodes));
1877 }
1878
1879 fn signing_key() -> SigningKey {
1882 SigningKey::from_bytes(&[9_u8; 32])
1883 }
1884
1885 fn base64url_encode(input: &[u8]) -> String {
1886 const ALPHABET: &[u8; 64] =
1887 b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_";
1888 let mut out = String::new();
1889 let mut i = 0;
1890 while i + 3 <= input.len() {
1891 let n = (u32::from(input[i]) << 16)
1892 | (u32::from(input[i + 1]) << 8)
1893 | u32::from(input[i + 2]);
1894 out.push(ALPHABET[((n >> 18) & 63) as usize] as char);
1895 out.push(ALPHABET[((n >> 12) & 63) as usize] as char);
1896 out.push(ALPHABET[((n >> 6) & 63) as usize] as char);
1897 out.push(ALPHABET[(n & 63) as usize] as char);
1898 i += 3;
1899 }
1900 let remainder = input.len() - i;
1901 if remainder == 1 {
1902 let n = u32::from(input[i]) << 16;
1903 out.push(ALPHABET[((n >> 18) & 63) as usize] as char);
1904 out.push(ALPHABET[((n >> 12) & 63) as usize] as char);
1905 } else if remainder == 2 {
1906 let n = (u32::from(input[i]) << 16) | (u32::from(input[i + 1]) << 8);
1907 out.push(ALPHABET[((n >> 18) & 63) as usize] as char);
1908 out.push(ALPHABET[((n >> 12) & 63) as usize] as char);
1909 out.push(ALPHABET[((n >> 6) & 63) as usize] as char);
1910 }
1911 out
1912 }
1913
1914 fn sign_token(payload: &Value, key: &SigningKey, key_id: &str) -> String {
1915 let header = base64url_encode(format!(r#"{{"alg":"EdDSA","kid":"{key_id}"}}"#).as_bytes());
1916 let payload_b64 = base64url_encode(payload.to_string().as_bytes());
1917 let signing_input = format!("{header}.{payload_b64}");
1918 let signature = key.sign(signing_input.as_bytes());
1919 let signature_b64 = base64url_encode(&signature.to_bytes());
1920 format!("{header}.{payload_b64}.{signature_b64}")
1921 }
1922
1923 fn valid_claims_payload() -> Value {
1924 json!({
1925 "jti": "token-001",
1926 "iss": "traverse-approval-service",
1927 "aud": "traverse-runtime",
1928 "sub": "principal-001",
1929 "workspace_id": "workspace-001",
1930 "proposal_digest": "digest-p",
1931 "snapshot_digest": "digest-s",
1932 "permitted_effects": ["external_effect"],
1933 "permitted_connectors": ["traverse.http"],
1934 "max_use_count": 1,
1935 "exp": 4_102_444_800_i64, })
1937 }
1938
1939 fn verification_context(
1940 keys: &HashMap<String, ed25519_dalek::VerifyingKey>,
1941 ) -> ApprovalTokenVerificationContext<'_> {
1942 ApprovalTokenVerificationContext {
1943 expected_issuer: "traverse-approval-service",
1944 expected_audience: "traverse-runtime",
1945 expected_workspace_id: "workspace-001",
1946 expected_proposal_digest: "digest-p",
1947 expected_snapshot_digest: "digest-s",
1948 verifying_keys_by_key_id: keys,
1949 }
1950 }
1951
1952 fn keys_with_signing_key() -> HashMap<String, ed25519_dalek::VerifyingKey> {
1953 let mut keys = HashMap::new();
1954 keys.insert("key-1".to_string(), signing_key().verifying_key());
1955 keys
1956 }
1957
1958 #[test]
1959 fn verifies_a_well_formed_approval_token() {
1960 let token = sign_token(&valid_claims_payload(), &signing_key(), "key-1");
1961 let keys = keys_with_signing_key();
1962 let claims = verify_approval_token(&token, &verification_context(&keys))
1963 .expect("well-formed token must verify");
1964 assert_eq!(claims.token_id, "token-001");
1965 assert_eq!(claims.principal, "principal-001");
1966 assert_eq!(claims.permitted_effects, vec![EffectClass::ExternalEffect]);
1967 }
1968
1969 #[test]
1970 fn rejects_malformed_token_shape() {
1971 let keys = keys_with_signing_key();
1972 let failure = verify_approval_token("not-a-token", &verification_context(&keys))
1973 .expect_err("malformed token must be rejected");
1974 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
1975 }
1976
1977 #[test]
1978 fn rejects_disallowed_algorithm() {
1979 let header = base64url_encode(br#"{"alg":"HS256","kid":"key-1"}"#);
1980 let payload = base64url_encode(valid_claims_payload().to_string().as_bytes());
1981 let token = format!("{header}.{payload}.sig");
1982 let keys = keys_with_signing_key();
1983 let failure = verify_approval_token(&token, &verification_context(&keys))
1984 .expect_err("disallowed alg must be rejected");
1985 assert_eq!(failure.code, ApprovalTokenErrorCode::AlgorithmNotAllowed);
1986 }
1987
1988 #[test]
1989 fn rejects_unknown_key_id() {
1990 let token = sign_token(&valid_claims_payload(), &signing_key(), "unknown-key");
1991 let keys = keys_with_signing_key();
1992 let failure = verify_approval_token(&token, &verification_context(&keys))
1993 .expect_err("unknown key id must be rejected");
1994 assert_eq!(failure.code, ApprovalTokenErrorCode::UnknownKeyId);
1995 }
1996
1997 #[test]
1998 fn rejects_bad_signature() {
1999 let other_key = SigningKey::from_bytes(&[3_u8; 32]);
2000 let token = sign_token(&valid_claims_payload(), &other_key, "key-1");
2001 let keys = keys_with_signing_key();
2002 let failure = verify_approval_token(&token, &verification_context(&keys))
2003 .expect_err("signature from the wrong key must be rejected");
2004 assert_eq!(
2005 failure.code,
2006 ApprovalTokenErrorCode::SignatureVerificationFailed
2007 );
2008 }
2009
2010 #[test]
2011 fn rejects_issuer_mismatch() {
2012 let mut payload = valid_claims_payload();
2013 payload["iss"] = json!("someone-else");
2014 let token = sign_token(&payload, &signing_key(), "key-1");
2015 let keys = keys_with_signing_key();
2016 let failure = verify_approval_token(&token, &verification_context(&keys))
2017 .expect_err("issuer mismatch must be rejected");
2018 assert_eq!(failure.code, ApprovalTokenErrorCode::IssuerMismatch);
2019 }
2020
2021 #[test]
2022 fn rejects_audience_mismatch() {
2023 let mut payload = valid_claims_payload();
2024 payload["aud"] = json!("someone-else");
2025 let token = sign_token(&payload, &signing_key(), "key-1");
2026 let keys = keys_with_signing_key();
2027 let failure = verify_approval_token(&token, &verification_context(&keys))
2028 .expect_err("audience mismatch must be rejected");
2029 assert_eq!(failure.code, ApprovalTokenErrorCode::AudienceMismatch);
2030 }
2031
2032 #[test]
2033 fn rejects_workspace_mismatch() {
2034 let mut payload = valid_claims_payload();
2035 payload["workspace_id"] = json!("someone-elses-workspace");
2036 let token = sign_token(&payload, &signing_key(), "key-1");
2037 let keys = keys_with_signing_key();
2038 let failure = verify_approval_token(&token, &verification_context(&keys))
2039 .expect_err("workspace mismatch must be rejected");
2040 assert_eq!(failure.code, ApprovalTokenErrorCode::WorkspaceMismatch);
2041 }
2042
2043 #[test]
2044 fn rejects_proposal_digest_mismatch() {
2045 let mut payload = valid_claims_payload();
2046 payload["proposal_digest"] = json!("different-digest");
2047 let token = sign_token(&payload, &signing_key(), "key-1");
2048 let keys = keys_with_signing_key();
2049 let failure = verify_approval_token(&token, &verification_context(&keys))
2050 .expect_err("proposal digest mismatch must be rejected");
2051 assert_eq!(failure.code, ApprovalTokenErrorCode::ProposalDigestMismatch);
2052 }
2053
2054 #[test]
2055 fn rejects_snapshot_digest_mismatch() {
2056 let mut payload = valid_claims_payload();
2057 payload["snapshot_digest"] = json!("different-digest");
2058 let token = sign_token(&payload, &signing_key(), "key-1");
2059 let keys = keys_with_signing_key();
2060 let failure = verify_approval_token(&token, &verification_context(&keys))
2061 .expect_err("snapshot digest mismatch must be rejected");
2062 assert_eq!(failure.code, ApprovalTokenErrorCode::SnapshotDigestMismatch);
2063 }
2064
2065 #[test]
2066 fn rejects_expired_token() {
2067 let mut payload = valid_claims_payload();
2068 payload["exp"] = json!(1); let token = sign_token(&payload, &signing_key(), "key-1");
2070 let keys = keys_with_signing_key();
2071 let failure = verify_approval_token(&token, &verification_context(&keys))
2072 .expect_err("expired token must be rejected");
2073 assert_eq!(failure.code, ApprovalTokenErrorCode::Expired);
2074 }
2075
2076 #[test]
2077 fn rejects_a_token_with_an_invalid_base64url_character() {
2078 let payload = base64url_encode(valid_claims_payload().to_string().as_bytes());
2079 let token = format!("not!valid!base64.{payload}.sig");
2080 let keys = keys_with_signing_key();
2081 let failure = verify_approval_token(&token, &verification_context(&keys))
2082 .expect_err("an invalid base64url character must be rejected");
2083 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2084 }
2085
2086 #[test]
2087 fn rejects_a_token_whose_header_is_not_valid_json() {
2088 let header = base64url_encode(b"not-json");
2089 let payload = base64url_encode(valid_claims_payload().to_string().as_bytes());
2090 let token = format!("{header}.{payload}.sig");
2091 let keys = keys_with_signing_key();
2092 let failure = verify_approval_token(&token, &verification_context(&keys))
2093 .expect_err("a non-JSON header must be rejected");
2094 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2095 }
2096
2097 #[test]
2098 fn rejects_a_token_with_a_wrong_length_signature() {
2099 let header = base64url_encode(br#"{"alg":"EdDSA","kid":"key-1"}"#);
2100 let payload_b64 = base64url_encode(valid_claims_payload().to_string().as_bytes());
2101 let short_signature = base64url_encode(b"too-short");
2102 let token = format!("{header}.{payload_b64}.{short_signature}");
2103 let keys = keys_with_signing_key();
2104 let failure = verify_approval_token(&token, &verification_context(&keys))
2105 .expect_err("a signature that is not 64 bytes must be rejected");
2106 assert_eq!(
2107 failure.code,
2108 ApprovalTokenErrorCode::SignatureVerificationFailed
2109 );
2110 }
2111
2112 #[test]
2113 fn rejects_a_token_whose_payload_is_not_valid_json() {
2114 let key = signing_key();
2115 let header = base64url_encode(br#"{"alg":"EdDSA","kid":"key-1"}"#);
2116 let payload_b64 = base64url_encode(b"not-json");
2117 let signing_input = format!("{header}.{payload_b64}");
2118 let signature = key.sign(signing_input.as_bytes());
2119 let signature_b64 = base64url_encode(&signature.to_bytes());
2120 let token = format!("{header}.{payload_b64}.{signature_b64}");
2121 let keys = keys_with_signing_key();
2122 let failure = verify_approval_token(&token, &verification_context(&keys))
2123 .expect_err("a non-JSON payload must be rejected");
2124 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2125 }
2126
2127 #[test]
2128 fn rejects_a_token_missing_a_required_string_claim() {
2129 let mut payload = valid_claims_payload();
2130 payload
2131 .as_object_mut()
2132 .expect("payload fixture is an object")
2133 .remove("jti");
2134 let token = sign_token(&payload, &signing_key(), "key-1");
2135 let keys = keys_with_signing_key();
2136 let failure = verify_approval_token(&token, &verification_context(&keys))
2137 .expect_err("a missing required string claim must be rejected");
2138 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2139 }
2140
2141 #[test]
2142 fn rejects_a_token_missing_max_use_count() {
2143 let mut payload = valid_claims_payload();
2144 payload
2145 .as_object_mut()
2146 .expect("payload fixture is an object")
2147 .remove("max_use_count");
2148 let token = sign_token(&payload, &signing_key(), "key-1");
2149 let keys = keys_with_signing_key();
2150 let failure = verify_approval_token(&token, &verification_context(&keys))
2151 .expect_err("a missing max_use_count claim must be rejected");
2152 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2153 }
2154
2155 #[test]
2156 fn rejects_a_token_missing_exp() {
2157 let mut payload = valid_claims_payload();
2158 payload
2159 .as_object_mut()
2160 .expect("payload fixture is an object")
2161 .remove("exp");
2162 let token = sign_token(&payload, &signing_key(), "key-1");
2163 let keys = keys_with_signing_key();
2164 let failure = verify_approval_token(&token, &verification_context(&keys))
2165 .expect_err("a missing exp claim must be rejected");
2166 assert_eq!(failure.code, ApprovalTokenErrorCode::Malformed);
2167 }
2168
2169 #[test]
2170 fn accepts_a_token_permitting_an_irreversible_effect() {
2171 let mut payload = valid_claims_payload();
2172 payload["permitted_effects"] = json!(["irreversible_effect"]);
2173 let token = sign_token(&payload, &signing_key(), "key-1");
2174 let keys = keys_with_signing_key();
2175 let claims = verify_approval_token(&token, &verification_context(&keys))
2176 .expect("a token permitting an irreversible effect must verify");
2177 assert_eq!(
2178 claims.permitted_effects,
2179 vec![EffectClass::IrreversibleEffect]
2180 );
2181 }
2182
2183 #[test]
2184 fn ignores_an_unrecognized_permitted_effect_string() {
2185 let mut payload = valid_claims_payload();
2186 payload["permitted_effects"] = json!(["not_a_real_effect", "external_effect"]);
2187 let token = sign_token(&payload, &signing_key(), "key-1");
2188 let keys = keys_with_signing_key();
2189 let claims = verify_approval_token(&token, &verification_context(&keys))
2190 .expect("an unrecognized effect string must be filtered out, not rejected");
2191 assert_eq!(claims.permitted_effects, vec![EffectClass::ExternalEffect]);
2192 }
2193
2194 #[test]
2195 fn approval_token_store_default_starts_with_no_recorded_uses() {
2196 let store = ApprovalTokenStore::default();
2197 store
2198 .check_and_record_use(&claims_with_use_count(1))
2199 .expect("a freshly defaulted store has no recorded uses");
2200 }
2201
2202 #[test]
2203 fn check_and_record_use_fails_closed_when_the_store_mutex_is_poisoned() {
2204 let store = ApprovalTokenStore::new();
2205
2206 let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
2207 let _guard = store
2208 .used
2209 .lock()
2210 .expect("lock must be acquirable to poison it");
2211 panic!("poison approval token store lock for test");
2212 }));
2213
2214 let failure = store
2215 .check_and_record_use(&claims_with_use_count(10))
2216 .expect_err("a poisoned store must fail closed");
2217 assert_eq!(failure.code, ApprovalTokenErrorCode::StoreUnavailable);
2218 }
2219
2220 fn claims_with_use_count(max_use_count: u32) -> ApprovalTokenClaims {
2223 ApprovalTokenClaims {
2224 token_id: "token-001".to_string(),
2225 issuer: "traverse-approval-service".to_string(),
2226 key_id: "key-1".to_string(),
2227 audience: "traverse-runtime".to_string(),
2228 principal: "principal-001".to_string(),
2229 workspace_id: "workspace-001".to_string(),
2230 proposal_digest: "digest-p".to_string(),
2231 snapshot_digest: "digest-s".to_string(),
2232 permitted_effects: Vec::new(),
2233 permitted_connectors: Vec::new(),
2234 max_use_count,
2235 expiry_unix: 4_102_444_800,
2236 }
2237 }
2238
2239 #[test]
2240 fn token_store_enforces_max_use_count() {
2241 let store = ApprovalTokenStore::new();
2242 let claims = claims_with_use_count(2);
2243 store
2244 .check_and_record_use(&claims)
2245 .expect("first use must succeed");
2246 store
2247 .check_and_record_use(&claims)
2248 .expect("second use must succeed");
2249 let failure = store
2250 .check_and_record_use(&claims)
2251 .expect_err("third use beyond max_use_count must be rejected");
2252 assert_eq!(failure.code, ApprovalTokenErrorCode::UseCountExhausted);
2253 }
2254
2255 #[test]
2256 fn token_store_denies_a_revoked_token() {
2257 let store = ApprovalTokenStore::new();
2258 let claims = claims_with_use_count(10);
2259 store.revoke(&claims.token_id);
2260 let failure = store
2261 .check_and_record_use(&claims)
2262 .expect_err("revoked token must be rejected");
2263 assert_eq!(failure.code, ApprovalTokenErrorCode::Revoked);
2264 }
2265
2266 #[test]
2269 fn quota_tracker_denies_when_principal_limit_reached() {
2270 let tracker = QuotaTracker::new();
2271 let limits = QuotaLimits {
2272 max_concurrent_per_principal: 1,
2273 max_concurrent_per_app: 10,
2274 max_concurrent_per_workspace: 10,
2275 };
2276 let _first = tracker
2277 .reserve("principal-1", "app-1", "workspace-1", &limits)
2278 .expect("first reservation must succeed");
2279 let denial = tracker
2280 .reserve("principal-1", "app-2", "workspace-2", &limits)
2281 .expect_err("second reservation for the same principal must be denied");
2282 assert_eq!(denial.scope, "principal");
2283 }
2284
2285 #[test]
2286 fn quota_tracker_denies_when_workspace_limit_reached() {
2287 let tracker = QuotaTracker::new();
2288 let limits = QuotaLimits {
2289 max_concurrent_per_principal: 10,
2290 max_concurrent_per_app: 10,
2291 max_concurrent_per_workspace: 1,
2292 };
2293 let _first = tracker
2294 .reserve("principal-1", "app-1", "workspace-1", &limits)
2295 .expect("first reservation must succeed");
2296 let denial = tracker
2297 .reserve("principal-2", "app-2", "workspace-1", &limits)
2298 .expect_err("second reservation for the same workspace must be denied");
2299 assert_eq!(denial.scope, "workspace");
2300 }
2301
2302 #[test]
2303 fn quota_tracker_releases_on_drop_allowing_reuse() {
2304 let tracker = QuotaTracker::new();
2305 let limits = QuotaLimits {
2306 max_concurrent_per_principal: 1,
2307 max_concurrent_per_app: 1,
2308 max_concurrent_per_workspace: 1,
2309 };
2310 {
2311 let _reservation = tracker
2312 .reserve("principal-1", "app-1", "workspace-1", &limits)
2313 .expect("first reservation must succeed");
2314 }
2315 tracker
2316 .reserve("principal-1", "app-1", "workspace-1", &limits)
2317 .expect("reservation must succeed again after the first is dropped");
2318 }
2319
2320 #[test]
2321 fn quota_tracker_denies_when_app_limit_reached_and_rolls_back_the_principal_reservation() {
2322 let tracker = QuotaTracker::new();
2323 let limits = QuotaLimits {
2324 max_concurrent_per_principal: 1,
2325 max_concurrent_per_app: 1,
2326 max_concurrent_per_workspace: 10,
2327 };
2328 let _first = tracker
2329 .reserve("principal-1", "app-1", "workspace-1", &limits)
2330 .expect("first reservation must succeed");
2331
2332 let denial = tracker
2333 .reserve("principal-2", "app-1", "workspace-2", &limits)
2334 .expect_err("second reservation against the same app must be denied");
2335 assert_eq!(denial.scope, "app");
2336
2337 tracker
2341 .reserve("principal-2", "app-2", "workspace-3", &limits)
2342 .expect("principal-2's slot must have been released by the app-limit rollback");
2343 }
2344
2345 #[test]
2346 fn quota_limits_default_matches_the_documented_defaults() {
2347 let limits = QuotaLimits::default();
2348 assert_eq!(
2349 limits.max_concurrent_per_principal,
2350 DEFAULT_MAX_CONCURRENT_PER_PRINCIPAL
2351 );
2352 assert_eq!(
2353 limits.max_concurrent_per_app,
2354 DEFAULT_MAX_CONCURRENT_PER_APP
2355 );
2356 assert_eq!(
2357 limits.max_concurrent_per_workspace,
2358 DEFAULT_MAX_CONCURRENT_PER_WORKSPACE
2359 );
2360 }
2361
2362 #[test]
2363 fn reserve_fails_closed_when_a_quota_dimension_mutex_is_poisoned() {
2364 let tracker = QuotaTracker::new();
2365
2366 let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
2367 let _guard = tracker
2368 .principal_slots
2369 .lock()
2370 .expect("lock must be acquirable to poison it");
2371 panic!("poison quota tracker principal_slots lock for test");
2372 }));
2373
2374 let limits = QuotaLimits {
2375 max_concurrent_per_principal: 10,
2376 max_concurrent_per_app: 10,
2377 max_concurrent_per_workspace: 10,
2378 };
2379 let denial = tracker
2380 .reserve("principal-1", "app-1", "workspace-1", &limits)
2381 .expect_err("a poisoned quota dimension must fail closed");
2382 assert_eq!(denial.scope, "store");
2383 }
2384
2385 #[derive(Default)]
2388 struct MappingAwareExecutor {
2389 consumer_saw_input: std::sync::Arc<Mutex<Option<Value>>>,
2390 }
2391
2392 impl LocalExecutor for MappingAwareExecutor {
2393 fn execute(
2394 &self,
2395 capability: &traverse_registry::ResolvedCapability,
2396 input: &Value,
2397 ) -> Result<LocalExecutionOutput, LocalExecutionFailure> {
2398 if capability.contract.id == "test.produce" {
2399 Ok(LocalExecutionOutput {
2400 value: json!({"value": "produced-by-a"}),
2401 emitted_events: Vec::new(),
2402 })
2403 } else {
2404 if let Ok(mut seen) = self.consumer_saw_input.lock() {
2405 *seen = Some(input.clone());
2406 }
2407 Ok(LocalExecutionOutput {
2408 value: json!({"received": input.clone()}),
2409 emitted_events: Vec::new(),
2410 })
2411 }
2412 }
2413 }
2414
2415 struct AlwaysFailingExecutor;
2416
2417 impl LocalExecutor for AlwaysFailingExecutor {
2418 fn execute(
2419 &self,
2420 _capability: &traverse_registry::ResolvedCapability,
2421 _input: &Value,
2422 ) -> Result<LocalExecutionOutput, LocalExecutionFailure> {
2423 Err(LocalExecutionFailure {
2424 code: LocalExecutionFailureCode::ExecutionFailed,
2425 message: "always fails".to_string(),
2426 })
2427 }
2428 }
2429
2430 fn resolved_nodes() -> Vec<ResolvedProposalNode> {
2431 vec![
2432 ResolvedProposalNode {
2433 node_id: "a".to_string(),
2434 contract: contract(
2435 "test.produce",
2436 "1.0.0",
2437 json!({}),
2438 json!({}),
2439 automatic_risk(),
2440 ),
2441 },
2442 ResolvedProposalNode {
2443 node_id: "b".to_string(),
2444 contract: contract(
2445 "test.consume",
2446 "1.0.0",
2447 json!({}),
2448 json!({}),
2449 automatic_risk(),
2450 ),
2451 },
2452 ]
2453 }
2454
2455 #[test]
2456 fn executes_a_linear_proposal_threading_mapped_data_between_nodes() {
2457 let registry = registry_with(vec![
2458 (
2459 contract(
2460 "test.produce",
2461 "1.0.0",
2462 json!({}),
2463 json!({}),
2464 automatic_risk(),
2465 ),
2466 artifact("digest-a"),
2467 ),
2468 (
2469 contract(
2470 "test.consume",
2471 "1.0.0",
2472 json!({}),
2473 json!({}),
2474 automatic_risk(),
2475 ),
2476 artifact("digest-b"),
2477 ),
2478 ]);
2479 let consumer_saw_input = std::sync::Arc::new(Mutex::new(None));
2480 let executor = MappingAwareExecutor {
2481 consumer_saw_input: consumer_saw_input.clone(),
2482 };
2483 let runtime = Runtime::new(registry, executor)
2484 .with_security_config(RuntimeSecurityConfig::development());
2485
2486 let canonical = linear_canonical();
2487 let digest = proposal_digest(&canonical.proposal);
2488 let trace = execute_proposal(
2489 &runtime,
2490 &canonical,
2491 &resolved_nodes(),
2492 AuthorizationSummary {
2493 automatic: true,
2494 approval_token_id: None,
2495 },
2496 &digest,
2497 "snapshot-digest",
2498 );
2499
2500 assert_eq!(trace.terminal_state, ProposalTerminalState::Succeeded);
2501 assert_eq!(trace.node_outcomes.len(), 2);
2502 assert_eq!(trace.node_outcomes[0].status, ProposalNodeStatus::Succeeded);
2503 assert_eq!(trace.node_outcomes[1].status, ProposalNodeStatus::Succeeded);
2504
2505 let seen_input = consumer_saw_input
2506 .lock()
2507 .expect("test mutex must not be poisoned")
2508 .clone()
2509 .expect("consumer node must have executed and recorded its input");
2510 assert_eq!(
2511 seen_input,
2512 json!({"value": "produced-by-a"}),
2513 "node b's input must be assembled solely from the declared mapping, not node a's full output"
2514 );
2515 }
2516
2517 #[test]
2518 fn stops_at_first_failure_and_skips_remaining_nodes() {
2519 let registry = registry_with(vec![
2520 (
2521 contract(
2522 "test.produce",
2523 "1.0.0",
2524 json!({}),
2525 json!({}),
2526 automatic_risk(),
2527 ),
2528 artifact("digest-a"),
2529 ),
2530 (
2531 contract(
2532 "test.consume",
2533 "1.0.0",
2534 json!({}),
2535 json!({}),
2536 automatic_risk(),
2537 ),
2538 artifact("digest-b"),
2539 ),
2540 ]);
2541 let runtime = Runtime::new(registry, AlwaysFailingExecutor)
2542 .with_security_config(RuntimeSecurityConfig::development());
2543
2544 let canonical = linear_canonical();
2545 let digest = proposal_digest(&canonical.proposal);
2546 let trace = execute_proposal(
2547 &runtime,
2548 &canonical,
2549 &resolved_nodes(),
2550 AuthorizationSummary {
2551 automatic: true,
2552 approval_token_id: None,
2553 },
2554 &digest,
2555 "snapshot-digest",
2556 );
2557
2558 assert_eq!(trace.terminal_state, ProposalTerminalState::Failed);
2559 assert_eq!(trace.node_outcomes[0].status, ProposalNodeStatus::Failed);
2560 assert_eq!(
2561 trace.node_outcomes[1].status,
2562 ProposalNodeStatus::SkippedAfterEarlierFailure
2563 );
2564 }
2565
2566 #[test]
2567 fn executes_a_mapping_sourced_from_initial_input_into_a_nested_target_path() {
2568 let registry = registry_with(vec![(
2569 contract(
2570 "test.consume",
2571 "1.0.0",
2572 json!({}),
2573 json!({}),
2574 automatic_risk(),
2575 ),
2576 artifact("digest-x"),
2577 )]);
2578 let consumer_saw_input = std::sync::Arc::new(Mutex::new(None));
2579 let executor = MappingAwareExecutor {
2580 consumer_saw_input: consumer_saw_input.clone(),
2581 };
2582 let runtime = Runtime::new(registry, executor)
2583 .with_security_config(RuntimeSecurityConfig::development());
2584
2585 let canonical = initial_input_canonical();
2586 let digest = proposal_digest(&canonical.proposal);
2587 let trace = execute_proposal(
2588 &runtime,
2589 &canonical,
2590 &[ResolvedProposalNode {
2591 node_id: "x".to_string(),
2592 contract: contract(
2593 "test.consume",
2594 "1.0.0",
2595 json!({}),
2596 json!({}),
2597 automatic_risk(),
2598 ),
2599 }],
2600 AuthorizationSummary {
2601 automatic: true,
2602 approval_token_id: None,
2603 },
2604 &digest,
2605 "snapshot-digest",
2606 );
2607
2608 assert_eq!(trace.terminal_state, ProposalTerminalState::Succeeded);
2609 let seen_input = consumer_saw_input
2610 .lock()
2611 .expect("test mutex must not be poisoned")
2612 .clone()
2613 .expect("consumer must have received input");
2614 assert_eq!(seen_input["nested"]["value"], json!("from-caller"));
2615 }
2616
2617 #[test]
2618 fn execute_proposal_skips_an_execution_order_entry_with_no_matching_node() {
2619 let registry = registry_with(vec![(
2625 contract(
2626 "test.produce",
2627 "1.0.0",
2628 json!({}),
2629 json!({}),
2630 automatic_risk(),
2631 ),
2632 artifact("digest-a"),
2633 )]);
2634 let runtime = Runtime::new(registry, MappingAwareExecutor::default())
2635 .with_security_config(RuntimeSecurityConfig::development());
2636
2637 let mut proposal = linear_proposal_source();
2638 proposal.nodes.truncate(1);
2639 proposal.edges.clear();
2640 proposal.mappings.clear();
2641 let canonical = CanonicalProposal {
2642 execution_order: vec!["a".to_string(), "ghost".to_string()],
2643 proposal,
2644 };
2645 let digest = proposal_digest(&canonical.proposal);
2646 let trace = execute_proposal(
2647 &runtime,
2648 &canonical,
2649 &[ResolvedProposalNode {
2650 node_id: "a".to_string(),
2651 contract: contract(
2652 "test.produce",
2653 "1.0.0",
2654 json!({}),
2655 json!({}),
2656 automatic_risk(),
2657 ),
2658 }],
2659 AuthorizationSummary {
2660 automatic: true,
2661 approval_token_id: None,
2662 },
2663 &digest,
2664 "snapshot-digest",
2665 );
2666
2667 assert_eq!(trace.terminal_state, ProposalTerminalState::Succeeded);
2668 assert_eq!(trace.node_outcomes.len(), 1);
2669 }
2670
2671 #[test]
2672 fn pointer_set_with_an_empty_pointer_replaces_the_entire_target() {
2673 let mut target = json!({"unused": true});
2674 pointer_set(&mut target, "", json!({"replaced": true}));
2675 assert_eq!(target, json!({"replaced": true}));
2676 }
2677}