1use std::{
13 collections::HashMap,
14 sync::{Arc, Mutex},
15 time::Instant,
16};
17
18use chrono::Utc;
19use serde_json::Value;
20use traverse_contracts::{CapabilityContract, ServiceType, ViolationRecord};
21
22use crate::{
23 events::types::{EventBroker, TraverseEvent},
24 executor::{ArtifactType, CapabilityExecutor, ExecutorCapability},
25 placement::{
26 PlacementConstraintEvaluator, PlacementDecision, PlacementError, PlacementRequest,
27 RuntimeSnapshot,
28 },
29 trace::{PrivateTraceEntry, PublicTraceEntry, TraceOutcome, TraceStore, new_trace_id_and_time},
30};
31
32use traverse_contracts::ExecutionTarget;
33
34pub type CapabilityExecutorRegistry = HashMap<ArtifactType, Box<dyn CapabilityExecutor>>;
40
41pub struct RouterRequest {
43 pub capability_id: String,
45 pub artifact_type: ArtifactType,
47 pub contract: CapabilityContract,
49 pub target_hint: Option<ExecutionTarget>,
51 pub runtime_snapshot: RuntimeSnapshot,
53 pub input: Value,
55 pub executor_capability: ExecutorCapability,
57 pub trace_id_override: Option<String>,
59}
60
61#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum RouterError {
64 PlacementFailed(PlacementError),
66 ExecutorNotFound(String),
68 ExecutionFailed(String),
70 ContractViolation(Vec<ViolationRecord>),
72 TraceLockPoisoned,
74}
75
76impl std::fmt::Display for RouterError {
77 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78 match self {
79 Self::PlacementFailed(e) => write!(f, "placement failed: {e:?}"),
80 Self::ExecutorNotFound(t) => write!(f, "no executor registered for artifact type: {t}"),
81 Self::ExecutionFailed(msg) => write!(f, "execution failed: {msg}"),
82 Self::ContractViolation(violations) => {
83 write!(f, "contract violation: {} violation(s)", violations.len())
84 }
85 Self::TraceLockPoisoned => write!(f, "trace store lock is poisoned"),
86 }
87 }
88}
89
90impl std::error::Error for RouterError {}
91
92#[derive(Debug)]
94pub struct RouterResponse {
95 pub output: Value,
97 pub emitted_events: Vec<TraverseEvent>,
101 pub trace_id: String,
103 pub placement_decision: PlacementDecision,
105}
106
107pub struct PlacementRouter {
116 evaluator: PlacementConstraintEvaluator,
117 executor_registry: CapabilityExecutorRegistry,
118 trace_store: Arc<Mutex<TraceStore>>,
119 event_broker: Arc<dyn EventBroker>,
120}
121
122impl PlacementRouter {
123 #[must_use]
125 pub fn new(
126 evaluator: PlacementConstraintEvaluator,
127 executor_registry: CapabilityExecutorRegistry,
128 trace_store: Arc<Mutex<TraceStore>>,
129 event_broker: Arc<dyn EventBroker>,
130 ) -> Self {
131 Self {
132 evaluator,
133 executor_registry,
134 trace_store,
135 event_broker,
136 }
137 }
138
139 pub fn execute(&self, request: RouterRequest) -> Result<RouterResponse, RouterError> {
152 let executor = self
153 .executor_registry
154 .get(&request.artifact_type)
155 .ok_or_else(|| RouterError::ExecutorNotFound(format!("{:?}", request.artifact_type)))?;
156 self.execute_with_executor(request, executor.as_ref())
157 }
158
159 pub fn execute_with_executor(
168 &self,
169 request: RouterRequest,
170 executor: &dyn CapabilityExecutor,
171 ) -> Result<RouterResponse, RouterError> {
172 let placement_req = PlacementRequest {
174 capability_id: request.capability_id.clone(),
175 target_hint: request.target_hint,
176 runtime_snapshot: request.runtime_snapshot,
177 };
178
179 let decision = self
180 .evaluator
181 .evaluate(&placement_req, &request.contract)
182 .map_err(RouterError::PlacementFailed)?;
183
184 let placement_target_str = format!("{:?}", decision.target);
185
186 let start = Instant::now();
193 let exec_result = executor.execute(&request.executor_capability, &request.input);
194 let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
195
196 let (output, emitted_events, outcome) = match exec_result {
197 Ok(exec_output) => (
198 exec_output.value,
199 exec_output.emitted_events,
200 TraceOutcome::Success,
201 ),
202 Err(e) => return Err(RouterError::ExecutionFailed(format!("{e}"))),
203 };
204
205 let (trace_id, time) = match request.trace_id_override {
207 Some(override_id) => (override_id, Utc::now().to_rfc3339()),
208 None => new_trace_id_and_time(),
209 };
210
211 let public_entry = PublicTraceEntry::new(
212 trace_id.clone(),
213 request.capability_id.clone(),
214 placement_target_str,
215 outcome,
216 duration_ms,
217 time,
218 );
219
220 let input_str = serde_json::to_string(&request.input).unwrap_or_default();
221 let output_str = serde_json::to_string(&output).unwrap_or_default();
222 let private_entry =
223 PrivateTraceEntry::new(trace_id.clone(), &input_str, &output_str, duration_ms);
224
225 {
226 let mut store = self
227 .trace_store
228 .lock()
229 .map_err(|_| RouterError::TraceLockPoisoned)?;
230 store.insert(public_entry, Some(private_entry));
231 }
232
233 let published_events = if request.contract.service_type == ServiceType::Subscribable {
235 for event in &emitted_events {
236 let _ = self.event_broker.publish(event.clone());
238 }
239 emitted_events
240 } else {
241 Vec::new()
242 };
243
244 Ok(RouterResponse {
245 output,
246 emitted_events: published_events,
247 trace_id,
248 placement_decision: decision,
249 })
250 }
251}