1use async_trait::async_trait;
6use everruns_core::atoms::{
7 ActAtom, ActInput, ActResult, Atom, InputAtom, InputAtomInput, InputAtomResult, ReasonAtom,
8 ReasonInput, ReasonResult,
9};
10use everruns_core::capabilities::{
11 SystemPromptContext, collect_capabilities_with_configs, report_result_tool_for_child_session,
12 report_task_progress_tool_for_child_session,
13};
14use everruns_core::events::{
15 EventContext, EventRequest, OutputMessageCompletedData, SessionActivatedData, SessionIdledData,
16 TurnCompletedData, TurnFailedData, TurnStartedData,
17};
18use everruns_core::message::{ContentPart, Message};
19use everruns_core::message_retriever::MessageRetriever;
20use everruns_core::platform_store::PlatformStore;
21use everruns_core::session::SessionStatus;
22use everruns_core::tools::Tool;
23use everruns_core::traits::{
24 AgentStore, BudgetChecker, EventEmitter, HarnessStore, ImageArtifactStore, ImageResolver,
25 LeasedResourceStore, PaymentAuthority, ProviderCredentialStore, ProviderStore, ResolvedModel,
26 SessionCreationAuthority, SessionFileSystem, SessionMutator, SessionResourceRegistry,
27 SessionScheduleStore, SessionSqlDbStoreRef, SessionStorageStore, SessionStore,
28 ToolContextServices, UserConnectionResolver,
29};
30use everruns_core::typed_id::{AgentId, HarnessId, MessageId, SessionId, TurnId};
31use everruns_core::vector_store::KnowledgeIndexSearch;
32use everruns_core::{
33 Agent, CapabilityRegistry, CapabilityStatus, DependencyBlocker, DriverRegistry, EgressService,
34 ErrorDisclosure, Harness, Session, TokenUsage, ToolDefinition, ToolRegistry, UserFacingError,
35 UtilityLlmService, assemble_turn_context, org_public_id_from_internal,
36 resolve_runtime_capabilities,
37};
38use std::sync::Arc;
39use tracing::warn;
40
41#[derive(Debug, Clone)]
43pub struct RuntimeHostTurnContext {
44 pub agent: Option<Agent>,
45 pub session: Session,
46 pub messages: Vec<Message>,
47 pub model: Option<ResolvedModel>,
48 pub mcp_tool_definitions: Vec<ToolDefinition>,
49}
50
51#[async_trait]
62pub trait RuntimeHostAdapter: Send + Sync + Clone + 'static {
63 async fn get_agent(
64 &self,
65 org_id: i64,
66 agent_id: AgentId,
67 ) -> everruns_core::error::Result<Option<Agent>>;
68
69 async fn get_harness(
70 &self,
71 org_id: i64,
72 harness_id: HarnessId,
73 ) -> everruns_core::error::Result<Option<Harness>>;
74
75 async fn set_session_status(
76 &self,
77 org_id: i64,
78 session_id: SessionId,
79 status: SessionStatus,
80 ) -> everruns_core::error::Result<Session>;
81
82 async fn load_turn_context(
83 &self,
84 org_id: i64,
85 session_id: SessionId,
86 ) -> everruns_core::error::Result<RuntimeHostTurnContext>;
87
88 fn capability_registry(&self) -> CapabilityRegistry;
89
90 fn driver_registry(&self) -> DriverRegistry;
91
92 fn harness_store(&self, org_id: i64) -> Arc<dyn HarnessStore>;
93
94 fn agent_store(&self, org_id: i64) -> Arc<dyn AgentStore>;
95
96 fn session_store(&self, org_id: i64) -> Arc<dyn SessionStore>;
97
98 fn session_mutator(&self, org_id: i64) -> Arc<dyn SessionMutator>;
99
100 fn provider_store(&self, org_id: i64) -> Arc<dyn ProviderStore>;
101
102 fn message_store(&self) -> Arc<dyn MessageRetriever>;
103
104 fn event_emitter(&self) -> Arc<dyn EventEmitter>;
105
106 fn file_store(&self) -> Arc<dyn SessionFileSystem>;
107
108 fn image_resolver(&self, _org_id: i64) -> Option<Arc<dyn ImageResolver>> {
109 None
110 }
111
112 fn image_artifact_store(&self, _org_id: i64) -> Option<Arc<dyn ImageArtifactStore>> {
113 None
114 }
115
116 fn provider_credential_store(&self, _org_id: i64) -> Option<Arc<dyn ProviderCredentialStore>> {
117 None
118 }
119
120 fn utility_llm_service(&self) -> Option<Arc<dyn UtilityLlmService>> {
121 None
122 }
123
124 fn egress_service(&self) -> Option<Arc<dyn EgressService>> {
125 None
126 }
127
128 fn storage_store(&self) -> Option<Arc<dyn SessionStorageStore>> {
129 None
130 }
131
132 fn knowledge_store(&self) -> Option<Arc<dyn everruns_core::traits::KnowledgeStore>> {
134 None
135 }
136
137 fn connection_resolver(&self) -> Option<Arc<dyn UserConnectionResolver>> {
138 None
139 }
140
141 fn sqldb_store(&self) -> Option<SessionSqlDbStoreRef> {
142 None
143 }
144
145 fn leased_resource_store(&self) -> Option<Arc<dyn LeasedResourceStore>> {
146 None
147 }
148
149 fn session_resource_registry(&self) -> Option<Arc<dyn SessionResourceRegistry>> {
150 None
151 }
152
153 fn session_task_registry(
154 &self,
155 ) -> Option<Arc<dyn everruns_core::session_task::SessionTaskRegistry>> {
156 None
157 }
158
159 fn schedule_store(&self, _org_id: i64) -> Option<Arc<dyn SessionScheduleStore>> {
160 None
161 }
162
163 fn platform_store(
164 &self,
165 _org_id: i64,
166 _session_id: SessionId,
167 ) -> Option<Arc<dyn PlatformStore>> {
168 None
169 }
170
171 fn knowledge_index_search(&self, _org_id: i64) -> Option<Arc<dyn KnowledgeIndexSearch>> {
175 None
176 }
177
178 fn budget_checker(
179 &self,
180 _org_id: i64,
181 _agent_id: Option<AgentId>,
182 ) -> Option<Arc<dyn BudgetChecker>> {
183 None
184 }
185
186 fn payment_authority(
187 &self,
188 _org_id: i64,
189 _agent_id: Option<AgentId>,
190 ) -> Option<Arc<dyn PaymentAuthority>> {
191 None
192 }
193
194 fn session_creation_authority(
195 &self,
196 _org_id: i64,
197 _session_id: SessionId,
198 ) -> Option<Arc<dyn SessionCreationAuthority>> {
199 None
200 }
201
202 fn outbound_tool_rate_limiter(
205 &self,
206 _org_id: i64,
207 ) -> Option<Arc<dyn everruns_core::OutboundToolRateLimiter>> {
208 None
209 }
210
211 fn durable_tool_result_store(&self) -> Option<Arc<dyn everruns_core::DurableToolResultStore>> {
214 None
215 }
216
217 fn subagent_spawn_store(&self) -> Option<Arc<dyn everruns_core::SubagentSpawnStore>> {
220 None
221 }
222
223 fn stream_heartbeater(&self) -> Option<Arc<dyn everruns_core::StreamHeartbeater>> {
226 None
227 }
228
229 fn partial_stream_store(&self) -> Option<Arc<dyn everruns_core::PartialStreamStore>> {
232 None
233 }
234
235 fn reasoning_effort_handle(
244 &self,
245 _session_id: SessionId,
246 ) -> Option<everruns_core::ReasoningEffortHandle> {
247 None
248 }
249
250 fn provider_stall_timeout(&self) -> Option<std::time::Duration> {
253 None
254 }
255
256 async fn mcp_executor(
260 &self,
261 _org_id: i64,
262 _session_id: SessionId,
263 ) -> Option<Arc<everruns_mcp::McpExecutor>> {
264 None
265 }
266}
267
268struct RuntimeExecutionCapabilities {
269 tool_registry: ToolRegistry,
270 post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>>,
271 pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>>,
272 tool_call_hooks: Vec<Arc<dyn everruns_core::ToolCallHook>>,
273 subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
274}
275
276fn subagent_nesting_policy_from_configs(
277 resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
278) -> everruns_core::SubagentNestingPolicy {
279 let subagents_config = resolved_capability_configs.iter().find(|config| {
280 config.capability_id() == everruns_core::capabilities::SUBAGENTS_CAPABILITY_ID
281 });
282
283 let configured_depth = subagents_config
284 .and_then(|config| {
285 config
286 .config
287 .get("max_subagent_depth")
288 .or_else(|| config.config.get("max_depth"))
289 })
290 .and_then(|value| value.as_u64())
291 .and_then(|value| u32::try_from(value).ok());
292 let configured_max_active = subagents_config
293 .and_then(|config| {
294 config
295 .config
296 .get("max_active_descendant_tasks")
297 .or_else(|| config.config.get("max_concurrent_descendant_tasks"))
298 })
299 .and_then(|value| value.as_u64())
300 .and_then(|value| u32::try_from(value).ok());
301 let configured_max_total = subagents_config
302 .and_then(|config| config.config.get("max_total_descendant_tasks"))
303 .and_then(|value| value.as_u64())
304 .and_then(|value| u32::try_from(value).ok());
305 let configured_max_active_detached = subagents_config
306 .and_then(|config| config.config.get("max_active_detached_tasks"))
307 .and_then(|value| value.as_u64())
308 .and_then(|value| u32::try_from(value).ok());
309 let configured_max_total_detached = subagents_config
310 .and_then(|config| config.config.get("max_total_detached_tasks"))
311 .and_then(|value| value.as_u64())
312 .and_then(|value| u32::try_from(value).ok());
313
314 everruns_core::SubagentNestingPolicy::default()
315 .with_agent_override(configured_depth)
316 .with_agent_task_caps_override(configured_max_active, configured_max_total)
317 .with_agent_detached_task_caps_override(
318 configured_max_active_detached,
319 configured_max_total_detached,
320 )
321}
322
323fn finalize_specs_from_configs(
333 resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
334 capability_registry: &CapabilityRegistry,
335) -> Vec<everruns_core::user_hook_types::UserHookSpec> {
336 let mut hook_contributions: Vec<(String, Vec<everruns_core::user_hook_types::UserHookSpec>)> =
337 Vec::new();
338 let mut disabled_contributions: Vec<String> = Vec::new();
339 for config in resolved_capability_configs {
340 let Some(capability) = capability_registry.get(config.capability_id()) else {
341 continue;
342 };
343 let specs = capability.user_hooks_with_config(&config.config);
344 if !specs.is_empty() {
345 hook_contributions.push((config.capability_id().to_string(), specs));
346 }
347 if config.capability_id() == "user_hooks" {
348 disabled_contributions.extend(
349 everruns_core::capabilities::user_hooks::disabled_contributions(&config.config),
350 );
351 }
352 }
353 everruns_core::hook_adapter::finalize_hook_specs(hook_contributions, &disabled_contributions)
354}
355
356async fn collect_lifecycle_hook_specs<A: RuntimeHostAdapter>(
361 adapter: &A,
362 org_id: i64,
363 session_id: SessionId,
364 harness_id: HarnessId,
365 agent_id: Option<AgentId>,
366) -> everruns_core::error::Result<(
367 Vec<everruns_core::user_hook_types::UserHookSpec>,
368 Arc<dyn everruns_core::hook_executor::BashHookDispatcher>,
369)> {
370 let capability_registry = adapter.capability_registry();
371 let harness_chain = adapter
372 .harness_store(org_id)
373 .get_harness_chain(harness_id)
374 .await?;
375 if harness_chain.is_empty() {
376 return Err(everruns_core::error::AgentLoopError::harness_not_found(
377 harness_id,
378 ));
379 }
380 let session = adapter
381 .session_store(org_id)
382 .get_session(session_id)
383 .await?
384 .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
385 let agent = match agent_id {
386 Some(agent_id) => adapter.agent_store(org_id).get_agent(agent_id).await?,
387 None => None,
388 };
389 let resolved = resolve_runtime_capabilities(
390 &harness_chain,
391 agent.as_ref(),
392 &session,
393 &capability_registry,
394 );
395 let specs =
396 finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
397 let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
398 everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
399 );
400 Ok((specs, dispatcher))
401}
402
403async fn load_execution_capabilities<A: RuntimeHostAdapter>(
404 adapter: &A,
405 org_id: i64,
406 session_id: SessionId,
407 harness_id: HarnessId,
408 agent_id: Option<AgentId>,
409 locale: Option<String>,
410 blueprint_id: Option<&str>,
411) -> everruns_core::error::Result<RuntimeExecutionCapabilities> {
412 let capability_registry = adapter.capability_registry();
413 if let Some(blueprint_id) = blueprint_id {
414 let mut registry = ToolRegistry::with_defaults();
415 let blueprint = capability_registry.blueprint(blueprint_id).ok_or_else(|| {
416 everruns_core::error::AgentLoopError::config(format!(
417 "Blueprint \"{blueprint_id}\" not found in registry"
418 ))
419 })?;
420 for tool in blueprint.tools {
421 registry.register_boxed(tool);
422 }
423 return Ok(RuntimeExecutionCapabilities {
424 tool_registry: registry,
425 post_tool_hooks: Vec::new(),
426 pre_tool_hooks: Vec::new(),
427 tool_call_hooks: Vec::new(),
428 subagent_nesting_policy: everruns_core::SubagentNestingPolicy::default(),
429 });
430 }
431
432 let harness_chain = adapter
433 .harness_store(org_id)
434 .get_harness_chain(harness_id)
435 .await?;
436 if harness_chain.is_empty() {
437 return Err(everruns_core::error::AgentLoopError::harness_not_found(
438 harness_id,
439 ));
440 }
441
442 let session = adapter
443 .session_store(org_id)
444 .get_session(session_id)
445 .await?
446 .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
447
448 let agent_store = adapter.agent_store(org_id);
449 let agent = match agent_id {
450 Some(agent_id) => Some(
451 agent_store
452 .get_agent(agent_id)
453 .await?
454 .ok_or_else(|| everruns_core::error::AgentLoopError::agent_not_found(agent_id))?,
455 ),
456 None => None,
457 };
458
459 let resolved = resolve_runtime_capabilities(
460 &harness_chain,
461 agent.as_ref(),
462 &session,
463 &capability_registry,
464 );
465 let prompt_ctx = SystemPromptContext {
472 session_id,
473 locale: locale.or(session.locale.clone()),
474 file_store: Some(everruns_core::scoped_prompt_file_store(
481 adapter.file_store(),
482 session.workspace_id,
483 )),
484 model: None,
485 };
486 let collected = collect_capabilities_with_configs(
487 &resolved.resolved_capability_configs,
488 &capability_registry,
489 &prompt_ctx,
490 )
491 .await;
492
493 let mut registry = ToolRegistry::with_defaults();
494 for tool in collected.tools {
495 registry.register_boxed(tool);
496 }
497
498 let mut post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>> = resolved
503 .resolved_capability_configs
504 .iter()
505 .flat_map(|config| {
506 capability_registry
507 .get(config.capability_id())
508 .filter(|capability| capability.status() == CapabilityStatus::Available)
509 .map(|capability| capability.post_tool_exec_hooks_with_config(&config.config))
510 .unwrap_or_default()
511 })
512 .collect();
513 post_tool_hooks.sort_by_key(|hook| hook.priority());
516
517 let user_hook_specs =
524 finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
525 let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>> = resolved
528 .resolved_capability_configs
529 .iter()
530 .flat_map(|config| {
531 capability_registry
532 .get(config.capability_id())
533 .filter(|capability| capability.status() == CapabilityStatus::Available)
534 .map(|capability| capability.pre_tool_use_hooks_with_config(&config.config))
535 .unwrap_or_default()
536 })
537 .collect();
538 if !user_hook_specs.is_empty() {
539 let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
540 everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
541 );
542 post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
543 &user_hook_specs,
544 dispatcher.clone(),
545 ));
546 pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
547 &user_hook_specs,
548 dispatcher,
549 ));
550 }
551
552 let tool_call_hooks = collected.tool_call_hooks;
563
564 Ok(RuntimeExecutionCapabilities {
565 tool_registry: registry,
566 post_tool_hooks,
567 pre_tool_hooks,
568 tool_call_hooks,
569 subagent_nesting_policy: subagent_nesting_policy_from_configs(
570 &resolved.resolved_capability_configs,
571 ),
572 })
573}
574
575fn runtime_tool_context_services<A: RuntimeHostAdapter>(
576 adapter: &A,
577 org_id: i64,
578 session_id: SessionId,
579 agent_id: Option<AgentId>,
580 tool_registry: Option<Arc<ToolRegistry>>,
581 mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
582 subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
583) -> ToolContextServices {
584 ToolContextServices {
585 file_store: Some(adapter.file_store()),
586 storage_store: adapter.storage_store(),
587 image_store: adapter.image_artifact_store(org_id),
588 provider_credential_store: adapter.provider_credential_store(org_id),
589 utility_llm_service: adapter.utility_llm_service(),
590 mcp_invoker,
591 egress_service: adapter.egress_service(),
592 sqldb_store: adapter.sqldb_store(),
593 message_retriever: Some(adapter.message_store()),
594 session_store: Some(adapter.session_store(org_id)),
595 session_mutator: Some(adapter.session_mutator(org_id)),
596 agent_store: Some(adapter.agent_store(org_id)),
597 connection_resolver: adapter.connection_resolver(),
598 schedule_store: adapter.schedule_store(org_id),
599 platform_store: adapter.platform_store(org_id, session_id),
600 knowledge_store: adapter.knowledge_store(),
601 knowledge_index_search: adapter.knowledge_index_search(org_id),
602 leased_resource_store: adapter.leased_resource_store(),
603 session_resource_registry: adapter.session_resource_registry(),
604 session_task_registry: adapter.session_task_registry(),
605 event_emitter: Some(adapter.event_emitter()),
606 capability_registry: Some(adapter.capability_registry()),
607 tool_registry,
608 org_id: Some(
609 org_public_id_from_internal(org_id)
610 .parse()
611 .expect("internal org id converts to valid public org id"),
612 ),
613 network_access: None,
614 budget_checker: adapter.budget_checker(org_id, agent_id),
615 payment_authority: adapter.payment_authority(org_id, agent_id),
616 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
617 subagent_spawn_store: adapter.subagent_spawn_store(),
618 subagent_nesting_policy,
619 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
620 }
621}
622
623pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
625 adapter: A,
626 org_id: i64,
627 session_id: SessionId,
628}
629
630impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
631 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
632 Self {
633 adapter,
634 org_id,
635 session_id,
636 }
637 }
638
639 async fn set_session_status(&self, status: SessionStatus, action: &'static str) {
640 if let Err(error) = self
641 .adapter
642 .set_session_status(self.org_id, self.session_id, status)
643 .await
644 {
645 warn!(
646 session_id = %self.session_id,
647 org_id = self.org_id,
648 action,
649 %error,
650 "runtime host lifecycle status update failed"
651 );
652 }
653 }
654
655 async fn emit_event(&self, request: EventRequest) {
656 let event_type = request.event_type.clone();
657 if let Err(error) = self.adapter.event_emitter().emit(request).await {
658 warn!(
659 session_id = %self.session_id,
660 org_id = self.org_id,
661 event_type,
662 %error,
663 "runtime host lifecycle event emission failed"
664 );
665 }
666 }
667
668 pub async fn turn_started(&self, turn_id: TurnId, input_message_id: MessageId) {
669 let input_content = self
670 .adapter
671 .message_store()
672 .get(self.session_id, input_message_id)
673 .await
674 .ok()
675 .flatten()
676 .map(|message| message.content_to_llm_string());
677
678 self.set_session_status(SessionStatus::Active, "turn_started")
679 .await;
680
681 self.emit_event(EventRequest::new(
682 self.session_id,
683 EventContext::turn(turn_id, input_message_id),
684 SessionActivatedData {
685 turn_id,
686 input_message_id,
687 },
688 ))
689 .await;
690
691 self.emit_event(EventRequest::new(
692 self.session_id,
693 EventContext::turn(turn_id, input_message_id),
694 TurnStartedData {
695 turn_id,
696 input_message_id,
697 input_content,
698 },
699 ))
700 .await;
701 }
702
703 pub async fn emit_turn_completed(&self, input_message_id: MessageId, data: TurnCompletedData) {
704 let turn_id = data.turn_id;
705 self.emit_event(EventRequest::new(
706 self.session_id,
707 EventContext::turn(turn_id, input_message_id),
708 data,
709 ))
710 .await;
711 }
712
713 pub async fn emit_session_idled(
714 &self,
715 turn_id: TurnId,
716 input_message_id: MessageId,
717 iterations: Option<u32>,
718 usage: Option<TokenUsage>,
719 ) {
720 self.set_session_status(SessionStatus::Idle, "emit_session_idled")
721 .await;
722
723 self.emit_event(EventRequest::new(
724 self.session_id,
725 EventContext::turn(turn_id, input_message_id),
726 SessionIdledData {
727 turn_id,
728 iterations,
729 usage,
730 },
731 ))
732 .await;
733 }
734
735 pub async fn turn_completed(
736 &self,
737 turn_id: TurnId,
738 input_message_id: MessageId,
739 iterations: u32,
740 usage: Option<TokenUsage>,
741 input_content: Option<String>,
742 ) {
743 self.emit_turn_completed(
744 input_message_id,
745 TurnCompletedData {
746 turn_id,
747 iterations,
748 duration_ms: None,
749 usage: usage.clone(),
750 input_content,
751 final_message_id: None,
752 final_answer_preview: None,
753 time_to_first_token_ms: None,
754 tool_call_count: None,
755 llm_call_count: None,
756 status: Some("completed".to_string()),
757 },
758 )
759 .await;
760 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
761 .await;
762 }
763
764 pub async fn turn_sealed(
771 &self,
772 turn_id: TurnId,
773 input_message_id: MessageId,
774 reason: &str,
775 iterations: u32,
776 usage: Option<TokenUsage>,
777 ) {
778 let context = EventContext::turn(turn_id, input_message_id);
779
780 self.emit_event(EventRequest::new(
781 self.session_id,
782 context.clone(),
783 everruns_core::events::TurnSealedData {
784 turn_id,
785 reason: reason.to_string(),
786 detail: None,
787 iterations: Some(iterations),
788 usage: usage.clone(),
789 },
790 ))
791 .await;
792
793 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
794 .await;
795 }
796
797 pub async fn fire_turn_end_hooks(
801 &self,
802 harness_id: HarnessId,
803 agent_id: Option<AgentId>,
804 turn_id: TurnId,
805 success: bool,
806 ) {
807 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
808 &self.adapter,
809 self.org_id,
810 self.session_id,
811 harness_id,
812 agent_id,
813 )
814 .await
815 {
816 Ok(pair) => pair,
817 Err(error) => {
818 warn!(
819 session_id = %self.session_id,
820 %error,
821 "failed to collect turn_end hook specs; skipping"
822 );
823 return;
824 }
825 };
826 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
827 &specs,
828 everruns_core::user_hook_types::HookEvent::TurnEnd,
829 dispatcher,
830 );
831 if hooks.is_empty() {
832 return;
833 }
834 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
835 session_id: self.session_id,
836 turn_id: Some(turn_id),
837 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
838 agent_id: agent_id.map(|a| a.to_string()),
839 };
840 everruns_core::lifecycle_hooks::run_turn_end_hooks(
841 &hooks,
842 &ctx,
843 serde_json::json!({ "success": success }),
844 )
845 .await;
846 }
847
848 pub async fn user_prompt_blocked(
853 &self,
854 turn_id: TurnId,
855 input_message_id: MessageId,
856 reason: &str,
857 user_message: Option<&str>,
858 ) {
859 let user_error =
860 UserFacingError::new(everruns_core::user_facing_error_codes::BLOCKED_BY_HOOK);
861 let shown = user_message.unwrap_or(reason);
862 let mut error_message = Message::assistant(shown);
863 let mut metadata = std::collections::HashMap::new();
864 user_error.apply_to_message_metadata(&mut metadata);
865 error_message.metadata = Some(metadata);
866
867 self.emit_event(EventRequest::new(
868 self.session_id,
869 EventContext::turn(turn_id, input_message_id),
870 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
871 ))
872 .await;
873
874 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
875 .await;
876 }
877
878 pub async fn turn_failed(
879 &self,
880 turn_id: TurnId,
881 input_message_id: MessageId,
882 error: &str,
883 user_error: Option<&UserFacingError>,
884 ) {
885 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
886 .await;
887 }
888
889 pub async fn turn_failed_with_disclosure(
893 &self,
894 turn_id: TurnId,
895 input_message_id: MessageId,
896 error: &str,
897 user_error: Option<&UserFacingError>,
898 disclosure: Option<ErrorDisclosure>,
899 ) {
900 self.set_session_status(SessionStatus::Idle, "turn_failed")
901 .await;
902
903 self.emit_event(EventRequest::new(
904 self.session_id,
905 EventContext::turn(turn_id, input_message_id),
906 {
907 let mut data = TurnFailedData {
908 turn_id,
909 error: error.to_string(),
910 error_code: None,
911 error_fields: None,
912 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
913 };
914 if let Some(user_error) = user_error {
915 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
916 }
917 data
918 },
919 ))
920 .await;
921
922 self.emit_event(EventRequest::new(
923 self.session_id,
924 EventContext::turn(turn_id, input_message_id),
925 SessionIdledData {
926 turn_id,
927 iterations: None,
928 usage: None,
929 },
930 ))
931 .await;
932 }
933
934 pub async fn waiting_for_tool_results(&self) {
935 self.set_session_status(
936 SessionStatus::WaitingForToolResults,
937 "waiting_for_tool_results",
938 )
939 .await;
940 }
941
942 pub async fn dependency_blocked(
943 &self,
944 turn_id: TurnId,
945 input_message_id: MessageId,
946 blocker: DependencyBlocker,
947 ) {
948 let user_error = UserFacingError::new(blocker.error_code())
949 .with_field(
950 "dependency",
951 match blocker {
952 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
953 "harness"
954 }
955 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
956 },
957 )
958 .with_field(
959 "state",
960 match blocker {
961 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
962 "archived"
963 }
964 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
965 "deleted"
966 }
967 },
968 );
969 let mut error_message = Message::assistant(blocker.message());
970 let mut metadata = std::collections::HashMap::new();
971 user_error.apply_to_message_metadata(&mut metadata);
972 error_message.metadata = Some(metadata);
973
974 self.emit_event(EventRequest::new(
975 self.session_id,
976 EventContext::turn(turn_id, input_message_id),
977 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
978 ))
979 .await;
980
981 self.turn_failed(
982 turn_id,
983 input_message_id,
984 blocker.message(),
985 Some(&user_error),
986 )
987 .await;
988 }
989}
990
991pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
992 adapter: &A,
993 org_id: i64,
994 harness_id: HarnessId,
995 agent_id: Option<AgentId>,
996) -> everruns_core::error::Result<Option<DependencyBlocker>> {
997 let harness_store = adapter.harness_store(org_id);
998 let agent_store = adapter.agent_store(org_id);
999 everruns_core::detect_dependency_blocker(
1000 harness_store.as_ref(),
1001 agent_store.as_ref(),
1002 harness_id,
1003 agent_id,
1004 )
1005 .await
1006}
1007
1008pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1009 adapter: &A,
1010 org_id: i64,
1011 input: InputAtomInput,
1012) -> everruns_core::error::Result<InputAtomResult> {
1013 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1017 handle.set(None);
1018 }
1019
1020 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1021 .turn_started(input.context.turn_id, input.context.input_message_id)
1022 .await;
1023
1024 let atom = InputAtom::new(adapter.message_store());
1025 atom.execute(input).await
1026}
1027
1028struct UserPromptHookResult {
1035 decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1036 original_message: String,
1037}
1038
1039async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1040 adapter: &A,
1041 org_id: i64,
1042 input: &ReasonInput,
1043) -> everruns_core::error::Result<Option<UserPromptHookResult>> {
1044 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1045 adapter,
1046 org_id,
1047 input.context.session_id,
1048 input.harness_id,
1049 input.agent_id,
1050 )
1051 .await
1052 {
1053 Ok(pair) => pair,
1054 Err(error) => {
1055 warn!(
1056 session_id = %input.context.session_id,
1057 %error,
1058 "failed to collect user_prompt_submit hook specs; continuing without them"
1059 );
1060 return Ok(None);
1061 }
1062 };
1063 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1064 &specs,
1065 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1066 dispatcher,
1067 );
1068 if hooks.is_empty() {
1069 return Ok(None);
1070 }
1071
1072 let message_text = adapter
1073 .message_store()
1074 .get(input.context.session_id, input.context.input_message_id)
1075 .await
1076 .ok()
1077 .flatten()
1078 .map(|m| m.content_to_llm_string())
1079 .unwrap_or_default();
1080
1081 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1082 session_id: input.context.session_id,
1083 turn_id: Some(input.context.turn_id),
1084 org_id: org_public_id_from_internal(org_id).parse().ok(),
1085 agent_id: input.agent_id.map(|a| a.to_string()),
1086 };
1087 let original_message = message_text.clone();
1088 let decision =
1089 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1090 .await;
1091 Ok(Some(UserPromptHookResult {
1092 decision,
1093 original_message,
1094 }))
1095}
1096
1097pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1098 adapter: &A,
1099 org_id: i64,
1100 input: ReasonInput,
1101) -> everruns_core::error::Result<ReasonResult> {
1102 if let Some(blocker) =
1103 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1104 {
1105 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1106 .dependency_blocked(
1107 input.context.turn_id,
1108 input.context.input_message_id,
1109 blocker,
1110 )
1111 .await;
1112 return Ok(ReasonResult {
1113 success: false,
1114 text: blocker.message().to_string(),
1115 tool_calls: vec![],
1116 has_tool_calls: false,
1117 tool_definitions: vec![],
1118 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1119 error: Some("dependency_unavailable".to_string()),
1120 user_facing_error: None,
1121 error_disclosure: None,
1122 usage: None,
1123 output_message_id: None,
1124 time_to_first_token_ms: None,
1125 response_id: None,
1126 locale: None,
1127 network_access: None,
1128 parallel_tool_calls: None,
1129 });
1130 }
1131
1132 let mut user_prompt_message_override = None;
1140 if input.iteration <= 1
1141 && let Some(hook_result) = run_user_prompt_submit_for_turn(adapter, org_id, &input).await?
1142 {
1143 match hook_result.decision {
1144 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1145 reason,
1146 user_message,
1147 } => {
1148 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1149 .user_prompt_blocked(
1150 input.context.turn_id,
1151 input.context.input_message_id,
1152 &reason,
1153 user_message.as_deref(),
1154 )
1155 .await;
1156 return Ok(ReasonResult {
1157 success: false,
1158 text: user_message.unwrap_or_else(|| reason.clone()),
1159 tool_calls: vec![],
1160 has_tool_calls: false,
1161 tool_definitions: vec![],
1162 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1163 error: Some("blocked_by_user_prompt_hook".to_string()),
1164 user_facing_error: None,
1165 error_disclosure: None,
1166 usage: None,
1167 output_message_id: None,
1168 time_to_first_token_ms: None,
1169 response_id: None,
1170 locale: None,
1171 network_access: None,
1172 parallel_tool_calls: None,
1173 });
1174 }
1175 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1176 if message != hook_result.original_message {
1177 user_prompt_message_override = Some(message);
1178 }
1179 }
1180 }
1181 }
1182
1183 let validation_session = adapter
1187 .session_store(org_id)
1188 .get_session(input.context.session_id)
1189 .await?
1190 .ok_or_else(|| {
1191 everruns_core::error::AgentLoopError::session_not_found(input.context.session_id)
1192 })?;
1193 let validation_capabilities = load_execution_capabilities(
1194 adapter,
1195 org_id,
1196 input.context.session_id,
1197 input.harness_id,
1198 input.agent_id,
1199 validation_session.locale.clone(),
1200 validation_session.blueprint_id.as_deref(),
1201 )
1202 .await?;
1203 let validation_services = runtime_tool_context_services(
1204 adapter,
1205 org_id,
1206 input.context.session_id,
1207 input.agent_id,
1208 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1209 None,
1210 validation_capabilities.subagent_nesting_policy,
1211 );
1212 validation_capabilities
1213 .tool_registry
1214 .validate_context_services(&validation_services)?;
1215
1216 let mut turn_context = adapter
1217 .load_turn_context(org_id, input.context.session_id)
1218 .await?;
1219 if let Some(registry) = adapter.session_task_registry() {
1220 let session_store = adapter.session_store(org_id);
1221 if let Some(tool) = report_result_tool_for_child_session(
1222 input.context.session_id,
1223 session_store.as_ref(),
1224 registry.as_ref(),
1225 )
1226 .await?
1227 {
1228 turn_context.mcp_tool_definitions.push(tool.to_definition());
1229 }
1230 if let Some(tool) = report_task_progress_tool_for_child_session(
1231 input.context.session_id,
1232 session_store.as_ref(),
1233 registry.as_ref(),
1234 )
1235 .await?
1236 {
1237 turn_context.mcp_tool_definitions.push(tool.to_definition());
1238 }
1239 }
1240
1241 let mut atom = ReasonAtom::new(
1242 adapter.harness_store(org_id),
1243 adapter.agent_store(org_id),
1244 adapter.session_store(org_id),
1245 adapter.message_store(),
1246 adapter.provider_store(org_id),
1247 adapter.capability_registry(),
1248 adapter.driver_registry(),
1249 adapter.event_emitter(),
1250 )
1251 .with_file_store(adapter.file_store());
1252 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1253 atom = atom.with_image_resolver(image_resolver);
1254 }
1255 if let Some(hb) = adapter.stream_heartbeater() {
1256 atom = atom.with_stream_heartbeater(hb);
1257 }
1258 if let Some(timeout) = adapter.provider_stall_timeout() {
1259 atom = atom.with_provider_stall_timeout(timeout);
1260 }
1261 if let Some(store) = adapter.partial_stream_store() {
1262 atom = atom.with_partial_stream_store(store);
1263 }
1264 if let Some(store) = adapter.durable_tool_result_store() {
1265 atom = atom.with_durable_tool_result_store(store);
1266 }
1267 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1268 atom = atom.with_reasoning_effort_handle(handle);
1269 }
1270 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1271 atom = atom.with_utility_llm_service(utility_llm_service);
1272 }
1273 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1276 atom = atom.with_schedule_store(schedule_store);
1277 }
1278
1279 let input = ReasonInput {
1280 mcp_tool_definitions: turn_context.mcp_tool_definitions,
1281 ..input
1282 };
1283
1284 if let Some(message_override) = user_prompt_message_override {
1285 let mut assembled = assemble_turn_context(
1286 adapter.harness_store(org_id).as_ref(),
1287 adapter.agent_store(org_id).as_ref(),
1288 adapter.session_store(org_id).as_ref(),
1289 adapter.message_store().as_ref(),
1290 adapter.provider_store(org_id).as_ref(),
1291 &adapter.capability_registry(),
1292 input.context.session_id,
1293 input.harness_id,
1294 input.agent_id,
1295 &input.mcp_tool_definitions,
1296 Some(adapter.file_store()),
1297 )
1298 .await?;
1299
1300 let message = assembled
1301 .messages
1302 .iter_mut()
1303 .find(|message| message.id == input.context.input_message_id)
1304 .ok_or_else(|| {
1305 everruns_core::error::AgentLoopError::config(
1306 "user_prompt_submit mutation: input message not found in assembled context",
1307 )
1308 })?;
1309
1310 message
1315 .content
1316 .retain(|part| !matches!(part, ContentPart::Text(_)));
1317 message
1318 .content
1319 .insert(0, ContentPart::text(message_override));
1320
1321 return atom.execute_with_assembled_context(input, assembled).await;
1322 }
1323
1324 atom.execute(input).await
1325}
1326
1327pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1328 adapter: &A,
1329 input: ActInput,
1330) -> everruns_core::error::Result<ActResult> {
1331 let org_id = input.org_id.ok_or_else(|| {
1332 everruns_core::error::AgentLoopError::config(
1333 "ActInput.org_id must be set for runtime host execution",
1334 )
1335 })?;
1336
1337 if let Some(blocker) =
1338 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1339 {
1340 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1341 .dependency_blocked(
1342 input.context.turn_id,
1343 input.context.input_message_id,
1344 blocker,
1345 )
1346 .await;
1347 return Ok(ActResult {
1348 results: vec![],
1349 completed: true,
1350 success_count: 0,
1351 error_count: 1,
1352 waiting_for_tool_results: false,
1353 blocked: true,
1354 client_tool_calls: vec![],
1355 client_tool_definitions: vec![],
1356 });
1357 }
1358
1359 let execution_capabilities = load_execution_capabilities(
1360 adapter,
1361 org_id,
1362 input.context.session_id,
1363 input.harness_id,
1364 input.agent_id,
1365 input.locale.clone(),
1366 input.blueprint_id.as_deref(),
1367 )
1368 .await?;
1369 let mut tool_registry = execution_capabilities.tool_registry;
1370
1371 if input
1372 .tool_definitions
1373 .iter()
1374 .any(|definition| definition.name() == "report_result")
1375 && let Some(registry) = adapter.session_task_registry()
1376 && let Some(tool) = report_result_tool_for_child_session(
1377 input.context.session_id,
1378 adapter.session_store(org_id).as_ref(),
1379 registry.as_ref(),
1380 )
1381 .await?
1382 {
1383 tool_registry.register_boxed(Box::new(tool.with_file_store(adapter.file_store())));
1384 }
1385 if input
1386 .tool_definitions
1387 .iter()
1388 .any(|definition| definition.name() == "report_task_progress")
1389 && let Some(registry) = adapter.session_task_registry()
1390 && let Some(tool) = report_task_progress_tool_for_child_session(
1391 input.context.session_id,
1392 adapter.session_store(org_id).as_ref(),
1393 registry.as_ref(),
1394 )
1395 .await?
1396 {
1397 tool_registry.register_boxed(Box::new(tool));
1398 }
1399
1400 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1410 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1411 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1412 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1413 tool_registry.register_boxed(tool);
1414 }
1415 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1416 &input.tool_definitions,
1417 invoker,
1418 )));
1419 }
1420
1421 let builtin_tool_registry = Arc::new(tool_registry.clone());
1422 let context_services = runtime_tool_context_services(
1423 adapter,
1424 org_id,
1425 input.context.session_id,
1426 input.agent_id,
1427 Some(builtin_tool_registry),
1428 mcp_invoker,
1429 execution_capabilities.subagent_nesting_policy,
1430 );
1431 tool_registry.validate_context_services(&context_services)?;
1432 let executor: Arc<dyn everruns_core::traits::ToolExecutor> = Arc::new(tool_registry);
1433
1434 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1435 .with_context_services(context_services)
1436 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1437 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1438 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1439
1440 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1441 atom = atom.with_outbound_tool_rate_limiter(limiter);
1442 }
1443 if let Some(store) = adapter.durable_tool_result_store() {
1444 atom = atom.with_durable_tool_result_store(store);
1445 }
1446
1447 atom.execute(input).await
1448}