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