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 async fn mcp_executor(
266 &self,
267 _org_id: i64,
268 _session_id: SessionId,
269 ) -> Option<Arc<everruns_mcp::McpExecutor>> {
270 None
271 }
272}
273
274struct RuntimeExecutionCapabilities {
275 tool_registry: ToolRegistry,
276 post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>>,
277 pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>>,
278 tool_call_hooks: Vec<Arc<dyn everruns_core::ToolCallHook>>,
279 subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
280}
281
282fn subagent_nesting_policy_from_configs(
283 resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
284) -> everruns_core::SubagentNestingPolicy {
285 let subagents_config = resolved_capability_configs.iter().find(|config| {
286 config.capability_id() == everruns_core::capabilities::SUBAGENTS_CAPABILITY_ID
287 });
288
289 let configured_depth = subagents_config
290 .and_then(|config| {
291 config
292 .config
293 .get("max_subagent_depth")
294 .or_else(|| config.config.get("max_depth"))
295 })
296 .and_then(|value| value.as_u64())
297 .and_then(|value| u32::try_from(value).ok());
298 let configured_max_active = subagents_config
299 .and_then(|config| {
300 config
301 .config
302 .get("max_active_descendant_tasks")
303 .or_else(|| config.config.get("max_concurrent_descendant_tasks"))
304 })
305 .and_then(|value| value.as_u64())
306 .and_then(|value| u32::try_from(value).ok());
307 let configured_max_total = subagents_config
308 .and_then(|config| config.config.get("max_total_descendant_tasks"))
309 .and_then(|value| value.as_u64())
310 .and_then(|value| u32::try_from(value).ok());
311 let configured_max_active_detached = subagents_config
312 .and_then(|config| config.config.get("max_active_detached_tasks"))
313 .and_then(|value| value.as_u64())
314 .and_then(|value| u32::try_from(value).ok());
315 let configured_max_total_detached = subagents_config
316 .and_then(|config| config.config.get("max_total_detached_tasks"))
317 .and_then(|value| value.as_u64())
318 .and_then(|value| u32::try_from(value).ok());
319
320 everruns_core::SubagentNestingPolicy::default()
321 .with_agent_override(configured_depth)
322 .with_agent_task_caps_override(configured_max_active, configured_max_total)
323 .with_agent_detached_task_caps_override(
324 configured_max_active_detached,
325 configured_max_total_detached,
326 )
327}
328
329fn finalize_specs_from_configs(
339 resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
340 capability_registry: &CapabilityRegistry,
341) -> Vec<everruns_core::user_hook_types::UserHookSpec> {
342 let mut hook_contributions: Vec<(String, Vec<everruns_core::user_hook_types::UserHookSpec>)> =
343 Vec::new();
344 let mut disabled_contributions: Vec<String> = Vec::new();
345 for config in resolved_capability_configs {
346 let Some(capability) = capability_registry.get(config.capability_id()) else {
347 continue;
348 };
349 let specs = capability.user_hooks_with_config(&config.config);
350 if !specs.is_empty() {
351 hook_contributions.push((config.capability_id().to_string(), specs));
352 }
353 if config.capability_id() == "user_hooks" {
354 disabled_contributions.extend(
355 everruns_core::capabilities::user_hooks::disabled_contributions(&config.config),
356 );
357 }
358 }
359 everruns_core::hook_adapter::finalize_hook_specs(hook_contributions, &disabled_contributions)
360}
361
362async fn collect_lifecycle_hook_specs<A: RuntimeHostAdapter>(
367 adapter: &A,
368 org_id: i64,
369 session_id: SessionId,
370 harness_id: HarnessId,
371 agent_id: Option<AgentId>,
372) -> everruns_core::error::Result<(
373 Vec<everruns_core::user_hook_types::UserHookSpec>,
374 Arc<dyn everruns_core::hook_executor::BashHookDispatcher>,
375)> {
376 let capability_registry = adapter.capability_registry();
377 let harness_chain = adapter
378 .harness_store(org_id)
379 .get_harness_chain(harness_id)
380 .await?;
381 if harness_chain.is_empty() {
382 return Err(everruns_core::error::AgentLoopError::harness_not_found(
383 harness_id,
384 ));
385 }
386 let session = adapter
387 .session_store(org_id)
388 .get_session(session_id)
389 .await?
390 .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
391 let agent = match agent_id {
392 Some(agent_id) => adapter.agent_store(org_id).get_agent(agent_id).await?,
393 None => None,
394 };
395 let resolved = resolve_runtime_capabilities(
396 &harness_chain,
397 agent.as_ref(),
398 &session,
399 &capability_registry,
400 );
401 let specs =
402 finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
403 let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
404 everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
405 );
406 Ok((specs, dispatcher))
407}
408
409async fn load_execution_capabilities<A: RuntimeHostAdapter>(
410 adapter: &A,
411 org_id: i64,
412 session_id: SessionId,
413 harness_id: HarnessId,
414 agent_id: Option<AgentId>,
415 locale: Option<String>,
416 blueprint_id: Option<&str>,
417) -> everruns_core::error::Result<RuntimeExecutionCapabilities> {
418 let capability_registry = adapter.capability_registry();
419 if let Some(blueprint_id) = blueprint_id {
420 let mut registry = ToolRegistry::with_defaults();
421 let blueprint = capability_registry.blueprint(blueprint_id).ok_or_else(|| {
422 everruns_core::error::AgentLoopError::config(format!(
423 "Blueprint \"{blueprint_id}\" not found in registry"
424 ))
425 })?;
426 for tool in blueprint.tools {
427 registry.register_boxed(tool);
428 }
429 return Ok(RuntimeExecutionCapabilities {
430 tool_registry: registry,
431 post_tool_hooks: Vec::new(),
432 pre_tool_hooks: Vec::new(),
433 tool_call_hooks: Vec::new(),
434 subagent_nesting_policy: everruns_core::SubagentNestingPolicy::default(),
435 });
436 }
437
438 let harness_chain = adapter
439 .harness_store(org_id)
440 .get_harness_chain(harness_id)
441 .await?;
442 if harness_chain.is_empty() {
443 return Err(everruns_core::error::AgentLoopError::harness_not_found(
444 harness_id,
445 ));
446 }
447
448 let session = adapter
449 .session_store(org_id)
450 .get_session(session_id)
451 .await?
452 .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
453
454 let agent_store = adapter.agent_store(org_id);
455 let agent = match agent_id {
456 Some(agent_id) => Some(
457 agent_store
458 .get_agent(agent_id)
459 .await?
460 .ok_or_else(|| everruns_core::error::AgentLoopError::agent_not_found(agent_id))?,
461 ),
462 None => None,
463 };
464
465 let resolved = resolve_runtime_capabilities(
466 &harness_chain,
467 agent.as_ref(),
468 &session,
469 &capability_registry,
470 );
471 let prompt_ctx = SystemPromptContext {
478 session_id,
479 locale: locale.or(session.locale.clone()),
480 file_store: Some(everruns_core::scoped_prompt_file_store(
487 adapter.file_store(),
488 session.workspace_id,
489 )),
490 model: None,
491 };
492 let collected = collect_capabilities_with_configs(
493 &resolved.resolved_capability_configs,
494 &capability_registry,
495 &prompt_ctx,
496 )
497 .await;
498
499 let mut registry = ToolRegistry::with_defaults();
500 for tool in collected.tools {
501 registry.register_boxed(tool);
502 }
503
504 let mut post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>> = resolved
509 .resolved_capability_configs
510 .iter()
511 .flat_map(|config| {
512 capability_registry
513 .get(config.capability_id())
514 .filter(|capability| capability.status() == CapabilityStatus::Available)
515 .map(|capability| capability.post_tool_exec_hooks_with_config(&config.config))
516 .unwrap_or_default()
517 })
518 .collect();
519 post_tool_hooks.sort_by_key(|hook| hook.priority());
522
523 let user_hook_specs =
530 finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
531 let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>> = resolved
534 .resolved_capability_configs
535 .iter()
536 .flat_map(|config| {
537 capability_registry
538 .get(config.capability_id())
539 .filter(|capability| capability.status() == CapabilityStatus::Available)
540 .map(|capability| capability.pre_tool_use_hooks_with_config(&config.config))
541 .unwrap_or_default()
542 })
543 .collect();
544 if !user_hook_specs.is_empty() {
545 let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
546 everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
547 );
548 post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
549 &user_hook_specs,
550 dispatcher.clone(),
551 ));
552 pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
553 &user_hook_specs,
554 dispatcher,
555 ));
556 }
557
558 let tool_call_hooks = collected.tool_call_hooks;
569
570 Ok(RuntimeExecutionCapabilities {
571 tool_registry: registry,
572 post_tool_hooks,
573 pre_tool_hooks,
574 tool_call_hooks,
575 subagent_nesting_policy: subagent_nesting_policy_from_configs(
576 &resolved.resolved_capability_configs,
577 ),
578 })
579}
580
581fn runtime_tool_context_services<A: RuntimeHostAdapter>(
582 adapter: &A,
583 org_id: i64,
584 session_id: SessionId,
585 agent_id: Option<AgentId>,
586 tool_registry: Option<Arc<ToolRegistry>>,
587 mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
588 subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
589) -> ToolContextServices {
590 ToolContextServices {
591 file_store: Some(adapter.file_store()),
592 storage_store: adapter.storage_store(),
593 image_store: adapter.image_artifact_store(org_id),
594 provider_credential_store: adapter.provider_credential_store(org_id),
595 utility_llm_service: adapter.utility_llm_service(),
596 mcp_invoker,
597 egress_service: adapter.egress_service(),
598 sqldb_store: adapter.sqldb_store(),
599 message_retriever: Some(adapter.message_store()),
600 session_store: Some(adapter.session_store(org_id)),
601 session_mutator: Some(adapter.session_mutator(org_id)),
602 agent_store: Some(adapter.agent_store(org_id)),
603 connection_resolver: adapter.connection_resolver(),
604 schedule_store: adapter.schedule_store(org_id),
605 platform_store: adapter.platform_store(org_id, session_id),
606 knowledge_store: adapter.knowledge_store(),
607 knowledge_index_search: adapter.knowledge_index_search(org_id),
608 leased_resource_store: adapter.leased_resource_store(),
609 session_resource_registry: adapter.session_resource_registry(),
610 session_task_registry: adapter.session_task_registry(),
611 event_emitter: Some(adapter.event_emitter()),
612 capability_registry: Some(adapter.capability_registry()),
613 tool_registry,
614 org_id: Some(
615 org_public_id_from_internal(org_id)
616 .parse()
617 .expect("internal org id converts to valid public org id"),
618 ),
619 network_access: None,
620 budget_checker: adapter.budget_checker(org_id, agent_id),
621 payment_authority: adapter.payment_authority(org_id, agent_id),
622 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
623 subagent_spawn_store: adapter.subagent_spawn_store(),
624 subagent_nesting_policy,
625 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
626 }
627}
628
629pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
631 adapter: A,
632 org_id: i64,
633 session_id: SessionId,
634}
635
636impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
637 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
638 Self {
639 adapter,
640 org_id,
641 session_id,
642 }
643 }
644
645 async fn set_session_status(&self, status: SessionStatus, action: &'static str) {
646 if let Err(error) = self
647 .adapter
648 .set_session_status(self.org_id, self.session_id, status)
649 .await
650 {
651 warn!(
652 session_id = %self.session_id,
653 org_id = self.org_id,
654 action,
655 %error,
656 "runtime host lifecycle status update failed"
657 );
658 }
659 }
660
661 async fn emit_event(&self, request: EventRequest) {
662 let event_type = request.event_type.clone();
663 if let Err(error) = self.adapter.event_emitter().emit(request).await {
664 warn!(
665 session_id = %self.session_id,
666 org_id = self.org_id,
667 event_type,
668 %error,
669 "runtime host lifecycle event emission failed"
670 );
671 }
672 }
673
674 pub async fn turn_started(&self, turn_id: TurnId, input_message_id: MessageId) {
675 let input_content = self
676 .adapter
677 .message_store()
678 .get(self.session_id, input_message_id)
679 .await
680 .ok()
681 .flatten()
682 .map(|message| message.content_to_llm_string());
683
684 self.set_session_status(SessionStatus::Active, "turn_started")
685 .await;
686
687 self.emit_event(EventRequest::new(
688 self.session_id,
689 EventContext::turn(turn_id, input_message_id),
690 SessionActivatedData {
691 turn_id,
692 input_message_id,
693 },
694 ))
695 .await;
696
697 self.emit_event(EventRequest::new(
698 self.session_id,
699 EventContext::turn(turn_id, input_message_id),
700 TurnStartedData {
701 turn_id,
702 input_message_id,
703 input_content,
704 },
705 ))
706 .await;
707 }
708
709 pub async fn emit_turn_completed(&self, input_message_id: MessageId, data: TurnCompletedData) {
710 let turn_id = data.turn_id;
711 self.emit_event(EventRequest::new(
712 self.session_id,
713 EventContext::turn(turn_id, input_message_id),
714 data,
715 ))
716 .await;
717 }
718
719 pub async fn emit_session_idled(
720 &self,
721 turn_id: TurnId,
722 input_message_id: MessageId,
723 iterations: Option<u32>,
724 usage: Option<TokenUsage>,
725 ) {
726 self.set_session_status(SessionStatus::Idle, "emit_session_idled")
727 .await;
728
729 self.emit_event(EventRequest::new(
730 self.session_id,
731 EventContext::turn(turn_id, input_message_id),
732 SessionIdledData {
733 turn_id,
734 iterations,
735 usage,
736 },
737 ))
738 .await;
739 }
740
741 pub async fn turn_completed(
742 &self,
743 turn_id: TurnId,
744 input_message_id: MessageId,
745 iterations: u32,
746 usage: Option<TokenUsage>,
747 input_content: Option<String>,
748 ) {
749 self.emit_turn_completed(
750 input_message_id,
751 TurnCompletedData {
752 turn_id,
753 iterations,
754 duration_ms: None,
755 usage: usage.clone(),
756 input_content,
757 final_message_id: None,
758 final_answer_preview: None,
759 time_to_first_token_ms: None,
760 tool_call_count: None,
761 llm_call_count: None,
762 status: Some("completed".to_string()),
763 },
764 )
765 .await;
766 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
767 .await;
768 }
769
770 pub async fn turn_sealed(
777 &self,
778 turn_id: TurnId,
779 input_message_id: MessageId,
780 reason: &str,
781 iterations: u32,
782 usage: Option<TokenUsage>,
783 ) {
784 let context = EventContext::turn(turn_id, input_message_id);
785
786 self.emit_event(EventRequest::new(
787 self.session_id,
788 context.clone(),
789 everruns_core::events::TurnSealedData {
790 turn_id,
791 reason: reason.to_string(),
792 detail: None,
793 iterations: Some(iterations),
794 usage: usage.clone(),
795 },
796 ))
797 .await;
798
799 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
800 .await;
801 }
802
803 pub async fn fire_turn_end_hooks(
807 &self,
808 harness_id: HarnessId,
809 agent_id: Option<AgentId>,
810 turn_id: TurnId,
811 success: bool,
812 ) {
813 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
814 &self.adapter,
815 self.org_id,
816 self.session_id,
817 harness_id,
818 agent_id,
819 )
820 .await
821 {
822 Ok(pair) => pair,
823 Err(error) => {
824 warn!(
825 session_id = %self.session_id,
826 %error,
827 "failed to collect turn_end hook specs; skipping"
828 );
829 return;
830 }
831 };
832 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
833 &specs,
834 everruns_core::user_hook_types::HookEvent::TurnEnd,
835 dispatcher,
836 );
837 if hooks.is_empty() {
838 return;
839 }
840 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
841 session_id: self.session_id,
842 turn_id: Some(turn_id),
843 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
844 agent_id: agent_id.map(|a| a.to_string()),
845 };
846 everruns_core::lifecycle_hooks::run_turn_end_hooks(
847 &hooks,
848 &ctx,
849 serde_json::json!({ "success": success }),
850 )
851 .await;
852 }
853
854 pub async fn user_prompt_blocked(
859 &self,
860 turn_id: TurnId,
861 input_message_id: MessageId,
862 reason: &str,
863 user_message: Option<&str>,
864 ) {
865 let user_error =
866 UserFacingError::new(everruns_core::user_facing_error_codes::BLOCKED_BY_HOOK);
867 let shown = user_message.unwrap_or(reason);
868 let mut error_message = Message::assistant(shown);
869 let mut metadata = std::collections::HashMap::new();
870 user_error.apply_to_message_metadata(&mut metadata);
871 error_message.metadata = Some(metadata);
872
873 self.emit_event(EventRequest::new(
874 self.session_id,
875 EventContext::turn(turn_id, input_message_id),
876 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
877 ))
878 .await;
879
880 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
881 .await;
882 }
883
884 pub async fn turn_failed(
885 &self,
886 turn_id: TurnId,
887 input_message_id: MessageId,
888 error: &str,
889 user_error: Option<&UserFacingError>,
890 ) {
891 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
892 .await;
893 }
894
895 pub async fn turn_failed_with_disclosure(
899 &self,
900 turn_id: TurnId,
901 input_message_id: MessageId,
902 error: &str,
903 user_error: Option<&UserFacingError>,
904 disclosure: Option<ErrorDisclosure>,
905 ) {
906 self.set_session_status(SessionStatus::Idle, "turn_failed")
907 .await;
908
909 self.emit_event(EventRequest::new(
910 self.session_id,
911 EventContext::turn(turn_id, input_message_id),
912 {
913 let mut data = TurnFailedData {
914 turn_id,
915 error: error.to_string(),
916 error_code: None,
917 error_fields: None,
918 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
919 };
920 if let Some(user_error) = user_error {
921 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
922 }
923 data
924 },
925 ))
926 .await;
927
928 self.emit_event(EventRequest::new(
929 self.session_id,
930 EventContext::turn(turn_id, input_message_id),
931 SessionIdledData {
932 turn_id,
933 iterations: None,
934 usage: None,
935 },
936 ))
937 .await;
938 }
939
940 pub async fn waiting_for_tool_results(&self) {
941 self.set_session_status(
942 SessionStatus::WaitingForToolResults,
943 "waiting_for_tool_results",
944 )
945 .await;
946 }
947
948 pub async fn dependency_blocked(
949 &self,
950 turn_id: TurnId,
951 input_message_id: MessageId,
952 blocker: DependencyBlocker,
953 ) {
954 let user_error = UserFacingError::new(blocker.error_code())
955 .with_field(
956 "dependency",
957 match blocker {
958 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
959 "harness"
960 }
961 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
962 },
963 )
964 .with_field(
965 "state",
966 match blocker {
967 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
968 "archived"
969 }
970 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
971 "deleted"
972 }
973 },
974 );
975 let mut error_message = Message::assistant(blocker.message());
976 let mut metadata = std::collections::HashMap::new();
977 user_error.apply_to_message_metadata(&mut metadata);
978 error_message.metadata = Some(metadata);
979
980 self.emit_event(EventRequest::new(
981 self.session_id,
982 EventContext::turn(turn_id, input_message_id),
983 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
984 ))
985 .await;
986
987 self.turn_failed(
988 turn_id,
989 input_message_id,
990 blocker.message(),
991 Some(&user_error),
992 )
993 .await;
994 }
995}
996
997pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
998 adapter: &A,
999 org_id: i64,
1000 harness_id: HarnessId,
1001 agent_id: Option<AgentId>,
1002) -> everruns_core::error::Result<Option<DependencyBlocker>> {
1003 let harness_store = adapter.harness_store(org_id);
1004 let agent_store = adapter.agent_store(org_id);
1005 everruns_core::detect_dependency_blocker(
1006 harness_store.as_ref(),
1007 agent_store.as_ref(),
1008 harness_id,
1009 agent_id,
1010 )
1011 .await
1012}
1013
1014pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1015 adapter: &A,
1016 org_id: i64,
1017 input: InputAtomInput,
1018) -> everruns_core::error::Result<InputAtomResult> {
1019 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1023 handle.set(None);
1024 }
1025
1026 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1027 .turn_started(input.context.turn_id, input.context.input_message_id)
1028 .await;
1029
1030 let atom = InputAtom::new(adapter.message_store());
1031 atom.execute(input).await
1032}
1033
1034struct UserPromptHookResult {
1041 decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1042 original_message: String,
1043}
1044
1045async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1046 adapter: &A,
1047 org_id: i64,
1048 input: &ReasonInput,
1049) -> everruns_core::error::Result<Option<UserPromptHookResult>> {
1050 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1051 adapter,
1052 org_id,
1053 input.context.session_id,
1054 input.harness_id,
1055 input.agent_id,
1056 )
1057 .await
1058 {
1059 Ok(pair) => pair,
1060 Err(error) => {
1061 warn!(
1062 session_id = %input.context.session_id,
1063 %error,
1064 "failed to collect user_prompt_submit hook specs; continuing without them"
1065 );
1066 return Ok(None);
1067 }
1068 };
1069 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1070 &specs,
1071 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1072 dispatcher,
1073 );
1074 if hooks.is_empty() {
1075 return Ok(None);
1076 }
1077
1078 let message_text = adapter
1079 .message_store()
1080 .get(input.context.session_id, input.context.input_message_id)
1081 .await
1082 .ok()
1083 .flatten()
1084 .map(|m| m.content_to_llm_string())
1085 .unwrap_or_default();
1086
1087 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1088 session_id: input.context.session_id,
1089 turn_id: Some(input.context.turn_id),
1090 org_id: org_public_id_from_internal(org_id).parse().ok(),
1091 agent_id: input.agent_id.map(|a| a.to_string()),
1092 };
1093 let original_message = message_text.clone();
1094 let decision =
1095 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1096 .await;
1097 Ok(Some(UserPromptHookResult {
1098 decision,
1099 original_message,
1100 }))
1101}
1102
1103pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1104 adapter: &A,
1105 org_id: i64,
1106 input: ReasonInput,
1107) -> everruns_core::error::Result<ReasonResult> {
1108 if let Some(blocker) =
1109 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1110 {
1111 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1112 .dependency_blocked(
1113 input.context.turn_id,
1114 input.context.input_message_id,
1115 blocker,
1116 )
1117 .await;
1118 return Ok(ReasonResult {
1119 success: false,
1120 text: blocker.message().to_string(),
1121 tool_calls: vec![],
1122 has_tool_calls: false,
1123 tool_definitions: vec![],
1124 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1125 error: Some("dependency_unavailable".to_string()),
1126 user_facing_error: None,
1127 error_disclosure: None,
1128 usage: None,
1129 output_message_id: None,
1130 time_to_first_token_ms: None,
1131 response_id: None,
1132 finish_reason: None,
1133 locale: None,
1134 network_access: None,
1135 parallel_tool_calls: None,
1136 });
1137 }
1138
1139 let mut user_prompt_message_override = None;
1147 if input.iteration <= 1
1148 && let Some(hook_result) = run_user_prompt_submit_for_turn(adapter, org_id, &input).await?
1149 {
1150 match hook_result.decision {
1151 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1152 reason,
1153 user_message,
1154 } => {
1155 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1156 .user_prompt_blocked(
1157 input.context.turn_id,
1158 input.context.input_message_id,
1159 &reason,
1160 user_message.as_deref(),
1161 )
1162 .await;
1163 return Ok(ReasonResult {
1164 success: false,
1165 text: user_message.unwrap_or_else(|| reason.clone()),
1166 tool_calls: vec![],
1167 has_tool_calls: false,
1168 tool_definitions: vec![],
1169 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1170 error: Some("blocked_by_user_prompt_hook".to_string()),
1171 user_facing_error: None,
1172 error_disclosure: None,
1173 usage: None,
1174 output_message_id: None,
1175 time_to_first_token_ms: None,
1176 response_id: None,
1177 finish_reason: None,
1178 locale: None,
1179 network_access: None,
1180 parallel_tool_calls: None,
1181 });
1182 }
1183 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1184 if message != hook_result.original_message {
1185 user_prompt_message_override = Some(message);
1186 }
1187 }
1188 }
1189 }
1190
1191 let validation_session = adapter
1195 .session_store(org_id)
1196 .get_session(input.context.session_id)
1197 .await?
1198 .ok_or_else(|| {
1199 everruns_core::error::AgentLoopError::session_not_found(input.context.session_id)
1200 })?;
1201 let validation_capabilities = load_execution_capabilities(
1202 adapter,
1203 org_id,
1204 input.context.session_id,
1205 input.harness_id,
1206 input.agent_id,
1207 validation_session.locale.clone(),
1208 validation_session.blueprint_id.as_deref(),
1209 )
1210 .await?;
1211 let validation_services = runtime_tool_context_services(
1212 adapter,
1213 org_id,
1214 input.context.session_id,
1215 input.agent_id,
1216 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1217 None,
1218 validation_capabilities.subagent_nesting_policy,
1219 );
1220 validation_capabilities
1221 .tool_registry
1222 .validate_context_services(&validation_services)?;
1223
1224 let mut turn_context = adapter
1225 .load_turn_context(org_id, input.context.session_id)
1226 .await?;
1227 if let Some(registry) = adapter.session_task_registry() {
1228 let session_store = adapter.session_store(org_id);
1229 if let Some(tool) = report_result_tool_for_child_session(
1230 input.context.session_id,
1231 session_store.as_ref(),
1232 registry.as_ref(),
1233 )
1234 .await?
1235 {
1236 turn_context.mcp_tool_definitions.push(tool.to_definition());
1237 }
1238 if let Some(tool) = report_task_progress_tool_for_child_session(
1239 input.context.session_id,
1240 session_store.as_ref(),
1241 registry.as_ref(),
1242 )
1243 .await?
1244 {
1245 turn_context.mcp_tool_definitions.push(tool.to_definition());
1246 }
1247 }
1248
1249 let mut atom = ReasonAtom::new(
1250 adapter.harness_store(org_id),
1251 adapter.agent_store(org_id),
1252 adapter.session_store(org_id),
1253 adapter.message_store(),
1254 adapter.provider_store(org_id),
1255 adapter.capability_registry(),
1256 adapter.driver_registry(),
1257 adapter.event_emitter(),
1258 )
1259 .with_file_store(adapter.file_store());
1260 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1261 atom = atom.with_image_resolver(image_resolver);
1262 }
1263 if let Some(hb) = adapter.stream_heartbeater() {
1264 atom = atom.with_stream_heartbeater(hb);
1265 }
1266 if let Some(timeout) = adapter.provider_stall_timeout() {
1267 atom = atom.with_provider_stall_timeout(timeout);
1268 }
1269 if let Some(store) = adapter.partial_stream_store() {
1270 atom = atom.with_partial_stream_store(store);
1271 }
1272 if let Some(store) = adapter.durable_tool_result_store() {
1273 atom = atom.with_durable_tool_result_store(store);
1274 }
1275 if let Some(store) = adapter.compaction_checkpoint_store() {
1276 atom = atom.with_compaction_checkpoint_store(store);
1277 }
1278 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1279 atom = atom.with_reasoning_effort_handle(handle);
1280 }
1281 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1282 atom = atom.with_utility_llm_service(utility_llm_service);
1283 }
1284 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1287 atom = atom.with_schedule_store(schedule_store);
1288 }
1289
1290 let input = ReasonInput {
1291 mcp_tool_definitions: turn_context.mcp_tool_definitions,
1292 ..input
1293 };
1294
1295 if let Some(message_override) = user_prompt_message_override {
1296 let mut assembled = assemble_turn_context(
1297 adapter.harness_store(org_id).as_ref(),
1298 adapter.agent_store(org_id).as_ref(),
1299 adapter.session_store(org_id).as_ref(),
1300 adapter.message_store().as_ref(),
1301 adapter.provider_store(org_id).as_ref(),
1302 &adapter.capability_registry(),
1303 input.context.session_id,
1304 input.harness_id,
1305 input.agent_id,
1306 &input.mcp_tool_definitions,
1307 Some(adapter.file_store()),
1308 )
1309 .await?;
1310
1311 let message = assembled
1312 .messages
1313 .iter_mut()
1314 .find(|message| message.id == input.context.input_message_id)
1315 .ok_or_else(|| {
1316 everruns_core::error::AgentLoopError::config(
1317 "user_prompt_submit mutation: input message not found in assembled context",
1318 )
1319 })?;
1320
1321 message
1326 .content
1327 .retain(|part| !matches!(part, ContentPart::Text(_)));
1328 message
1329 .content
1330 .insert(0, ContentPart::text(message_override));
1331
1332 return atom.execute_with_assembled_context(input, assembled).await;
1333 }
1334
1335 atom.execute(input).await
1336}
1337
1338pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1339 adapter: &A,
1340 input: ActInput,
1341) -> everruns_core::error::Result<ActResult> {
1342 let org_id = input.org_id.ok_or_else(|| {
1343 everruns_core::error::AgentLoopError::config(
1344 "ActInput.org_id must be set for runtime host execution",
1345 )
1346 })?;
1347
1348 if let Some(blocker) =
1349 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1350 {
1351 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1352 .dependency_blocked(
1353 input.context.turn_id,
1354 input.context.input_message_id,
1355 blocker,
1356 )
1357 .await;
1358 return Ok(ActResult {
1359 results: vec![],
1360 completed: true,
1361 success_count: 0,
1362 error_count: 1,
1363 waiting_for_tool_results: false,
1364 blocked: true,
1365 client_tool_calls: vec![],
1366 client_tool_definitions: vec![],
1367 });
1368 }
1369
1370 let execution_capabilities = load_execution_capabilities(
1371 adapter,
1372 org_id,
1373 input.context.session_id,
1374 input.harness_id,
1375 input.agent_id,
1376 input.locale.clone(),
1377 input.blueprint_id.as_deref(),
1378 )
1379 .await?;
1380 let mut tool_registry = execution_capabilities.tool_registry;
1381
1382 if input
1383 .tool_definitions
1384 .iter()
1385 .any(|definition| definition.name() == "report_result")
1386 && let Some(registry) = adapter.session_task_registry()
1387 && let Some(tool) = report_result_tool_for_child_session(
1388 input.context.session_id,
1389 adapter.session_store(org_id).as_ref(),
1390 registry.as_ref(),
1391 )
1392 .await?
1393 {
1394 tool_registry.register_boxed(Box::new(tool.with_file_store(adapter.file_store())));
1395 }
1396 if input
1397 .tool_definitions
1398 .iter()
1399 .any(|definition| definition.name() == "report_task_progress")
1400 && let Some(registry) = adapter.session_task_registry()
1401 && let Some(tool) = report_task_progress_tool_for_child_session(
1402 input.context.session_id,
1403 adapter.session_store(org_id).as_ref(),
1404 registry.as_ref(),
1405 )
1406 .await?
1407 {
1408 tool_registry.register_boxed(Box::new(tool));
1409 }
1410
1411 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1421 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1422 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1423 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1424 tool_registry.register_boxed(tool);
1425 }
1426 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1427 &input.tool_definitions,
1428 invoker,
1429 )));
1430 }
1431
1432 let builtin_tool_registry = Arc::new(tool_registry.clone());
1433 let context_services = runtime_tool_context_services(
1434 adapter,
1435 org_id,
1436 input.context.session_id,
1437 input.agent_id,
1438 Some(builtin_tool_registry),
1439 mcp_invoker,
1440 execution_capabilities.subagent_nesting_policy,
1441 );
1442 tool_registry.validate_context_services(&context_services)?;
1443 let executor: Arc<dyn everruns_core::traits::ToolExecutor> = Arc::new(tool_registry);
1444
1445 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1446 .with_context_services(context_services)
1447 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1448 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1449 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1450
1451 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1452 atom = atom.with_outbound_tool_rate_limiter(limiter);
1453 }
1454 if let Some(store) = adapter.durable_tool_result_store() {
1455 atom = atom.with_durable_tool_result_store(store);
1456 }
1457
1458 atom.execute(input).await
1459}