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 if user_hook_specs
536 .iter()
537 .any(|spec| spec.event == everruns_core::user_hook_types::HookEvent::UserPromptSubmit)
538 {
539 registry.unregister("query_history");
540 }
541 let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>> = resolved
544 .resolved_capability_configs
545 .iter()
546 .flat_map(|config| {
547 capability_registry
548 .get(config.capability_id())
549 .filter(|capability| capability.status() == CapabilityStatus::Available)
550 .map(|capability| capability.pre_tool_use_hooks_with_config(&config.config))
551 .unwrap_or_default()
552 })
553 .collect();
554 if !user_hook_specs.is_empty() {
555 let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
556 everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
557 );
558 post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
559 &user_hook_specs,
560 dispatcher.clone(),
561 ));
562 pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
563 &user_hook_specs,
564 dispatcher,
565 ));
566 }
567
568 let tool_call_hooks = collected.tool_call_hooks;
579
580 Ok(RuntimeExecutionCapabilities {
581 tool_registry: registry,
582 post_tool_hooks,
583 pre_tool_hooks,
584 tool_call_hooks,
585 subagent_nesting_policy: subagent_nesting_policy_from_configs(
586 &resolved.resolved_capability_configs,
587 ),
588 })
589}
590
591fn runtime_tool_context_services<A: RuntimeHostAdapter>(
592 adapter: &A,
593 org_id: i64,
594 session_id: SessionId,
595 agent_id: Option<AgentId>,
596 tool_registry: Option<Arc<ToolRegistry>>,
597 mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
598 subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
599) -> ToolContextServices {
600 ToolContextServices {
601 file_store: Some(adapter.file_store()),
602 storage_store: adapter.storage_store(),
603 image_store: adapter.image_artifact_store(org_id),
604 provider_credential_store: adapter.provider_credential_store(org_id),
605 utility_llm_service: adapter.utility_llm_service(),
606 mcp_invoker,
607 egress_service: adapter.egress_service(),
608 sqldb_store: adapter.sqldb_store(),
609 message_retriever: Some(adapter.message_store()),
610 session_store: Some(adapter.session_store(org_id)),
611 session_mutator: Some(adapter.session_mutator(org_id)),
612 agent_store: Some(adapter.agent_store(org_id)),
613 connection_resolver: adapter.connection_resolver(),
614 schedule_store: adapter.schedule_store(org_id),
615 platform_store: adapter.platform_store(org_id, session_id),
616 knowledge_store: adapter.knowledge_store(),
617 knowledge_index_search: adapter.knowledge_index_search(org_id),
618 leased_resource_store: adapter.leased_resource_store(),
619 session_resource_registry: adapter.session_resource_registry(),
620 session_task_registry: adapter.session_task_registry(),
621 event_emitter: Some(adapter.event_emitter()),
622 capability_registry: Some(adapter.capability_registry()),
623 tool_registry,
624 org_id: Some(
625 org_public_id_from_internal(org_id)
626 .parse()
627 .expect("internal org id converts to valid public org id"),
628 ),
629 network_access: None,
630 budget_checker: adapter.budget_checker(org_id, agent_id),
631 payment_authority: adapter.payment_authority(org_id, agent_id),
632 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
633 subagent_spawn_store: adapter.subagent_spawn_store(),
634 subagent_nesting_policy,
635 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
636 }
637}
638
639pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
641 adapter: A,
642 org_id: i64,
643 session_id: SessionId,
644}
645
646impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
647 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
648 Self {
649 adapter,
650 org_id,
651 session_id,
652 }
653 }
654
655 async fn set_session_status(&self, status: SessionStatus, action: &'static str) {
656 if let Err(error) = self
657 .adapter
658 .set_session_status(self.org_id, self.session_id, status)
659 .await
660 {
661 warn!(
662 session_id = %self.session_id,
663 org_id = self.org_id,
664 action,
665 %error,
666 "runtime host lifecycle status update failed"
667 );
668 }
669 }
670
671 async fn emit_event(&self, request: EventRequest) {
672 let event_type = request.event_type.clone();
673 if let Err(error) = self.adapter.event_emitter().emit(request).await {
674 warn!(
675 session_id = %self.session_id,
676 org_id = self.org_id,
677 event_type,
678 %error,
679 "runtime host lifecycle event emission failed"
680 );
681 }
682 }
683
684 pub async fn turn_started(&self, turn_id: TurnId, input_message_id: MessageId) {
685 let input_content = self
686 .adapter
687 .message_store()
688 .get(self.session_id, input_message_id)
689 .await
690 .ok()
691 .flatten()
692 .map(|message| message.content_to_llm_string());
693
694 self.set_session_status(SessionStatus::Active, "turn_started")
695 .await;
696
697 self.emit_event(EventRequest::new(
698 self.session_id,
699 EventContext::turn(turn_id, input_message_id),
700 SessionActivatedData {
701 turn_id,
702 input_message_id,
703 },
704 ))
705 .await;
706
707 self.emit_event(EventRequest::new(
708 self.session_id,
709 EventContext::turn(turn_id, input_message_id),
710 TurnStartedData {
711 turn_id,
712 input_message_id,
713 input_content,
714 },
715 ))
716 .await;
717 }
718
719 pub async fn emit_turn_completed(&self, input_message_id: MessageId, data: TurnCompletedData) {
720 let turn_id = data.turn_id;
721 self.emit_event(EventRequest::new(
722 self.session_id,
723 EventContext::turn(turn_id, input_message_id),
724 data,
725 ))
726 .await;
727 }
728
729 pub async fn emit_session_idled(
730 &self,
731 turn_id: TurnId,
732 input_message_id: MessageId,
733 iterations: Option<u32>,
734 usage: Option<TokenUsage>,
735 ) {
736 self.set_session_status(SessionStatus::Idle, "emit_session_idled")
737 .await;
738
739 self.emit_event(EventRequest::new(
740 self.session_id,
741 EventContext::turn(turn_id, input_message_id),
742 SessionIdledData {
743 turn_id,
744 iterations,
745 usage,
746 },
747 ))
748 .await;
749 }
750
751 pub async fn turn_completed(
752 &self,
753 turn_id: TurnId,
754 input_message_id: MessageId,
755 iterations: u32,
756 usage: Option<TokenUsage>,
757 input_content: Option<String>,
758 ) {
759 self.emit_turn_completed(
760 input_message_id,
761 TurnCompletedData {
762 turn_id,
763 iterations,
764 duration_ms: None,
765 usage: usage.clone(),
766 input_content,
767 final_message_id: None,
768 final_answer_preview: None,
769 time_to_first_token_ms: None,
770 tool_call_count: None,
771 llm_call_count: None,
772 status: Some("completed".to_string()),
773 },
774 )
775 .await;
776 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
777 .await;
778 }
779
780 pub async fn turn_sealed(
787 &self,
788 turn_id: TurnId,
789 input_message_id: MessageId,
790 reason: &str,
791 iterations: u32,
792 usage: Option<TokenUsage>,
793 ) {
794 let context = EventContext::turn(turn_id, input_message_id);
795
796 self.emit_event(EventRequest::new(
797 self.session_id,
798 context.clone(),
799 everruns_core::events::TurnSealedData {
800 turn_id,
801 reason: reason.to_string(),
802 detail: None,
803 iterations: Some(iterations),
804 usage: usage.clone(),
805 },
806 ))
807 .await;
808
809 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
810 .await;
811 }
812
813 pub async fn fire_turn_end_hooks(
817 &self,
818 harness_id: HarnessId,
819 agent_id: Option<AgentId>,
820 turn_id: TurnId,
821 success: bool,
822 ) {
823 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
824 &self.adapter,
825 self.org_id,
826 self.session_id,
827 harness_id,
828 agent_id,
829 )
830 .await
831 {
832 Ok(pair) => pair,
833 Err(error) => {
834 warn!(
835 session_id = %self.session_id,
836 %error,
837 "failed to collect turn_end hook specs; skipping"
838 );
839 return;
840 }
841 };
842 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
843 &specs,
844 everruns_core::user_hook_types::HookEvent::TurnEnd,
845 dispatcher,
846 );
847 if hooks.is_empty() {
848 return;
849 }
850 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
851 session_id: self.session_id,
852 turn_id: Some(turn_id),
853 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
854 agent_id: agent_id.map(|a| a.to_string()),
855 };
856 everruns_core::lifecycle_hooks::run_turn_end_hooks(
857 &hooks,
858 &ctx,
859 serde_json::json!({ "success": success }),
860 )
861 .await;
862 }
863
864 pub async fn user_prompt_blocked(
869 &self,
870 turn_id: TurnId,
871 input_message_id: MessageId,
872 reason: &str,
873 user_message: Option<&str>,
874 ) {
875 let user_error =
876 UserFacingError::new(everruns_core::user_facing_error_codes::BLOCKED_BY_HOOK);
877 let shown = user_message.unwrap_or(reason);
878 let mut error_message = Message::assistant(shown);
879 let mut metadata = std::collections::HashMap::new();
880 user_error.apply_to_message_metadata(&mut metadata);
881 error_message.metadata = Some(metadata);
882
883 self.emit_event(EventRequest::new(
884 self.session_id,
885 EventContext::turn(turn_id, input_message_id),
886 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
887 ))
888 .await;
889
890 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
891 .await;
892 }
893
894 pub async fn turn_failed(
895 &self,
896 turn_id: TurnId,
897 input_message_id: MessageId,
898 error: &str,
899 user_error: Option<&UserFacingError>,
900 ) {
901 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
902 .await;
903 }
904
905 pub async fn turn_failed_with_disclosure(
909 &self,
910 turn_id: TurnId,
911 input_message_id: MessageId,
912 error: &str,
913 user_error: Option<&UserFacingError>,
914 disclosure: Option<ErrorDisclosure>,
915 ) {
916 self.set_session_status(SessionStatus::Idle, "turn_failed")
917 .await;
918
919 self.emit_event(EventRequest::new(
920 self.session_id,
921 EventContext::turn(turn_id, input_message_id),
922 {
923 let mut data = TurnFailedData {
924 turn_id,
925 error: error.to_string(),
926 error_code: None,
927 error_fields: None,
928 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
929 };
930 if let Some(user_error) = user_error {
931 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
932 }
933 data
934 },
935 ))
936 .await;
937
938 self.emit_event(EventRequest::new(
939 self.session_id,
940 EventContext::turn(turn_id, input_message_id),
941 SessionIdledData {
942 turn_id,
943 iterations: None,
944 usage: None,
945 },
946 ))
947 .await;
948 }
949
950 pub async fn waiting_for_tool_results(&self) {
951 self.set_session_status(
952 SessionStatus::WaitingForToolResults,
953 "waiting_for_tool_results",
954 )
955 .await;
956 }
957
958 pub async fn dependency_blocked(
959 &self,
960 turn_id: TurnId,
961 input_message_id: MessageId,
962 blocker: DependencyBlocker,
963 ) {
964 let user_error = UserFacingError::new(blocker.error_code())
965 .with_field(
966 "dependency",
967 match blocker {
968 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
969 "harness"
970 }
971 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
972 },
973 )
974 .with_field(
975 "state",
976 match blocker {
977 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
978 "archived"
979 }
980 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
981 "deleted"
982 }
983 },
984 );
985 let mut error_message = Message::assistant(blocker.message());
986 let mut metadata = std::collections::HashMap::new();
987 user_error.apply_to_message_metadata(&mut metadata);
988 error_message.metadata = Some(metadata);
989
990 self.emit_event(EventRequest::new(
991 self.session_id,
992 EventContext::turn(turn_id, input_message_id),
993 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
994 ))
995 .await;
996
997 self.turn_failed(
998 turn_id,
999 input_message_id,
1000 blocker.message(),
1001 Some(&user_error),
1002 )
1003 .await;
1004 }
1005}
1006
1007pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1008 adapter: &A,
1009 org_id: i64,
1010 harness_id: HarnessId,
1011 agent_id: Option<AgentId>,
1012) -> everruns_core::error::Result<Option<DependencyBlocker>> {
1013 let harness_store = adapter.harness_store(org_id);
1014 let agent_store = adapter.agent_store(org_id);
1015 everruns_core::detect_dependency_blocker(
1016 harness_store.as_ref(),
1017 agent_store.as_ref(),
1018 harness_id,
1019 agent_id,
1020 )
1021 .await
1022}
1023
1024pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1025 adapter: &A,
1026 org_id: i64,
1027 input: InputAtomInput,
1028) -> everruns_core::error::Result<InputAtomResult> {
1029 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1033 handle.set(None);
1034 }
1035
1036 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1037 .turn_started(input.context.turn_id, input.context.input_message_id)
1038 .await;
1039
1040 let atom = InputAtom::new(adapter.message_store());
1041 atom.execute(input).await
1042}
1043
1044struct UserPromptHookResult {
1051 decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1052 original_message: String,
1053}
1054
1055async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1056 adapter: &A,
1057 org_id: i64,
1058 input: &ReasonInput,
1059) -> everruns_core::error::Result<Option<UserPromptHookResult>> {
1060 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1061 adapter,
1062 org_id,
1063 input.context.session_id,
1064 input.harness_id,
1065 input.agent_id,
1066 )
1067 .await
1068 {
1069 Ok(pair) => pair,
1070 Err(error) => {
1071 warn!(
1072 session_id = %input.context.session_id,
1073 %error,
1074 "failed to collect user_prompt_submit hook specs; continuing without them"
1075 );
1076 return Ok(None);
1077 }
1078 };
1079 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1080 &specs,
1081 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1082 dispatcher,
1083 );
1084 if hooks.is_empty() {
1085 return Ok(None);
1086 }
1087
1088 let message_text = adapter
1089 .message_store()
1090 .get(input.context.session_id, input.context.input_message_id)
1091 .await
1092 .ok()
1093 .flatten()
1094 .map(|m| m.content_to_llm_string())
1095 .unwrap_or_default();
1096
1097 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1098 session_id: input.context.session_id,
1099 turn_id: Some(input.context.turn_id),
1100 org_id: org_public_id_from_internal(org_id).parse().ok(),
1101 agent_id: input.agent_id.map(|a| a.to_string()),
1102 };
1103 let original_message = message_text.clone();
1104 let decision =
1105 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1106 .await;
1107 Ok(Some(UserPromptHookResult {
1108 decision,
1109 original_message,
1110 }))
1111}
1112
1113pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1114 adapter: &A,
1115 org_id: i64,
1116 input: ReasonInput,
1117) -> everruns_core::error::Result<ReasonResult> {
1118 if let Some(blocker) =
1119 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1120 {
1121 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1122 .dependency_blocked(
1123 input.context.turn_id,
1124 input.context.input_message_id,
1125 blocker,
1126 )
1127 .await;
1128 return Ok(ReasonResult {
1129 success: false,
1130 text: blocker.message().to_string(),
1131 tool_calls: vec![],
1132 has_tool_calls: false,
1133 tool_definitions: vec![],
1134 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1135 error: Some("dependency_unavailable".to_string()),
1136 user_facing_error: None,
1137 error_disclosure: None,
1138 usage: None,
1139 output_message_id: None,
1140 time_to_first_token_ms: None,
1141 response_id: None,
1142 finish_reason: None,
1143 locale: None,
1144 network_access: None,
1145 parallel_tool_calls: None,
1146 });
1147 }
1148
1149 let mut user_prompt_message_override = None;
1157 if input.iteration <= 1
1158 && let Some(hook_result) = run_user_prompt_submit_for_turn(adapter, org_id, &input).await?
1159 {
1160 match hook_result.decision {
1161 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1162 reason,
1163 user_message,
1164 } => {
1165 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1166 .user_prompt_blocked(
1167 input.context.turn_id,
1168 input.context.input_message_id,
1169 &reason,
1170 user_message.as_deref(),
1171 )
1172 .await;
1173 return Ok(ReasonResult {
1174 success: false,
1175 text: user_message.unwrap_or_else(|| reason.clone()),
1176 tool_calls: vec![],
1177 has_tool_calls: false,
1178 tool_definitions: vec![],
1179 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1180 error: Some("blocked_by_user_prompt_hook".to_string()),
1181 user_facing_error: None,
1182 error_disclosure: None,
1183 usage: None,
1184 output_message_id: None,
1185 time_to_first_token_ms: None,
1186 response_id: None,
1187 finish_reason: None,
1188 locale: None,
1189 network_access: None,
1190 parallel_tool_calls: None,
1191 });
1192 }
1193 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1194 if message != hook_result.original_message {
1195 user_prompt_message_override = Some(message);
1196 }
1197 }
1198 }
1199 }
1200
1201 let validation_session = adapter
1205 .session_store(org_id)
1206 .get_session(input.context.session_id)
1207 .await?
1208 .ok_or_else(|| {
1209 everruns_core::error::AgentLoopError::session_not_found(input.context.session_id)
1210 })?;
1211 let validation_capabilities = load_execution_capabilities(
1212 adapter,
1213 org_id,
1214 input.context.session_id,
1215 input.harness_id,
1216 input.agent_id,
1217 validation_session.locale.clone(),
1218 validation_session.blueprint_id.as_deref(),
1219 )
1220 .await?;
1221 let query_history_allowed = validation_capabilities
1222 .tool_registry
1223 .get("query_history")
1224 .is_some();
1225 let validation_services = runtime_tool_context_services(
1226 adapter,
1227 org_id,
1228 input.context.session_id,
1229 input.agent_id,
1230 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1231 None,
1232 validation_capabilities.subagent_nesting_policy,
1233 );
1234 validation_capabilities
1235 .tool_registry
1236 .validate_context_services(&validation_services)?;
1237
1238 let mut turn_context = adapter
1239 .load_turn_context(org_id, input.context.session_id)
1240 .await?;
1241 if let Some(registry) = adapter.session_task_registry() {
1242 let session_store = adapter.session_store(org_id);
1243 if let Some(tool) = report_result_tool_for_child_session(
1244 input.context.session_id,
1245 session_store.as_ref(),
1246 registry.as_ref(),
1247 )
1248 .await?
1249 {
1250 turn_context.mcp_tool_definitions.push(tool.to_definition());
1251 }
1252 if let Some(tool) = report_task_progress_tool_for_child_session(
1253 input.context.session_id,
1254 session_store.as_ref(),
1255 registry.as_ref(),
1256 )
1257 .await?
1258 {
1259 turn_context.mcp_tool_definitions.push(tool.to_definition());
1260 }
1261 }
1262
1263 let mut reason_capability_registry = adapter.capability_registry();
1264 if !query_history_allowed {
1265 reason_capability_registry
1269 .register(everruns_core::capabilities::InfinityContextFilterOnlyCapability);
1270 }
1271 let mut atom = ReasonAtom::new(
1272 adapter.harness_store(org_id),
1273 adapter.agent_store(org_id),
1274 adapter.session_store(org_id),
1275 adapter.message_store(),
1276 adapter.provider_store(org_id),
1277 reason_capability_registry.clone(),
1278 adapter.driver_registry(),
1279 adapter.event_emitter(),
1280 )
1281 .with_file_store(adapter.file_store());
1282 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1283 atom = atom.with_image_resolver(image_resolver);
1284 }
1285 if let Some(hb) = adapter.stream_heartbeater() {
1286 atom = atom.with_stream_heartbeater(hb);
1287 }
1288 if let Some(timeout) = adapter.provider_stall_timeout() {
1289 atom = atom.with_provider_stall_timeout(timeout);
1290 }
1291 if let Some(store) = adapter.partial_stream_store() {
1292 atom = atom.with_partial_stream_store(store);
1293 }
1294 if let Some(store) = adapter.durable_tool_result_store() {
1295 atom = atom.with_durable_tool_result_store(store);
1296 }
1297 if let Some(store) = adapter.compaction_checkpoint_store() {
1298 atom = atom.with_compaction_checkpoint_store(store);
1299 }
1300 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1301 atom = atom.with_reasoning_effort_handle(handle);
1302 }
1303 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1304 atom = atom.with_utility_llm_service(utility_llm_service);
1305 }
1306 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1309 atom = atom.with_schedule_store(schedule_store);
1310 }
1311
1312 let input = ReasonInput {
1313 mcp_tool_definitions: turn_context.mcp_tool_definitions,
1314 ..input
1315 };
1316
1317 if let Some(message_override) = user_prompt_message_override {
1318 let mut assembled = assemble_turn_context(
1319 adapter.harness_store(org_id).as_ref(),
1320 adapter.agent_store(org_id).as_ref(),
1321 adapter.session_store(org_id).as_ref(),
1322 adapter.message_store().as_ref(),
1323 adapter.provider_store(org_id).as_ref(),
1324 &reason_capability_registry,
1325 input.context.session_id,
1326 input.harness_id,
1327 input.agent_id,
1328 &input.mcp_tool_definitions,
1329 Some(adapter.file_store()),
1330 )
1331 .await?;
1332
1333 let message = assembled
1334 .messages
1335 .iter_mut()
1336 .find(|message| message.id == input.context.input_message_id)
1337 .ok_or_else(|| {
1338 everruns_core::error::AgentLoopError::config(
1339 "user_prompt_submit mutation: input message not found in assembled context",
1340 )
1341 })?;
1342
1343 message
1348 .content
1349 .retain(|part| !matches!(part, ContentPart::Text(_)));
1350 message
1351 .content
1352 .insert(0, ContentPart::text(message_override));
1353
1354 return atom.execute_with_assembled_context(input, assembled).await;
1355 }
1356
1357 atom.execute(input).await
1358}
1359
1360pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1361 adapter: &A,
1362 input: ActInput,
1363) -> everruns_core::error::Result<ActResult> {
1364 let org_id = input.org_id.ok_or_else(|| {
1365 everruns_core::error::AgentLoopError::config(
1366 "ActInput.org_id must be set for runtime host execution",
1367 )
1368 })?;
1369
1370 if let Some(blocker) =
1371 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1372 {
1373 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1374 .dependency_blocked(
1375 input.context.turn_id,
1376 input.context.input_message_id,
1377 blocker,
1378 )
1379 .await;
1380 return Ok(ActResult {
1381 results: vec![],
1382 completed: true,
1383 success_count: 0,
1384 error_count: 1,
1385 waiting_for_tool_results: false,
1386 blocked: true,
1387 client_tool_calls: vec![],
1388 client_tool_definitions: vec![],
1389 });
1390 }
1391
1392 let execution_capabilities = load_execution_capabilities(
1393 adapter,
1394 org_id,
1395 input.context.session_id,
1396 input.harness_id,
1397 input.agent_id,
1398 input.locale.clone(),
1399 input.blueprint_id.as_deref(),
1400 )
1401 .await?;
1402 let mut tool_registry = execution_capabilities.tool_registry;
1403
1404 if input
1405 .tool_definitions
1406 .iter()
1407 .any(|definition| definition.name() == "report_result")
1408 && let Some(registry) = adapter.session_task_registry()
1409 && let Some(tool) = report_result_tool_for_child_session(
1410 input.context.session_id,
1411 adapter.session_store(org_id).as_ref(),
1412 registry.as_ref(),
1413 )
1414 .await?
1415 {
1416 tool_registry.register_boxed(Box::new(tool.with_file_store(adapter.file_store())));
1417 }
1418 if input
1419 .tool_definitions
1420 .iter()
1421 .any(|definition| definition.name() == "report_task_progress")
1422 && let Some(registry) = adapter.session_task_registry()
1423 && let Some(tool) = report_task_progress_tool_for_child_session(
1424 input.context.session_id,
1425 adapter.session_store(org_id).as_ref(),
1426 registry.as_ref(),
1427 )
1428 .await?
1429 {
1430 tool_registry.register_boxed(Box::new(tool));
1431 }
1432
1433 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1443 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1444 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1445 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1446 tool_registry.register_boxed(tool);
1447 }
1448 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1449 &input.tool_definitions,
1450 invoker,
1451 )));
1452 }
1453
1454 let builtin_tool_registry = Arc::new(tool_registry.clone());
1455 let context_services = runtime_tool_context_services(
1456 adapter,
1457 org_id,
1458 input.context.session_id,
1459 input.agent_id,
1460 Some(builtin_tool_registry),
1461 mcp_invoker,
1462 execution_capabilities.subagent_nesting_policy,
1463 );
1464 tool_registry.validate_context_services(&context_services)?;
1465 let executor: Arc<dyn everruns_core::traits::ToolExecutor> = Arc::new(tool_registry);
1466
1467 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1468 .with_context_services(context_services)
1469 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1470 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1471 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1472
1473 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1474 atom = atom.with_outbound_tool_rate_limiter(limiter);
1475 }
1476 if let Some(store) = adapter.durable_tool_result_store() {
1477 atom = atom.with_durable_tool_result_store(store);
1478 }
1479
1480 atom.execute(input).await
1481}