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::session::SessionStatus;
21use everruns_core::tools::Tool;
22use everruns_core::traits::{
23 AgentStore, BudgetChecker, EventEmitter, HarnessStore, ImageArtifactStore, ImageResolver,
24 LeasedResourceStore, PaymentAuthority, ProviderCredentialStore, ProviderStore, ResolvedModel,
25 SessionCreationAuthority, SessionFileSystem, SessionMutator, SessionResourceRegistry,
26 SessionScheduleStore, SessionSqlDbStoreRef, SessionStorageStore, SessionStore,
27 ToolContextServices, UserConnectionResolver,
28};
29use everruns_core::typed_id::{AgentId, HarnessId, MessageId, SessionId, TurnId};
30use everruns_core::vector_store::KnowledgeIndexSearch;
31use everruns_core::{
32 Agent, CapabilityRegistry, CapabilityStatus, DependencyBlocker, DriverRegistry, EgressService,
33 ErrorDisclosure, Harness, Session, TokenUsage, ToolDefinition, ToolRegistry, UserFacingError,
34 UtilityLlmService, assemble_turn_context, org_public_id_from_internal,
35 resolve_runtime_capabilities,
36};
37use everruns_platform::PlatformStore;
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 let platform_store = adapter.platform_store(org_id, session_id);
611 let subagent_delegate = platform_store.clone().map(|store| {
612 Arc::new(everruns_platform::PlatformStoreSubagentDelegate(store))
613 as Arc<dyn everruns_core::subagent_delegation::SubagentSessionDelegate>
614 });
615 let extensions = {
616 let mut extensions = everruns_core::traits::ToolContextExtensions::default();
617 if let Some(store) = platform_store {
618 extensions.insert(Arc::new(everruns_platform::PlatformStoreExt(store)));
619 }
620 extensions
621 };
622 ToolContextServices {
623 file_store: Some(adapter.file_store()),
624 storage_store: adapter.storage_store(),
625 image_store: adapter.image_artifact_store(org_id),
626 provider_credential_store: adapter.provider_credential_store(org_id),
627 utility_llm_service: adapter.utility_llm_service(),
628 mcp_invoker,
629 egress_service: adapter.egress_service(),
630 sqldb_store: adapter.sqldb_store(),
631 message_retriever: Some(adapter.message_store()),
632 session_store: Some(adapter.session_store(org_id)),
633 session_mutator: Some(adapter.session_mutator(org_id)),
634 agent_store: Some(adapter.agent_store(org_id)),
635 connection_resolver: adapter.connection_resolver(),
636 schedule_store: adapter.schedule_store(org_id),
637 subagent_delegate,
638 extensions,
639 knowledge_store: adapter.knowledge_store(),
640 knowledge_index_search: adapter.knowledge_index_search(org_id),
641 leased_resource_store: adapter.leased_resource_store(),
642 session_resource_registry: adapter.session_resource_registry(),
643 session_task_registry: adapter.session_task_registry(),
644 event_emitter: Some(adapter.event_emitter()),
645 capability_registry: Some(adapter.capability_registry()),
646 tool_registry,
647 org_id: Some(
648 org_public_id_from_internal(org_id)
649 .parse()
650 .expect("internal org id converts to valid public org id"),
651 ),
652 network_access: None,
653 budget_checker: adapter.budget_checker(org_id, agent_id),
654 payment_authority: adapter.payment_authority(org_id, agent_id),
655 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
656 subagent_spawn_store: adapter.subagent_spawn_store(),
657 subagent_nesting_policy,
658 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
659 }
660}
661
662pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
664 adapter: A,
665 org_id: i64,
666 session_id: SessionId,
667}
668
669impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
670 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
671 Self {
672 adapter,
673 org_id,
674 session_id,
675 }
676 }
677
678 async fn set_session_status(&self, status: SessionStatus, action: &'static str) {
679 if let Err(error) = self
680 .adapter
681 .set_session_status(self.org_id, self.session_id, status)
682 .await
683 {
684 warn!(
685 session_id = %self.session_id,
686 org_id = self.org_id,
687 action,
688 %error,
689 "runtime host lifecycle status update failed"
690 );
691 }
692 }
693
694 async fn emit_event(&self, request: EventRequest) {
695 let event_type = request.event_type.clone();
696 if let Err(error) = self.adapter.event_emitter().emit(request).await {
697 warn!(
698 session_id = %self.session_id,
699 org_id = self.org_id,
700 event_type,
701 %error,
702 "runtime host lifecycle event emission failed"
703 );
704 }
705 }
706
707 pub async fn turn_started(&self, turn_id: TurnId, input_message_id: MessageId) {
708 let input_content = self
709 .adapter
710 .message_store()
711 .get(self.session_id, input_message_id)
712 .await
713 .ok()
714 .flatten()
715 .map(|message| message.content_to_llm_string());
716
717 self.set_session_status(SessionStatus::Active, "turn_started")
718 .await;
719
720 self.emit_event(EventRequest::new(
721 self.session_id,
722 EventContext::turn(turn_id, input_message_id),
723 SessionActivatedData {
724 turn_id,
725 input_message_id,
726 },
727 ))
728 .await;
729
730 self.emit_event(EventRequest::new(
731 self.session_id,
732 EventContext::turn(turn_id, input_message_id),
733 TurnStartedData {
734 turn_id,
735 input_message_id,
736 input_content,
737 },
738 ))
739 .await;
740 }
741
742 pub async fn emit_turn_completed(&self, input_message_id: MessageId, data: TurnCompletedData) {
743 let turn_id = data.turn_id;
744 self.emit_event(EventRequest::new(
745 self.session_id,
746 EventContext::turn(turn_id, input_message_id),
747 data,
748 ))
749 .await;
750 }
751
752 pub async fn emit_session_idled(
753 &self,
754 turn_id: TurnId,
755 input_message_id: MessageId,
756 iterations: Option<u32>,
757 usage: Option<TokenUsage>,
758 ) {
759 self.set_session_status(SessionStatus::Idle, "emit_session_idled")
760 .await;
761
762 self.emit_event(EventRequest::new(
763 self.session_id,
764 EventContext::turn(turn_id, input_message_id),
765 SessionIdledData {
766 turn_id,
767 iterations,
768 usage,
769 },
770 ))
771 .await;
772 }
773
774 pub async fn turn_completed(
775 &self,
776 turn_id: TurnId,
777 input_message_id: MessageId,
778 iterations: u32,
779 usage: Option<TokenUsage>,
780 input_content: Option<String>,
781 ) {
782 self.emit_turn_completed(
783 input_message_id,
784 TurnCompletedData {
785 turn_id,
786 iterations,
787 duration_ms: None,
788 usage: usage.clone(),
789 input_content,
790 final_message_id: None,
791 final_answer_preview: None,
792 time_to_first_token_ms: None,
793 tool_call_count: None,
794 llm_call_count: None,
795 status: Some("completed".to_string()),
796 },
797 )
798 .await;
799 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
800 .await;
801 }
802
803 pub async fn turn_sealed(
810 &self,
811 turn_id: TurnId,
812 input_message_id: MessageId,
813 reason: &str,
814 iterations: u32,
815 usage: Option<TokenUsage>,
816 ) {
817 let context = EventContext::turn(turn_id, input_message_id);
818
819 self.emit_event(EventRequest::new(
820 self.session_id,
821 context.clone(),
822 everruns_core::events::TurnSealedData {
823 turn_id,
824 reason: reason.to_string(),
825 detail: None,
826 iterations: Some(iterations),
827 usage: usage.clone(),
828 },
829 ))
830 .await;
831
832 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
833 .await;
834 }
835
836 pub async fn fire_turn_end_hooks(
840 &self,
841 harness_id: HarnessId,
842 agent_id: Option<AgentId>,
843 turn_id: TurnId,
844 success: bool,
845 ) {
846 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
847 &self.adapter,
848 self.org_id,
849 self.session_id,
850 harness_id,
851 agent_id,
852 )
853 .await
854 {
855 Ok(pair) => pair,
856 Err(error) => {
857 warn!(
858 session_id = %self.session_id,
859 %error,
860 "failed to collect turn_end hook specs; skipping"
861 );
862 return;
863 }
864 };
865 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
866 &specs,
867 everruns_core::user_hook_types::HookEvent::TurnEnd,
868 dispatcher,
869 );
870 if hooks.is_empty() {
871 return;
872 }
873 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
874 session_id: self.session_id,
875 turn_id: Some(turn_id),
876 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
877 agent_id: agent_id.map(|a| a.to_string()),
878 };
879 everruns_core::lifecycle_hooks::run_turn_end_hooks(
880 &hooks,
881 &ctx,
882 serde_json::json!({ "success": success }),
883 )
884 .await;
885 }
886
887 pub async fn user_prompt_blocked(
892 &self,
893 turn_id: TurnId,
894 input_message_id: MessageId,
895 reason: &str,
896 user_message: Option<&str>,
897 ) {
898 let user_error =
899 UserFacingError::new(everruns_core::user_facing_error_codes::BLOCKED_BY_HOOK);
900 let shown = user_message.unwrap_or(reason);
901 let mut error_message = Message::assistant(shown);
902 let mut metadata = std::collections::HashMap::new();
903 user_error.apply_to_message_metadata(&mut metadata);
904 error_message.metadata = Some(metadata);
905
906 self.emit_event(EventRequest::new(
907 self.session_id,
908 EventContext::turn(turn_id, input_message_id),
909 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
910 ))
911 .await;
912
913 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
914 .await;
915 }
916
917 pub async fn turn_failed(
918 &self,
919 turn_id: TurnId,
920 input_message_id: MessageId,
921 error: &str,
922 user_error: Option<&UserFacingError>,
923 ) {
924 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
925 .await;
926 }
927
928 pub async fn turn_failed_with_disclosure(
932 &self,
933 turn_id: TurnId,
934 input_message_id: MessageId,
935 error: &str,
936 user_error: Option<&UserFacingError>,
937 disclosure: Option<ErrorDisclosure>,
938 ) {
939 self.set_session_status(SessionStatus::Idle, "turn_failed")
940 .await;
941
942 self.emit_event(EventRequest::new(
943 self.session_id,
944 EventContext::turn(turn_id, input_message_id),
945 {
946 let mut data = TurnFailedData {
947 turn_id,
948 error: error.to_string(),
949 error_code: None,
950 error_fields: None,
951 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
952 };
953 if let Some(user_error) = user_error {
954 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
955 }
956 data
957 },
958 ))
959 .await;
960
961 self.emit_event(EventRequest::new(
962 self.session_id,
963 EventContext::turn(turn_id, input_message_id),
964 SessionIdledData {
965 turn_id,
966 iterations: None,
967 usage: None,
968 },
969 ))
970 .await;
971 }
972
973 pub async fn waiting_for_tool_results(&self) {
974 self.set_session_status(
975 SessionStatus::WaitingForToolResults,
976 "waiting_for_tool_results",
977 )
978 .await;
979 }
980
981 pub async fn dependency_blocked(
982 &self,
983 turn_id: TurnId,
984 input_message_id: MessageId,
985 blocker: DependencyBlocker,
986 ) {
987 let user_error = UserFacingError::new(blocker.error_code())
988 .with_field(
989 "dependency",
990 match blocker {
991 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
992 "harness"
993 }
994 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
995 },
996 )
997 .with_field(
998 "state",
999 match blocker {
1000 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
1001 "archived"
1002 }
1003 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
1004 "deleted"
1005 }
1006 },
1007 );
1008 let mut error_message = Message::assistant(blocker.message());
1009 let mut metadata = std::collections::HashMap::new();
1010 user_error.apply_to_message_metadata(&mut metadata);
1011 error_message.metadata = Some(metadata);
1012
1013 self.emit_event(EventRequest::new(
1014 self.session_id,
1015 EventContext::turn(turn_id, input_message_id),
1016 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1017 ))
1018 .await;
1019
1020 self.turn_failed(
1021 turn_id,
1022 input_message_id,
1023 blocker.message(),
1024 Some(&user_error),
1025 )
1026 .await;
1027 }
1028}
1029
1030pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1031 adapter: &A,
1032 org_id: i64,
1033 harness_id: HarnessId,
1034 agent_id: Option<AgentId>,
1035) -> everruns_core::error::Result<Option<DependencyBlocker>> {
1036 let harness_store = adapter.harness_store(org_id);
1037 let agent_store = adapter.agent_store(org_id);
1038 everruns_core::detect_dependency_blocker(
1039 harness_store.as_ref(),
1040 agent_store.as_ref(),
1041 harness_id,
1042 agent_id,
1043 )
1044 .await
1045}
1046
1047pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1048 adapter: &A,
1049 org_id: i64,
1050 input: InputAtomInput,
1051) -> everruns_core::error::Result<InputAtomResult> {
1052 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1056 handle.set(None);
1057 }
1058
1059 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1060 .turn_started(input.context.turn_id, input.context.input_message_id)
1061 .await;
1062
1063 let atom = InputAtom::new(adapter.message_store());
1064 atom.execute(input).await
1065}
1066
1067struct UserPromptHookResult {
1074 decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1075 original_message: String,
1076}
1077
1078async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1079 adapter: &A,
1080 org_id: i64,
1081 input: &ReasonInput,
1082) -> everruns_core::error::Result<Option<UserPromptHookResult>> {
1083 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1084 adapter,
1085 org_id,
1086 input.context.session_id,
1087 input.harness_id,
1088 input.agent_id,
1089 )
1090 .await
1091 {
1092 Ok(pair) => pair,
1093 Err(error) => {
1094 warn!(
1095 session_id = %input.context.session_id,
1096 %error,
1097 "failed to collect user_prompt_submit hook specs; continuing without them"
1098 );
1099 return Ok(None);
1100 }
1101 };
1102 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1103 &specs,
1104 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1105 dispatcher,
1106 );
1107 if hooks.is_empty() {
1108 return Ok(None);
1109 }
1110
1111 let message_text = adapter
1112 .message_store()
1113 .get(input.context.session_id, input.context.input_message_id)
1114 .await
1115 .ok()
1116 .flatten()
1117 .map(|m| m.content_to_llm_string())
1118 .unwrap_or_default();
1119
1120 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1121 session_id: input.context.session_id,
1122 turn_id: Some(input.context.turn_id),
1123 org_id: org_public_id_from_internal(org_id).parse().ok(),
1124 agent_id: input.agent_id.map(|a| a.to_string()),
1125 };
1126 let original_message = message_text.clone();
1127 let decision =
1128 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1129 .await;
1130 Ok(Some(UserPromptHookResult {
1131 decision,
1132 original_message,
1133 }))
1134}
1135
1136pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1137 adapter: &A,
1138 org_id: i64,
1139 input: ReasonInput,
1140) -> everruns_core::error::Result<ReasonResult> {
1141 let prompt_message_ids = (input.iteration <= 1)
1142 .then_some(input.context.input_message_id)
1143 .into_iter()
1144 .collect();
1145 execute_reason_activity_with_prompt_messages(adapter, org_id, input, prompt_message_ids).await
1146}
1147
1148pub async fn execute_reason_activity_with_prompt_messages<A: RuntimeHostAdapter>(
1153 adapter: &A,
1154 org_id: i64,
1155 input: ReasonInput,
1156 prompt_message_ids: Vec<MessageId>,
1157) -> everruns_core::error::Result<ReasonResult> {
1158 if let Some(blocker) =
1159 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1160 {
1161 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1162 .dependency_blocked(
1163 input.context.turn_id,
1164 input.context.input_message_id,
1165 blocker,
1166 )
1167 .await;
1168 return Ok(ReasonResult {
1169 success: false,
1170 text: blocker.message().to_string(),
1171 tool_calls: vec![],
1172 has_tool_calls: false,
1173 tool_definitions: vec![],
1174 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1175 error: Some("dependency_unavailable".to_string()),
1176 user_facing_error: None,
1177 error_disclosure: None,
1178 usage: None,
1179 output_message_id: None,
1180 time_to_first_token_ms: None,
1181 response_id: None,
1182 finish_reason: None,
1183 locale: None,
1184 network_access: None,
1185 parallel_tool_calls: None,
1186 });
1187 }
1188
1189 let mut user_prompt_message_overrides = Vec::new();
1194 for message_id in prompt_message_ids {
1195 let mut hook_input = input.clone();
1196 hook_input.context.input_message_id = message_id;
1197 let Some(hook_result) =
1198 run_user_prompt_submit_for_turn(adapter, org_id, &hook_input).await?
1199 else {
1200 continue;
1201 };
1202 match hook_result.decision {
1203 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1204 reason,
1205 user_message,
1206 } => {
1207 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1208 .user_prompt_blocked(
1209 input.context.turn_id,
1210 input.context.input_message_id,
1211 &reason,
1212 user_message.as_deref(),
1213 )
1214 .await;
1215 return Ok(ReasonResult {
1216 success: false,
1217 text: user_message.unwrap_or_else(|| reason.clone()),
1218 tool_calls: vec![],
1219 has_tool_calls: false,
1220 tool_definitions: vec![],
1221 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1222 error: Some("blocked_by_user_prompt_hook".to_string()),
1223 user_facing_error: None,
1224 error_disclosure: None,
1225 usage: None,
1226 output_message_id: None,
1227 time_to_first_token_ms: None,
1228 response_id: None,
1229 finish_reason: None,
1230 locale: None,
1231 network_access: None,
1232 parallel_tool_calls: None,
1233 });
1234 }
1235 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1236 if message != hook_result.original_message {
1237 user_prompt_message_overrides.push((message_id, message));
1238 }
1239 }
1240 }
1241 }
1242
1243 let validation_session = adapter
1247 .session_store(org_id)
1248 .get_session(input.context.session_id)
1249 .await?
1250 .ok_or_else(|| {
1251 everruns_core::error::AgentLoopError::session_not_found(input.context.session_id)
1252 })?;
1253 let validation_capabilities = load_execution_capabilities(
1254 adapter,
1255 org_id,
1256 input.context.session_id,
1257 input.harness_id,
1258 input.agent_id,
1259 validation_session.locale.clone(),
1260 validation_session.blueprint_id.as_deref(),
1261 )
1262 .await?;
1263 let query_history_allowed = validation_capabilities
1264 .tool_registry
1265 .get("query_history")
1266 .is_some();
1267 let validation_services = runtime_tool_context_services(
1268 adapter,
1269 org_id,
1270 input.context.session_id,
1271 input.agent_id,
1272 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1273 None,
1274 validation_capabilities.subagent_nesting_policy,
1275 );
1276 validation_capabilities
1277 .tool_registry
1278 .validate_context_services(&validation_services)?;
1279
1280 let mut turn_context = adapter
1281 .load_turn_context(org_id, input.context.session_id)
1282 .await?;
1283 if let Some(registry) = adapter.session_task_registry() {
1284 let session_store = adapter.session_store(org_id);
1285 if let Some(tool) = report_result_tool_for_child_session(
1286 input.context.session_id,
1287 session_store.as_ref(),
1288 registry.as_ref(),
1289 )
1290 .await?
1291 {
1292 turn_context.mcp_tool_definitions.push(tool.to_definition());
1293 }
1294 if let Some(tool) = report_task_progress_tool_for_child_session(
1295 input.context.session_id,
1296 session_store.as_ref(),
1297 registry.as_ref(),
1298 )
1299 .await?
1300 {
1301 turn_context.mcp_tool_definitions.push(tool.to_definition());
1302 }
1303 }
1304
1305 let mut reason_capability_registry = adapter.capability_registry();
1306 if !query_history_allowed {
1307 reason_capability_registry
1311 .register(everruns_core::capabilities::InfinityContextFilterOnlyCapability);
1312 }
1313 let mut atom = ReasonAtom::new(
1314 adapter.harness_store(org_id),
1315 adapter.agent_store(org_id),
1316 adapter.session_store(org_id),
1317 adapter.message_store(),
1318 adapter.provider_store(org_id),
1319 reason_capability_registry.clone(),
1320 adapter.driver_registry(),
1321 adapter.event_emitter(),
1322 )
1323 .with_file_store(adapter.file_store());
1324 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1325 atom = atom.with_image_resolver(image_resolver);
1326 }
1327 if let Some(hb) = adapter.stream_heartbeater() {
1328 atom = atom.with_stream_heartbeater(hb);
1329 }
1330 if let Some(timeout) = adapter.provider_stall_timeout() {
1331 atom = atom.with_provider_stall_timeout(timeout);
1332 }
1333 if let Some(config) = adapter.provider_retry_config() {
1334 atom = atom.with_provider_retry_config(config);
1335 }
1336 if let Some(store) = adapter.partial_stream_store() {
1337 atom = atom.with_partial_stream_store(store);
1338 }
1339 if let Some(store) = adapter.durable_tool_result_store() {
1340 atom = atom.with_durable_tool_result_store(store);
1341 }
1342 if let Some(store) = adapter.compaction_checkpoint_store() {
1343 atom = atom.with_compaction_checkpoint_store(store);
1344 }
1345 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1346 atom = atom.with_reasoning_effort_handle(handle);
1347 }
1348 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1349 atom = atom.with_utility_llm_service(utility_llm_service);
1350 }
1351 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1354 atom = atom.with_schedule_store(schedule_store);
1355 }
1356
1357 let input = ReasonInput {
1358 mcp_tool_definitions: turn_context.mcp_tool_definitions,
1359 ..input
1360 };
1361
1362 if !user_prompt_message_overrides.is_empty() {
1363 let mut assembled = assemble_turn_context(
1364 adapter.harness_store(org_id).as_ref(),
1365 adapter.agent_store(org_id).as_ref(),
1366 adapter.session_store(org_id).as_ref(),
1367 adapter.message_store().as_ref(),
1368 adapter.provider_store(org_id).as_ref(),
1369 &reason_capability_registry,
1370 input.context.session_id,
1371 input.harness_id,
1372 input.agent_id,
1373 &input.mcp_tool_definitions,
1374 Some(adapter.file_store()),
1375 )
1376 .await?;
1377
1378 for (message_id, message_override) in user_prompt_message_overrides {
1379 let message = assembled
1380 .messages
1381 .iter_mut()
1382 .find(|message| message.id == message_id)
1383 .ok_or_else(|| {
1384 everruns_core::error::AgentLoopError::config(
1385 "user_prompt_submit mutation: input message not found in assembled context",
1386 )
1387 })?;
1388
1389 message
1392 .content
1393 .retain(|part| !matches!(part, ContentPart::Text(_)));
1394 message
1395 .content
1396 .insert(0, ContentPart::text(message_override));
1397 }
1398
1399 return atom.execute_with_assembled_context(input, assembled).await;
1400 }
1401
1402 atom.execute(input).await
1403}
1404
1405pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1406 adapter: &A,
1407 input: ActInput,
1408) -> everruns_core::error::Result<ActResult> {
1409 let org_id = input.org_id.ok_or_else(|| {
1410 everruns_core::error::AgentLoopError::config(
1411 "ActInput.org_id must be set for runtime host execution",
1412 )
1413 })?;
1414
1415 if let Some(blocker) =
1416 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1417 {
1418 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1419 .dependency_blocked(
1420 input.context.turn_id,
1421 input.context.input_message_id,
1422 blocker,
1423 )
1424 .await;
1425 return Ok(ActResult {
1426 results: vec![],
1427 completed: true,
1428 success_count: 0,
1429 error_count: 1,
1430 waiting_for_tool_results: false,
1431 blocked: true,
1432 client_tool_calls: vec![],
1433 client_tool_definitions: vec![],
1434 });
1435 }
1436
1437 let execution_capabilities = load_execution_capabilities(
1438 adapter,
1439 org_id,
1440 input.context.session_id,
1441 input.harness_id,
1442 input.agent_id,
1443 input.locale.clone(),
1444 input.blueprint_id.as_deref(),
1445 )
1446 .await?;
1447 let mut tool_registry = execution_capabilities.tool_registry;
1448
1449 if input
1450 .tool_definitions
1451 .iter()
1452 .any(|definition| definition.name() == "report_result")
1453 && let Some(registry) = adapter.session_task_registry()
1454 && let Some(tool) = report_result_tool_for_child_session(
1455 input.context.session_id,
1456 adapter.session_store(org_id).as_ref(),
1457 registry.as_ref(),
1458 )
1459 .await?
1460 {
1461 tool_registry.register_boxed(Box::new(tool.with_file_store(adapter.file_store())));
1462 }
1463 if input
1464 .tool_definitions
1465 .iter()
1466 .any(|definition| definition.name() == "report_task_progress")
1467 && let Some(registry) = adapter.session_task_registry()
1468 && let Some(tool) = report_task_progress_tool_for_child_session(
1469 input.context.session_id,
1470 adapter.session_store(org_id).as_ref(),
1471 registry.as_ref(),
1472 )
1473 .await?
1474 {
1475 tool_registry.register_boxed(Box::new(tool));
1476 }
1477
1478 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1488 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1489 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1490 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1491 tool_registry.register_boxed(tool);
1492 }
1493 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1494 &input.tool_definitions,
1495 invoker,
1496 )));
1497 }
1498
1499 let builtin_tool_registry = Arc::new(tool_registry.clone());
1500 let context_services = runtime_tool_context_services(
1501 adapter,
1502 org_id,
1503 input.context.session_id,
1504 input.agent_id,
1505 Some(builtin_tool_registry),
1506 mcp_invoker,
1507 execution_capabilities.subagent_nesting_policy,
1508 );
1509 tool_registry.validate_context_services(&context_services)?;
1510 let executor: Arc<dyn everruns_core::traits::ToolExecutor> = Arc::new(tool_registry);
1511
1512 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1513 .with_context_services(context_services)
1514 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1515 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1516 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1517
1518 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1519 atom = atom.with_outbound_tool_rate_limiter(limiter);
1520 }
1521 if let Some(store) = adapter.durable_tool_result_store() {
1522 atom = atom.with_durable_tool_result_store(store);
1523 }
1524
1525 atom.execute(input).await
1526}