1use std::{
13 collections::HashMap,
14 sync::{Arc, Mutex},
15 time::Instant,
16};
17
18use serde_json::Value;
19use traverse_contracts::{CapabilityContract, ServiceType, ViolationRecord};
20
21use crate::{
22 events::types::{EventBroker, TraverseEvent},
23 executor::{ArtifactType, CapabilityExecutor, ExecutorCapability},
24 placement::{
25 PlacementConstraintEvaluator, PlacementDecision, PlacementError, PlacementRequest,
26 RuntimeSnapshot,
27 },
28 trace::{PrivateTraceEntry, PublicTraceEntry, TraceOutcome, TraceStore, new_trace_id_and_time},
29};
30
31use traverse_contracts::ExecutionTarget;
32
33pub type CapabilityExecutorRegistry = HashMap<ArtifactType, Box<dyn CapabilityExecutor>>;
39
40pub struct RouterRequest {
42 pub capability_id: String,
44 pub artifact_type: ArtifactType,
46 pub contract: CapabilityContract,
48 pub target_hint: Option<ExecutionTarget>,
50 pub runtime_snapshot: RuntimeSnapshot,
52 pub input: Value,
54 pub executor_capability: ExecutorCapability,
56 pub emitted_events: Vec<TraverseEvent>,
58}
59
60#[derive(Debug, Clone, PartialEq, Eq)]
62pub enum RouterError {
63 PlacementFailed(PlacementError),
65 ExecutorNotFound(String),
67 ExecutionFailed(String),
69 ContractViolation(Vec<ViolationRecord>),
71 TraceLockPoisoned,
73}
74
75impl std::fmt::Display for RouterError {
76 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77 match self {
78 Self::PlacementFailed(e) => write!(f, "placement failed: {e:?}"),
79 Self::ExecutorNotFound(t) => write!(f, "no executor registered for artifact type: {t}"),
80 Self::ExecutionFailed(msg) => write!(f, "execution failed: {msg}"),
81 Self::ContractViolation(violations) => {
82 write!(f, "contract violation: {} violation(s)", violations.len())
83 }
84 Self::TraceLockPoisoned => write!(f, "trace store lock is poisoned"),
85 }
86 }
87}
88
89impl std::error::Error for RouterError {}
90
91#[derive(Debug)]
93pub struct RouterResponse {
94 pub output: Value,
96 pub trace_id: String,
98 pub placement_decision: PlacementDecision,
100}
101
102pub struct PlacementRouter {
111 evaluator: PlacementConstraintEvaluator,
112 executor_registry: CapabilityExecutorRegistry,
113 trace_store: Arc<Mutex<TraceStore>>,
114 event_broker: Arc<dyn EventBroker>,
115}
116
117impl PlacementRouter {
118 #[must_use]
120 pub fn new(
121 evaluator: PlacementConstraintEvaluator,
122 executor_registry: CapabilityExecutorRegistry,
123 trace_store: Arc<Mutex<TraceStore>>,
124 event_broker: Arc<dyn EventBroker>,
125 ) -> Self {
126 Self {
127 evaluator,
128 executor_registry,
129 trace_store,
130 event_broker,
131 }
132 }
133
134 pub fn execute(&self, request: RouterRequest) -> Result<RouterResponse, RouterError> {
147 let placement_req = PlacementRequest {
149 capability_id: request.capability_id.clone(),
150 target_hint: request.target_hint,
151 runtime_snapshot: request.runtime_snapshot,
152 };
153
154 let decision = self
155 .evaluator
156 .evaluate(&placement_req, &request.contract)
157 .map_err(RouterError::PlacementFailed)?;
158
159 let placement_target_str = format!("{:?}", decision.target);
160
161 let executor = self
163 .executor_registry
164 .get(&request.artifact_type)
165 .ok_or_else(|| RouterError::ExecutorNotFound(format!("{:?}", request.artifact_type)))?;
166
167 let start = Instant::now();
169 let exec_result = executor.execute(&request.executor_capability, &request.input);
170 let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
171
172 let (output, outcome) = match exec_result {
173 Ok(v) => (v, TraceOutcome::Success),
174 Err(e) => return Err(RouterError::ExecutionFailed(format!("{e}"))),
175 };
176
177 let mut violations = Vec::new();
179 if request.contract.service_type == ServiceType::Subscribable
180 && !request.emitted_events.is_empty()
181 {
182 for event in &request.emitted_events {
183 let declared =
184 request.contract.emits.iter().any(|decl| {
185 decl.event_id == event.event_type && decl.version == event.version
186 });
187 if !declared {
188 violations.push(ViolationRecord::new(
189 "undeclared_event_emission",
190 &request.capability_id,
191 format!(
192 "capability emitted undeclared event {}@{}",
193 event.event_type, event.version
194 ),
195 ));
196 }
197 }
198 }
199
200 let outcome = if violations.is_empty() {
201 outcome
202 } else {
203 TraceOutcome::Failure
204 };
205
206 let (trace_id, time) = new_trace_id_and_time();
208
209 let mut public_entry = PublicTraceEntry::new(
210 trace_id.clone(),
211 request.capability_id.clone(),
212 placement_target_str,
213 outcome,
214 duration_ms,
215 time,
216 );
217 public_entry.violations.clone_from(&violations);
218
219 let input_str = serde_json::to_string(&request.input).unwrap_or_default();
220 let output_str = serde_json::to_string(&output).unwrap_or_default();
221 let private_entry =
222 PrivateTraceEntry::new(trace_id.clone(), &input_str, &output_str, duration_ms);
223
224 {
225 let mut store = self
226 .trace_store
227 .lock()
228 .map_err(|_| RouterError::TraceLockPoisoned)?;
229 store.insert(public_entry, Some(private_entry));
230 }
231
232 if !violations.is_empty() {
233 return Err(RouterError::ContractViolation(violations));
234 }
235
236 if request.contract.service_type == ServiceType::Subscribable {
238 for event in request.emitted_events {
239 let _ = self.event_broker.publish(event);
241 }
242 }
243
244 Ok(RouterResponse {
245 output,
246 trace_id,
247 placement_decision: decision,
248 })
249 }
250}