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