mod workflows;
pub use workflows::*;
mod artifact_router;
pub use artifact_router::*;
pub mod data_store;
pub mod events;
pub mod executor;
pub mod inference;
pub mod placement;
pub mod router;
pub mod security;
pub mod trace;
use chrono::Utc;
use events::{
EventBroker, EventCatalog, EventError, InProcessBroker, NoopRuntimeEventSink, RuntimeEventSink,
Subscription, SubscriptionPoll, TraverseEvent,
};
use executor::{
ArtifactType, CapabilityExecutor, ExecutorCapability, ExecutorError, ExecutorOutput,
};
use placement::{PlacementConstraintEvaluator, RuntimeSnapshot};
use router::{CapabilityExecutorRegistry, PlacementRouter, RouterError, RouterRequest};
use security::{
ArtifactVerificationFailure, ArtifactVerificationRecord, RuntimeIdentity,
RuntimeSecurityConfig, RuntimeWarning, derive_identity_from_jwt, verify_artifact,
};
use semver::Version;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value, json};
use std::fmt;
use std::fs;
use std::path::Path;
use std::sync::{Arc, Mutex};
use trace::TraceStore;
use traverse_contracts::{
EventReference, ExecutionTarget, HostApiAccess, Lifecycle, NetworkAccess, ServiceType,
ViolationRecord,
};
use traverse_registry::{
CapabilityRegistration, CapabilityRegistry, DiscoveryQuery, ImplementationKind, LookupScope,
ModelResolutionEvidence, RegistrationOutcome, RegistryFailure, RegistryScope, ResolutionError,
ResolvedCapability, WorkflowFailure, WorkflowRegistration, WorkflowRegistrationOutcome,
WorkflowRegistry, WorkspaceAppStateFailure, WorkspaceApplicationRegistration,
load_workspace_application_registries, resolve_dependencies, resolve_version_range,
};
use uuid::Uuid;
const RUNTIME_REQUEST_KIND: &str = "runtime_request";
const RUNTIME_RESULT_KIND: &str = "runtime_result";
const RUNTIME_STATE_EVENT_KIND: &str = "runtime_state_event";
const RUNTIME_TRACE_KIND: &str = "runtime_trace";
const RUNTIME_STATE_MACHINE_VALIDATION_KIND: &str = "runtime_state_machine_validation";
const BROWSER_SUBSCRIPTION_REQUEST_KIND: &str = "browser_runtime_subscription_request";
const BROWSER_SUBSCRIPTION_ERROR_KIND: &str = "browser_runtime_subscription_error";
const BROWSER_SUBSCRIPTION_LIFECYCLE_KIND: &str = "browser_runtime_subscription_lifecycle";
const BROWSER_SUBSCRIPTION_STATE_KIND: &str = "browser_runtime_subscription_state";
const BROWSER_SUBSCRIPTION_TRACE_KIND: &str = "browser_runtime_subscription_trace_artifact";
const BROWSER_SUBSCRIPTION_TERMINAL_KIND: &str = "browser_runtime_subscription_terminal";
const SUPPORTED_SCHEMA_VERSION: &str = "1.0.0";
const GOVERNING_SPEC: &str = "006-runtime-request-execution";
const STATE_MACHINE_GOVERNING_SPEC: &str = "010-runtime-state-machine";
const BROWSER_SUBSCRIPTION_GOVERNING_SPEC: &str = "013-browser-runtime-subscription";
const EXECUTION_PREFIX: &str = "exec_";
const TRACE_PREFIX: &str = "trace_";
const RUNTIME_EXECUTION_EVENT_TYPE: &str = "dev.traverse.runtime.execution.completed";
pub struct Runtime<E> {
registry: CapabilityRegistry,
workflow_registry: WorkflowRegistry,
applications: Vec<WorkspaceApplicationRegistration>,
executor: Arc<E>,
observability: RuntimeObservabilityConfig,
security: RuntimeSecurityConfig,
event_sink: Arc<dyn RuntimeEventSink>,
trace_store: Arc<Mutex<TraceStore>>,
event_broker: Arc<dyn EventBroker>,
}
impl<E: Clone> Clone for Runtime<E> {
fn clone(&self) -> Self {
Self {
registry: self.registry.clone(),
workflow_registry: self.workflow_registry.clone(),
applications: self.applications.clone(),
executor: Arc::clone(&self.executor),
observability: self.observability.clone(),
security: self.security.clone(),
event_sink: Arc::clone(&self.event_sink),
trace_store: Arc::clone(&self.trace_store),
event_broker: Arc::clone(&self.event_broker),
}
}
}
impl<E: fmt::Debug> fmt::Debug for Runtime<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Runtime")
.field("registry", &self.registry)
.field("workflow_registry", &self.workflow_registry)
.field("applications", &self.applications)
.field("executor", &self.executor)
.field("observability", &self.observability)
.field("security", &self.security)
.field("event_sink", &self.event_sink)
.field("trace_store", &self.trace_store)
.field("event_broker", &"Arc<dyn EventBroker>")
.finish()
}
}
impl<E> Runtime<E> {
#[must_use]
pub fn new(registry: CapabilityRegistry, executor: E) -> Self {
Self {
registry,
workflow_registry: WorkflowRegistry::new(),
applications: Vec::new(),
executor: Arc::new(executor),
observability: RuntimeObservabilityConfig::default(),
security: RuntimeSecurityConfig::default(),
event_sink: Arc::new(NoopRuntimeEventSink),
trace_store: Arc::new(Mutex::new(TraceStore::new())),
event_broker: default_event_broker(),
}
}
#[must_use]
pub fn with_workflow_registry(mut self, workflow_registry: WorkflowRegistry) -> Self {
self.workflow_registry = workflow_registry;
self
}
#[must_use]
pub fn with_workspace_applications(
mut self,
applications: Vec<WorkspaceApplicationRegistration>,
) -> Self {
self.applications = applications;
self
}
pub fn from_workspace_app_state(
workspace_root: &Path,
workspace_id: &str,
executor: E,
validator_version: &str,
) -> Result<Self, WorkspaceAppStateFailure> {
let loaded =
load_workspace_application_registries(workspace_root, workspace_id, validator_version)?;
Ok(Self::new(loaded.capability_registry, executor)
.with_workflow_registry(loaded.workflow_registry)
.with_workspace_applications(loaded.applications))
}
#[must_use]
pub fn with_observability_config(mut self, observability: RuntimeObservabilityConfig) -> Self {
self.observability = observability;
self
}
#[must_use]
pub fn observability_config(&self) -> &RuntimeObservabilityConfig {
&self.observability
}
#[must_use]
pub fn with_security_config(mut self, security: RuntimeSecurityConfig) -> Self {
self.security = security;
self
}
#[must_use]
pub fn with_event_sink(mut self, event_sink: Arc<dyn RuntimeEventSink>) -> Self {
self.event_sink = event_sink;
self
}
#[must_use]
pub fn with_event_broker(mut self, event_broker: Arc<dyn EventBroker>) -> Self {
self.event_broker = event_broker;
self
}
#[must_use]
pub fn with_trace_store(mut self, trace_store: Arc<Mutex<TraceStore>>) -> Self {
self.trace_store = trace_store;
self
}
#[must_use]
pub fn event_broker(&self) -> Arc<dyn EventBroker> {
Arc::clone(&self.event_broker)
}
#[must_use]
pub fn trace_store(&self) -> Arc<Mutex<TraceStore>> {
Arc::clone(&self.trace_store)
}
#[must_use]
pub fn security_config(&self) -> &RuntimeSecurityConfig {
&self.security
}
#[must_use]
pub fn capability_registry(&self) -> &CapabilityRegistry {
&self.registry
}
pub fn register_capability(
&mut self,
registration: CapabilityRegistration,
) -> Result<RegistrationOutcome, RegistryFailure> {
self.registry.register(registration)
}
#[must_use]
pub fn workflow_registry(&self) -> &WorkflowRegistry {
&self.workflow_registry
}
#[must_use]
pub fn workspace_applications(&self) -> &[WorkspaceApplicationRegistration] {
self.applications.as_slice()
}
#[must_use]
pub fn workflow_registry_mut(&mut self) -> &mut WorkflowRegistry {
&mut self.workflow_registry
}
pub fn register_workflow(
&mut self,
registration: WorkflowRegistration,
) -> Result<WorkflowRegistrationOutcome, WorkflowFailure> {
self.workflow_registry
.register(&self.registry, registration)
}
pub fn execute_governed_model_dependency(
&self,
app_id: &str,
app_version: &str,
request: &inference::GovernedModelExecutionRequest,
) -> Result<inference::GovernedModelExecutionOutcome, inference::GovernedModelExecutionError>
{
let Some(application) = self.applications.iter().find(|application| {
application.app_id == app_id && application.app_version == app_version
}) else {
return Err(inference::GovernedModelExecutionError::new(
inference::GovernedModelExecutionErrorCode::InterfaceNotDeclared,
"application registration was not loaded into this runtime",
));
};
let Some(dependency) = application
.model_dependencies
.iter()
.find(|dependency| dependency.interface_id == request.interface_id)
else {
return Err(inference::GovernedModelExecutionError::new(
inference::GovernedModelExecutionErrorCode::InterfaceNotDeclared,
"requested inference interface is not declared by this application",
));
};
inference::execute_governed_ollama_model_dependency(dependency, request)
}
}
pub trait LocalExecutor: Send + Sync {
fn execute(
&self,
capability: &ResolvedCapability,
input: &Value,
) -> Result<LocalExecutionOutput, LocalExecutionFailure>;
}
#[derive(Debug, Clone, PartialEq)]
pub struct LocalExecutionOutput {
pub value: Value,
pub emitted_events: Vec<TraverseEvent>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LocalExecutionFailure {
pub code: LocalExecutionFailureCode,
pub message: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LocalExecutionFailureCode {
ExecutionFailed,
Timeout,
InvalidInput,
ResourceExhausted,
ConstraintViolated,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeRequest {
pub kind: String,
pub schema_version: String,
pub request_id: String,
pub intent: RuntimeIntent,
pub input: Value,
pub lookup: RuntimeLookup,
pub context: RuntimeContext,
pub governing_spec: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeIntent {
#[serde(default)]
pub capability_id: Option<String>,
#[serde(default)]
pub capability_version: Option<String>,
#[serde(default)]
pub version_range: Option<String>,
#[serde(default)]
pub intent_key: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeLookup {
pub scope: RuntimeLookupScope,
pub allow_ambiguity: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeLookupScope {
PublicOnly,
PreferPrivate,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeContext {
pub requested_target: PlacementTarget,
#[serde(default)]
pub correlation_id: Option<String>,
#[serde(default)]
pub caller: Option<String>,
#[serde(default)]
pub traceparent: Option<String>,
#[serde(default)]
pub tracestate: Option<String>,
#[serde(default)]
pub metadata: Option<Value>,
#[serde(default)]
pub identity: Option<RuntimeIdentity>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeObservabilityConfig {
pub signals: OTelSignalConfig,
pub exporter: OTelExporterConfig,
pub deterministic_ids: bool,
#[serde(default)]
pub deterministic_seed: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelSignalConfig {
pub traces_enabled: bool,
pub logs_enabled: bool,
pub metrics_enabled: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelExporterConfig {
#[serde(default)]
pub endpoint: Option<String>,
pub protocol: OtlpProtocol,
}
impl Default for RuntimeObservabilityConfig {
fn default() -> Self {
Self {
signals: OTelSignalConfig {
traces_enabled: true,
logs_enabled: false,
metrics_enabled: false,
},
exporter: OTelExporterConfig {
endpoint: None,
protocol: OtlpProtocol::Http,
},
deterministic_ids: false,
deterministic_seed: None,
}
}
}
impl RuntimeObservabilityConfig {
#[must_use]
pub fn deterministic_test(seed: &str) -> Self {
Self {
deterministic_ids: true,
deterministic_seed: Some(seed.to_string()),
..Self::default()
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OtlpProtocol {
Http,
Grpc,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PlacementTarget {
Local,
Browser,
Edge,
Cloud,
Worker,
Device,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlacementDecisionRecord {
pub requested_target: PlacementTarget,
#[serde(default)]
pub selected_target: Option<PlacementTarget>,
pub status: PlacementDecisionStatus,
pub reason: PlacementDecisionReason,
pub supported_executor_targets: Vec<PlacementTarget>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PlacementDecisionStatus {
NotAttempted,
Selected,
Unsupported,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PlacementDecisionReason {
SelectionNotReached,
RequestedTargetSelected,
RequestedTargetUnsupported,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeStateEvent {
pub kind: String,
pub schema_version: String,
pub event_id: String,
pub execution_id: String,
pub request_id: String,
pub state: RuntimeState,
pub entered_at: String,
pub details: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeState {
Idle,
LoadingRegistry,
Ready,
Discovering,
EvaluatingConstraints,
Selecting,
Executing,
EmittingEvents,
Completed,
Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeTransitionReasonCode {
RuntimeInitializationStarted,
RegistryLoaded,
RegistryLoadFailed,
RequestStarted,
CandidatesCollected,
NoMatch,
ConstraintsEvaluated,
ConstraintValidationFailed,
CandidateSelected,
SelectionFailed,
ExecutionSucceededWithEvents,
ExecutionSucceeded,
ExecutionFailed,
EventsEmitted,
EventEmissionFailed,
ExecutionClosed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeTransitionRecord {
pub from_state: RuntimeState,
pub to_state: RuntimeState,
pub reason_code: RuntimeTransitionReasonCode,
pub occurred_at: String,
#[serde(default)]
pub request_id: Option<String>,
#[serde(default)]
pub execution_id: Option<String>,
#[serde(default)]
pub details: Option<Value>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeStateMachineValidationEvidence {
pub kind: String,
pub schema_version: String,
pub governing_spec: String,
pub validated_at: String,
pub status: RuntimeStateMachineValidationStatus,
pub checked_states: Vec<RuntimeState>,
pub checked_transitions: Vec<String>,
pub violations: Vec<Value>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeStateMachineValidationStatus {
Passed,
Failed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeTrace {
pub kind: String,
pub schema_version: String,
pub trace_id: String,
pub execution_id: String,
pub request_id: String,
pub governing_spec: String,
pub request: RuntimeRequest,
pub decision_evidence: TraceDecisionEvidence,
pub state_progression: TraceStateProgression,
pub terminal_outcome: TraceTerminalOutcome,
pub emitted_events: Vec<traverse_contracts::EventReference>,
#[serde(default)]
pub workflow_evidence: Option<WorkflowTraversalEvidence>,
#[serde(default)]
pub model_resolution: Vec<ModelResolutionEvidence>,
pub state_transitions: Vec<RuntimeTransitionRecord>,
pub state_machine_validation: RuntimeStateMachineValidationEvidence,
pub candidate_collection: CandidateCollectionRecord,
pub selection: SelectionRecord,
pub execution: ExecutionRecord,
pub result: TraceResultRecord,
pub otel_trace: OTelTraceRecord,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelTraceRecord {
pub trace_id: String,
#[serde(default)]
pub parent_traceparent: Option<String>,
#[serde(default)]
pub tracestate: Option<String>,
pub exporter: OTelExporterRecord,
pub spans: Vec<OTelSpanRecord>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelExporterRecord {
pub enabled: bool,
#[serde(default)]
pub endpoint: Option<String>,
pub protocol: OtlpProtocol,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelSpanRecord {
pub trace_id: String,
pub span_id: String,
#[serde(default)]
pub parent_span_id: Option<String>,
pub name: String,
pub kind: OTelSpanKind,
pub status: OTelSpanStatus,
pub started_at: String,
pub ended_at: String,
pub attributes: Vec<OTelAttribute>,
#[serde(default)]
pub events: Vec<OTelSpanEvent>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OTelSpanKind {
Internal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum OTelSpanStatus {
Ok,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelAttribute {
pub key: String,
pub value: Value,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OTelSpanEvent {
pub name: String,
pub timestamp: String,
pub attributes: Vec<OTelAttribute>,
}
impl RuntimeTrace {
#[must_use]
pub fn with_model_resolution(mut self, evidence: Vec<ModelResolutionEvidence>) -> Self {
self.decision_evidence
.model_resolution
.clone_from(&evidence);
self.model_resolution = evidence;
self
}
#[must_use]
pub fn selected_capability_id(&self) -> Option<&str> {
self.selection.selected_capability_id.as_deref()
}
#[must_use]
pub fn errors(&self) -> Option<&RuntimeError> {
self.terminal_outcome.error.as_ref()
}
#[must_use]
pub fn emitted_events(&self) -> &[traverse_contracts::EventReference] {
self.emitted_events.as_slice()
}
#[must_use]
pub fn output(&self) -> Option<&serde_json::Value> {
self.result.output.as_ref()
}
#[must_use]
pub fn is_success(&self) -> bool {
self.terminal_outcome.runtime_status == RuntimeResultStatus::Completed
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TraceDecisionEvidence {
pub candidate_collection: CandidateCollectionRecord,
pub selection: SelectionRecord,
#[serde(default)]
pub model_resolution: Vec<ModelResolutionEvidence>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TraceStateProgression {
pub state_events: Vec<RuntimeStateEvent>,
pub transitions: Vec<RuntimeTransitionRecord>,
pub validation: RuntimeStateMachineValidationEvidence,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TraceTerminalOutcome {
pub runtime_status: RuntimeResultStatus,
pub execution_status: ExecutionStatus,
#[serde(default)]
pub failure_reason: Option<ExecutionFailureReason>,
#[serde(default)]
pub error: Option<RuntimeError>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CandidateCollectionRecord {
pub lookup_scope: RuntimeLookupScope,
pub candidates: Vec<RuntimeCandidate>,
pub rejected_candidates: Vec<RejectedRuntimeCandidate>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeCandidate {
pub scope: RuntimeRegistryScope,
pub capability_id: String,
pub capability_version: String,
pub artifact_ref: String,
pub implementation_kind: RuntimeImplementationKind,
pub lifecycle: RuntimeLifecycle,
pub reason: CandidateReason,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RejectedRuntimeCandidate {
pub capability_id: String,
pub capability_version: String,
pub scope: RuntimeRegistryScope,
pub reason: RejectedCandidateReason,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CandidateReason {
ExactMatch,
IntentMatch,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RejectedCandidateReason {
WrongScope,
NotRunnableLocally,
LifecycleNotRunnable,
InputContractInvalid,
ArtifactMissing,
SupersededByPrivateOverlay,
NotSelectedAfterOrdering,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SelectionRecord {
pub status: SelectionStatus,
#[serde(default)]
pub selected_capability_id: Option<String>,
#[serde(default)]
pub selected_capability_version: Option<String>,
#[serde(default)]
pub failure_reason: Option<SelectionFailureReason>,
#[serde(default)]
pub remaining_candidates: Vec<RuntimeCandidate>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SelectionStatus {
Selected,
NoMatch,
Ambiguous,
InvalidRequest,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SelectionFailureReason {
InvalidRequest,
NoMatch,
Ambiguous,
NotRunnable,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ExecutionRecord {
pub placement: PlacementDecisionRecord,
pub placement_target: PlacementTarget,
pub status: ExecutionStatus,
#[serde(default)]
pub artifact_ref: Option<String>,
#[serde(default)]
pub started_at: Option<String>,
#[serde(default)]
pub completed_at: Option<String>,
#[serde(default)]
pub output_digest: Option<String>,
#[serde(default)]
pub failure_reason: Option<ExecutionFailureReason>,
#[serde(default)]
pub artifact_verification: Option<ArtifactVerificationRecord>,
#[serde(default)]
pub identity: Option<RuntimeIdentity>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionStatus {
NotStarted,
Succeeded,
Failed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionFailureReason {
ContractInputInvalid,
ArtifactMissing,
ArtifactNotRunnable,
PlacementUnsupported,
ExecutionFailed,
ContractOutputInvalid,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TraceResultRecord {
pub status: RuntimeResultStatus,
#[serde(default)]
pub output: Option<serde_json::Value>,
#[serde(default)]
pub error: Option<RuntimeError>,
#[serde(default)]
pub warnings: Vec<RuntimeWarning>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeResult {
pub kind: String,
pub schema_version: String,
pub execution_id: String,
pub request_id: String,
pub status: RuntimeResultStatus,
pub trace_ref: String,
#[serde(default)]
pub output: Option<Value>,
#[serde(default)]
pub error: Option<RuntimeError>,
#[serde(default)]
pub warnings: Vec<RuntimeWarning>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeResultStatus {
Completed,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeError {
pub code: RuntimeErrorCode,
pub message: String,
pub details: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeErrorCode {
RequestInvalid,
CapabilityNotFound,
CapabilityAmbiguous,
CapabilityNotRunnable,
PlacementUnsupported,
ArtifactMissing,
ExecutionFailed,
OutputValidationFailed,
ContractViolation,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RuntimeExecutionOutcome {
pub result: RuntimeResult,
pub trace: RuntimeTrace,
pub state_events: Vec<RuntimeStateEvent>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionRequest {
pub kind: String,
pub schema_version: String,
pub governing_spec: String,
#[serde(default)]
pub request_id: Option<String>,
#[serde(default)]
pub execution_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionErrorMessage {
pub kind: String,
pub schema_version: String,
pub sequence: u64,
pub code: BrowserRuntimeSubscriptionErrorCode,
pub message: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BrowserRuntimeSubscriptionErrorCode {
InvalidRequest,
NotFound,
UnsupportedOperation,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionLifecycleMessage {
pub kind: String,
pub schema_version: String,
pub sequence: u64,
pub request_id: String,
pub execution_id: String,
pub status: BrowserRuntimeSubscriptionLifecycleStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BrowserRuntimeSubscriptionLifecycleStatus {
SubscriptionEstablished,
StreamCompleted,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionStateMessage {
pub kind: String,
pub schema_version: String,
pub sequence: u64,
pub state_event: RuntimeStateEvent,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionTraceArtifactMessage {
pub kind: String,
pub schema_version: String,
pub sequence: u64,
pub trace: RuntimeTrace,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BrowserRuntimeSubscriptionTerminalMessage {
pub kind: String,
pub schema_version: String,
pub sequence: u64,
pub result: RuntimeResult,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum BrowserRuntimeSubscriptionMessage {
Error(BrowserRuntimeSubscriptionErrorMessage),
Lifecycle(Box<BrowserRuntimeSubscriptionLifecycleMessage>),
State(Box<BrowserRuntimeSubscriptionStateMessage>),
TraceArtifact(Box<BrowserRuntimeSubscriptionTraceArtifactMessage>),
StreamTerminal(Box<BrowserRuntimeSubscriptionTerminalMessage>),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeRegistryScope {
Public,
Private,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeImplementationKind {
Executable,
Workflow,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RuntimeLifecycle {
Draft,
Active,
Deprecated,
Retired,
Archived,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RequestParseFailure {
pub message: String,
}
impl fmt::Display for RequestParseFailure {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.message)
}
}
impl std::error::Error for RequestParseFailure {}
pub fn parse_runtime_request(json: &str) -> Result<RuntimeRequest, RequestParseFailure> {
serde_json::from_str::<RuntimeRequest>(json).map_err(|error| RequestParseFailure {
message: error.to_string(),
})
}
#[must_use]
pub fn browser_subscription_messages(
request: &BrowserRuntimeSubscriptionRequest,
outcome: &RuntimeExecutionOutcome,
) -> Vec<BrowserRuntimeSubscriptionMessage> {
if let Some(error) = validate_browser_subscription_request(request) {
return vec![BrowserRuntimeSubscriptionMessage::Error(error)];
}
if !subscription_targets_outcome(request, outcome) {
return vec![BrowserRuntimeSubscriptionMessage::Error(
BrowserRuntimeSubscriptionErrorMessage {
kind: BROWSER_SUBSCRIPTION_ERROR_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence: 0,
code: BrowserRuntimeSubscriptionErrorCode::NotFound,
message: "subscription target did not match the supplied execution outcome"
.to_string(),
},
)];
}
let mut sequence = 0_u64;
let mut messages = Vec::new();
messages.push(BrowserRuntimeSubscriptionMessage::Lifecycle(Box::new(
BrowserRuntimeSubscriptionLifecycleMessage {
kind: BROWSER_SUBSCRIPTION_LIFECYCLE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence,
request_id: outcome.result.request_id.clone(),
execution_id: outcome.result.execution_id.clone(),
status: BrowserRuntimeSubscriptionLifecycleStatus::SubscriptionEstablished,
},
)));
sequence += 1;
for state_event in &outcome.state_events {
messages.push(BrowserRuntimeSubscriptionMessage::State(Box::new(
BrowserRuntimeSubscriptionStateMessage {
kind: BROWSER_SUBSCRIPTION_STATE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence,
state_event: state_event.clone(),
},
)));
sequence += 1;
}
messages.push(BrowserRuntimeSubscriptionMessage::TraceArtifact(Box::new(
BrowserRuntimeSubscriptionTraceArtifactMessage {
kind: BROWSER_SUBSCRIPTION_TRACE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence,
trace: outcome.trace.clone(),
},
)));
sequence += 1;
messages.push(BrowserRuntimeSubscriptionMessage::StreamTerminal(Box::new(
BrowserRuntimeSubscriptionTerminalMessage {
kind: BROWSER_SUBSCRIPTION_TERMINAL_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence,
result: outcome.result.clone(),
},
)));
sequence += 1;
messages.push(BrowserRuntimeSubscriptionMessage::Lifecycle(Box::new(
BrowserRuntimeSubscriptionLifecycleMessage {
kind: BROWSER_SUBSCRIPTION_LIFECYCLE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence,
request_id: outcome.result.request_id.clone(),
execution_id: outcome.result.execution_id.clone(),
status: BrowserRuntimeSubscriptionLifecycleStatus::StreamCompleted,
},
)));
messages
}
fn validate_browser_subscription_request(
request: &BrowserRuntimeSubscriptionRequest,
) -> Option<BrowserRuntimeSubscriptionErrorMessage> {
if request.kind != BROWSER_SUBSCRIPTION_REQUEST_KIND {
return Some(browser_subscription_error(
BrowserRuntimeSubscriptionErrorCode::InvalidRequest,
"kind must equal browser_runtime_subscription_request",
));
}
if request.schema_version != SUPPORTED_SCHEMA_VERSION {
return Some(browser_subscription_error(
BrowserRuntimeSubscriptionErrorCode::InvalidRequest,
"schema_version must equal 1.0.0",
));
}
if request.governing_spec != BROWSER_SUBSCRIPTION_GOVERNING_SPEC {
return Some(browser_subscription_error(
BrowserRuntimeSubscriptionErrorCode::InvalidRequest,
"governing_spec must equal 013-browser-runtime-subscription",
));
}
match (&request.request_id, &request.execution_id) {
(Some(request_id), None) if non_empty(request_id) => None,
(None, Some(execution_id)) if non_empty(execution_id) => None,
(Some(_), Some(_)) => Some(browser_subscription_error(
BrowserRuntimeSubscriptionErrorCode::InvalidRequest,
"exactly one target selector must be supplied",
)),
_ => Some(browser_subscription_error(
BrowserRuntimeSubscriptionErrorCode::InvalidRequest,
"subscription request must include request_id or execution_id",
)),
}
}
fn subscription_targets_outcome(
request: &BrowserRuntimeSubscriptionRequest,
outcome: &RuntimeExecutionOutcome,
) -> bool {
match (&request.request_id, &request.execution_id) {
(Some(request_id), None) => request_id == &outcome.result.request_id,
(None, Some(execution_id)) => execution_id == &outcome.result.execution_id,
_ => false,
}
}
fn browser_subscription_error(
code: BrowserRuntimeSubscriptionErrorCode,
message: &str,
) -> BrowserRuntimeSubscriptionErrorMessage {
BrowserRuntimeSubscriptionErrorMessage {
kind: BROWSER_SUBSCRIPTION_ERROR_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
sequence: 0,
code,
message: message.to_string(),
}
}
impl<E> Runtime<E>
where
E: LocalExecutor,
{
#[must_use]
pub fn execute(&self, request: RuntimeRequest) -> RuntimeExecutionOutcome {
let identity = request.context.identity.clone();
let (attempt, mut emitter) = begin_attempt(request, self.observability.clone());
emitter.push(
RuntimeState::Discovering,
RuntimeTransitionReasonCode::RequestStarted,
json!({
"lookup_scope": attempt.request.lookup.scope,
"identity": attempt.request.context.identity,
}),
);
let mut outcome = if let Some(error) = validate_request(&attempt.request) {
invalid_request_outcome(attempt, emitter, error)
} else {
let resolution = self.resolve_candidates(&attempt.request, &mut emitter);
if resolution.eligible.is_empty() {
no_eligible_outcome(attempt, emitter, resolution.collection)
} else if resolution.eligible.len() > 1 {
ambiguous_outcome(attempt, emitter, resolution)
} else {
let mut eligible = resolution.eligible;
let selected = eligible.remove(0);
let selection = SelectionRecord {
status: SelectionStatus::Selected,
selected_capability_id: Some(selected.record.id.clone()),
selected_capability_version: Some(selected.record.version.clone()),
failure_reason: None,
remaining_candidates: Vec::new(),
};
self.execute_selected(
attempt,
emitter,
resolution.collection,
selection,
&selected,
)
}
};
self.emit_execution_lifecycle_event(&mut outcome, identity.as_ref());
outcome
}
fn emit_execution_lifecycle_event(
&self,
outcome: &mut RuntimeExecutionOutcome,
identity: Option<&RuntimeIdentity>,
) {
let event = TraverseEvent {
id: Uuid::new_v4().to_string(),
source: "traverse-runtime".to_string(),
event_type: RUNTIME_EXECUTION_EVENT_TYPE.to_string(),
datacontenttype: "application/json".to_string(),
time: Utc::now().to_rfc3339(),
data: json!({
"execution_id": outcome.result.execution_id,
"request_id": outcome.result.request_id,
"status": outcome.result.status,
"trace_ref": outcome.result.trace_ref,
}),
owner: "traverse-runtime".to_string(),
version: SUPPORTED_SCHEMA_VERSION.to_string(),
lifecycle_status: events::LifecycleStatus::Active,
deduplication_id: Some(outcome.result.execution_id.clone()),
ordering_scope: Some(outcome.result.request_id.clone()),
correlation_id: Some(outcome.result.request_id.clone()),
causation_id: Some(outcome.result.request_id.clone()),
subject_id: identity.map(|identity| identity.subject_id.clone()),
actor_id: identity.and_then(|identity| identity.actor_id.clone()),
};
if let Err(error) = self.event_sink.emit(event) {
let warning = RuntimeWarning {
code: "runtime_event_sink_delivery_failed".to_string(),
message: error.to_string(),
};
outcome.result.warnings.push(warning.clone());
outcome.trace.result.warnings.push(warning);
}
}
fn collect_candidates(
&self,
request: &RuntimeRequest,
_reason: CandidateReason,
) -> Vec<ResolvedCapability> {
let lookup_scope = map_lookup_scope(request.lookup.scope);
if is_exact_target(&request.intent) {
return request
.intent
.capability_id
.as_deref()
.zip(request.intent.capability_version.as_deref())
.and_then(|(id, version)| self.registry.find_exact(lookup_scope, id, version))
.into_iter()
.collect();
}
if let (Some(capability_id), Some(range_str)) = (
request.intent.capability_id.as_deref(),
request.intent.version_range.as_deref(),
) && non_empty(capability_id)
&& non_empty(range_str)
{
return match resolve_version_range(
&self.registry,
capability_id,
range_str,
lookup_scope,
) {
Ok(resolved) => {
let entry_lookup = match resolved.scope {
RegistryScope::Public => LookupScope::PublicOnly,
RegistryScope::Private => LookupScope::PreferPrivate,
};
self.registry
.find_exact(entry_lookup, &resolved.capability_id, &resolved.version)
.into_iter()
.collect()
}
Err(_) => Vec::new(),
};
}
let target = request
.intent
.capability_id
.as_deref()
.or(request.intent.intent_key.as_deref())
.unwrap_or_default();
self.registry
.discover(lookup_scope, &DiscoveryQuery::default())
.into_iter()
.filter(|entry| entry.id == target)
.filter_map(|entry| {
let scope = match entry.scope {
traverse_registry::RegistryScope::Public => LookupScope::PublicOnly,
traverse_registry::RegistryScope::Private => LookupScope::PreferPrivate,
};
self.registry.find_exact(scope, &entry.id, &entry.version)
})
.collect()
}
fn resolve_candidates(
&self,
request: &RuntimeRequest,
emitter: &mut StateEmitter,
) -> CandidateResolution {
let candidate_reason = if is_exact_target(&request.intent) {
CandidateReason::ExactMatch
} else {
CandidateReason::IntentMatch
};
let discovered = self.collect_candidates(request, candidate_reason);
if !discovered.is_empty() {
emitter.push(
RuntimeState::EvaluatingConstraints,
RuntimeTransitionReasonCode::CandidatesCollected,
json!({"candidate_count": discovered.len()}),
);
}
let mut eligible = Vec::new();
let mut rejected = Vec::new();
for candidate in discovered {
match evaluate_candidate(candidate) {
CandidateEvaluation::Eligible(capability) => eligible.push(capability),
CandidateEvaluation::Rejected(candidate, reason) => {
rejected.push(RejectedRuntimeCandidate {
capability_id: candidate.record.id.clone(),
capability_version: candidate.record.version.clone(),
scope: map_registry_scope(candidate.record.scope),
reason,
});
}
}
}
if !eligible.is_empty() {
emitter.push(
RuntimeState::Selecting,
RuntimeTransitionReasonCode::ConstraintsEvaluated,
json!({
"eligible_candidates": eligible.len(),
"rejected_candidates": rejected.len()
}),
);
}
CandidateResolution {
eligible: eligible.clone(),
collection: CandidateCollectionRecord {
lookup_scope: request.lookup.scope,
candidates: eligible
.iter()
.map(|capability| runtime_candidate(capability, candidate_reason))
.collect(),
rejected_candidates: rejected,
},
candidate_reason,
}
}
#[allow(clippy::too_many_lines)]
fn execute_selected(
&self,
attempt: AttemptContext,
emitter: StateEmitter,
candidate_collection: CandidateCollectionRecord,
selection: SelectionRecord,
selected: &ResolvedCapability,
) -> RuntimeExecutionOutcome {
let mut context = ExecutionContext {
attempt,
emitter,
candidate_collection,
selection,
};
let requested_target = context.attempt.request.context.requested_target;
if contains_drafts_segment(&selected.record.contract_path) {
let violation = ViolationRecord::new(
"draft_artifact_not_executable",
selected.record.contract_path.clone(),
"draft artifacts are quarantined under drafts/ and must not be executable",
);
let error = runtime_error(
RuntimeErrorCode::ContractViolation,
"draft artifacts are not executable",
json!({"violations": [violation]}),
);
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::ArtifactNotRunnable,
placement: placement_not_attempted(
requested_target,
PlacementDecisionReason::SelectionNotReached,
),
error,
artifact_verification: None,
},
);
}
let placement = match resolve_placement(requested_target) {
Ok(placement) => placement,
Err(error) => {
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::PlacementUnsupported,
placement: placement_not_attempted(
requested_target,
PlacementDecisionReason::RequestedTargetUnsupported,
),
error,
artifact_verification: None,
},
);
}
};
let artifact_bytes = match load_artifact_bytes_for_verification(selected) {
Ok(bytes) => bytes,
Err(error) => {
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::ArtifactMissing,
placement,
error,
artifact_verification: None,
},
);
}
};
match verify_artifact(selected, &artifact_bytes, &self.security) {
Ok(record) => {
if record.warning_code.is_some() {
context.attempt.warnings.push(RuntimeWarning {
code: record.warning_code.clone().unwrap_or_default(),
message: "unsigned local/dev artifact allowed by development security mode"
.to_string(),
});
}
context.attempt.artifact_verification = Some(record);
}
Err(error) => {
let record = error.record().clone();
let runtime_error = artifact_verification_runtime_error(&error);
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::ArtifactNotRunnable,
placement,
error: runtime_error,
artifact_verification: Some(record),
},
);
}
}
let lookup_scope = map_lookup_scope(context.attempt.request.lookup.scope);
if let Err(dep_error) = resolve_dependencies(
&self.registry,
&selected.record.id,
&selected.contract.dependencies,
lookup_scope,
) {
let (detail_id, detail_version) = match &dep_error {
ResolutionError::MissingDependency {
capability_id,
required_version,
} => (capability_id.clone(), required_version.clone()),
ResolutionError::CircularDependency { cycle } => {
(cycle.join(" -> "), String::new())
}
ResolutionError::MaxTransitiveDepthExceeded { depth, chain } => {
(format!("depth={depth}"), chain.join(" -> "))
}
};
let error = runtime_error(
RuntimeErrorCode::CapabilityNotFound,
"dependency resolution failed before execution",
serde_json::json!({
"dependency_id": detail_id,
"required_version": detail_version,
}),
);
let artifact_verification = context.attempt.artifact_verification.clone();
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::ArtifactMissing,
placement,
error,
artifact_verification,
},
);
}
if let Err(error) = validate_payload_against_contract(
&context.attempt.request.input,
&selected.contract.inputs.schema,
RuntimeErrorCode::RequestInvalid,
"runtime request input does not satisfy the selected capability input contract",
) {
let artifact_verification = context.attempt.artifact_verification.clone();
return pre_execution_failure_outcome(
context,
PreExecutionFailure {
artifact_ref: Some(selected.record.artifact_ref.clone()),
failure_reason: ExecutionFailureReason::ContractInputInvalid,
placement,
error,
artifact_verification,
},
);
}
self.execute_started_selection(context, selected, placement)
}
fn execute_started_selection(
&self,
mut context: ExecutionContext,
selected: &ResolvedCapability,
placement: PlacementDecisionRecord,
) -> RuntimeExecutionOutcome {
let identity = context.attempt.request.context.identity.clone();
let started_execution =
start_selected_execution(&mut context.emitter, selected, placement, identity.as_ref());
if selected.record.implementation_kind == ImplementationKind::Workflow {
return self.execute_workflow_capability(context, selected, started_execution);
}
let artifact_type = artifact_type_for(selected);
let executor_capability = executor_capability_for(selected, artifact_type.clone());
let bridge = BoundLocalExecutor {
executor: Arc::clone(&self.executor),
selected: selected.clone(),
};
let router = PlacementRouter::new(
PlacementConstraintEvaluator,
CapabilityExecutorRegistry::new(),
Arc::clone(&self.trace_store),
Arc::clone(&self.event_broker),
);
let target_hint = execution_target_from_placement(
started_execution
.placement
.selected_target
.unwrap_or(started_execution.placement.requested_target),
);
let router_request = RouterRequest {
capability_id: selected.record.id.clone(),
artifact_type,
contract: selected.contract.clone(),
target_hint: Some(target_hint),
runtime_snapshot: idle_runtime_snapshot(),
input: context.attempt.request.input.clone(),
executor_capability,
trace_id_override: Some(context.attempt.trace_id.clone()),
};
let router_result = router.execute_with_executor(router_request, &bridge);
match router_result {
Ok(response) => {
if let Err(error) = validate_payload_against_contract(
&response.output,
&selected.contract.outputs.schema,
RuntimeErrorCode::OutputValidationFailed,
"executor output does not satisfy the selected capability output contract",
) {
return execution_failure_outcome(
context,
ExecutionFailureState {
artifact_ref: selected.record.artifact_ref.clone(),
started_at: started_execution.started_at,
placement: started_execution.placement,
failure_reason: ExecutionFailureReason::ContractOutputInvalid,
},
error,
Vec::new(),
None,
);
}
let emitted_events: Vec<EventReference> = response
.emitted_events
.iter()
.map(|event| EventReference {
event_id: event.event_type.clone(),
version: event.version.clone(),
})
.collect();
successful_execution_outcome(
context,
selected,
started_execution,
response.output,
emitted_events,
None,
)
}
Err(error) => {
let (runtime_error, failure_reason, emitted_events) = map_router_error(&error);
execution_failure_outcome(
context,
ExecutionFailureState {
artifact_ref: selected.record.artifact_ref.clone(),
started_at: started_execution.started_at,
placement: started_execution.placement,
failure_reason,
},
runtime_error,
emitted_events,
None,
)
}
}
}
}
struct BoundLocalExecutor<E> {
executor: Arc<E>,
selected: ResolvedCapability,
}
impl<E> CapabilityExecutor for BoundLocalExecutor<E>
where
E: LocalExecutor,
{
fn execute(
&self,
_capability: &ExecutorCapability,
input: &Value,
) -> Result<ExecutorOutput, ExecutorError> {
let output = self
.executor
.execute(&self.selected, input)
.map_err(|failure| ExecutorError::ExecutionFailed(failure.message))?;
validate_natively_emitted_events(&self.selected.contract, &output.emitted_events)
.map_err(ExecutorError::ExecutionFailed)?;
Ok(ExecutorOutput {
value: output.value,
emitted_events: output.emitted_events,
})
}
}
fn validate_natively_emitted_events(
contract: &traverse_contracts::CapabilityContract,
emitted_events: &[TraverseEvent],
) -> Result<(), String> {
if emitted_events.is_empty() {
return Ok(());
}
if contract.service_type != ServiceType::Subscribable {
return Err(format!(
"capability '{}' emitted events but its service_type is not Subscribable",
contract.id
));
}
for event in emitted_events {
let declared = contract
.emits
.iter()
.any(|decl| decl.event_id == event.event_type && decl.version == event.version);
if !declared {
return Err(format!(
"capability '{}' emitted an undeclared event {}@{}",
contract.id, event.event_type, event.version
));
}
}
Ok(())
}
#[derive(Debug, Default)]
struct DiscardEventBroker;
impl EventBroker for DiscardEventBroker {
fn publish(&self, _event: TraverseEvent) -> Result<(), EventError> {
Ok(())
}
fn subscribe(&self, event_type: &str, _from_cursor: &str) -> Result<Subscription, EventError> {
Err(EventError::UnregisteredEventType(event_type.to_string()))
}
fn subscribe_for_subject(
&self,
event_type: &str,
_from_cursor: &str,
_subject_id: Option<&str>,
) -> Result<Subscription, EventError> {
Err(EventError::UnregisteredEventType(event_type.to_string()))
}
fn poll(
&self,
subscription_id: &str,
_max_events: usize,
) -> Result<SubscriptionPoll, EventError> {
Err(EventError::SubscriptionNotFound(
subscription_id.to_string(),
))
}
fn cancel(&self, subscription_id: &str) -> Result<(), EventError> {
Err(EventError::SubscriptionNotFound(
subscription_id.to_string(),
))
}
}
fn default_event_broker() -> Arc<dyn EventBroker> {
event_broker_or_discard(InProcessBroker::new(Arc::new(EventCatalog::new())))
}
fn event_broker_or_discard(result: Result<InProcessBroker, EventError>) -> Arc<dyn EventBroker> {
match result {
Ok(broker) => Arc::new(broker),
Err(_) => Arc::new(DiscardEventBroker),
}
}
fn idle_runtime_snapshot() -> RuntimeSnapshot {
RuntimeSnapshot {
target_loads: [
(ExecutionTarget::Local, 0.0),
(ExecutionTarget::Browser, 0.0),
(ExecutionTarget::Edge, 0.0),
(ExecutionTarget::Cloud, 0.0),
(ExecutionTarget::Worker, 0.0),
(ExecutionTarget::Device, 0.0),
]
.into_iter()
.collect(),
}
}
fn artifact_type_for(selected: &ResolvedCapability) -> ArtifactType {
if selected.artifact.binary.is_some() {
ArtifactType::Wasm
} else {
ArtifactType::Native
}
}
fn executor_capability_for(
selected: &ResolvedCapability,
artifact_type: ArtifactType,
) -> ExecutorCapability {
let binary = selected.artifact.binary.as_ref();
ExecutorCapability {
capability_id: selected.contract.id.clone(),
artifact_type,
wasm_binary_path: binary.map(|binary| binary.location.clone()),
wasm_checksum: selected
.artifact
.digests
.binary_digest
.as_deref()
.and_then(|digest| digest.strip_prefix("sha256:"))
.map(str::to_string),
host_abi_version: None,
emits: selected.contract.emits.clone(),
service_type: selected.contract.service_type.clone(),
}
}
fn execution_target_from_placement(target: PlacementTarget) -> ExecutionTarget {
match target {
PlacementTarget::Local => ExecutionTarget::Local,
PlacementTarget::Browser => ExecutionTarget::Browser,
PlacementTarget::Edge => ExecutionTarget::Edge,
PlacementTarget::Cloud => ExecutionTarget::Cloud,
PlacementTarget::Worker => ExecutionTarget::Worker,
PlacementTarget::Device => ExecutionTarget::Device,
}
}
fn map_router_error(
error: &RouterError,
) -> (RuntimeError, ExecutionFailureReason, Vec<EventReference>) {
match error {
RouterError::PlacementFailed(placement_error) => (
runtime_error(
RuntimeErrorCode::PlacementUnsupported,
"placement constraints rejected the capability execution request",
json!({"placement_error": format!("{placement_error:?}")}),
),
ExecutionFailureReason::PlacementUnsupported,
Vec::new(),
),
RouterError::ExecutorNotFound(artifact_type) => (
runtime_error(
RuntimeErrorCode::CapabilityNotRunnable,
"no executor is registered for the capability artifact type",
json!({"artifact_type": artifact_type}),
),
ExecutionFailureReason::ArtifactNotRunnable,
Vec::new(),
),
RouterError::ExecutionFailed(message) => (
runtime_error(
RuntimeErrorCode::ExecutionFailed,
message,
json!({"code": "execution_failed"}),
),
ExecutionFailureReason::ExecutionFailed,
Vec::new(),
),
RouterError::ContractViolation(violations) => (
runtime_error(
RuntimeErrorCode::ContractViolation,
"capability execution violated its governed contract",
json!({"violations": violations}),
),
ExecutionFailureReason::ExecutionFailed,
Vec::new(),
),
RouterError::TraceLockPoisoned => (
runtime_error(
RuntimeErrorCode::ExecutionFailed,
"trace store lock is poisoned",
json!({"code": "trace_lock_poisoned"}),
),
ExecutionFailureReason::ExecutionFailed,
Vec::new(),
),
}
}
fn terminal_failure(context: FailureContext) -> RuntimeExecutionOutcome {
let result_record = TraceResultRecord {
status: RuntimeResultStatus::Error,
output: None,
error: Some(context.error.clone()),
warnings: context.attempt.warnings.clone(),
};
let otel_trace = otel_trace_record(
&context.attempt,
&context.state_transitions,
&context.selection,
&context.execution,
&result_record,
);
let trace = RuntimeTrace {
kind: RUNTIME_TRACE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
trace_id: context.attempt.trace_id.clone(),
execution_id: context.attempt.execution_id.clone(),
request_id: context.attempt.request.request_id.clone(),
governing_spec: GOVERNING_SPEC.to_string(),
request: context.attempt.request.clone(),
decision_evidence: TraceDecisionEvidence {
candidate_collection: context.candidate_collection.clone(),
selection: context.selection.clone(),
model_resolution: Vec::new(),
},
state_progression: TraceStateProgression {
state_events: context.state_events.clone(),
transitions: context.state_transitions.clone(),
validation: context.state_machine_validation.clone(),
},
terminal_outcome: TraceTerminalOutcome {
runtime_status: RuntimeResultStatus::Error,
execution_status: context.execution.status,
failure_reason: context.execution.failure_reason,
error: Some(context.error.clone()),
},
emitted_events: context.emitted_events,
workflow_evidence: context.workflow_evidence,
model_resolution: Vec::new(),
state_transitions: context.state_transitions,
state_machine_validation: context.state_machine_validation,
candidate_collection: context.candidate_collection,
selection: context.selection,
execution: context.execution,
result: result_record,
otel_trace,
};
let result = RuntimeResult {
kind: RUNTIME_RESULT_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
execution_id: context.attempt.execution_id,
request_id: context.attempt.request.request_id,
status: RuntimeResultStatus::Error,
trace_ref: context.attempt.trace_id,
output: None,
error: Some(context.error),
warnings: context.attempt.warnings,
};
RuntimeExecutionOutcome {
result,
trace,
state_events: context.state_events,
}
}
fn supported_executor_targets() -> Vec<PlacementTarget> {
vec![PlacementTarget::Local]
}
fn placement_not_attempted(
requested_target: PlacementTarget,
reason: PlacementDecisionReason,
) -> PlacementDecisionRecord {
PlacementDecisionRecord {
requested_target,
selected_target: None,
status: PlacementDecisionStatus::NotAttempted,
reason,
supported_executor_targets: supported_executor_targets(),
}
}
fn resolve_placement(
requested_target: PlacementTarget,
) -> Result<PlacementDecisionRecord, RuntimeError> {
if requested_target == PlacementTarget::Local {
return Ok(PlacementDecisionRecord {
requested_target,
selected_target: Some(PlacementTarget::Local),
status: PlacementDecisionStatus::Selected,
reason: PlacementDecisionReason::RequestedTargetSelected,
supported_executor_targets: supported_executor_targets(),
});
}
Err(runtime_error(
RuntimeErrorCode::PlacementUnsupported,
"requested placement target is not supported by the available executor set",
json!({
"requested_target": requested_target,
"supported_executor_targets": supported_executor_targets(),
}),
))
}
fn sanitized_request(mut request: RuntimeRequest) -> RuntimeRequest {
if let Some(token) = request.context.caller.take() {
if let Some(identity) = derive_identity_from_jwt(&token) {
request.context.identity = Some(identity);
} else {
request.context.caller = Some(token);
}
}
request
}
fn load_artifact_bytes_for_verification(
selected: &ResolvedCapability,
) -> Result<Vec<u8>, RuntimeError> {
let Some(binary) = selected.artifact.binary.as_ref() else {
return Ok(Vec::new());
};
match fs::read(&binary.location) {
Ok(bytes) => Ok(bytes),
Err(_) if binary.signature.is_none() => Ok(selected
.artifact
.digests
.binary_digest
.clone()
.unwrap_or_else(|| selected.record.artifact_ref.clone())
.into_bytes()),
Err(error) => Err(runtime_error(
RuntimeErrorCode::ArtifactMissing,
"artifact bytes could not be loaded for signature verification",
json!({
"artifact_ref": selected.record.artifact_ref,
"location": binary.location,
"code": "artifact_load_failed",
"message": error.to_string(),
}),
)),
}
}
fn artifact_verification_runtime_error(error: &ArtifactVerificationFailure) -> RuntimeError {
runtime_error(
RuntimeErrorCode::ContractViolation,
"artifact signature verification failed before execution",
json!({
"code": error.code(),
"artifact_verification": error.record(),
}),
)
}
fn begin_attempt(
request: RuntimeRequest,
observability: RuntimeObservabilityConfig,
) -> (AttemptContext, StateEmitter) {
let request = sanitized_request(request);
let request_id = request.request_id.clone();
let execution_id = format!("{EXECUTION_PREFIX}{request_id}");
let trace_id = format!("{TRACE_PREFIX}{execution_id}");
let mut emitter = StateEmitter::new(&execution_id, &request_id);
emitter.push(
RuntimeState::LoadingRegistry,
RuntimeTransitionReasonCode::RuntimeInitializationStarted,
json!({
"registry_status": "available",
"identity": request.context.identity,
}),
);
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::RegistryLoaded,
json!({"governing_spec": GOVERNING_SPEC}),
);
(
AttemptContext {
request,
execution_id,
trace_id,
observability,
artifact_verification: None,
warnings: Vec::new(),
},
emitter,
)
}
fn invalid_request_outcome(
attempt: AttemptContext,
mut emitter: StateEmitter,
error: RuntimeError,
) -> RuntimeExecutionOutcome {
let placement = placement_not_attempted(
attempt.request.context.requested_target,
PlacementDecisionReason::SelectionNotReached,
);
emitter.push(
RuntimeState::EvaluatingConstraints,
RuntimeTransitionReasonCode::CandidatesCollected,
json!({"candidate_count": 0}),
);
emitter.push(
RuntimeState::Error,
RuntimeTransitionReasonCode::ConstraintValidationFailed,
json!({"code": error.code, "message": error.message}),
);
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Error}),
);
let finished = emitter.finish();
let identity = attempt.request.context.identity.clone();
terminal_failure(FailureContext {
attempt,
state_events: finished.events,
state_transitions: finished.transitions,
state_machine_validation: finished.validation,
candidate_collection: CandidateCollectionRecord {
lookup_scope: RuntimeLookupScope::PreferPrivate,
candidates: Vec::new(),
rejected_candidates: Vec::new(),
},
selection: SelectionRecord {
status: SelectionStatus::InvalidRequest,
selected_capability_id: None,
selected_capability_version: None,
failure_reason: Some(SelectionFailureReason::InvalidRequest),
remaining_candidates: Vec::new(),
},
execution: ExecutionRecord {
placement: placement.clone(),
placement_target: placement.requested_target,
status: ExecutionStatus::NotStarted,
artifact_ref: None,
started_at: None,
completed_at: None,
output_digest: None,
failure_reason: Some(ExecutionFailureReason::ContractInputInvalid),
artifact_verification: None,
identity,
},
error,
emitted_events: Vec::new(),
workflow_evidence: None,
})
}
fn no_eligible_outcome(
attempt: AttemptContext,
mut emitter: StateEmitter,
candidate_collection: CandidateCollectionRecord,
) -> RuntimeExecutionOutcome {
let placement = placement_not_attempted(
attempt.request.context.requested_target,
PlacementDecisionReason::SelectionNotReached,
);
let error = if candidate_collection.rejected_candidates.is_empty() {
runtime_error(
RuntimeErrorCode::CapabilityNotFound,
"no eligible capability matched the runtime request",
json!({"request_id": attempt.request.request_id}),
)
} else {
runtime_error(
RuntimeErrorCode::CapabilityNotRunnable,
"matching capabilities were found but none were runnable locally",
json!({"rejected_candidates": candidate_collection.rejected_candidates}),
)
};
let reason = if candidate_collection.rejected_candidates.is_empty() {
RuntimeTransitionReasonCode::NoMatch
} else {
RuntimeTransitionReasonCode::ConstraintValidationFailed
};
emitter.push(RuntimeState::Error, reason, json!({"code": error.code}));
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Error}),
);
let failure_reason = if error.code == RuntimeErrorCode::CapabilityNotFound {
SelectionFailureReason::NoMatch
} else {
SelectionFailureReason::NotRunnable
};
let finished = emitter.finish();
let identity = attempt.request.context.identity.clone();
terminal_failure(FailureContext {
attempt,
state_events: finished.events,
state_transitions: finished.transitions,
state_machine_validation: finished.validation,
candidate_collection,
selection: SelectionRecord {
status: SelectionStatus::NoMatch,
selected_capability_id: None,
selected_capability_version: None,
failure_reason: Some(failure_reason),
remaining_candidates: Vec::new(),
},
execution: ExecutionRecord {
placement: placement.clone(),
placement_target: placement.requested_target,
status: ExecutionStatus::NotStarted,
artifact_ref: None,
started_at: None,
completed_at: None,
output_digest: None,
failure_reason: Some(ExecutionFailureReason::ArtifactNotRunnable),
artifact_verification: None,
identity,
},
error,
emitted_events: Vec::new(),
workflow_evidence: None,
})
}
fn ambiguous_outcome(
attempt: AttemptContext,
mut emitter: StateEmitter,
resolution: CandidateResolution,
) -> RuntimeExecutionOutcome {
let placement = placement_not_attempted(
attempt.request.context.requested_target,
PlacementDecisionReason::SelectionNotReached,
);
let remaining_candidates = resolution
.eligible
.iter()
.map(|candidate| runtime_candidate(candidate, resolution.candidate_reason))
.collect::<Vec<_>>();
let error = runtime_error(
RuntimeErrorCode::CapabilityAmbiguous,
"runtime request matched more than one eligible capability",
json!({"remaining_candidates": remaining_candidates}),
);
emitter.push(
RuntimeState::Error,
RuntimeTransitionReasonCode::SelectionFailed,
json!({"code": error.code}),
);
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Error}),
);
let finished = emitter.finish();
let identity = attempt.request.context.identity.clone();
terminal_failure(FailureContext {
attempt,
state_events: finished.events,
state_transitions: finished.transitions,
state_machine_validation: finished.validation,
candidate_collection: resolution.collection,
selection: SelectionRecord {
status: SelectionStatus::Ambiguous,
selected_capability_id: None,
selected_capability_version: None,
failure_reason: Some(SelectionFailureReason::Ambiguous),
remaining_candidates,
},
execution: ExecutionRecord {
placement: placement.clone(),
placement_target: placement.requested_target,
status: ExecutionStatus::NotStarted,
artifact_ref: None,
started_at: None,
completed_at: None,
output_digest: None,
failure_reason: Some(ExecutionFailureReason::ArtifactNotRunnable),
artifact_verification: None,
identity,
},
error,
emitted_events: Vec::new(),
workflow_evidence: None,
})
}
fn pre_execution_failure_outcome(
context: ExecutionContext,
failure: PreExecutionFailure,
) -> RuntimeExecutionOutcome {
let ExecutionContext {
attempt,
mut emitter,
candidate_collection,
selection,
} = context;
let reason = if emitter.current_state == RuntimeState::Selecting {
RuntimeTransitionReasonCode::SelectionFailed
} else {
RuntimeTransitionReasonCode::ConstraintValidationFailed
};
emitter.push(
RuntimeState::Error,
reason,
json!({"code": failure.error.code, "details": failure.error.details}),
);
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Error}),
);
let finished = emitter.finish();
let identity = attempt.request.context.identity.clone();
terminal_failure(FailureContext {
attempt,
state_events: finished.events,
state_transitions: finished.transitions,
state_machine_validation: finished.validation,
candidate_collection,
selection,
execution: ExecutionRecord {
placement: failure.placement.clone(),
placement_target: failure
.placement
.selected_target
.unwrap_or(failure.placement.requested_target),
status: ExecutionStatus::NotStarted,
artifact_ref: failure.artifact_ref,
started_at: None,
completed_at: None,
output_digest: None,
failure_reason: Some(failure.failure_reason),
artifact_verification: failure.artifact_verification,
identity,
},
error: failure.error,
emitted_events: Vec::new(),
workflow_evidence: None,
})
}
#[allow(clippy::too_many_arguments)]
fn execution_failure_outcome(
context: ExecutionContext,
failure: ExecutionFailureState,
error: RuntimeError,
emitted_events: Vec<traverse_contracts::EventReference>,
workflow_evidence: Option<WorkflowTraversalEvidence>,
) -> RuntimeExecutionOutcome {
let ExecutionContext {
attempt,
mut emitter,
candidate_collection,
selection,
} = context;
emitter.push(
RuntimeState::Error,
RuntimeTransitionReasonCode::ExecutionFailed,
json!({"code": error.code, "details": error.details}),
);
let completed_at = emitter.next_timestamp();
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Error}),
);
let finished = emitter.finish();
let identity = attempt.request.context.identity.clone();
let artifact_verification = attempt.artifact_verification.clone();
terminal_failure(FailureContext {
attempt,
state_events: finished.events,
state_transitions: finished.transitions,
state_machine_validation: finished.validation,
candidate_collection,
selection,
execution: ExecutionRecord {
placement: failure.placement.clone(),
placement_target: failure
.placement
.selected_target
.unwrap_or(failure.placement.requested_target),
status: ExecutionStatus::Failed,
artifact_ref: Some(failure.artifact_ref),
started_at: Some(failure.started_at),
completed_at: Some(completed_at),
output_digest: None,
failure_reason: Some(failure.failure_reason),
artifact_verification,
identity,
},
error,
emitted_events,
workflow_evidence,
})
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
fn successful_execution_outcome(
context: ExecutionContext,
selected: &ResolvedCapability,
started_execution: StartedExecution,
execution_output: Value,
emitted_events: Vec<traverse_contracts::EventReference>,
workflow_evidence: Option<WorkflowTraversalEvidence>,
) -> RuntimeExecutionOutcome {
let ExecutionContext {
attempt,
mut emitter,
candidate_collection,
selection,
} = context;
let completed_at = emitter.next_timestamp();
let emits_events = selected.record.implementation_kind == ImplementationKind::Workflow
|| !selected.contract.emits.is_empty();
if emits_events {
emitter.push(
RuntimeState::EmittingEvents,
RuntimeTransitionReasonCode::ExecutionSucceededWithEvents,
json!({
"capability_id": selected.record.id,
"capability_version": selected.record.version,
"declared_event_count": selected.contract.emits.len(),
}),
);
emitter.push(
RuntimeState::Completed,
RuntimeTransitionReasonCode::EventsEmitted,
json!({
"capability_id": selected.record.id,
"capability_version": selected.record.version,
}),
);
} else {
emitter.push(
RuntimeState::Completed,
RuntimeTransitionReasonCode::ExecutionSucceeded,
json!({
"capability_id": selected.record.id,
"capability_version": selected.record.version,
}),
);
}
emitter.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed,
json!({"terminal_state": RuntimeState::Completed}),
);
let finished = emitter.finish();
let execution = ExecutionRecord {
placement: started_execution.placement.clone(),
placement_target: started_execution
.placement
.selected_target
.unwrap_or(started_execution.placement.requested_target),
status: ExecutionStatus::Succeeded,
artifact_ref: Some(selected.record.artifact_ref.clone()),
started_at: Some(started_execution.started_at),
completed_at: Some(completed_at),
output_digest: Some(content_digest(&execution_output)),
failure_reason: None,
artifact_verification: attempt.artifact_verification.clone(),
identity: attempt.request.context.identity.clone(),
};
let result_record = TraceResultRecord {
status: RuntimeResultStatus::Completed,
output: Some(execution_output.clone()),
error: None,
warnings: attempt.warnings.clone(),
};
let otel_trace = otel_trace_record(
&attempt,
&finished.transitions,
&selection,
&execution,
&result_record,
);
let trace = RuntimeTrace {
kind: RUNTIME_TRACE_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
trace_id: attempt.trace_id.clone(),
execution_id: attempt.execution_id.clone(),
request_id: attempt.request.request_id.clone(),
governing_spec: GOVERNING_SPEC.to_string(),
request: attempt.request.clone(),
decision_evidence: TraceDecisionEvidence {
candidate_collection: candidate_collection.clone(),
selection: selection.clone(),
model_resolution: Vec::new(),
},
state_progression: TraceStateProgression {
state_events: finished.events.clone(),
transitions: finished.transitions.clone(),
validation: finished.validation.clone(),
},
terminal_outcome: TraceTerminalOutcome {
runtime_status: RuntimeResultStatus::Completed,
execution_status: execution.status,
failure_reason: None,
error: None,
},
emitted_events,
workflow_evidence,
model_resolution: Vec::new(),
state_transitions: finished.transitions.clone(),
state_machine_validation: finished.validation.clone(),
candidate_collection,
selection,
execution,
result: result_record,
otel_trace,
};
let result = RuntimeResult {
kind: RUNTIME_RESULT_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
execution_id: attempt.execution_id,
request_id: attempt.request.request_id,
status: RuntimeResultStatus::Completed,
trace_ref: attempt.trace_id,
output: Some(execution_output),
error: None,
warnings: attempt.warnings,
};
RuntimeExecutionOutcome {
result,
trace,
state_events: finished.events,
}
}
fn validate_request(request: &RuntimeRequest) -> Option<RuntimeError> {
if request.kind != RUNTIME_REQUEST_KIND {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"kind must equal runtime_request",
json!({"path": "$.kind"}),
));
}
if request.schema_version != SUPPORTED_SCHEMA_VERSION {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"schema_version must equal 1.0.0",
json!({"path": "$.schema_version"}),
));
}
if request.governing_spec != GOVERNING_SPEC {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"governing_spec must equal 006-runtime-request-execution",
json!({"path": "$.governing_spec"}),
));
}
if request.request_id.trim().is_empty() {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"request_id must be non-empty",
json!({"path": "$.request_id"}),
));
}
if request.lookup.allow_ambiguity {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"allow_ambiguity must be false in this runtime slice",
json!({"path": "$.lookup.allow_ambiguity"}),
));
}
if request
.intent
.capability_version
.as_deref()
.is_some_and(|version| Version::parse(version).is_err())
{
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"capability_version must be valid semantic versioning",
json!({"path": "$.intent.capability_version"}),
));
}
let exact_id = request
.intent
.capability_id
.as_deref()
.is_some_and(non_empty);
let exact_version = request
.intent
.capability_version
.as_deref()
.is_some_and(non_empty);
let intent_key = request.intent.intent_key.as_deref().is_some_and(non_empty);
if !(exact_id || intent_key) {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"runtime intent must include capability_id or intent_key",
json!({"path": "$.intent"}),
));
}
if exact_version && !exact_id {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"capability_version requires capability_id",
json!({"path": "$.intent.capability_version"}),
));
}
let has_version_range = request
.intent
.version_range
.as_deref()
.is_some_and(non_empty);
if has_version_range && !exact_id {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"version_range requires capability_id",
json!({"path": "$.intent.version_range"}),
));
}
if has_version_range && exact_version {
return Some(runtime_error(
RuntimeErrorCode::RequestInvalid,
"version_range and capability_version are mutually exclusive",
json!({"path": "$.intent.version_range"}),
));
}
None
}
fn is_exact_target(intent: &RuntimeIntent) -> bool {
intent.capability_id.as_deref().is_some_and(non_empty)
&& intent.capability_version.as_deref().is_some_and(non_empty)
}
fn non_empty(value: &str) -> bool {
!value.trim().is_empty()
}
fn map_lookup_scope(scope: RuntimeLookupScope) -> LookupScope {
match scope {
RuntimeLookupScope::PublicOnly => LookupScope::PublicOnly,
RuntimeLookupScope::PreferPrivate => LookupScope::PreferPrivate,
}
}
fn evaluate_candidate(candidate: ResolvedCapability) -> CandidateEvaluation {
if !candidate.contract.lifecycle.is_runtime_eligible() {
return CandidateEvaluation::Rejected(
candidate,
RejectedCandidateReason::LifecycleNotRunnable,
);
}
if candidate.record.implementation_kind == ImplementationKind::Workflow {
if candidate.artifact.workflow_ref.is_some() {
return CandidateEvaluation::Eligible(candidate);
}
return CandidateEvaluation::Rejected(candidate, RejectedCandidateReason::ArtifactMissing);
}
let execution = &candidate.contract.execution;
if !execution
.preferred_targets
.contains(&ExecutionTarget::Local)
|| execution.constraints.host_api_access != HostApiAccess::None
|| execution.constraints.network_access != NetworkAccess::Forbidden
{
return CandidateEvaluation::Rejected(
candidate,
RejectedCandidateReason::NotRunnableLocally,
);
}
match candidate.artifact.binary.as_ref() {
Some(binary) if binary.location.trim().is_empty() => {
CandidateEvaluation::Rejected(candidate, RejectedCandidateReason::ArtifactMissing)
}
None | Some(_) => CandidateEvaluation::Eligible(candidate),
}
}
fn validate_payload_against_contract(
payload: &Value,
schema: &Value,
code: RuntimeErrorCode,
message: &str,
) -> Result<(), RuntimeError> {
let mut errors = Vec::new();
validate_value_against_schema(payload, schema, "$", &mut errors);
if errors.is_empty() {
Ok(())
} else {
Err(runtime_error(
code,
message,
json!({ "violations": errors }),
))
}
}
pub(crate) fn validate_value_against_schema(
value: &Value,
schema: &Value,
path: &str,
errors: &mut Vec<Value>,
) {
let Some(schema_object) = schema.as_object() else {
errors.push(json!({
"path": path,
"message": "schema must be an object"
}));
return;
};
if let Some(schema_type) = schema_object.get("type").and_then(Value::as_str) {
match schema_type {
"object" => {
let Some(instance) = value.as_object() else {
errors.push(type_error(path, "object"));
return;
};
validate_required(instance, schema_object, path, errors);
validate_properties(instance, schema_object, path, errors);
}
"array" => {
let Some(items) = value.as_array() else {
errors.push(type_error(path, "array"));
return;
};
if let Some(item_schema) = schema_object.get("items") {
for (index, item) in items.iter().enumerate() {
validate_value_against_schema(
item,
item_schema,
&format!("{path}[{index}]"),
errors,
);
}
}
}
"string" if !value.is_string() => errors.push(type_error(path, "string")),
"integer" if value.as_i64().is_none() && value.as_u64().is_none() => {
errors.push(type_error(path, "integer"));
}
"number" if !value.is_number() => errors.push(type_error(path, "number")),
"boolean" if !value.is_boolean() => errors.push(type_error(path, "boolean")),
"null" if !value.is_null() => errors.push(type_error(path, "null")),
_ => {}
}
}
}
fn validate_required(
instance: &Map<String, Value>,
schema_object: &Map<String, Value>,
path: &str,
errors: &mut Vec<Value>,
) {
let Some(required) = schema_object.get("required").and_then(Value::as_array) else {
return;
};
for required_field in required.iter().filter_map(Value::as_str) {
if !instance.contains_key(required_field) {
errors.push(json!({
"path": format!("{path}.{required_field}"),
"message": "required property is missing"
}));
}
}
}
fn validate_properties(
instance: &Map<String, Value>,
schema_object: &Map<String, Value>,
path: &str,
errors: &mut Vec<Value>,
) {
let Some(properties) = schema_object.get("properties").and_then(Value::as_object) else {
return;
};
for (key, value) in instance {
if let Some(property_schema) = properties.get(key) {
validate_value_against_schema(value, property_schema, &format!("{path}.{key}"), errors);
}
}
}
fn type_error(path: &str, expected: &str) -> Value {
json!({
"path": path,
"message": format!("expected {expected}")
})
}
fn runtime_candidate(capability: &ResolvedCapability, reason: CandidateReason) -> RuntimeCandidate {
RuntimeCandidate {
scope: map_registry_scope(capability.record.scope),
capability_id: capability.record.id.clone(),
capability_version: capability.record.version.clone(),
artifact_ref: capability.record.artifact_ref.clone(),
implementation_kind: map_implementation_kind(capability.record.implementation_kind),
lifecycle: map_lifecycle(&capability.record.lifecycle),
reason,
}
}
fn map_registry_scope(scope: RegistryScope) -> RuntimeRegistryScope {
match scope {
RegistryScope::Public => RuntimeRegistryScope::Public,
RegistryScope::Private => RuntimeRegistryScope::Private,
}
}
fn map_implementation_kind(kind: ImplementationKind) -> RuntimeImplementationKind {
match kind {
ImplementationKind::Executable => RuntimeImplementationKind::Executable,
ImplementationKind::Workflow => RuntimeImplementationKind::Workflow,
}
}
fn map_lifecycle(lifecycle: &Lifecycle) -> RuntimeLifecycle {
match lifecycle {
Lifecycle::Draft => RuntimeLifecycle::Draft,
Lifecycle::Active => RuntimeLifecycle::Active,
Lifecycle::Deprecated => RuntimeLifecycle::Deprecated,
Lifecycle::Retired => RuntimeLifecycle::Retired,
Lifecycle::Archived => RuntimeLifecycle::Archived,
}
}
fn runtime_error(code: RuntimeErrorCode, message: &str, details: Value) -> RuntimeError {
RuntimeError {
code,
message: message.to_string(),
details,
}
}
fn content_digest(value: &Value) -> String {
let json = value.to_string();
let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
for byte in json.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x0000_0001_0000_01b3);
}
format!("0.1.0:{hash:016x}")
}
fn otel_trace_record(
attempt: &AttemptContext,
state_transitions: &[RuntimeTransitionRecord],
selection: &SelectionRecord,
execution: &ExecutionRecord,
result: &TraceResultRecord,
) -> OTelTraceRecord {
let trace_id = otel_trace_id(attempt);
let root_span_id = otel_span_id(attempt, "runtime.request", 0);
let mut spans = vec![otel_span(OTelSpanInput {
trace_id: &trace_id,
span_id: &root_span_id,
parent_span_id: None,
name: "traverse.runtime.request",
status: span_status(result.status),
started_at: first_transition_time(state_transitions),
ended_at: last_transition_time(state_transitions),
attributes: base_otel_attributes(attempt, selection, execution),
events: error_events(result),
})];
for (index, phase) in otel_phase_names().iter().enumerate() {
spans.push(otel_span(OTelSpanInput {
trace_id: &trace_id,
span_id: &otel_span_id(attempt, phase, index + 1),
parent_span_id: Some(root_span_id.clone()),
name: phase,
status: phase_status(phase, result.status),
started_at: phase_started_at(state_transitions, index),
ended_at: phase_ended_at(state_transitions, index),
attributes: base_otel_attributes(attempt, selection, execution),
events: if result.status == RuntimeResultStatus::Error
&& *phase == "traverse.trace.assembly"
{
error_events(result)
} else {
Vec::new()
},
}));
}
OTelTraceRecord {
trace_id,
parent_traceparent: attempt.request.context.traceparent.clone(),
tracestate: attempt.request.context.tracestate.clone(),
exporter: OTelExporterRecord {
enabled: attempt.observability.exporter.endpoint.is_some(),
endpoint: attempt.observability.exporter.endpoint.clone(),
protocol: attempt.observability.exporter.protocol,
},
spans,
}
}
fn otel_phase_names() -> [&'static str; 5] {
[
"traverse.request.intake",
"traverse.registry.lookup",
"traverse.contract.validation",
"traverse.capability.execution",
"traverse.trace.assembly",
]
}
fn otel_trace_id(attempt: &AttemptContext) -> String {
if attempt.observability.deterministic_ids {
let seed = attempt
.observability
.deterministic_seed
.as_deref()
.unwrap_or("traverse-test");
return deterministic_hex(seed, &attempt.trace_id, 32);
}
deterministic_hex("traverse-runtime", &attempt.trace_id, 32)
}
fn otel_span_id(attempt: &AttemptContext, name: &str, index: usize) -> String {
let seed = attempt
.observability
.deterministic_seed
.as_deref()
.unwrap_or("traverse-runtime");
deterministic_hex(
seed,
&format!("{}:{name}:{index}", attempt.execution_id),
16,
)
}
fn deterministic_hex(seed: &str, value: &str, len: usize) -> String {
let mut hash: u128 = 0x6c62_272e_07bb_0142_62b8_2175_6295_c58d;
for byte in seed.as_bytes().iter().chain(value.as_bytes()) {
hash ^= u128::from(*byte);
hash = hash.wrapping_mul(0x0000_0000_0100_0000_0000_0000_0000_013b);
}
format!("{hash:032x}").chars().take(len).collect()
}
struct OTelSpanInput<'a> {
trace_id: &'a str,
span_id: &'a str,
parent_span_id: Option<String>,
name: &'a str,
status: OTelSpanStatus,
started_at: String,
ended_at: String,
attributes: Vec<OTelAttribute>,
events: Vec<OTelSpanEvent>,
}
fn otel_span(input: OTelSpanInput<'_>) -> OTelSpanRecord {
OTelSpanRecord {
trace_id: input.trace_id.to_string(),
span_id: input.span_id.to_string(),
parent_span_id: input.parent_span_id,
name: input.name.to_string(),
kind: OTelSpanKind::Internal,
status: input.status,
started_at: input.started_at,
ended_at: input.ended_at,
attributes: input.attributes,
events: input.events,
}
}
fn base_otel_attributes(
attempt: &AttemptContext,
selection: &SelectionRecord,
execution: &ExecutionRecord,
) -> Vec<OTelAttribute> {
let mut attributes = vec![
otel_attr("traverse.request.id", json!(attempt.request.request_id)),
otel_attr("traverse.execution.id", json!(attempt.execution_id)),
otel_attr("traverse.lookup.scope", json!(attempt.request.lookup.scope)),
otel_attr(
"traverse.runtime.placement.target",
json!(execution.placement_target),
),
];
if let Some(correlation_id) = &attempt.request.context.correlation_id {
attributes.push(otel_attr("traverse.correlation.id", json!(correlation_id)));
}
if let Some(capability_id) = &selection.selected_capability_id {
attributes.push(otel_attr("traverse.capability.id", json!(capability_id)));
}
if let Some(capability_version) = &selection.selected_capability_version {
attributes.push(otel_attr(
"traverse.capability.version",
json!(capability_version),
));
}
attributes
}
fn otel_attr(key: &str, value: Value) -> OTelAttribute {
OTelAttribute {
key: key.to_string(),
value,
}
}
fn error_events(result: &TraceResultRecord) -> Vec<OTelSpanEvent> {
result
.error
.as_ref()
.map(|error| {
vec![OTelSpanEvent {
name: "exception".to_string(),
timestamp: "1970-01-01T00:00:00Z".to_string(),
attributes: vec![
otel_attr("traverse.error.classification", json!(error.code)),
otel_attr("traverse.error.message", json!(error.message)),
],
}]
})
.unwrap_or_default()
}
fn span_status(status: RuntimeResultStatus) -> OTelSpanStatus {
match status {
RuntimeResultStatus::Completed => OTelSpanStatus::Ok,
RuntimeResultStatus::Error => OTelSpanStatus::Error,
}
}
fn phase_status(phase: &str, status: RuntimeResultStatus) -> OTelSpanStatus {
if status == RuntimeResultStatus::Error && phase == "traverse.trace.assembly" {
OTelSpanStatus::Error
} else {
OTelSpanStatus::Ok
}
}
fn first_transition_time(transitions: &[RuntimeTransitionRecord]) -> String {
transitions.first().map_or_else(
|| "1970-01-01T00:00:00Z".to_string(),
|transition| transition.occurred_at.clone(),
)
}
fn last_transition_time(transitions: &[RuntimeTransitionRecord]) -> String {
transitions.last().map_or_else(
|| "1970-01-01T00:00:00Z".to_string(),
|transition| transition.occurred_at.clone(),
)
}
fn phase_started_at(transitions: &[RuntimeTransitionRecord], index: usize) -> String {
transitions.get(index).map_or_else(
|| first_transition_time(transitions),
|transition| transition.occurred_at.clone(),
)
}
fn phase_ended_at(transitions: &[RuntimeTransitionRecord], index: usize) -> String {
transitions.get(index + 1).map_or_else(
|| last_transition_time(transitions),
|transition| transition.occurred_at.clone(),
)
}
fn contains_drafts_segment(path: &str) -> bool {
path.replace('\\', "/")
.split('/')
.any(|segment| segment == "drafts")
}
struct AttemptContext {
request: RuntimeRequest,
execution_id: String,
trace_id: String,
observability: RuntimeObservabilityConfig,
artifact_verification: Option<ArtifactVerificationRecord>,
warnings: Vec<RuntimeWarning>,
}
struct CandidateResolution {
eligible: Vec<ResolvedCapability>,
collection: CandidateCollectionRecord,
candidate_reason: CandidateReason,
}
struct FailureContext {
attempt: AttemptContext,
state_events: Vec<RuntimeStateEvent>,
state_transitions: Vec<RuntimeTransitionRecord>,
state_machine_validation: RuntimeStateMachineValidationEvidence,
candidate_collection: CandidateCollectionRecord,
selection: SelectionRecord,
execution: ExecutionRecord,
error: RuntimeError,
emitted_events: Vec<traverse_contracts::EventReference>,
workflow_evidence: Option<WorkflowTraversalEvidence>,
}
struct ExecutionFailureState {
artifact_ref: String,
started_at: String,
placement: PlacementDecisionRecord,
failure_reason: ExecutionFailureReason,
}
struct ExecutionContext {
attempt: AttemptContext,
emitter: StateEmitter,
candidate_collection: CandidateCollectionRecord,
selection: SelectionRecord,
}
struct StartedExecution {
started_at: String,
placement: PlacementDecisionRecord,
}
struct PreExecutionFailure {
artifact_ref: Option<String>,
failure_reason: ExecutionFailureReason,
placement: PlacementDecisionRecord,
error: RuntimeError,
artifact_verification: Option<ArtifactVerificationRecord>,
}
enum CandidateEvaluation {
Eligible(ResolvedCapability),
Rejected(ResolvedCapability, RejectedCandidateReason),
}
struct StateEmitter {
execution_id: String,
request_id: String,
next_second: u32,
next_event_index: u32,
current_state: RuntimeState,
events: Vec<RuntimeStateEvent>,
transitions: Vec<RuntimeTransitionRecord>,
violations: Vec<Value>,
}
struct FinishedStateMachineArtifacts {
events: Vec<RuntimeStateEvent>,
transitions: Vec<RuntimeTransitionRecord>,
validation: RuntimeStateMachineValidationEvidence,
}
fn start_selected_execution(
emitter: &mut StateEmitter,
selected: &ResolvedCapability,
placement: PlacementDecisionRecord,
identity: Option<&RuntimeIdentity>,
) -> StartedExecution {
let started_at = emitter.next_timestamp();
emitter.push(
RuntimeState::Executing,
RuntimeTransitionReasonCode::CandidateSelected,
json!({
"capability_id": selected.record.id,
"capability_version": selected.record.version,
"artifact_ref": selected.record.artifact_ref,
"requested_target": placement.requested_target,
"selected_target": placement.selected_target,
"placement_status": placement.status,
"placement_reason": placement.reason,
"identity": identity,
}),
);
StartedExecution {
started_at,
placement,
}
}
impl StateEmitter {
fn new(execution_id: &str, request_id: &str) -> Self {
Self {
execution_id: execution_id.to_string(),
request_id: request_id.to_string(),
next_second: 0,
next_event_index: 0,
current_state: RuntimeState::Idle,
events: Vec::new(),
transitions: Vec::new(),
violations: Vec::new(),
}
}
fn push(&mut self, state: RuntimeState, reason: RuntimeTransitionReasonCode, details: Value) {
let transitioned = self.try_push(state, reason, details);
debug_assert!(transitioned, "runtime state transition must be spec-valid");
}
fn try_push(
&mut self,
state: RuntimeState,
reason: RuntimeTransitionReasonCode,
details: Value,
) -> bool {
let from_state = self.current_state;
if !is_allowed_transition(from_state, state, reason) {
self.violations.push(json!({
"from_state": from_state,
"to_state": state,
"reason_code": reason,
"message": "unexpected runtime state transition"
}));
return false;
}
let entered_at = self.next_timestamp();
let mut event_details = detail_object(details);
event_details.insert(
"transition_reason".to_string(),
serde_json::to_value(reason)
.unwrap_or_else(|_| Value::String("serialization_failed".to_string())),
);
let event = RuntimeStateEvent {
kind: RUNTIME_STATE_EVENT_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
event_id: format!("rse_{}_{:04}", self.execution_id, self.next_event_index),
execution_id: self.execution_id.clone(),
request_id: self.request_id.clone(),
state,
entered_at: entered_at.clone(),
details: Value::Object(event_details.clone()),
};
self.next_event_index += 1;
self.events.push(event);
self.transitions.push(RuntimeTransitionRecord {
from_state,
to_state: state,
reason_code: reason,
occurred_at: entered_at,
request_id: Some(self.request_id.clone()),
execution_id: Some(self.execution_id.clone()),
details: Some(Value::Object(event_details)),
});
self.current_state = state;
true
}
fn next_timestamp(&mut self) -> String {
let timestamp = format!("1970-01-01T00:00:{:02}Z", self.next_second);
self.next_second += 1;
timestamp
}
fn finish(self) -> FinishedStateMachineArtifacts {
let checked_states = vec![
RuntimeState::Idle,
RuntimeState::LoadingRegistry,
RuntimeState::Ready,
RuntimeState::Discovering,
RuntimeState::EvaluatingConstraints,
RuntimeState::Selecting,
RuntimeState::Executing,
RuntimeState::EmittingEvents,
RuntimeState::Completed,
RuntimeState::Error,
];
let checked_transitions = self
.transitions
.iter()
.map(|transition| {
format!(
"{}->{}",
runtime_state_name(transition.from_state),
runtime_state_name(transition.to_state)
)
})
.collect();
let validation = RuntimeStateMachineValidationEvidence {
kind: RUNTIME_STATE_MACHINE_VALIDATION_KIND.to_string(),
schema_version: SUPPORTED_SCHEMA_VERSION.to_string(),
governing_spec: STATE_MACHINE_GOVERNING_SPEC.to_string(),
validated_at: format!(
"1970-01-01T00:00:{:02}Z",
self.next_second.saturating_sub(1)
),
status: if self.violations.is_empty() {
RuntimeStateMachineValidationStatus::Passed
} else {
RuntimeStateMachineValidationStatus::Failed
},
checked_states,
checked_transitions,
violations: self.violations,
};
FinishedStateMachineArtifacts {
events: self.events,
transitions: self.transitions,
validation,
}
}
}
fn is_allowed_transition(
from: RuntimeState,
to: RuntimeState,
reason: RuntimeTransitionReasonCode,
) -> bool {
matches!(
(from, to, reason),
(
RuntimeState::Idle,
RuntimeState::LoadingRegistry,
RuntimeTransitionReasonCode::RuntimeInitializationStarted
) | (
RuntimeState::LoadingRegistry,
RuntimeState::Ready,
RuntimeTransitionReasonCode::RegistryLoaded
) | (
RuntimeState::LoadingRegistry,
RuntimeState::Error,
RuntimeTransitionReasonCode::RegistryLoadFailed
) | (
RuntimeState::Ready,
RuntimeState::Discovering,
RuntimeTransitionReasonCode::RequestStarted
) | (
RuntimeState::Discovering,
RuntimeState::EvaluatingConstraints,
RuntimeTransitionReasonCode::CandidatesCollected
) | (
RuntimeState::Discovering,
RuntimeState::Error,
RuntimeTransitionReasonCode::NoMatch
) | (
RuntimeState::EvaluatingConstraints,
RuntimeState::Selecting,
RuntimeTransitionReasonCode::ConstraintsEvaluated
) | (
RuntimeState::EvaluatingConstraints,
RuntimeState::Error,
RuntimeTransitionReasonCode::ConstraintValidationFailed
) | (
RuntimeState::Selecting,
RuntimeState::Executing,
RuntimeTransitionReasonCode::CandidateSelected
) | (
RuntimeState::Selecting,
RuntimeState::Error,
RuntimeTransitionReasonCode::SelectionFailed
) | (
RuntimeState::Executing,
RuntimeState::EmittingEvents,
RuntimeTransitionReasonCode::ExecutionSucceededWithEvents
) | (
RuntimeState::Executing,
RuntimeState::Completed,
RuntimeTransitionReasonCode::ExecutionSucceeded
) | (
RuntimeState::Executing,
RuntimeState::Error,
RuntimeTransitionReasonCode::ExecutionFailed
) | (
RuntimeState::EmittingEvents,
RuntimeState::Completed,
RuntimeTransitionReasonCode::EventsEmitted
) | (
RuntimeState::EmittingEvents,
RuntimeState::Error,
RuntimeTransitionReasonCode::EventEmissionFailed
) | (
RuntimeState::Completed | RuntimeState::Error,
RuntimeState::Ready,
RuntimeTransitionReasonCode::ExecutionClosed
)
)
}
fn detail_object(details: Value) -> Map<String, Value> {
match details {
Value::Object(map) => map,
other => {
let mut map = Map::new();
map.insert("value".to_string(), other);
map
}
}
}
fn runtime_state_name(state: RuntimeState) -> &'static str {
match state {
RuntimeState::Idle => "idle",
RuntimeState::LoadingRegistry => "loading_registry",
RuntimeState::Ready => "ready",
RuntimeState::Discovering => "discovering",
RuntimeState::EvaluatingConstraints => "evaluating_constraints",
RuntimeState::Selecting => "selecting",
RuntimeState::Executing => "executing",
RuntimeState::EmittingEvents => "emitting_events",
RuntimeState::Completed => "completed",
RuntimeState::Error => "error",
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::expect_used)]
use std::fmt::Write as _;
use super::security::{
ArtifactVerificationFailure, ArtifactVerificationScheme, ArtifactVerificationStatus,
RuntimeSecurityConfig, derive_identity_from_jwt, verify_artifact,
};
use super::{
BrowserRuntimeSubscriptionErrorCode, BrowserRuntimeSubscriptionMessage,
BrowserRuntimeSubscriptionRequest, CandidateEvaluation, CandidateReason, LocalExecutor,
PlacementTarget, RejectedCandidateReason, Runtime, RuntimeContext, RuntimeIntent,
RuntimeLookup, RuntimeLookupScope, RuntimeLookupScope::*, RuntimeRequest,
RuntimeResultStatus, RuntimeState, RuntimeTransitionReasonCode,
browser_subscription_messages, evaluate_candidate, map_implementation_kind, map_lifecycle,
map_registry_scope, parse_runtime_request, runtime_candidate, subscription_targets_outcome,
validate_browser_subscription_request, validate_payload_against_contract, validate_request,
};
use ed25519_dalek::{Signer, SigningKey};
use serde_json::json;
use sha2::{Digest, Sha256};
use std::collections::BTreeMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use traverse_contracts::{
BinaryFormat as ContractBinaryFormat, Entrypoint, EntrypointKind, Execution,
ExecutionConstraints, ExecutionTarget, FilesystemAccess, HostApiAccess, Lifecycle,
NetworkAccess, Owner, Provenance, ProvenanceSource, SchemaContainer, ServiceType,
};
use traverse_registry::{
ArtifactDigests, ArtifactSignature, ArtifactSignatureScheme, BinaryFormat, BinaryReference,
CapabilityArtifactRecord, CapabilityRegistration, CapabilityRegistry,
CapabilityRegistryRecord, ComposabilityMetadata, CompositionKind, CompositionPattern,
DiscoveryIndexEntry, ImplementationKind, ModelCandidateReadiness,
ModelCandidateRejectionCode, ModelResolutionEvidence, ModelResolutionPhase,
RegistryProvenance, RegistryScope, ResolvedCapability, SelectedModelCandidate, SourceKind,
SourceReference, WorkspaceAppStateErrorCode,
};
const HEX_TABLE: &[u8; 16] = b"0123456789abcdef";
static TEMP_COUNTER: AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Default)]
struct RecordingEventSink {
events: Mutex<Vec<super::events::TraverseEvent>>,
}
impl super::events::RuntimeEventSink for RecordingEventSink {
fn emit(
&self,
event: super::events::TraverseEvent,
) -> Result<(), super::events::EventError> {
self.events
.lock()
.expect("recording sink lock must not be poisoned")
.push(event);
Ok(())
}
}
#[derive(Debug)]
struct FailingEventSink;
impl super::events::RuntimeEventSink for FailingEventSink {
fn emit(
&self,
_event: super::events::TraverseEvent,
) -> Result<(), super::events::EventError> {
Err(super::events::EventError::JournalWrite(
"sink unavailable".to_string(),
))
}
}
#[test]
fn missing_binary_metadata_is_eligible_for_native_host_execution() {
let capability = resolved_capability(None, Lifecycle::Active);
let evaluation = evaluate_candidate(capability);
assert!(matches!(evaluation, CandidateEvaluation::Eligible(_)));
}
#[test]
fn live_wiring_helpers_support_native_and_wasm_artifact_types() {
use super::executor::ArtifactType;
use super::placement::PlacementError;
use super::router::RouterError;
use super::{ExecutionFailureReason, RuntimeErrorCode};
let native = resolved_capability(None, Lifecycle::Active);
assert_eq!(super::artifact_type_for(&native), ArtifactType::Native);
assert_eq!(
super::executor_capability_for(&native, ArtifactType::Native).artifact_type,
ArtifactType::Native
);
let wasm = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
assert_eq!(super::artifact_type_for(&wasm), ArtifactType::Wasm);
let (code, reason, _) =
super::map_router_error(&RouterError::ExecutorNotFound("Native".to_string()));
assert_eq!(code.code, RuntimeErrorCode::CapabilityNotRunnable);
assert_eq!(reason, ExecutionFailureReason::ArtifactNotRunnable);
let (code, reason, _) = super::map_router_error(&RouterError::TraceLockPoisoned);
assert_eq!(code.code, RuntimeErrorCode::ExecutionFailed);
assert_eq!(reason, ExecutionFailureReason::ExecutionFailed);
let (code, reason, _) = super::map_router_error(&RouterError::PlacementFailed(
PlacementError::NoEligibleTarget,
));
assert_eq!(code.code, RuntimeErrorCode::PlacementUnsupported);
assert_eq!(reason, ExecutionFailureReason::PlacementUnsupported);
let (code, reason, _) =
super::map_router_error(&RouterError::ExecutionFailed("boom".to_string()));
assert_eq!(code.code, RuntimeErrorCode::ExecutionFailed);
assert_eq!(reason, ExecutionFailureReason::ExecutionFailed);
let (code, reason, _) = super::map_router_error(&RouterError::ContractViolation(vec![
traverse_contracts::ViolationRecord::new(
"undeclared_event_emission",
"test.cap",
"test message",
),
]));
assert_eq!(code.code, RuntimeErrorCode::ContractViolation);
assert_eq!(reason, ExecutionFailureReason::ExecutionFailed);
}
#[test]
fn live_wiring_runtime_surfaces_and_fallback_broker_are_covered() {
use super::PlacementTarget;
use super::events::{EventError, LifecycleStatus, TraverseEvent};
use traverse_contracts::ExecutionTarget;
let runtime = Runtime::new(CapabilityRegistry::new(), NoopExecutor);
let _ = runtime.clone();
let debug = format!("{runtime:?}");
assert!(debug.contains("Runtime"));
assert!(Arc::ptr_eq(
&runtime.event_broker(),
&runtime.event_broker()
));
assert!(Arc::ptr_eq(&runtime.trace_store(), &runtime.trace_store()));
let mapped = [
(PlacementTarget::Local, ExecutionTarget::Local),
(PlacementTarget::Browser, ExecutionTarget::Browser),
(PlacementTarget::Edge, ExecutionTarget::Edge),
(PlacementTarget::Cloud, ExecutionTarget::Cloud),
(PlacementTarget::Worker, ExecutionTarget::Worker),
(PlacementTarget::Device, ExecutionTarget::Device),
];
for (placement, expected) in mapped {
assert_eq!(super::execution_target_from_placement(placement), expected);
}
let discard = super::event_broker_or_discard(Err(EventError::InvalidRetentionWindow(
"forced".to_string(),
)));
assert!(
discard
.publish(TraverseEvent {
id: "evt".to_string(),
source: "test".to_string(),
event_type: "dev.traverse.discard".to_string(),
datacontenttype: "application/json".to_string(),
time: "2026-08-06T00:00:00Z".to_string(),
data: json!({}),
owner: "test".to_string(),
version: "1.0.0".to_string(),
lifecycle_status: LifecycleStatus::Active,
deduplication_id: None,
ordering_scope: None,
correlation_id: None,
causation_id: None,
subject_id: None,
actor_id: None,
})
.is_ok()
);
assert!(discard.subscribe("dev.traverse.discard", "0").is_err());
assert!(
discard
.subscribe_for_subject("dev.traverse.discard", "0", None)
.is_err()
);
assert!(discard.poll("missing", 1).is_err());
assert!(discard.cancel("missing").is_err());
}
fn run_native_executor_through_placement_router<E: super::LocalExecutor + 'static>(
executor: E,
) -> (
Result<super::router::RouterResponse, super::router::RouterError>,
Arc<super::events::InProcessBroker>,
String,
) {
use super::events::{
EventBroker, EventCatalog, EventCatalogEntry, InProcessBroker, LifecycleStatus,
};
use super::executor::ArtifactType;
use super::placement::PlacementConstraintEvaluator;
use super::router::{CapabilityExecutorRegistry, PlacementRouter, RouterRequest};
use super::trace::TraceStore;
let event_type = "dev.traverse.native.live-emitted";
let catalog = Arc::new(EventCatalog::new());
catalog
.register(EventCatalogEntry {
event_type: event_type.to_string(),
owner: "native.live".to_string(),
version: "1.0.0".to_string(),
lifecycle_status: LifecycleStatus::Active,
consumer_count: 0,
})
.expect("catalog entry should register");
let broker = Arc::new(InProcessBroker::new(catalog).expect("broker should construct"));
let subscription = broker
.subscribe(event_type, "0")
.expect("subscribe should succeed");
let trace_store = Arc::new(Mutex::new(TraceStore::new()));
let mut selected = resolved_capability(None, Lifecycle::Active);
selected.contract.service_type = ServiceType::Subscribable;
selected.contract.event_trigger = Some("dev.traverse.native.trigger".to_string());
selected.contract.emits = vec![traverse_contracts::EventReference {
event_id: event_type.to_string(),
version: "1.0.0".to_string(),
}];
selected.contract.permitted_targets = vec![ExecutionTarget::Local, ExecutionTarget::Cloud];
let mut registry = CapabilityExecutorRegistry::new();
registry.insert(
ArtifactType::Native,
Box::new(super::BoundLocalExecutor {
executor: Arc::new(executor),
selected: selected.clone(),
}),
);
let router = PlacementRouter::new(
PlacementConstraintEvaluator,
registry,
Arc::clone(&trace_store),
broker.clone(),
);
let response = router.execute(RouterRequest {
capability_id: selected.record.id.clone(),
artifact_type: ArtifactType::Native,
contract: selected.contract.clone(),
target_hint: Some(ExecutionTarget::Local),
runtime_snapshot: super::idle_runtime_snapshot(),
input: json!({}),
executor_capability: super::executor_capability_for(&selected, ArtifactType::Native),
trace_id_override: Some("trace_native_live".to_string()),
});
(response, broker, subscription.subscription_id)
}
fn native_traverse_event(event_type: &str) -> super::events::TraverseEvent {
super::events::TraverseEvent {
id: "native-event-1".to_string(),
source: "traverse-runtime/test.native".to_string(),
event_type: event_type.to_string(),
datacontenttype: "application/json".to_string(),
time: "2026-01-01T00:00:00Z".to_string(),
data: json!({}),
owner: "test.native".to_string(),
version: "1.0.0".to_string(),
lifecycle_status: super::events::LifecycleStatus::Active,
deduplication_id: Some("native-event-1".to_string()),
ordering_scope: Some("test.native".to_string()),
correlation_id: None,
causation_id: None,
subject_id: None,
actor_id: None,
}
}
#[test]
fn bound_local_executor_publishes_declared_native_events_through_placement_router() {
use super::events::EventBroker;
use serde_json::Value;
struct NativeEmitExecutor;
impl super::LocalExecutor for NativeEmitExecutor {
fn execute(
&self,
_capability: &ResolvedCapability,
_input: &Value,
) -> Result<super::LocalExecutionOutput, super::LocalExecutionFailure> {
Ok(super::LocalExecutionOutput {
value: json!({ "draft_id": "native-1" }),
emitted_events: vec![native_traverse_event("dev.traverse.native.live-emitted")],
})
}
}
let (response, broker, subscription_id) =
run_native_executor_through_placement_router(NativeEmitExecutor);
let response = response.expect("native placement router execution should succeed");
assert_eq!(response.trace_id, "trace_native_live");
assert_eq!(response.emitted_events.len(), 1);
let poll = broker
.poll(&subscription_id, 10)
.expect("poll should succeed");
assert_eq!(poll.events.len(), 1);
assert_eq!(
poll.events[0].event.event_type,
"dev.traverse.native.live-emitted"
);
}
#[test]
fn bound_local_executor_rejects_undeclared_native_event() {
use super::events::EventBroker;
use serde_json::Value;
struct UndeclaredEmitExecutor;
impl super::LocalExecutor for UndeclaredEmitExecutor {
fn execute(
&self,
_capability: &ResolvedCapability,
_input: &Value,
) -> Result<super::LocalExecutionOutput, super::LocalExecutionFailure> {
Ok(super::LocalExecutionOutput {
value: json!({ "draft_id": "native-1" }),
emitted_events: vec![native_traverse_event("dev.traverse.native.undeclared")],
})
}
}
let (response, broker, subscription_id) =
run_native_executor_through_placement_router(UndeclaredEmitExecutor);
assert!(response.is_err());
let poll = broker
.poll(&subscription_id, 10)
.expect("poll should succeed");
assert!(poll.events.is_empty());
}
#[test]
fn bound_local_executor_rejects_native_event_from_non_subscribable_capability() {
use serde_json::Value;
struct NonSubscribableEmitExecutor;
impl super::LocalExecutor for NonSubscribableEmitExecutor {
fn execute(
&self,
_capability: &ResolvedCapability,
_input: &Value,
) -> Result<super::LocalExecutionOutput, super::LocalExecutionFailure> {
Ok(super::LocalExecutionOutput {
value: json!({ "draft_id": "native-1" }),
emitted_events: vec![native_traverse_event("dev.traverse.native.live-emitted")],
})
}
}
use super::events::{
EventBroker, EventCatalog, EventCatalogEntry, InProcessBroker, LifecycleStatus,
};
use super::executor::ArtifactType;
use super::placement::PlacementConstraintEvaluator;
use super::router::{CapabilityExecutorRegistry, PlacementRouter, RouterRequest};
use super::trace::TraceStore;
let event_type = "dev.traverse.native.live-emitted";
let catalog = Arc::new(EventCatalog::new());
catalog
.register(EventCatalogEntry {
event_type: event_type.to_string(),
owner: "native.live".to_string(),
version: "1.0.0".to_string(),
lifecycle_status: LifecycleStatus::Active,
consumer_count: 0,
})
.expect("catalog entry should register");
let broker = Arc::new(InProcessBroker::new(catalog).expect("broker should construct"));
let subscription = broker
.subscribe(event_type, "0")
.expect("subscribe should succeed");
let trace_store = Arc::new(Mutex::new(TraceStore::new()));
let mut selected = resolved_capability(None, Lifecycle::Active);
selected.contract.service_type = ServiceType::Stateless;
selected.contract.emits = vec![traverse_contracts::EventReference {
event_id: event_type.to_string(),
version: "1.0.0".to_string(),
}];
selected.contract.permitted_targets = vec![ExecutionTarget::Local, ExecutionTarget::Cloud];
let mut registry = CapabilityExecutorRegistry::new();
registry.insert(
ArtifactType::Native,
Box::new(super::BoundLocalExecutor {
executor: Arc::new(NonSubscribableEmitExecutor),
selected: selected.clone(),
}),
);
let router = PlacementRouter::new(
PlacementConstraintEvaluator,
registry,
Arc::clone(&trace_store),
broker.clone(),
);
let response = router.execute(RouterRequest {
capability_id: selected.record.id.clone(),
artifact_type: ArtifactType::Native,
contract: selected.contract.clone(),
target_hint: Some(ExecutionTarget::Local),
runtime_snapshot: super::idle_runtime_snapshot(),
input: json!({}),
executor_capability: super::executor_capability_for(&selected, ArtifactType::Native),
trace_id_override: Some("trace_native_live_non_subscribable".to_string()),
});
assert!(response.is_err());
let poll = broker
.poll(&subscription.subscription_id, 10)
.expect("poll should succeed");
assert!(poll.events.is_empty());
}
#[test]
fn invalid_json_request_reports_parse_error_text() {
let error = parse_runtime_request("{invalid").err();
assert!(error.is_some());
let message = error.map(|item| item.to_string()).unwrap_or_default();
assert!(!message.is_empty());
}
#[test]
fn runtime_loads_durable_workspace_app_state() {
let workspace_root = unique_workspace_state_dir();
write_runtime_workspace_app_state_fixture(&workspace_root, "local");
let runtime = Runtime::from_workspace_app_state(
&workspace_root,
"local",
NoopExecutor,
"test-runtime",
)
.expect("workspace app state should load");
assert!(
runtime
.capability_registry()
.find_exact(
traverse_registry::LookupScope::PreferPrivate,
"expedition.planning.validate-team-readiness",
"1.0.0"
)
.is_some()
);
assert!(
runtime
.workflow_registry()
.find_exact(
traverse_registry::LookupScope::PreferPrivate,
"expedition.planning.plan-expedition",
"1.0.0"
)
.is_some()
);
assert_eq!(
runtime.workspace_applications()[0].model_dependencies[0].interface_id,
"traverse.inference.generate"
);
}
#[test]
fn governed_model_execution_resolves_from_loaded_app_declaration() {
let workspace_root = unique_workspace_state_dir();
write_runtime_workspace_app_state_fixture(&workspace_root, "local");
let runtime = Runtime::from_workspace_app_state(
&workspace_root,
"local",
NoopExecutor,
"test-runtime",
)
.expect("workspace app state should load");
let mut provider_configs = BTreeMap::new();
provider_configs.insert(
"ollama.local.generate".to_string(),
crate::inference::OllamaProviderConfig {
base_url: "http://127.0.0.1:9".to_string(),
request_timeout_ms: Some(50),
max_response_bytes: None,
},
);
let error = runtime
.execute_governed_model_dependency(
"expedition.readiness",
"1.0.0",
&crate::inference::GovernedModelExecutionRequest {
interface_id: "traverse.inference.generate".to_string(),
prompt: "Summarize readiness.".to_string(),
system_prompt: None,
options: json!({}),
requested_placement: ExecutionTarget::Local,
provider_configs,
},
)
.expect_err("unavailable local provider should fail before output");
assert_eq!(
error.code,
crate::inference::GovernedModelExecutionErrorCode::ModelDependencyUnsatisfied
);
let evidence = error
.model_resolution
.expect("failed model execution should include resolution evidence");
assert_eq!(
evidence.requested_interface_id,
"traverse.inference.generate"
);
assert_eq!(
evidence.machine_failure_code(),
Some("model_dependency_unsatisfied")
);
}
#[test]
fn governed_model_execution_rejects_missing_app_or_interface() {
let workspace_root = unique_workspace_state_dir();
write_runtime_workspace_app_state_fixture(&workspace_root, "local");
let runtime = Runtime::from_workspace_app_state(
&workspace_root,
"local",
NoopExecutor,
"test-runtime",
)
.expect("workspace app state should load");
let request = crate::inference::GovernedModelExecutionRequest {
interface_id: "traverse.inference.embed".to_string(),
prompt: "Summarize readiness.".to_string(),
system_prompt: None,
options: json!({}),
requested_placement: ExecutionTarget::Local,
provider_configs: BTreeMap::new(),
};
let missing_app = runtime
.execute_governed_model_dependency("missing.app", "1.0.0", &request)
.expect_err("unknown app should fail");
assert_eq!(
missing_app.code,
crate::inference::GovernedModelExecutionErrorCode::InterfaceNotDeclared
);
let missing_interface = runtime
.execute_governed_model_dependency("expedition.readiness", "1.0.0", &request)
.expect_err("undeclared interface should fail");
assert_eq!(
missing_interface.code,
crate::inference::GovernedModelExecutionErrorCode::InterfaceNotDeclared
);
}
#[test]
fn runtime_reports_missing_workspace_app_state() {
let workspace_root = unique_workspace_state_dir();
let failure = Runtime::from_workspace_app_state(
&workspace_root,
"local",
NoopExecutor,
"test-runtime",
)
.expect_err("missing workspace app state should fail");
assert_eq!(
failure.errors[0].code,
WorkspaceAppStateErrorCode::MissingWorkspaceState
);
}
#[test]
fn request_validation_rejects_all_invalid_request_guards() {
let mut request = valid_request();
request.kind = "wrong".to_string();
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.schema_version = "9.9.9".to_string();
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.governing_spec = "wrong-spec".to_string();
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.request_id.clear();
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.lookup.allow_ambiguity = true;
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.context.requested_target = PlacementTarget::Local;
request.intent.capability_version = Some("bad".to_string());
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.intent.capability_id = None;
request.intent.intent_key = None;
request.intent.capability_version = None;
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
let mut request = valid_request();
request.intent.capability_id = None;
request.intent.capability_version = Some("1.0.0".to_string());
assert_eq!(
validate_request(&request).map(|error| error.code),
Some(super::RuntimeErrorCode::RequestInvalid)
);
}
#[test]
fn candidate_evaluation_covers_local_runnability_branches() {
let mut capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
capability.record.implementation_kind = ImplementationKind::Workflow;
assert!(matches!(
evaluate_candidate(capability.clone()),
CandidateEvaluation::Rejected(_, RejectedCandidateReason::ArtifactMissing)
));
capability.artifact.workflow_ref = Some(traverse_registry::WorkflowReference {
workflow_id: "workflow".to_string(),
workflow_version: "1.0.0".to_string(),
});
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Eligible(_)
));
let capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: String::new(),
signature: None,
}),
Lifecycle::Active,
);
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Rejected(_, RejectedCandidateReason::ArtifactMissing)
));
let mut capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
capability.contract.execution.preferred_targets = vec![ExecutionTarget::Cloud];
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Rejected(_, RejectedCandidateReason::NotRunnableLocally)
));
let mut capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
capability.contract.execution.constraints.host_api_access =
HostApiAccess::ExceptionRequired;
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Rejected(_, RejectedCandidateReason::NotRunnableLocally)
));
let mut capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
capability.contract.execution.constraints.network_access = NetworkAccess::Required;
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Rejected(_, RejectedCandidateReason::NotRunnableLocally)
));
let capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
assert!(matches!(
evaluate_candidate(capability),
CandidateEvaluation::Eligible(_)
));
}
#[test]
fn payload_validation_covers_schema_branches() {
let invalid_schema = validate_payload_against_contract(
&json!({"field": "value"}),
&json!("bad-schema"),
super::RuntimeErrorCode::RequestInvalid,
"invalid schema",
);
assert!(invalid_schema.is_err());
let wrong_object = validate_payload_against_contract(
&json!("value"),
&json!({"type": "object"}),
super::RuntimeErrorCode::RequestInvalid,
"wrong object",
);
assert!(wrong_object.is_err());
let wrong_array = validate_payload_against_contract(
&json!("value"),
&json!({"type": "array"}),
super::RuntimeErrorCode::RequestInvalid,
"wrong array",
);
assert!(wrong_array.is_err());
let typed_array = validate_payload_against_contract(
&json!(["value", 2]),
&json!({"type": "array", "items": {"type": "string"}}),
super::RuntimeErrorCode::RequestInvalid,
"typed array",
);
assert!(typed_array.is_err());
for (value, schema) in [
(json!("value"), json!({"type": "integer"})),
(json!("value"), json!({"type": "number"})),
(json!("value"), json!({"type": "boolean"})),
(json!("value"), json!({"type": "null"})),
] {
let result = validate_payload_against_contract(
&value,
&schema,
super::RuntimeErrorCode::RequestInvalid,
"typed validation",
);
assert!(result.is_err());
}
let missing_required = validate_payload_against_contract(
&json!({}),
&json!({"type": "object", "required": ["draft_id"]}),
super::RuntimeErrorCode::RequestInvalid,
"required field",
);
assert!(missing_required.is_err());
let property_mismatch = validate_payload_against_contract(
&json!({"draft_id": 3}),
&json!({"type": "object", "properties": {"draft_id": {"type": "string"}}}),
super::RuntimeErrorCode::RequestInvalid,
"property mismatch",
);
assert!(property_mismatch.is_err());
let array_without_item_schema = validate_payload_against_contract(
&json!(["draft-1"]),
&json!({"type": "array"}),
super::RuntimeErrorCode::RequestInvalid,
"array without item schema",
);
assert!(array_without_item_schema.is_ok());
let object_without_type = validate_payload_against_contract(
&json!({"draft_id": "draft-1"}),
&json!({}),
super::RuntimeErrorCode::RequestInvalid,
"object without type",
);
assert!(object_without_type.is_ok());
}
#[test]
fn runtime_mapping_helpers_cover_all_variants() {
assert_eq!(
map_registry_scope(RegistryScope::Public),
super::RuntimeRegistryScope::Public
);
assert_eq!(
map_registry_scope(RegistryScope::Private),
super::RuntimeRegistryScope::Private
);
assert_eq!(
map_implementation_kind(ImplementationKind::Executable),
super::RuntimeImplementationKind::Executable
);
assert_eq!(
map_implementation_kind(ImplementationKind::Workflow),
super::RuntimeImplementationKind::Workflow
);
assert_eq!(
map_lifecycle(&Lifecycle::Draft),
super::RuntimeLifecycle::Draft
);
assert_eq!(
map_lifecycle(&Lifecycle::Active),
super::RuntimeLifecycle::Active
);
assert_eq!(
map_lifecycle(&Lifecycle::Deprecated),
super::RuntimeLifecycle::Deprecated
);
assert_eq!(
map_lifecycle(&Lifecycle::Retired),
super::RuntimeLifecycle::Retired
);
assert_eq!(
map_lifecycle(&Lifecycle::Archived),
super::RuntimeLifecycle::Archived
);
}
#[test]
fn runtime_candidate_helper_copies_registry_shape() {
let capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Deprecated,
);
let candidate = runtime_candidate(&capability, CandidateReason::IntentMatch);
assert_eq!(candidate.reason, CandidateReason::IntentMatch);
assert_eq!(candidate.lifecycle, super::RuntimeLifecycle::Deprecated);
assert_eq!(
candidate.implementation_kind,
super::RuntimeImplementationKind::Executable
);
}
#[test]
fn successful_runtime_execution_reports_completed_result_status() {
let mut events = super::StateEmitter::new("exec_1", "req_1");
events.push(
RuntimeState::LoadingRegistry,
RuntimeTransitionReasonCode::RuntimeInitializationStarted,
json!({}),
);
events.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::RegistryLoaded,
json!({}),
);
events.push(
RuntimeState::Discovering,
RuntimeTransitionReasonCode::RequestStarted,
json!({}),
);
events.push(
RuntimeState::EvaluatingConstraints,
RuntimeTransitionReasonCode::CandidatesCollected,
json!({"candidate_count": 1}),
);
events.push(
RuntimeState::Selecting,
RuntimeTransitionReasonCode::ConstraintsEvaluated,
json!({"eligible_candidates": 1}),
);
events.push(
RuntimeState::Executing,
RuntimeTransitionReasonCode::CandidateSelected,
json!({"capability_id": "content.comments.create-comment-draft"}),
);
let attempt = super::AttemptContext {
request: valid_request(),
execution_id: "exec_1".to_string(),
trace_id: "trace_exec_1".to_string(),
observability: super::RuntimeObservabilityConfig::default(),
artifact_verification: None,
warnings: Vec::new(),
};
let capability = resolved_capability(
Some(traverse_registry::BinaryReference {
format: traverse_registry::BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
let outcome = super::successful_execution_outcome(
super::ExecutionContext {
attempt,
emitter: events,
candidate_collection: super::CandidateCollectionRecord {
lookup_scope: PreferPrivate,
candidates: vec![runtime_candidate(&capability, CandidateReason::ExactMatch)],
rejected_candidates: Vec::new(),
},
selection: super::SelectionRecord {
status: super::SelectionStatus::Selected,
selected_capability_id: Some(capability.record.id.clone()),
selected_capability_version: Some(capability.record.version.clone()),
failure_reason: None,
remaining_candidates: Vec::new(),
},
},
&capability,
super::StartedExecution {
started_at: "1970-01-01T00:00:00Z".to_string(),
placement: super::resolve_placement(PlacementTarget::Local)
.unwrap_or_else(|_| unreachable!("local placement should resolve")),
},
json!({"draft_id": "draft-1"}),
capability.contract.emits.clone(),
None,
);
assert_eq!(outcome.result.status, RuntimeResultStatus::Completed);
assert_eq!(
outcome.state_events.last().map(|event| event.state),
Some(RuntimeState::Ready)
);
assert_eq!(
outcome.trace.decision_evidence.selection.status,
super::SelectionStatus::Selected
);
assert_eq!(
outcome.trace.state_progression.state_events,
outcome.state_events
);
assert_eq!(
outcome.trace.terminal_outcome.runtime_status,
RuntimeResultStatus::Completed
);
assert_eq!(outcome.trace.emitted_events, capability.contract.emits);
assert_eq!(
outcome.trace.state_machine_validation.status,
super::RuntimeStateMachineValidationStatus::Passed
);
}
#[test]
fn runtime_emits_a_token_free_terminal_event_for_invalid_requests() {
let sink = Arc::new(RecordingEventSink::default());
let mut request = valid_request();
request.kind = "unsupported_runtime_request".to_string();
request.context.identity = Some(super::security::RuntimeIdentity {
subject_id: "subject_123".to_string(),
actor_id: Some("actor_456".to_string()),
token_reference_hash: "must-not-leak".to_string(),
});
let outcome = Runtime::new(CapabilityRegistry::new(), NoopExecutor)
.with_event_sink(sink.clone())
.execute(request);
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
let events = sink
.events
.lock()
.expect("recording sink lock must not be poisoned");
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, super::RUNTIME_EXECUTION_EVENT_TYPE);
assert_eq!(events[0].subject_id.as_deref(), Some("subject_123"));
assert_eq!(events[0].actor_id.as_deref(), Some("actor_456"));
let serialized = serde_json::to_string(&events[0]).expect("event must serialize");
assert!(!serialized.contains("must-not-leak"));
}
#[test]
fn runtime_records_sink_delivery_failures_without_changing_execution_status() {
let outcome = Runtime::new(CapabilityRegistry::new(), NoopExecutor)
.with_event_sink(Arc::new(FailingEventSink))
.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(outcome.result.warnings.len(), 1);
assert_eq!(
outcome.result.warnings[0].code,
"runtime_event_sink_delivery_failed"
);
assert_eq!(outcome.trace.result.warnings, outcome.result.warnings);
}
#[test]
fn runtime_execution_produces_otel_phase_spans() {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let outcome = runtime.execute(valid_request());
let spans = &outcome.trace.otel_trace.spans;
let names: Vec<&str> = spans.iter().map(|span| span.name.as_str()).collect();
assert_eq!(spans.len(), 6);
assert!(names.contains(&"traverse.runtime.request"));
assert!(names.contains(&"traverse.request.intake"));
assert!(names.contains(&"traverse.registry.lookup"));
assert!(names.contains(&"traverse.contract.validation"));
assert!(names.contains(&"traverse.capability.execution"));
assert!(names.contains(&"traverse.trace.assembly"));
assert!(
spans
.iter()
.all(|span| span.status == super::OTelSpanStatus::Ok)
);
assert!(spans.iter().all(|span| {
span.attributes
.iter()
.all(|attr| attr.key.starts_with("traverse.") || attr.key == "service.name")
}));
}
#[test]
fn runtime_otel_trace_propagates_w3c_context_and_exporter_config() {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development())
.with_observability_config(super::RuntimeObservabilityConfig {
exporter: super::OTelExporterConfig {
endpoint: Some("http://collector:4318".to_string()),
protocol: super::OtlpProtocol::Http,
},
..super::RuntimeObservabilityConfig::deterministic_test("seed-1")
});
let mut request = valid_request();
request.context.traceparent =
Some("00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01".to_string());
request.context.tracestate = Some("vendor=value".to_string());
let first = runtime.execute(request.clone()).trace.otel_trace;
let second = runtime.execute(request).trace.otel_trace;
assert_eq!(first.trace_id, second.trace_id);
assert_eq!(first.spans[0].span_id, second.spans[0].span_id);
assert_eq!(
first.parent_traceparent.as_deref(),
Some("00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01")
);
assert_eq!(first.tracestate.as_deref(), Some("vendor=value"));
assert!(first.exporter.enabled);
assert_eq!(
first.exporter.endpoint.as_deref(),
Some("http://collector:4318")
);
assert_eq!(
runtime.observability_config().exporter.endpoint.as_deref(),
Some("http://collector:4318")
);
}
#[test]
fn governed_artifact_with_valid_ed25519_signature_executes() {
let artifact_bytes = b"governed wasm bytes";
let path = temp_artifact_path("ed25519-valid");
assert!(fs::write(&path, artifact_bytes).is_ok());
let signature = ed25519_signature_for(artifact_bytes);
let mut registry = CapabilityRegistry::new();
assert!(
registry
.register(governed_registration(&path, Some(signature)))
.is_ok()
);
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Completed);
assert_eq!(
outcome
.trace
.execution
.artifact_verification
.as_ref()
.map(|record| record.status),
Some(ArtifactVerificationStatus::Verified)
);
assert_eq!(
outcome
.trace
.execution
.artifact_verification
.as_ref()
.and_then(|record| record.scheme),
Some(ArtifactVerificationScheme::Ed25519)
);
}
#[test]
fn governed_artifact_without_signature_is_rejected_in_production() {
let path = temp_artifact_path("missing-signature");
assert!(fs::write(&path, b"unsigned governed bytes").is_ok());
let mut registry = CapabilityRegistry::new();
assert!(
registry
.register(governed_registration(&path, None))
.is_ok()
);
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("missing_signature")
);
assert_eq!(
outcome
.trace
.execution
.artifact_verification
.as_ref()
.and_then(|record| record.error_code.as_deref()),
Some("missing_signature")
);
}
#[test]
fn governed_artifact_without_checksum_is_rejected_before_execution() {
let bytes = b"governed checksum missing";
let path = temp_artifact_path("missing-checksum");
assert!(fs::write(&path, bytes).is_ok());
let mut registration = governed_registration(&path, Some(ed25519_signature_for(bytes)));
registration.artifact.digests.binary_digest = None;
let mut registry = CapabilityRegistry::new();
assert!(registry.register(registration).is_ok());
let outcome = Runtime::new(registry, NoopExecutor).execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("missing_checksum")
);
}
#[test]
fn governed_artifact_with_mismatched_checksum_is_rejected_before_execution() {
let bytes = b"governed checksum mismatch";
let path = temp_artifact_path("checksum-mismatch");
assert!(fs::write(&path, bytes).is_ok());
let mut registration = governed_registration(&path, Some(ed25519_signature_for(bytes)));
registration.artifact.digests.binary_digest = Some(
"sha256:0000000000000000000000000000000000000000000000000000000000000000".to_string(),
);
let mut registry = CapabilityRegistry::new();
assert!(registry.register(registration).is_ok());
let outcome = Runtime::new(registry, NoopExecutor).execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("checksum_mismatch")
);
}
#[test]
fn local_artifact_checksum_does_not_override_local_development_policy() {
let bytes = b"local checksum is advisory";
let path = temp_artifact_path("local-checksum");
assert!(fs::write(&path, bytes).is_ok());
let mut registration = governed_registration(&path, Some(ed25519_signature_for(bytes)));
registration.artifact.source.kind = SourceKind::Local;
registration.artifact.digests.binary_digest = Some(
"sha256:0000000000000000000000000000000000000000000000000000000000000000".to_string(),
);
let mut registry = CapabilityRegistry::new();
assert!(registry.register(registration).is_ok());
let outcome = Runtime::new(registry, NoopExecutor).execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Completed);
}
#[test]
fn local_dev_unsigned_artifact_warns_and_executes_in_development_mode() {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Completed);
assert_eq!(
outcome
.result
.warnings
.first()
.map(|warning| warning.code.as_str()),
Some("unsigned_local_dev_artifact")
);
assert_eq!(
outcome
.trace
.execution
.artifact_verification
.as_ref()
.map(|record| record.status),
Some(ArtifactVerificationStatus::Warning)
);
}
#[test]
fn default_security_config_is_production() {
assert_eq!(
RuntimeSecurityConfig::default(),
RuntimeSecurityConfig::production()
);
assert_ne!(
RuntimeSecurityConfig::default(),
RuntimeSecurityConfig::development()
);
}
#[test]
fn unsigned_local_artifact_rejected_under_default_security_config() {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor);
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("missing_signature")
);
}
#[test]
fn governed_artifact_rejects_placeholder_sigstore_bundle() {
let path = temp_artifact_path("sigstore-valid");
assert!(fs::write(&path, b"sigstore governed bytes").is_ok());
let signature = ArtifactSignature {
scheme: ArtifactSignatureScheme::Sigstore,
public_key_hex: None,
signature_hex: None,
sigstore_bundle_ref: Some("verified://bundle/comment-draft".to_string()),
};
let mut registry = CapabilityRegistry::new();
assert!(
registry
.register(governed_registration(&path, Some(signature)))
.is_ok()
);
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::production());
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("sigstore_unreachable")
);
}
#[test]
fn jwt_identity_is_derived_and_raw_token_is_not_traced() {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let mut request = valid_request();
let token = make_jwt_with_actor("alice", "workflow-agent");
request.context.caller = Some(token.clone());
let outcome = runtime.execute(request);
let trace_json = serde_json::to_string(&outcome.trace).unwrap_or_default();
assert_eq!(outcome.result.status, RuntimeResultStatus::Completed);
assert!(!trace_json.contains(&token));
assert_eq!(
outcome
.trace
.request
.context
.identity
.as_ref()
.map(|identity| identity.subject_id.as_str()),
Some("alice")
);
assert_eq!(
outcome
.trace
.execution
.identity
.as_ref()
.and_then(|identity| identity.actor_id.as_deref()),
Some("workflow-agent")
);
assert!(
outcome
.state_events
.iter()
.any(|event| event.details.to_string().contains("alice"))
);
}
#[test]
#[allow(clippy::too_many_lines)]
fn security_branch_guards_cover_malformed_signatures_sigstore_and_jwts() {
let capability = governed_resolved_capability(None);
let missing_public_key = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: None,
signature_hex: Some("00".to_string()),
sigstore_bundle_ref: None,
};
let missing_signature = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("00".to_string()),
signature_hex: None,
sigstore_bundle_ref: None,
};
let bad_public_hex = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("abc".to_string()),
signature_hex: Some("00".to_string()),
sigstore_bundle_ref: None,
};
let bad_signature_hex = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("00".to_string()),
signature_hex: Some("zz".to_string()),
sigstore_bundle_ref: None,
};
let short_public_key = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("00".to_string()),
signature_hex: Some("00".repeat(64)),
sigstore_bundle_ref: None,
};
let short_signature = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("00".repeat(32)),
signature_hex: Some("00".to_string()),
sigstore_bundle_ref: None,
};
let invalid_public_key = ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some("ff".repeat(32)),
signature_hex: Some("00".repeat(64)),
sigstore_bundle_ref: None,
};
let mismatch = ed25519_signature_for(b"other bytes");
for signature in [
missing_public_key,
missing_signature,
bad_public_hex,
bad_signature_hex,
short_public_key,
short_signature,
invalid_public_key,
mismatch,
] {
let mut capability = capability.clone();
capability.artifact.binary = Some(BinaryReference {
format: BinaryFormat::Wasm,
location: "unused.wasm".to_string(),
signature: Some(signature),
});
let error = verify_artifact(
&capability,
b"artifact bytes",
&RuntimeSecurityConfig::production(),
);
assert!(matches!(
error,
Err(ArtifactVerificationFailure::SignatureVerificationFailed(_))
));
let failure = error.err();
assert_eq!(
failure.as_ref().map(ArtifactVerificationFailure::code),
Some("signature_verification_failed")
);
assert_eq!(
failure
.as_ref()
.and_then(|item| item.record().error_code.as_deref()),
Some("signature_verification_failed")
);
}
let mut capability = capability.clone();
capability.artifact.binary = Some(BinaryReference {
format: BinaryFormat::Wasm,
location: "unused.wasm".to_string(),
signature: Some(ArtifactSignature {
scheme: ArtifactSignatureScheme::Sigstore,
public_key_hex: None,
signature_hex: None,
sigstore_bundle_ref: Some("https://rekor.example/bundle".to_string()),
}),
});
let error = verify_artifact(
&capability,
b"artifact bytes",
&RuntimeSecurityConfig::production(),
);
assert!(matches!(
error,
Err(ArtifactVerificationFailure::SigstoreUnreachable(_))
));
assert_eq!(
error.err().as_ref().map(ArtifactVerificationFailure::code),
Some("sigstore_unreachable")
);
let bad_payload = base64url_encode(b"{");
let no_subject = base64url_encode(b"{}");
assert!(derive_identity_from_jwt("not-a-jwt").is_none());
assert!(derive_identity_from_jwt("a.b.c.d").is_none());
assert!(derive_identity_from_jwt("a.abc=.c").is_none());
assert!(derive_identity_from_jwt("a.*.c").is_none());
assert!(derive_identity_from_jwt("a.a.c").is_none());
assert!(derive_identity_from_jwt("a.-___.c").is_none());
assert!(derive_identity_from_jwt(&format!("a.{bad_payload}.c")).is_none());
assert!(derive_identity_from_jwt(&format!("a.{no_subject}.c")).is_none());
assert_eq!(base64url_encode(b""), "");
}
#[test]
fn runtime_security_config_accessor_returns_current_config() {
let runtime = Runtime::new(CapabilityRegistry::new(), NoopExecutor)
.with_security_config(RuntimeSecurityConfig::production());
assert_eq!(
runtime.security_config(),
&RuntimeSecurityConfig::production()
);
}
#[test]
fn signed_artifact_missing_from_disk_fails_before_execution() {
let path = temp_artifact_path("missing-from-disk");
let signature = ed25519_signature_for(b"governed wasm bytes");
let mut registry = CapabilityRegistry::new();
assert!(
registry
.register(governed_registration(&path, Some(signature)))
.is_ok()
);
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::production());
let outcome = runtime.execute(valid_request());
assert_eq!(outcome.result.status, RuntimeResultStatus::Error);
assert_eq!(
outcome
.result
.error
.as_ref()
.and_then(|error| error.details.get("code"))
.and_then(serde_json::Value::as_str),
Some("artifact_load_failed")
);
}
#[test]
fn otel_timestamp_helpers_default_without_transitions() {
let transitions = Vec::new();
assert_eq!(
super::first_transition_time(&transitions),
"1970-01-01T00:00:00Z"
);
assert_eq!(
super::last_transition_time(&transitions),
"1970-01-01T00:00:00Z"
);
assert_eq!(
super::phase_started_at(&transitions, 0),
"1970-01-01T00:00:00Z"
);
assert_eq!(
super::phase_ended_at(&transitions, 0),
"1970-01-01T00:00:00Z"
);
}
#[test]
fn state_emitter_records_transition_validation_and_rejects_invalid_moves() {
let mut events = super::StateEmitter::new("exec_1", "req_1");
assert!(events.try_push(
RuntimeState::LoadingRegistry,
RuntimeTransitionReasonCode::RuntimeInitializationStarted,
json!({})
));
assert!(!events.try_push(
RuntimeState::Completed,
RuntimeTransitionReasonCode::ExecutionSucceeded,
json!({})
));
let finished = events.finish();
assert_eq!(finished.events.len(), 1);
assert_eq!(finished.transitions.len(), 1);
assert_eq!(
finished.validation.status,
super::RuntimeStateMachineValidationStatus::Failed
);
assert_eq!(finished.validation.violations.len(), 1);
}
#[test]
fn pre_execution_failure_from_constraint_phase_uses_constraint_reason() {
let mut events = super::StateEmitter::new("exec_1", "req_1");
events.push(
RuntimeState::LoadingRegistry,
RuntimeTransitionReasonCode::RuntimeInitializationStarted,
json!({}),
);
events.push(
RuntimeState::Ready,
RuntimeTransitionReasonCode::RegistryLoaded,
json!({}),
);
events.push(
RuntimeState::Discovering,
RuntimeTransitionReasonCode::RequestStarted,
json!({}),
);
events.push(
RuntimeState::EvaluatingConstraints,
RuntimeTransitionReasonCode::CandidatesCollected,
json!({"candidate_count": 1}),
);
let outcome = super::pre_execution_failure_outcome(
super::ExecutionContext {
attempt: super::AttemptContext {
request: valid_request(),
execution_id: "exec_1".to_string(),
trace_id: "trace_exec_1".to_string(),
observability: super::RuntimeObservabilityConfig::default(),
artifact_verification: None,
warnings: Vec::new(),
},
emitter: events,
candidate_collection: super::CandidateCollectionRecord {
lookup_scope: PreferPrivate,
candidates: Vec::new(),
rejected_candidates: Vec::new(),
},
selection: super::SelectionRecord {
status: super::SelectionStatus::NoMatch,
selected_capability_id: None,
selected_capability_version: None,
failure_reason: Some(super::SelectionFailureReason::NotRunnable),
remaining_candidates: Vec::new(),
},
},
super::PreExecutionFailure {
artifact_ref: None,
failure_reason: super::ExecutionFailureReason::ArtifactMissing,
placement: super::placement_not_attempted(
PlacementTarget::Local,
super::PlacementDecisionReason::SelectionNotReached,
),
error: super::runtime_error(
super::RuntimeErrorCode::CapabilityNotRunnable,
"not runnable",
json!({}),
),
artifact_verification: None,
},
);
assert_eq!(
outcome.trace.state_transitions[4].reason_code,
RuntimeTransitionReasonCode::ConstraintValidationFailed
);
}
#[test]
fn detail_object_wraps_non_object_values() {
let wrapped = super::detail_object(json!("value"));
assert_eq!(wrapped.get("value"), Some(&json!("value")));
}
#[test]
fn collect_candidates_handles_missing_target_and_public_discovery() {
let runtime = super::Runtime::new(CapabilityRegistry::new(), NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let mut request = valid_request();
request.intent.capability_id = None;
request.intent.capability_version = None;
request.intent.intent_key = None;
assert!(
runtime
.collect_candidates(&request, CandidateReason::IntentMatch)
.is_empty()
);
let mut registry = CapabilityRegistry::new();
let outcome = registry.register(public_registration());
assert!(outcome.is_ok());
let runtime = super::Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let mut request = valid_request();
request.lookup.scope = PublicOnly;
request.intent.capability_id = None;
request.intent.capability_version = None;
request.intent.intent_key = Some("content.comments.create-comment-draft".to_string());
let candidates = runtime.collect_candidates(&request, CandidateReason::IntentMatch);
assert_eq!(candidates.len(), 1);
assert_eq!(candidates[0].record.scope, RegistryScope::Public);
}
#[test]
fn noop_executor_returns_structured_output() {
let executor = NoopExecutor;
let capability = resolved_capability(
Some(BinaryReference {
format: BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
}),
Lifecycle::Active,
);
let result = executor.execute(&capability, &json!({}));
assert_eq!(
result,
Ok(super::LocalExecutionOutput {
value: json!({"draft_id": "draft"}),
emitted_events: Vec::new(),
})
);
}
#[test]
fn browser_subscription_validation_covers_guard_branches() {
let mut request = valid_browser_subscription_request();
request.kind = "wrong".to_string();
assert_eq!(
validate_browser_subscription_request(&request).map(|error| error.code),
Some(BrowserRuntimeSubscriptionErrorCode::InvalidRequest)
);
let mut request = valid_browser_subscription_request();
request.schema_version = "9.9.9".to_string();
assert_eq!(
validate_browser_subscription_request(&request).map(|error| error.code),
Some(BrowserRuntimeSubscriptionErrorCode::InvalidRequest)
);
let mut request = valid_browser_subscription_request();
request.governing_spec = "wrong-spec".to_string();
assert_eq!(
validate_browser_subscription_request(&request).map(|error| error.code),
Some(BrowserRuntimeSubscriptionErrorCode::InvalidRequest)
);
}
#[test]
fn browser_subscription_reports_not_found_for_mismatched_target() {
let outcome = runtime_outcome_for_browser_subscription();
let request = BrowserRuntimeSubscriptionRequest {
request_id: Some("req_other".to_string()),
execution_id: None,
..valid_browser_subscription_request()
};
let messages = browser_subscription_messages(&request, &outcome);
assert_eq!(
messages,
vec![BrowserRuntimeSubscriptionMessage::Error(
super::BrowserRuntimeSubscriptionErrorMessage {
kind: "browser_runtime_subscription_error".to_string(),
schema_version: "1.0.0".to_string(),
sequence: 0,
code: BrowserRuntimeSubscriptionErrorCode::NotFound,
message: "subscription target did not match the supplied execution outcome"
.to_string(),
}
)]
);
}
#[test]
fn browser_subscription_target_helper_covers_fallback_branch() {
let outcome = runtime_outcome_for_browser_subscription();
let invalid_request = BrowserRuntimeSubscriptionRequest {
request_id: Some("req_123".to_string()),
execution_id: Some(outcome.result.execution_id.clone()),
..valid_browser_subscription_request()
};
assert!(!subscription_targets_outcome(&invalid_request, &outcome));
}
fn valid_request() -> RuntimeRequest {
RuntimeRequest {
kind: "runtime_request".to_string(),
schema_version: "1.0.0".to_string(),
request_id: "req_123".to_string(),
intent: RuntimeIntent {
capability_id: Some("content.comments.create-comment-draft".to_string()),
capability_version: Some("1.0.0".to_string()),
version_range: None,
intent_key: Some("content.comments.create-comment-draft".to_string()),
},
input: json!({"comment_text": "Hello", "resource_id": "res-1"}),
lookup: RuntimeLookup {
scope: RuntimeLookupScope::PreferPrivate,
allow_ambiguity: false,
},
context: RuntimeContext {
requested_target: PlacementTarget::Local,
correlation_id: None,
caller: None,
traceparent: None,
tracestate: None,
metadata: None,
identity: None,
},
governing_spec: "006-runtime-request-execution".to_string(),
}
}
fn valid_browser_subscription_request() -> BrowserRuntimeSubscriptionRequest {
BrowserRuntimeSubscriptionRequest {
kind: "browser_runtime_subscription_request".to_string(),
schema_version: "1.0.0".to_string(),
governing_spec: "013-browser-runtime-subscription".to_string(),
request_id: Some("req_123".to_string()),
execution_id: None,
}
}
fn runtime_outcome_for_browser_subscription() -> super::RuntimeExecutionOutcome {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
runtime.execute(valid_request())
}
fn governed_registration(
path: &std::path::Path,
signature: Option<ArtifactSignature>,
) -> CapabilityRegistration {
let mut registration = public_registration();
registration.contract_path = "contracts/approved/comment-draft.json".to_string();
registration.artifact.source = SourceReference {
kind: SourceKind::Git,
location: "https://github.com/enricopiovesan/Traverse".to_string(),
};
registration.artifact.binary = Some(BinaryReference {
format: BinaryFormat::Wasm,
location: path.display().to_string(),
signature,
});
let bytes = fs::read(path).unwrap_or_default();
let hex_digest = Sha256::digest(bytes)
.iter()
.fold(String::new(), |mut acc, byte| {
let _ = write!(acc, "{byte:02x}");
acc
});
registration.artifact.digests.binary_digest = Some(format!("sha256:{hex_digest}"));
registration
}
fn governed_resolved_capability(signature: Option<ArtifactSignature>) -> ResolvedCapability {
let mut capability = resolved_capability(
Some(BinaryReference {
format: BinaryFormat::Wasm,
location: "unused.wasm".to_string(),
signature,
}),
Lifecycle::Active,
);
capability.record.contract_path = "contracts/approved/comment-draft.json".to_string();
capability.artifact.source = SourceReference {
kind: SourceKind::Git,
location: "https://github.com/enricopiovesan/Traverse".to_string(),
};
capability
}
fn ed25519_signature_for(bytes: &[u8]) -> ArtifactSignature {
let signing_key = SigningKey::from_bytes(&[7_u8; 32]);
let signature = signing_key.sign(bytes);
ArtifactSignature {
scheme: ArtifactSignatureScheme::Ed25519,
public_key_hex: Some(hex_encode(signing_key.verifying_key().as_bytes())),
signature_hex: Some(hex_encode(&signature.to_bytes())),
sigstore_bundle_ref: None,
}
}
fn temp_artifact_path(name: &str) -> std::path::PathBuf {
std::env::temp_dir().join(format!(
"traverse-runtime-{name}-{}-{}.wasm",
std::process::id(),
"req_123"
))
}
fn make_jwt_with_actor(subject_id: &str, actor_id: &str) -> String {
let header = base64url_encode(br#"{"alg":"none","typ":"JWT"}"#);
let payload = serde_json::json!({
"sub": subject_id,
"act": {"sub": actor_id}
});
format!(
"{}.{}.signature",
header,
base64url_encode(payload.to_string().as_bytes())
)
}
fn hex_encode(bytes: &[u8]) -> String {
let mut output = String::with_capacity(bytes.len() * 2);
for byte in bytes {
output.push(char::from(HEX_TABLE[(byte >> 4) as usize]));
output.push(char::from(HEX_TABLE[(byte & 0x0f) as usize]));
}
output
}
fn base64url_encode(bytes: &[u8]) -> String {
const TABLE: &[u8; 64] =
b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_";
let mut out = String::new();
let mut index = 0;
while index + 3 <= bytes.len() {
let chunk = &bytes[index..index + 3];
let n = (u32::from(chunk[0]) << 16) | (u32::from(chunk[1]) << 8) | u32::from(chunk[2]);
out.push(char::from(TABLE[((n >> 18) & 0x3f) as usize]));
out.push(char::from(TABLE[((n >> 12) & 0x3f) as usize]));
out.push(char::from(TABLE[((n >> 6) & 0x3f) as usize]));
out.push(char::from(TABLE[(n & 0x3f) as usize]));
index += 3;
}
match bytes.len() - index {
1 => {
let n = u32::from(bytes[index]) << 16;
out.push(char::from(TABLE[((n >> 18) & 0x3f) as usize]));
out.push(char::from(TABLE[((n >> 12) & 0x3f) as usize]));
}
2 => {
let n = (u32::from(bytes[index]) << 16) | (u32::from(bytes[index + 1]) << 8);
out.push(char::from(TABLE[((n >> 18) & 0x3f) as usize]));
out.push(char::from(TABLE[((n >> 12) & 0x3f) as usize]));
out.push(char::from(TABLE[((n >> 6) & 0x3f) as usize]));
}
_ => {}
}
out
}
fn public_registration() -> CapabilityRegistration {
CapabilityRegistration {
scope: RegistryScope::Public,
contract: test_contract(Lifecycle::Active),
contract_path: "registry/contract.json".to_string(),
artifact: test_artifact(Some(BinaryReference {
format: BinaryFormat::Wasm,
location: "artifact.wasm".to_string(),
signature: None,
})),
registered_at: "2026-03-27T00:00:00Z".to_string(),
tags: vec!["comments".to_string()],
composability: ComposabilityMetadata {
kind: CompositionKind::Atomic,
patterns: vec![CompositionPattern::Sequential],
provides: vec!["draft".to_string()],
requires: vec!["authenticated-user".to_string()],
},
governing_spec: "005-capability-registry".to_string(),
validator_version: "0.1.0".to_string(),
}
}
fn resolved_capability(
binary: Option<traverse_registry::BinaryReference>,
lifecycle: Lifecycle,
) -> ResolvedCapability {
ResolvedCapability {
contract: test_contract(lifecycle.clone()),
record: test_record(lifecycle.clone()),
artifact: test_artifact(binary),
index_entry: test_index_entry(lifecycle),
}
}
fn test_contract(lifecycle: Lifecycle) -> traverse_contracts::CapabilityContract {
traverse_contracts::CapabilityContract {
kind: "capability_contract".to_string(),
schema_version: "1.0.0".to_string(),
id: "content.comments.create-comment-draft".to_string(),
namespace: "content.comments".to_string(),
name: "create-comment-draft".to_string(),
version: "1.0.0".to_string(),
lifecycle,
owner: Owner {
team: "comments".to_string(),
contact: "comments@example.com".to_string(),
},
summary: "Create a comment draft for a resource".to_string(),
description: "Creates a draft comment and returns the generated draft identifier."
.to_string(),
inputs: SchemaContainer {
schema: json!({"type": "object"}),
},
outputs: SchemaContainer {
schema: json!({"type": "object"}),
},
preconditions: Vec::new(),
postconditions: Vec::new(),
side_effects: vec![traverse_contracts::SideEffect {
kind: traverse_contracts::SideEffectKind::MemoryOnly,
description: "Produces a draft representation in memory.".to_string(),
}],
emits: Vec::new(),
consumes: Vec::new(),
permissions: Vec::new(),
execution: Execution {
binary_format: ContractBinaryFormat::Wasm,
entrypoint: Entrypoint {
kind: EntrypointKind::WasiCommand,
command: "run".to_string(),
},
preferred_targets: vec![ExecutionTarget::Local],
constraints: ExecutionConstraints {
host_api_access: HostApiAccess::None,
network_access: NetworkAccess::Forbidden,
filesystem_access: FilesystemAccess::None,
},
},
policies: Vec::new(),
dependencies: Vec::new(),
provenance: Provenance {
source: ProvenanceSource::Greenfield,
author: "Enrico Piovesan".to_string(),
created_at: "2026-03-27T00:00:00Z".to_string(),
spec_ref: Some("006-runtime-request-execution".to_string()),
adr_refs: Vec::new(),
exception_refs: Vec::new(),
},
evidence: Vec::new(),
service_type: ServiceType::Stateless,
permitted_targets: vec![
ExecutionTarget::Local,
ExecutionTarget::Cloud,
ExecutionTarget::Edge,
ExecutionTarget::Device,
],
event_trigger: None,
connector_requirements: Vec::new(),
state_schema: None,
}
}
fn test_record(lifecycle: Lifecycle) -> CapabilityRegistryRecord {
CapabilityRegistryRecord {
scope: RegistryScope::Private,
id: "content.comments.create-comment-draft".to_string(),
version: "1.0.0".to_string(),
lifecycle,
owner: Owner {
team: "comments".to_string(),
contact: "comments@example.com".to_string(),
},
contract_path: "registry/contract.json".to_string(),
contract_digest: "digest".to_string(),
implementation_kind: ImplementationKind::Executable,
artifact_ref: "artifact:content.comments.create-comment-draft:1.0.0".to_string(),
registered_at: "2026-03-27T00:00:00Z".to_string(),
provenance: RegistryProvenance {
source: "test".to_string(),
author: "Enrico Piovesan".to_string(),
created_at: "2026-03-27T00:00:00Z".to_string(),
},
evidence: traverse_registry::RegistrationEvidence {
evidence_id: "evidence".to_string(),
artifact_ref: "artifact:content.comments.create-comment-draft:1.0.0".to_string(),
capability_id: "content.comments.create-comment-draft".to_string(),
capability_version: "1.0.0".to_string(),
scope: RegistryScope::Private,
governing_spec: "005-capability-registry".to_string(),
validator_version: "0.1.0".to_string(),
produced_at: "2026-03-27T00:00:00Z".to_string(),
result: traverse_registry::RegistrationResult::Passed,
},
}
}
fn test_artifact(
binary: Option<traverse_registry::BinaryReference>,
) -> CapabilityArtifactRecord {
CapabilityArtifactRecord {
artifact_ref: "artifact:content.comments.create-comment-draft:1.0.0".to_string(),
implementation_kind: ImplementationKind::Executable,
source: SourceReference {
kind: SourceKind::Git,
location: "https://github.com/enricopiovesan/cogolo".to_string(),
},
binary,
workflow_ref: None,
digests: ArtifactDigests {
source_digest: "src-digest".to_string(),
binary_digest: Some("bin-digest".to_string()),
},
provenance: RegistryProvenance {
source: "test".to_string(),
author: "Enrico Piovesan".to_string(),
created_at: "2026-03-27T00:00:00Z".to_string(),
},
}
}
fn test_index_entry(lifecycle: Lifecycle) -> DiscoveryIndexEntry {
DiscoveryIndexEntry {
scope: RegistryScope::Private,
id: "content.comments.create-comment-draft".to_string(),
version: "1.0.0".to_string(),
lifecycle,
owner: Owner {
team: "comments".to_string(),
contact: "comments@example.com".to_string(),
},
summary: "Create a comment draft for a resource".to_string(),
tags: vec!["comments".to_string()],
permissions: Vec::new(),
emits: Vec::new(),
consumes: Vec::new(),
implementation_kind: ImplementationKind::Executable,
composability: traverse_registry::ComposabilityMetadata {
kind: traverse_registry::CompositionKind::Atomic,
patterns: vec![traverse_registry::CompositionPattern::Sequential],
provides: vec!["draft".to_string()],
requires: vec!["authenticated-user".to_string()],
},
artifact_ref: "artifact:content.comments.create-comment-draft:1.0.0".to_string(),
registered_at: "2026-03-27T00:00:00Z".to_string(),
}
}
fn write_runtime_workspace_app_state_fixture(workspace_root: &Path, workspace_id: &str) {
let repo = repo_root();
let state_path = workspace_root
.join(".traverse/workspaces")
.join(workspace_id)
.join("apps/expedition.readiness/1.0.0/registration.json");
fs::create_dir_all(state_path.parent().expect("state path must have parent"))
.expect("workspace state parent should create");
fs::write(
state_path,
serde_json::to_string_pretty(&serde_json::json!({
"status": "registered",
"workspace_id": workspace_id,
"app_id": "expedition.readiness",
"app_version": "1.0.0",
"schema_version": "1.0.0",
"manifest_path": repo.join("examples/applications/expedition-readiness/app.manifest.json").display().to_string(),
"manifest_digest": "sha256:test-manifest",
"bundle_digest": "sha256:test-bundle",
"component_ids": [
"expedition.readiness.capture-expedition-objective-component",
"expedition.readiness.interpret-expedition-intent-component",
"expedition.readiness.assess-conditions-summary-component",
"expedition.readiness.validate-team-readiness-component",
"expedition.readiness.assemble-expedition-plan-component"
],
"workflow_ids": ["expedition.planning.plan-expedition"],
"components": runtime_workspace_components_json(&repo),
"workflows": [{
"workflow_id": "expedition.planning.plan-expedition",
"workflow_version": "1.0.0",
"workflow_digest": "sha256:test-workflow",
"path": repo.join("workflows/examples/expedition/plan-expedition/workflow.json").display().to_string()
}],
"model_dependencies": [{
"interface_id": "traverse.inference.generate",
"version_range": "^1.0",
"selection_policy": {
"strategy": "priority",
"allow_fallback": true
},
"required_capabilities": ["text_generation"],
"minimum_context_window": 8192,
"candidates": [{
"candidate_id": "ollama-llama-3-2-readiness",
"provider_capability_id": "traverse.inference.generate",
"provider_implementation_id": "ollama.local.generate",
"model_identifier": "llama3.2:3b",
"placement_target": "local",
"priority": 10,
"required_provider_config_keys": ["ollama_base_url"],
"metadata": {
"implementation_kind": "real_local_provider",
"provider": "ollama",
"model_context_window": 8192
}
}]
}],
"effective_config": {
"values": {
"workspace_id": "expedition-local",
"readiness_mode": "deterministic"
},
"redacted_secret_keys": []
},
"state_scope": "workspace_persisted",
"registration_fingerprint": {
"app_id": "expedition.readiness",
"app_version": "1.0.0",
"manifest_digest": "sha256:test-manifest"
}
}))
.expect("workspace app state should serialize"),
)
.expect("workspace app state should write");
}
fn runtime_workspace_components_json(repo: &Path) -> Vec<serde_json::Value> {
[
(
"capture-expedition-objective",
"expedition.planning.capture-expedition-objective",
),
(
"interpret-expedition-intent",
"expedition.planning.interpret-expedition-intent",
),
(
"assess-conditions-summary",
"expedition.planning.assess-conditions-summary",
),
(
"validate-team-readiness",
"expedition.planning.validate-team-readiness",
),
(
"assemble-expedition-plan",
"expedition.planning.assemble-expedition-plan",
),
]
.into_iter()
.map(|(leaf, capability_id)| {
serde_json::json!({
"component_id": format!("expedition.readiness.{leaf}-component"),
"component_version": "1.0.0",
"capability_id": capability_id,
"capability_version": "1.0.0",
"wasm_digest": "sha256:5647c39a1d25d8728350f9619025292a62e78a602068a2ad9b6f075751c93d99",
"manifest_path": repo.join("examples/applications/expedition-readiness/components/validate-team-readiness/component.manifest.json").display().to_string(),
"contract_path": repo.join(format!("contracts/examples/expedition/capabilities/{leaf}/contract.json")).display().to_string(),
"artifact_ref": repo.join("examples/capabilities/team-readiness-agent/artifacts/validate-team-readiness-agent.wasm").display().to_string()
})
})
.collect()
}
fn repo_root() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../..")
}
fn unique_workspace_state_dir() -> PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let counter = TEMP_COUNTER.fetch_add(1, Ordering::SeqCst);
let path = std::env::temp_dir().join(format!(
"traverse-runtime-workspace-state-test-{}-{nanos}-{counter}",
std::process::id()
));
fs::create_dir_all(&path).expect("temporary workspace should create");
path
}
#[derive(Debug, Clone)]
struct NoopExecutor;
impl super::LocalExecutor for NoopExecutor {
fn execute(
&self,
_capability: &ResolvedCapability,
_input: &serde_json::Value,
) -> Result<super::LocalExecutionOutput, super::LocalExecutionFailure> {
Ok(super::LocalExecutionOutput {
value: json!({"draft_id": "draft"}),
emitted_events: Vec::new(),
})
}
}
struct FailingExecutor;
impl super::LocalExecutor for FailingExecutor {
fn execute(
&self,
_capability: &ResolvedCapability,
_input: &serde_json::Value,
) -> Result<super::LocalExecutionOutput, super::LocalExecutionFailure> {
Err(super::LocalExecutionFailure {
code: super::LocalExecutionFailureCode::ExecutionFailed,
message: "forced failure".to_string(),
})
}
}
fn successful_trace() -> super::RuntimeTrace {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
runtime.execute(valid_request()).trace
}
fn failed_trace() -> super::RuntimeTrace {
let mut registry = CapabilityRegistry::new();
assert!(registry.register(public_registration()).is_ok());
let runtime = Runtime::new(registry, FailingExecutor)
.with_security_config(RuntimeSecurityConfig::development());
runtime.execute(valid_request()).trace
}
#[test]
fn selected_capability_id_returns_id_on_success() {
let trace = successful_trace();
assert_eq!(
trace.selected_capability_id(),
Some("content.comments.create-comment-draft")
);
}
#[test]
fn selected_capability_id_returns_none_when_no_selection() {
let registry = CapabilityRegistry::new();
let runtime = Runtime::new(registry, NoopExecutor)
.with_security_config(RuntimeSecurityConfig::development());
let trace = runtime.execute(valid_request()).trace;
assert!(trace.selected_capability_id().is_none());
}
#[test]
fn errors_returns_none_on_success() {
let trace = successful_trace();
assert!(trace.errors().is_none());
}
#[test]
fn errors_returns_error_on_failure() {
let trace = failed_trace();
assert!(trace.errors().is_some());
}
#[test]
fn emitted_events_returns_slice() {
let trace = successful_trace();
let _ = trace.emitted_events();
}
#[test]
fn runtime_trace_exposes_non_sensitive_model_resolution_evidence() {
let trace = successful_trace().with_model_resolution(vec![model_resolution_evidence()]);
let serialized = serde_json::to_string(&trace).unwrap_or_default();
assert_eq!(trace.model_resolution.len(), 1);
assert_eq!(
trace.decision_evidence.model_resolution,
trace.model_resolution
);
assert!(serialized.contains("model_resolution"));
assert!(serialized.contains("ollama.local.generate"));
assert!(serialized.contains("llama3.2:3b"));
assert!(!serialized.contains("private prompt"));
assert!(!serialized.contains("raw source text"));
assert!(!serialized.contains("sk-local-secret"));
}
#[test]
fn output_returns_value_on_success() {
let trace = successful_trace();
assert_eq!(trace.output(), Some(&json!({"draft_id": "draft"})));
}
#[test]
fn output_returns_none_on_failure() {
let trace = failed_trace();
assert!(trace.output().is_none());
}
#[test]
fn is_success_true_on_completed() {
let trace = successful_trace();
assert!(trace.is_success());
}
#[test]
fn is_success_false_on_error() {
let trace = failed_trace();
assert!(!trace.is_success());
}
fn model_resolution_evidence() -> ModelResolutionEvidence {
ModelResolutionEvidence {
phase: ModelResolutionPhase::Execution,
interface_id: "traverse.inference.generate".to_string(),
requested_interface_id: "traverse.inference.generate".to_string(),
requested_placement: ExecutionTarget::Local,
selected: Some(SelectedModelCandidate {
candidate_id: "ollama-llama-3-2".to_string(),
provider_capability_id: "traverse.inference.generate".to_string(),
provider_implementation_id: "ollama.local.generate".to_string(),
model_identifier: "llama3.2:3b".to_string(),
placement_target: ExecutionTarget::Local,
priority: 10,
selection_reason: "selected highest-priority passing candidate".to_string(),
}),
candidates: vec![traverse_registry::ModelCandidateEvaluation {
candidate_id: "ollama-llama-3-2".to_string(),
provider_capability_id: "traverse.inference.generate".to_string(),
provider_implementation_id: "ollama.local.generate".to_string(),
model_identifier: "llama3.2:3b".to_string(),
placement_target: ExecutionTarget::Local,
priority: 10,
readiness: ModelCandidateReadiness::Ready,
rejection_code: Option::<ModelCandidateRejectionCode>::None,
reason: "candidate passed availability, interface, placement, and context checks"
.to_string(),
manifest_order: 0,
}],
failure_code: Option::<ModelCandidateRejectionCode>::None,
}
}
}