1use async_trait::async_trait;
4use awaken_runtime_contract::registry_spec::AgentSpec;
5use thiserror::Error;
6
7use crate::RuntimeError;
8use crate::backend::ExecutionBackendError;
9
10mod capabilities;
11mod local_registry;
12mod types;
13
14pub use capabilities::{
15 BackendProfile, BackendRequirements, CancellationCapability, CapabilityDecision,
16 CapabilityMismatch, ContinuationCapability, DecisionCapability, FrontendToolCapability,
17 OutputCapability, OverrideCapability, PersistenceCapability, TranscriptCapability,
18 WaitCapability,
19};
20pub use local_registry::LocalRegistryResolver;
21pub use types::{
22 DelegatePersistence, ExecutionPlan, ExecutionRole, HandoffTranscriptRef, LiveOnlyScope,
23 PersistenceRequirement, RegistryResolutionScope, ReplayableResolvedRun, ReplayableScope,
24 ResolutionArtifact, ResolutionPolicy, ResolutionRequest, ResolutionTarget,
25 ResolvedModelBinding, ResolvedRun, ResolvedRunPlan, ResolvedTool, RootScopeKind, RunFeatureSet,
26};
27
28#[derive(Debug, Error)]
29pub enum ResolveError {
30 #[error("runtime resolve failed: {0}")]
31 Runtime(String),
32 #[error("unsupported resolution target: {0}")]
33 UnsupportedTarget(String),
34 #[error("unsupported persistence: {0}")]
35 UnsupportedPersistence(String),
36 #[error("backend capability mismatch: {0:?}")]
37 CapabilityMismatch(Vec<CapabilityMismatch>),
38 #[error("nested resolution scope mismatch: {0}")]
39 NestedScopeMismatch(String),
40}
41
42impl From<RuntimeError> for ResolveError {
43 fn from(error: RuntimeError) -> Self {
44 Self::Runtime(error.to_string())
45 }
46}
47
48impl From<ExecutionBackendError> for ResolveError {
49 fn from(error: ExecutionBackendError) -> Self {
50 Self::Runtime(error.to_string())
51 }
52}
53
54#[async_trait]
55pub trait Resolver: Send + Sync {
56 async fn resolve(&self, request: ResolutionRequest) -> Result<ResolvedRunPlan, ResolveError>;
57}
58
59#[async_trait]
60pub trait AgentSpecLookup: Send + Sync {
61 async fn resolve_spec(&self, agent_id: &str) -> Result<AgentSpec, ResolveError>;
62}
63
64#[cfg(test)]
65mod tests {
66 use std::sync::Arc;
67
68 use super::*;
69 use awaken_runtime_contract::contract::identity::RunIdentity;
70 use awaken_runtime_contract::contract::message::Message;
71
72 use crate::registry::{AgentResolver, ResolvedAgent};
73 use crate::run::RunActivation;
74
75 fn resume_decision() -> awaken_runtime_contract::contract::suspension::ToolCallResume {
76 awaken_runtime_contract::contract::suspension::ToolCallResume {
77 decision_id: "d1".into(),
78 action: awaken_runtime_contract::contract::suspension::ResumeDecisionAction::Resume,
79 result: serde_json::Value::Null,
80 reason: None,
81 updated_at: 0,
82 }
83 }
84
85 #[test]
86 fn features_are_derived_from_activation_and_policy() {
87 let activation = RunActivation::new("thread", vec![Message::user("hi")])
88 .with_decisions(vec![("call".into(), resume_decision())])
89 .with_hitl_resume_run_id("run-1");
90 let features =
91 RunFeatureSet::from_activation(&activation, ResolutionPolicy::PersistentServer);
92 assert!(features.has_seeded_decisions);
93 assert!(features.is_human_resume);
94 assert_eq!(
95 features.requested_persistence,
96 PersistenceRequirement::CheckpointRequired
97 );
98 }
99
100 #[test]
101 fn requirements_are_derived_from_features() {
102 let features = RunFeatureSet {
103 has_seeded_decisions: true,
104 has_live_decision_channel: true,
105 has_overrides: true,
106 has_frontend_tools: true,
107 is_continuation: true,
108 requested_persistence: PersistenceRequirement::CheckpointRequired,
109 ..Default::default()
110 };
111 let req = BackendRequirements::from_features(&features);
112 assert_eq!(req.decisions, Some(DecisionCapability::LiveAndDurable));
113 assert_eq!(req.overrides, Some(OverrideCapability::InferenceParams));
114 assert_eq!(
115 req.frontend_tools,
116 Some(FrontendToolCapability::DescriptorsOnly)
117 );
118 assert_eq!(req.persistence, Some(PersistenceCapability::Checkpoint));
119 assert_eq!(
120 req.continuation,
121 Some(ContinuationCapability::InProcessState)
122 );
123 }
124
125 #[test]
126 fn profile_check_returns_all_mismatches() {
127 let profile = BackendProfile {
128 decisions: DecisionCapability::None,
129 overrides: OverrideCapability::None,
130 frontend_tools: FrontendToolCapability::None,
131 persistence: PersistenceCapability::Ephemeral,
132 ..BackendProfile::full_local()
133 };
134 let req = BackendRequirements {
135 decisions: Some(DecisionCapability::DurableResume),
136 overrides: Some(OverrideCapability::InferenceParams),
137 frontend_tools: Some(FrontendToolCapability::DescriptorsOnly),
138 persistence: Some(PersistenceCapability::Checkpoint),
139 cancellation: None,
140 continuation: None,
141 waits: None,
142 transcript: None,
143 output: None,
144 };
145 let CapabilityDecision::Unsupported(mismatches) = profile.check(&req) else {
146 panic!("expected mismatches")
147 };
148 assert_eq!(mismatches.len(), 4);
149 }
150
151 #[test]
152 fn full_local_matches_contract_table() {
153 let profile = BackendProfile::full_local();
154 assert_eq!(profile.decisions, DecisionCapability::LiveAndDurable);
155 assert_eq!(profile.overrides, OverrideCapability::ModelAndParams);
156 assert_eq!(profile.frontend_tools, FrontendToolCapability::Executable);
157 assert_eq!(profile.persistence, PersistenceCapability::Checkpoint);
158 }
159
160 #[test]
161 fn remote_state_transcript_satisfies_full_transcript_requirement() {
162 let profile = BackendProfile {
163 transcript: TranscriptCapability::IncrementalUserMessagesWithRemoteState,
164 ..BackendProfile::full_local()
165 };
166 let req = BackendRequirements {
167 transcript: Some(TranscriptCapability::FullTranscript),
168 cancellation: None,
169 continuation: None,
170 decisions: None,
171 overrides: None,
172 frontend_tools: None,
173 persistence: None,
174 waits: None,
175 output: None,
176 };
177 assert!(matches!(profile.check(&req), CapabilityDecision::Supported));
178 }
179
180 struct MockResolver;
181
182 #[async_trait::async_trait]
183 impl awaken_runtime_contract::contract::executor::LlmExecutor for MockResolver {
184 async fn execute(
185 &self,
186 _request: awaken_runtime_contract::contract::executor::InferenceRequest,
187 ) -> Result<
188 awaken_runtime_contract::contract::inference::StreamResult,
189 awaken_runtime_contract::contract::executor::InferenceExecutionError,
190 > {
191 Ok(awaken_runtime_contract::contract::inference::StreamResult {
192 content: vec![],
193 tool_calls: vec![],
194 usage: None,
195 stop_reason: None,
196 has_incomplete_tool_calls: false,
197 })
198 }
199
200 fn name(&self) -> &str {
201 "mock"
202 }
203 }
204
205 impl AgentResolver for MockResolver {
206 fn resolve(&self, agent_id: &str) -> Result<ResolvedAgent, RuntimeError> {
207 Ok(ResolvedAgent::new(
208 agent_id,
209 "model",
210 "system",
211 Arc::new(MockResolver),
212 ))
213 }
214 }
215
216 #[tokio::test]
217 async fn legacy_adapter_rejects_pinned_and_non_root_requests() {
218 let adapter = LocalRegistryResolver::new(Arc::new(MockResolver));
219 let activation =
220 RunActivation::new("thread", vec![Message::user("hi")]).with_agent_id("agent-a");
221 let mut request =
222 ResolutionRequest::from_activation(&activation, ResolutionPolicy::LiveOnlyEmbedded);
223 request.resolution_scope = RegistryResolutionScope::Pinned("resolution-1".to_string());
224 assert!(matches!(
225 adapter.resolve(request).await,
226 Err(ResolveError::UnsupportedPersistence(_))
227 ));
228
229 let mut request =
230 ResolutionRequest::from_activation(&activation, ResolutionPolicy::LiveOnlyEmbedded);
231 request.target = ResolutionTarget::Delegate {
232 agent_id: "agent-a".into(),
233 parent_run: RunIdentity::for_thread("thread"),
234 persistence: DelegatePersistence::Ephemeral,
235 };
236 assert!(matches!(
237 adapter.resolve(request).await,
238 Err(ResolveError::UnsupportedTarget(_))
239 ));
240 }
241
242 #[tokio::test]
243 async fn runtime_rejects_live_only_plan_for_persistent_policy() {
244 let runtime = crate::AgentRuntime::new(Arc::new(MockResolver));
245 let activation =
246 RunActivation::new("thread", vec![Message::user("hi")]).with_agent_id("agent-a");
247 assert!(matches!(
248 runtime
249 .resolve_activation(&activation, ResolutionPolicy::PersistentServer)
250 .await,
251 Err(ResolveError::UnsupportedPersistence(_))
252 ));
253 }
254
255 struct LiveOnlyEverywhereResolver;
259
260 #[async_trait::async_trait]
261 impl Resolver for LiveOnlyEverywhereResolver {
262 async fn resolve(&self, req: ResolutionRequest) -> Result<ResolvedRunPlan, ResolveError> {
263 let agent_id = match &req.target {
264 ResolutionTarget::Root { agent_id, .. } => agent_id.clone(),
265 ResolutionTarget::Delegate { agent_id, .. } => agent_id.clone(),
266 ResolutionTarget::Handoff { agent_id, .. } => agent_id.clone(),
267 };
268 let role = match &req.target {
269 ResolutionTarget::Root { .. } => ExecutionRole::Root,
270 ResolutionTarget::Delegate { .. } => ExecutionRole::Delegate,
271 ResolutionTarget::Handoff { .. } => ExecutionRole::Handoff,
272 };
273 let requirements = BackendRequirements::from_features(&req.features);
274 let agent = ResolvedAgent::new(&agent_id, "model", "system", Arc::new(MockResolver));
275 Ok(ResolvedRunPlan::LiveOnly(ResolvedRun {
276 agent_spec: (*agent.spec).clone(),
277 role,
278 execution: ExecutionPlan::from_resolved_agent(&agent),
279 model: ResolvedModelBinding {
280 upstream_model: agent.upstream_model.clone(),
281 },
282 tools: Vec::new(),
283 overrides: req.overrides,
284 backend_profile: BackendProfile::full_local(),
285 requirements,
286 scope: LiveOnlyScope,
287 }))
288 }
289 }
290
291 struct ReplayableEverywhereResolver;
292
293 #[async_trait::async_trait]
294 impl Resolver for ReplayableEverywhereResolver {
295 async fn resolve(&self, req: ResolutionRequest) -> Result<ResolvedRunPlan, ResolveError> {
296 let agent_id = match &req.target {
297 ResolutionTarget::Root { agent_id, .. } => agent_id.clone(),
298 ResolutionTarget::Delegate { agent_id, .. } => agent_id.clone(),
299 ResolutionTarget::Handoff { agent_id, .. } => agent_id.clone(),
300 };
301 let role = match &req.target {
302 ResolutionTarget::Root { .. } => ExecutionRole::Root,
303 ResolutionTarget::Delegate { .. } => ExecutionRole::Delegate,
304 ResolutionTarget::Handoff { .. } => ExecutionRole::Handoff,
305 };
306 let requirements = BackendRequirements::from_features(&req.features);
307 let agent = ResolvedAgent::new(&agent_id, "model", "system", Arc::new(MockResolver));
308 Ok(ResolvedRunPlan::Replayable(ReplayableResolvedRun {
309 execution: ResolvedRun {
310 agent_spec: (*agent.spec).clone(),
311 role,
312 execution: ExecutionPlan::from_resolved_agent(&agent),
313 model: ResolvedModelBinding {
314 upstream_model: agent.upstream_model.clone(),
315 },
316 tools: Vec::new(),
317 overrides: req.overrides,
318 backend_profile: BackendProfile::full_local(),
319 requirements,
320 scope: ReplayableScope,
321 },
322 artifact: ResolutionArtifact {
323 resolution_id: "override-publication".to_string(),
324 },
325 }))
326 }
327 }
328
329 fn runtime_with_resolver<R: Resolver + 'static>(r: R) -> crate::AgentRuntime {
330 let runtime = crate::AgentRuntime::new(Arc::new(MockResolver));
331 runtime.set_run_resolver(Arc::new(r));
332 runtime
333 }
334
335 #[tokio::test]
336 async fn nested_resolve_rejects_root_target() {
337 let runtime = runtime_with_resolver(LiveOnlyEverywhereResolver);
338 let sub =
339 RunActivation::new("thread", vec![Message::user("delegate")]).with_agent_id("agent-a");
340 let target = ResolutionTarget::Root {
341 agent_id: "agent-a".into(),
342 thread_id: "thread".into(),
343 };
344 assert!(matches!(
345 runtime
346 .resolve_nested(RootScopeKind::LiveOnly, &sub, target)
347 .await,
348 Err(ResolveError::UnsupportedTarget(_))
349 ));
350 }
351
352 #[tokio::test]
353 async fn nested_resolve_replayable_parent_rejects_live_only_sub() {
354 let runtime = runtime_with_resolver(LiveOnlyEverywhereResolver);
355 let sub =
356 RunActivation::new("thread", vec![Message::user("delegate")]).with_agent_id("agent-a");
357 let target = ResolutionTarget::Delegate {
358 agent_id: "agent-a".into(),
359 parent_run: RunIdentity::for_thread("thread"),
360 persistence: DelegatePersistence::Ephemeral,
361 };
362 assert!(matches!(
363 runtime
364 .resolve_nested(RootScopeKind::Replayable, &sub, target)
365 .await,
366 Err(ResolveError::NestedScopeMismatch(_))
367 ));
368 }
369
370 #[tokio::test]
371 async fn nested_resolve_live_only_parent_accepts_live_only_sub() {
372 let runtime = runtime_with_resolver(LiveOnlyEverywhereResolver);
373 let sub =
374 RunActivation::new("thread", vec![Message::user("delegate")]).with_agent_id("agent-a");
375 let target = ResolutionTarget::Delegate {
376 agent_id: "agent-a".into(),
377 parent_run: RunIdentity::for_thread("thread"),
378 persistence: DelegatePersistence::Ephemeral,
379 };
380 let plan = runtime
381 .resolve_nested(RootScopeKind::LiveOnly, &sub, target)
382 .await
383 .expect("live-only parent accepts live-only sub");
384 assert!(matches!(plan, ResolvedRunPlan::LiveOnly(_)));
385 }
386
387 #[tokio::test]
388 async fn activation_scoped_resolver_overrides_runtime_default() {
389 let runtime = crate::AgentRuntime::new(Arc::new(MockResolver));
390 let activation = RunActivation::new("thread", vec![Message::user("hi")])
391 .with_agent_id("agent-a")
392 .with_run_resolver(Arc::new(ReplayableEverywhereResolver));
393
394 let plan = runtime
395 .resolve_activation(&activation, ResolutionPolicy::PersistentServer)
396 .await
397 .expect("activation resolver supplies replayable plan");
398
399 assert_eq!(plan.resolution_id(), Some("override-publication"));
400 assert_eq!(plan.agent_spec().id, "agent-a");
402 assert_eq!(plan.role(), ExecutionRole::Root);
403 let _ = plan.execution();
404 let _ = plan.backend_profile();
405 let replayable = plan.into_replayable().expect("plan is replayable");
406 assert_eq!(replayable.artifact.resolution_id, "override-publication");
407 }
408}