1use crate::{SessionMutator, SessionMutatorExt};
7use async_trait::async_trait;
8use everruns_core::capabilities::{
9 Capability, SystemPromptContext, collect_capabilities_with_configs,
10};
11use everruns_core::events::{
12 EventContext, EventRequest, OutputMessageCompletedData, SessionActivatedData, SessionIdledData,
13 SessionModelChangedData, TurnCompletedData, TurnFailedData, TurnStartedData,
14};
15use everruns_core::message::{ContentPart, Message, MessageRole};
16use everruns_core::message_retriever::MessageRetriever;
17use everruns_core::runtime_context::AssembledTurnContext;
18use everruns_core::session::SessionExecutionState;
19use everruns_core::{
20 CapabilityRegistry, CapabilityStatus, DependencyBlocker, EgressService,
21 ResolvedExecutionSnapshot, TokenUsage, ToolRegistry, UtilityLlmService,
22 org_public_id_from_internal, resolve_runtime_capabilities,
23};
24use everruns_core::{
25 connection_services::ProviderCredentialStore, connection_services::UserConnectionResolver,
26 delegation_services::SessionCreationAuthority, event_emitter::EventEmitter,
27 execution_loading::AgentStore, execution_loading::HarnessStore,
28 execution_loading::SessionStore, file_services::FileResolver,
29 image_services::ImageArtifactStore, image_services::ImageResolver,
30 provider_resolution::ProviderStore, session_files::SessionFileSystem,
31 session_services::LeasedResourceStore, session_services::SessionResourceRegistry,
32 session_services::SessionScheduleStore, session_services::SessionStorageStore,
33 tool_context::ToolContextServices, tool_execution::BudgetChecker,
34 tool_execution::PaymentAuthority,
35};
36use everruns_engine::{
37 ActAtom, ActInput, ActResult, InputAtom, InputAtomInput, InputAtomResult, ReasonAtom,
38 ReasonInput, ReasonResult,
39};
40use everruns_provider::driver_registry::DriverRegistry;
41use everruns_provider::tool_types::ToolDefinition;
42use everruns_provider::typed_id::{AgentId, HarnessId, MessageId, ModelId, SessionId, TurnId};
43use everruns_provider::user_facing_error::{ErrorDisclosure, UserFacingError};
44use std::sync::Arc;
45use tracing::warn;
46
47struct MessageFilterOnlyCapability(Arc<dyn Capability>);
50
51impl Capability for MessageFilterOnlyCapability {
52 fn id(&self) -> &str {
53 self.0.id()
54 }
55
56 fn aliases(&self) -> Vec<&'static str> {
57 self.0.aliases()
58 }
59
60 fn name(&self) -> &str {
61 self.0.name()
62 }
63
64 fn description(&self) -> &str {
65 self.0.description()
66 }
67
68 fn status(&self) -> CapabilityStatus {
69 self.0.status()
70 }
71
72 fn message_filter_provider(
73 &self,
74 ) -> Option<Arc<dyn everruns_core::message_filter::MessageFilterProvider>> {
75 self.0.message_filter_provider()
76 }
77
78 fn message_filter_config(
79 &self,
80 config: &serde_json::Value,
81 compaction_enabled: bool,
82 ) -> serde_json::Value {
83 self.0.message_filter_config(config, compaction_enabled)
84 }
85}
86
87#[cfg(feature = "bashkit")]
88fn bash_hook_dispatcher(
89 file_store: Arc<dyn SessionFileSystem>,
90) -> Arc<dyn everruns_core::hook_executor::BashHookDispatcher> {
91 Arc::new(everruns_integrations_bashkit::BashkitShellHookDispatcher::new(file_store))
92}
93
94#[cfg(not(feature = "bashkit"))]
95fn bash_hook_dispatcher(
96 _file_store: Arc<dyn SessionFileSystem>,
97) -> Arc<dyn everruns_core::hook_executor::BashHookDispatcher> {
98 struct DisabledDispatcher;
99
100 #[async_trait]
101 impl everruns_core::hook_executor::BashHookDispatcher for DisabledDispatcher {
102 async fn dispatch(
103 &self,
104 _payload: &everruns_core::hook_executor::HookPayload,
105 _command: &str,
106 _extra_env: &std::collections::BTreeMap<String, String>,
107 _opts: &everruns_core::hook_executor::ExecutorOpts,
108 ) -> std::result::Result<everruns_core::hook_executor::BashExecOutput, String> {
109 Err("bash hooks require the everruns-host `bashkit` feature".to_string())
110 }
111 }
112
113 Arc::new(DisabledDispatcher)
114}
115
116#[derive(Debug, Clone)]
123pub struct ResolvedTurnInputs {
124 pub snapshot: ResolvedExecutionSnapshot,
126 pub messages: Vec<Message>,
128 pub mcp_tool_definitions: Vec<ToolDefinition>,
130}
131
132#[async_trait]
143pub trait RuntimeHostAdapter: Send + Sync + Clone + 'static {
144 fn turn_cancellation(&self) -> Option<tokio::sync::watch::Receiver<bool>> {
146 None
147 }
148 async fn set_session_status(
151 &self,
152 org_id: i64,
153 session_id: SessionId,
154 status: SessionExecutionState,
155 ) -> everruns_provider::error::Result<()>;
156
157 async fn load_resolved_turn(
165 &self,
166 org_id: i64,
167 session_id: SessionId,
168 ) -> everruns_provider::error::Result<ResolvedTurnInputs>;
169
170 fn capability_registry(&self) -> CapabilityRegistry;
171
172 fn driver_registry(&self) -> DriverRegistry;
173
174 fn harness_store(&self, org_id: i64) -> Arc<dyn HarnessStore>;
175
176 fn agent_store(&self, org_id: i64) -> Arc<dyn AgentStore>;
177
178 fn session_store(&self, org_id: i64) -> Arc<dyn SessionStore>;
179
180 fn session_mutator(&self, org_id: i64) -> Arc<dyn SessionMutator>;
181
182 fn provider_store(&self, org_id: i64) -> Arc<dyn ProviderStore>;
183
184 fn message_store(&self) -> Arc<dyn MessageRetriever>;
185
186 fn native_async_store(
187 &self,
188 ) -> Option<Arc<dyn everruns_core::native_async_store::NativeAsyncStore>> {
189 None
190 }
191
192 fn compaction_checkpoint_store(
193 &self,
194 ) -> Option<Arc<dyn everruns_core::CompactionCheckpointStore>> {
195 None
196 }
197
198 fn event_emitter(&self) -> Arc<dyn EventEmitter>;
199
200 fn file_store(&self) -> Arc<dyn SessionFileSystem>;
201
202 fn image_resolver(&self, _org_id: i64) -> Option<Arc<dyn ImageResolver>> {
203 None
204 }
205
206 fn file_resolver(&self, _org_id: i64) -> Option<Arc<dyn FileResolver>> {
207 None
208 }
209
210 fn image_artifact_store(&self, _org_id: i64) -> Option<Arc<dyn ImageArtifactStore>> {
211 None
212 }
213
214 fn provider_credential_store(&self, _org_id: i64) -> Option<Arc<dyn ProviderCredentialStore>> {
215 None
216 }
217
218 fn utility_llm_service(&self) -> Option<Arc<dyn UtilityLlmService>> {
219 None
220 }
221
222 fn egress_service(&self) -> Option<Arc<dyn EgressService>> {
223 None
224 }
225
226 fn storage_store(&self) -> Option<Arc<dyn SessionStorageStore>> {
227 None
228 }
229
230 fn connection_resolver(&self) -> Option<Arc<dyn UserConnectionResolver>> {
231 None
232 }
233
234 fn tool_context_extensions(
236 &self,
237 _org_id: i64,
238 _session_id: SessionId,
239 ) -> everruns_core::tool_context::ToolContextExtensions {
240 Default::default()
241 }
242
243 fn subagent_delegate(
245 &self,
246 _org_id: i64,
247 _session_id: SessionId,
248 ) -> Option<Arc<dyn everruns_core::subagent_delegation::SubagentSessionDelegate>> {
249 None
250 }
251
252 fn tool_augmentor(&self) -> Option<Arc<dyn crate::HostToolAugmentor>> {
254 None
255 }
256
257 fn leased_resource_store(&self) -> Option<Arc<dyn LeasedResourceStore>> {
258 None
259 }
260
261 fn session_resource_registry(&self) -> Option<Arc<dyn SessionResourceRegistry>> {
262 None
263 }
264
265 fn session_task_registry(
266 &self,
267 ) -> Option<Arc<dyn everruns_core::session_task::SessionTaskRegistry>> {
268 None
269 }
270
271 fn schedule_store(&self, _org_id: i64) -> Option<Arc<dyn SessionScheduleStore>> {
272 None
273 }
274
275 fn budget_checker(
276 &self,
277 _org_id: i64,
278 _agent_id: Option<AgentId>,
279 ) -> Option<Arc<dyn BudgetChecker>> {
280 None
281 }
282
283 fn payment_authority(
284 &self,
285 _org_id: i64,
286 _agent_id: Option<AgentId>,
287 ) -> Option<Arc<dyn PaymentAuthority>> {
288 None
289 }
290
291 fn session_creation_authority(
292 &self,
293 _org_id: i64,
294 _session_id: SessionId,
295 ) -> Option<Arc<dyn SessionCreationAuthority>> {
296 None
297 }
298
299 fn outbound_tool_rate_limiter(
302 &self,
303 _org_id: i64,
304 ) -> Option<Arc<dyn everruns_core::tool_execution::OutboundToolRateLimiter>> {
305 None
306 }
307
308 fn durable_tool_result_store(
311 &self,
312 ) -> Option<Arc<dyn everruns_core::durability::DurableToolResultStore>> {
313 None
314 }
315
316 fn subagent_spawn_store(
319 &self,
320 ) -> Option<Arc<dyn everruns_core::delegation_services::SubagentSpawnStore>> {
321 None
322 }
323
324 fn stream_heartbeater(&self) -> Option<Arc<dyn everruns_core::durability::StreamHeartbeater>> {
327 None
328 }
329
330 fn partial_stream_store(
333 &self,
334 ) -> Option<Arc<dyn everruns_core::durability::PartialStreamStore>> {
335 None
336 }
337
338 fn reasoning_effort_handle(
347 &self,
348 _session_id: SessionId,
349 ) -> Option<everruns_core::tool_context::ReasoningEffortHandle> {
350 None
351 }
352
353 fn provider_stall_timeout(&self) -> Option<std::time::Duration> {
356 None
357 }
358
359 fn provider_retry_config(&self) -> Option<everruns_provider::llm_retry::LlmRetryConfig> {
362 None
363 }
364
365 async fn mcp_executor(
369 &self,
370 _org_id: i64,
371 _session_id: SessionId,
372 ) -> Option<Arc<dyn everruns_core::McpToolInvoker>> {
373 None
374 }
375}
376
377struct RuntimeExecutionCapabilities {
378 tool_registry: ToolRegistry,
379 post_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PostToolExecHook>>,
380 pre_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>>,
381 tool_call_hooks: Vec<Arc<dyn everruns_core::ToolCallHook>>,
382 subagent_nesting_policy: everruns_core::delegation_services::SubagentNestingPolicy,
383}
384
385fn subagent_nesting_policy_from_configs(
386 resolved_capability_configs: &[everruns_capability::CapabilityRef],
387) -> everruns_core::delegation_services::SubagentNestingPolicy {
388 let subagents_config = resolved_capability_configs
389 .iter()
390 .find(|config| config.capability_id() == "subagents");
391
392 let configured_depth = subagents_config
393 .and_then(|config| {
394 config
395 .config_value()
396 .get("max_subagent_depth")
397 .or_else(|| config.config_value().get("max_depth"))
398 })
399 .and_then(|value| value.as_u64())
400 .and_then(|value| u32::try_from(value).ok());
401 let configured_max_active = subagents_config
402 .and_then(|config| {
403 config
404 .config_value()
405 .get("max_active_descendant_tasks")
406 .or_else(|| config.config_value().get("max_concurrent_descendant_tasks"))
407 })
408 .and_then(|value| value.as_u64())
409 .and_then(|value| u32::try_from(value).ok());
410 let configured_max_total = subagents_config
411 .and_then(|config| config.config_value().get("max_total_descendant_tasks"))
412 .and_then(|value| value.as_u64())
413 .and_then(|value| u32::try_from(value).ok());
414 let configured_max_active_detached = subagents_config
415 .and_then(|config| config.config_value().get("max_active_detached_tasks"))
416 .and_then(|value| value.as_u64())
417 .and_then(|value| u32::try_from(value).ok());
418 let configured_max_total_detached = subagents_config
419 .and_then(|config| config.config_value().get("max_total_detached_tasks"))
420 .and_then(|value| value.as_u64())
421 .and_then(|value| u32::try_from(value).ok());
422
423 everruns_core::delegation_services::SubagentNestingPolicy::default()
424 .with_agent_override(configured_depth)
425 .with_agent_task_caps_override(configured_max_active, configured_max_total)
426 .with_agent_detached_task_caps_override(
427 configured_max_active_detached,
428 configured_max_total_detached,
429 )
430}
431
432fn finalize_specs_from_configs(
442 resolved_capability_configs: &[everruns_capability::CapabilityRef],
443 capability_registry: &CapabilityRegistry,
444 tool_augmentor: Option<&dyn crate::HostToolAugmentor>,
445) -> Vec<everruns_core::user_hook_types::UserHookSpec> {
446 let mut hook_contributions: Vec<(String, Vec<everruns_core::user_hook_types::UserHookSpec>)> =
447 Vec::new();
448 let mut disabled_contributions: Vec<String> = Vec::new();
449 for config in resolved_capability_configs {
450 let Some(capability) = capability_registry.get(config.capability_id()) else {
451 continue;
452 };
453 let specs = capability.user_hooks_with_config(config.config_value());
454 if !specs.is_empty() {
455 hook_contributions.push((config.capability_id().to_string(), specs));
456 }
457 if let Some(augmentor) = tool_augmentor {
458 disabled_contributions.extend(
459 augmentor
460 .disabled_hook_contributions(config.capability_id(), config.config_value()),
461 );
462 }
463 }
464 everruns_core::hook_adapter::finalize_hook_specs(hook_contributions, &disabled_contributions)
465}
466
467async fn collect_lifecycle_hook_specs<A: RuntimeHostAdapter>(
472 adapter: &A,
473 org_id: i64,
474 session_id: SessionId,
475 harness_id: HarnessId,
476 agent_id: Option<AgentId>,
477) -> everruns_provider::error::Result<(
478 Vec<everruns_core::user_hook_types::UserHookSpec>,
479 Arc<dyn everruns_core::hook_executor::BashHookDispatcher>,
480)> {
481 let capability_registry = adapter.capability_registry();
482 let harness = adapter
483 .harness_store(org_id)
484 .get_harness(harness_id)
485 .await?
486 .ok_or_else(|| everruns_provider::error::AgentLoopError::harness_not_found(harness_id))?;
487 let session = adapter
488 .session_store(org_id)
489 .get_session(session_id)
490 .await?
491 .ok_or_else(|| everruns_provider::error::AgentLoopError::session_not_found(session_id))?;
492 let agent = match agent_id {
493 Some(agent_id) => adapter.agent_store(org_id).get_agent(agent_id).await?,
494 None => None,
495 };
496 let resolved =
497 resolve_runtime_capabilities(&harness, agent.as_ref(), &session, &capability_registry);
498 let tool_augmentor = adapter.tool_augmentor();
499 let specs = finalize_specs_from_configs(
500 &resolved.resolved_capability_configs,
501 &capability_registry,
502 tool_augmentor.as_deref(),
503 );
504 let dispatcher = bash_hook_dispatcher(adapter.file_store());
505 Ok((specs, dispatcher))
506}
507
508async fn load_execution_capabilities<A: RuntimeHostAdapter>(
509 adapter: &A,
510 org_id: i64,
511 session_id: SessionId,
512 harness_id: HarnessId,
513 agent_id: Option<AgentId>,
514 locale: Option<String>,
515 blueprint_id: Option<&str>,
516) -> everruns_provider::error::Result<RuntimeExecutionCapabilities> {
517 let capability_registry = adapter.capability_registry();
518 if let Some(blueprint_id) = blueprint_id {
519 let mut registry = ToolRegistry::with_defaults();
520 #[cfg(feature = "builtins")]
521 everruns_builtins::register_default_tools(&mut registry);
522 let blueprint = capability_registry.blueprint(blueprint_id).ok_or_else(|| {
523 everruns_provider::error::AgentLoopError::config(format!(
524 "Blueprint \"{blueprint_id}\" not found in registry"
525 ))
526 })?;
527 for tool in blueprint.tools {
528 registry.register_boxed(tool);
529 }
530 return Ok(RuntimeExecutionCapabilities {
531 tool_registry: registry,
532 post_tool_hooks: Vec::new(),
533 pre_tool_hooks: Vec::new(),
534 tool_call_hooks: Vec::new(),
535 subagent_nesting_policy:
536 everruns_core::delegation_services::SubagentNestingPolicy::default(),
537 });
538 }
539
540 let harness = adapter
541 .harness_store(org_id)
542 .get_harness(harness_id)
543 .await?
544 .ok_or_else(|| everruns_provider::error::AgentLoopError::harness_not_found(harness_id))?;
545
546 let session = adapter
547 .session_store(org_id)
548 .get_session(session_id)
549 .await?
550 .ok_or_else(|| everruns_provider::error::AgentLoopError::session_not_found(session_id))?;
551
552 let agent_store = adapter.agent_store(org_id);
553 let agent =
554 match agent_id {
555 Some(agent_id) => Some(agent_store.get_agent(agent_id).await?.ok_or_else(|| {
556 everruns_provider::error::AgentLoopError::agent_not_found(agent_id)
557 })?),
558 None => None,
559 };
560
561 let resolved =
562 resolve_runtime_capabilities(&harness, agent.as_ref(), &session, &capability_registry);
563 let prompt_ctx = SystemPromptContext {
570 session_id,
571 locale: locale.or(session.locale.clone()),
572 file_store: Some(everruns_core::scoped_prompt_file_store(
579 adapter.file_store(),
580 session.workspace_id,
581 )),
582 model: None,
583 };
584 let collected = collect_capabilities_with_configs(
585 &resolved.resolved_capability_configs,
586 &capability_registry,
587 &prompt_ctx,
588 )
589 .await;
590
591 let mut registry = ToolRegistry::with_defaults();
592 #[cfg(feature = "builtins")]
593 everruns_builtins::register_default_tools(&mut registry);
594 for tool in collected.tools {
595 registry.register_boxed(tool);
596 }
597
598 let mut post_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PostToolExecHook>> = resolved
603 .resolved_capability_configs
604 .iter()
605 .flat_map(|config| {
606 capability_registry
607 .get(config.capability_id())
608 .filter(|capability| capability.status().is_active())
609 .map(|capability| {
610 capability.post_tool_exec_hooks_with_config(config.config_value())
611 })
612 .unwrap_or_default()
613 })
614 .collect();
615 post_tool_hooks.sort_by_key(|hook| hook.priority());
618
619 let tool_augmentor = adapter.tool_augmentor();
626 let user_hook_specs = finalize_specs_from_configs(
627 &resolved.resolved_capability_configs,
628 &capability_registry,
629 tool_augmentor.as_deref(),
630 );
631 if user_hook_specs
636 .iter()
637 .any(|spec| spec.event == everruns_core::user_hook_types::HookEvent::UserPromptSubmit)
638 {
639 registry.unregister("query_history");
640 }
641 let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>> = resolved
644 .resolved_capability_configs
645 .iter()
646 .flat_map(|config| {
647 capability_registry
648 .get(config.capability_id())
649 .filter(|capability| capability.status().is_active())
650 .map(|capability| capability.pre_tool_use_hooks_with_config(config.config_value()))
651 .unwrap_or_default()
652 })
653 .collect();
654 if !user_hook_specs.is_empty() {
655 let dispatcher = bash_hook_dispatcher(adapter.file_store());
656 post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
657 &user_hook_specs,
658 dispatcher.clone(),
659 ));
660 pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
661 &user_hook_specs,
662 dispatcher,
663 ));
664 }
665
666 let tool_call_hooks = collected.tool_call_hooks;
677
678 Ok(RuntimeExecutionCapabilities {
679 tool_registry: registry,
680 post_tool_hooks,
681 pre_tool_hooks,
682 tool_call_hooks,
683 subagent_nesting_policy: subagent_nesting_policy_from_configs(
684 &resolved.resolved_capability_configs,
685 ),
686 })
687}
688
689fn runtime_tool_context_services<A: RuntimeHostAdapter>(
690 adapter: &A,
691 org_id: i64,
692 session_id: SessionId,
693 agent_id: Option<AgentId>,
694 tool_registry: Option<Arc<ToolRegistry>>,
695 mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
696 subagent_nesting_policy: everruns_core::delegation_services::SubagentNestingPolicy,
697) -> ToolContextServices {
698 let extensions = {
699 let mut extensions = adapter.tool_context_extensions(org_id, session_id);
700 extensions.insert(Arc::new(SessionMutatorExt(adapter.session_mutator(org_id))));
701 extensions
702 };
703 ToolContextServices {
704 file_store: Some(adapter.file_store()),
705 storage_store: adapter.storage_store(),
706 image_store: adapter.image_artifact_store(org_id),
707 provider_credential_store: adapter.provider_credential_store(org_id),
708 utility_llm_service: adapter.utility_llm_service(),
709 mcp_invoker,
710 egress_service: adapter.egress_service(),
711 message_retriever: Some(adapter.message_store()),
712 session_store: Some(adapter.session_store(org_id)),
713 agent_store: Some(adapter.agent_store(org_id)),
714 connection_resolver: adapter.connection_resolver(),
715 schedule_store: adapter.schedule_store(org_id),
716 subagent_delegate: adapter.subagent_delegate(org_id, session_id),
717 extensions,
718 leased_resource_store: adapter.leased_resource_store(),
719 session_resource_registry: adapter.session_resource_registry(),
720 session_task_registry: adapter.session_task_registry(),
721 event_emitter: Some(adapter.event_emitter()),
722 capability_registry: Some(adapter.capability_registry()),
723 tool_registry,
724 org_id: Some(
725 org_public_id_from_internal(org_id)
726 .parse()
727 .expect("internal org id converts to valid public org id"),
728 ),
729 network_access: None,
730 budget_checker: adapter.budget_checker(org_id, agent_id),
731 payment_authority: adapter.payment_authority(org_id, agent_id),
732 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
733 subagent_spawn_store: adapter.subagent_spawn_store(),
734 subagent_nesting_policy,
735 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
736 }
737}
738
739#[derive(Debug, Default)]
742struct TurnAgentIdentity {
743 id: Option<AgentId>,
744 name: Option<String>,
745 description: Option<String>,
746}
747
748pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
749 adapter: A,
750 org_id: i64,
751 session_id: SessionId,
752}
753
754impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
755 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
756 Self {
757 adapter,
758 org_id,
759 session_id,
760 }
761 }
762
763 async fn set_session_status(
764 &self,
765 status: SessionExecutionState,
766 _action: &'static str,
767 ) -> everruns_provider::error::Result<()> {
768 self.adapter
769 .set_session_status(self.org_id, self.session_id, status)
770 .await
771 }
772
773 async fn emit_event(&self, request: EventRequest) -> everruns_provider::error::Result<()> {
774 self.adapter.event_emitter().emit(request).await.map(|_| ())
775 }
776
777 pub async fn turn_started(
778 &self,
779 turn_id: TurnId,
780 input_message_id: MessageId,
781 ) -> everruns_provider::error::Result<()> {
782 let input_content = self
783 .adapter
784 .message_store()
785 .get(self.session_id, input_message_id)
786 .await
787 .ok()
788 .flatten()
789 .map(|message| message.content_to_llm_string());
790
791 self.set_session_status(SessionExecutionState::Active, "turn_started")
792 .await?;
793
794 self.emit_event(EventRequest::new(
795 self.session_id,
796 EventContext::turn(turn_id, input_message_id),
797 SessionActivatedData {
798 turn_id,
799 input_message_id,
800 },
801 ))
802 .await?;
803
804 let agent = self.agent_identity().await;
805 self.emit_event(EventRequest::new(
806 self.session_id,
807 EventContext::turn(turn_id, input_message_id),
808 TurnStartedData {
809 turn_id,
810 input_message_id,
811 input_content,
812 agent_id: agent.id,
813 agent_name: agent.name,
814 agent_description: agent.description,
815 },
816 ))
817 .await?;
818 Ok(())
819 }
820
821 async fn agent_identity(&self) -> TurnAgentIdentity {
825 let agent_id = self
826 .adapter
827 .session_store(self.org_id)
828 .get_session(self.session_id)
829 .await
830 .ok()
831 .flatten()
832 .and_then(|session| session.agent_id);
833 let Some(agent_id) = agent_id else {
834 return TurnAgentIdentity::default();
835 };
836 let agent = self
837 .adapter
838 .agent_store(self.org_id)
839 .get_agent(agent_id)
840 .await
841 .ok()
842 .flatten();
843 TurnAgentIdentity {
844 id: Some(agent_id),
845 name: agent
848 .as_ref()
849 .map(|a| a.display_name.clone().unwrap_or_else(|| a.name.clone())),
850 description: agent.and_then(|a| a.description),
851 }
852 }
853
854 pub async fn emit_turn_completed(
855 &self,
856 input_message_id: MessageId,
857 data: TurnCompletedData,
858 ) -> everruns_provider::error::Result<()> {
859 let turn_id = data.turn_id;
860 self.emit_event(EventRequest::new(
861 self.session_id,
862 EventContext::turn(turn_id, input_message_id),
863 data,
864 ))
865 .await
866 }
867
868 pub async fn emit_session_idled(
869 &self,
870 turn_id: TurnId,
871 input_message_id: MessageId,
872 iterations: Option<u32>,
873 usage: Option<TokenUsage>,
874 ) -> everruns_provider::error::Result<()> {
875 self.set_session_status(SessionExecutionState::Idle, "emit_session_idled")
876 .await?;
877
878 self.emit_event(EventRequest::new(
879 self.session_id,
880 EventContext::turn(turn_id, input_message_id),
881 SessionIdledData {
882 turn_id,
883 iterations,
884 usage,
885 },
886 ))
887 .await
888 }
889
890 pub async fn turn_completed(
891 &self,
892 turn_id: TurnId,
893 input_message_id: MessageId,
894 iterations: u32,
895 usage: Option<TokenUsage>,
896 input_content: Option<String>,
897 ) -> everruns_provider::error::Result<()> {
898 self.emit_turn_completed(
899 input_message_id,
900 TurnCompletedData {
901 turn_id,
902 iterations,
903 duration_ms: None,
904 usage: usage.clone(),
905 input_content,
906 final_message_id: None,
907 final_answer_preview: None,
908 time_to_first_token_ms: None,
909 tool_call_count: None,
910 llm_call_count: None,
911 status: Some("completed".to_string()),
912 },
913 )
914 .await?;
915 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
916 .await
917 }
918
919 pub async fn turn_sealed(
926 &self,
927 turn_id: TurnId,
928 input_message_id: MessageId,
929 reason: &str,
930 iterations: u32,
931 usage: Option<TokenUsage>,
932 ) -> everruns_provider::error::Result<()> {
933 let context = EventContext::turn(turn_id, input_message_id);
934
935 self.emit_event(EventRequest::new(
936 self.session_id,
937 context.clone(),
938 everruns_core::events::TurnSealedData {
939 turn_id,
940 reason: reason.to_string(),
941 detail: None,
942 iterations: Some(iterations),
943 usage: usage.clone(),
944 },
945 ))
946 .await?;
947
948 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
949 .await
950 }
951
952 pub async fn fire_turn_end_hooks(
956 &self,
957 harness_id: HarnessId,
958 agent_id: Option<AgentId>,
959 turn_id: TurnId,
960 success: bool,
961 ) {
962 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
963 &self.adapter,
964 self.org_id,
965 self.session_id,
966 harness_id,
967 agent_id,
968 )
969 .await
970 {
971 Ok(pair) => pair,
972 Err(error) => {
973 warn!(
974 session_id = %self.session_id,
975 %error,
976 "failed to collect turn_end hook specs; skipping"
977 );
978 return;
979 }
980 };
981 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
982 &specs,
983 everruns_core::user_hook_types::HookEvent::TurnEnd,
984 dispatcher,
985 );
986 if hooks.is_empty() {
987 return;
988 }
989 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
990 session_id: self.session_id,
991 turn_id: Some(turn_id),
992 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
993 agent_id: agent_id.map(|a| a.to_string()),
994 };
995 everruns_core::lifecycle_hooks::run_turn_end_hooks(
996 &hooks,
997 &ctx,
998 serde_json::json!({ "success": success }),
999 )
1000 .await;
1001 }
1002
1003 pub async fn user_prompt_blocked(
1008 &self,
1009 turn_id: TurnId,
1010 input_message_id: MessageId,
1011 reason: &str,
1012 user_message: Option<&str>,
1013 ) -> everruns_provider::error::Result<()> {
1014 let user_error =
1015 UserFacingError::new(everruns_provider::user_facing_error::codes::BLOCKED_BY_HOOK);
1016 let shown = user_message.unwrap_or(reason);
1017 let mut error_message = Message::assistant(shown);
1018 let mut metadata = std::collections::HashMap::new();
1019 user_error.apply_to_message_metadata(&mut metadata);
1020 error_message.metadata = Some(metadata);
1021
1022 self.emit_event(EventRequest::new(
1023 self.session_id,
1024 EventContext::turn(turn_id, input_message_id),
1025 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1026 ))
1027 .await?;
1028
1029 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
1030 .await
1031 }
1032
1033 pub async fn turn_failed(
1034 &self,
1035 turn_id: TurnId,
1036 input_message_id: MessageId,
1037 error: &str,
1038 user_error: Option<&UserFacingError>,
1039 ) -> everruns_provider::error::Result<()> {
1040 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
1041 .await
1042 }
1043
1044 pub async fn turn_failed_with_disclosure(
1048 &self,
1049 turn_id: TurnId,
1050 input_message_id: MessageId,
1051 error: &str,
1052 user_error: Option<&UserFacingError>,
1053 disclosure: Option<ErrorDisclosure>,
1054 ) -> everruns_provider::error::Result<()> {
1055 self.set_session_status(SessionExecutionState::Idle, "turn_failed")
1056 .await?;
1057
1058 self.emit_event(EventRequest::new(
1059 self.session_id,
1060 EventContext::turn(turn_id, input_message_id),
1061 {
1062 let mut data = TurnFailedData {
1063 turn_id,
1064 error: error.to_string(),
1065 error_code: None,
1066 error_fields: None,
1067 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
1068 };
1069 if let Some(user_error) = user_error {
1070 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
1071 }
1072 data
1073 },
1074 ))
1075 .await?;
1076
1077 self.emit_event(EventRequest::new(
1078 self.session_id,
1079 EventContext::turn(turn_id, input_message_id),
1080 SessionIdledData {
1081 turn_id,
1082 iterations: None,
1083 usage: None,
1084 },
1085 ))
1086 .await
1087 }
1088
1089 pub async fn waiting_for_tool_results(&self) -> everruns_provider::error::Result<()> {
1090 self.set_session_status(
1091 SessionExecutionState::WaitingForToolResults,
1092 "waiting_for_tool_results",
1093 )
1094 .await
1095 }
1096
1097 pub async fn dependency_blocked(
1098 &self,
1099 turn_id: TurnId,
1100 input_message_id: MessageId,
1101 blocker: DependencyBlocker,
1102 ) -> everruns_provider::error::Result<()> {
1103 let user_error = UserFacingError::new(blocker.error_code())
1104 .with_field(
1105 "dependency",
1106 match blocker {
1107 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
1108 "harness"
1109 }
1110 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
1111 },
1112 )
1113 .with_field(
1114 "state",
1115 match blocker {
1116 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
1117 "archived"
1118 }
1119 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
1120 "deleted"
1121 }
1122 },
1123 );
1124 let mut error_message = Message::assistant(blocker.message());
1125 let mut metadata = std::collections::HashMap::new();
1126 user_error.apply_to_message_metadata(&mut metadata);
1127 error_message.metadata = Some(metadata);
1128
1129 self.emit_event(EventRequest::new(
1130 self.session_id,
1131 EventContext::turn(turn_id, input_message_id),
1132 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1133 ))
1134 .await?;
1135
1136 self.turn_failed(
1137 turn_id,
1138 input_message_id,
1139 blocker.message(),
1140 Some(&user_error),
1141 )
1142 .await
1143 }
1144}
1145
1146pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1147 adapter: &A,
1148 org_id: i64,
1149 harness_id: HarnessId,
1150 agent_id: Option<AgentId>,
1151) -> everruns_provider::error::Result<Option<DependencyBlocker>> {
1152 let harness_store = adapter.harness_store(org_id);
1153 let agent_store = adapter.agent_store(org_id);
1154 if let Some(blocker) = harness_store.get_harness_blocker(harness_id).await? {
1155 return Ok(Some(blocker));
1156 }
1157 if let Some(agent_id) = agent_id
1158 && let Some(blocker) = agent_store.get_agent_blocker(agent_id).await?
1159 {
1160 return Ok(Some(blocker));
1161 }
1162 Ok(None)
1163}
1164
1165pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1166 adapter: &A,
1167 org_id: i64,
1168 input: InputAtomInput,
1169) -> everruns_provider::error::Result<InputAtomResult> {
1170 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1174 handle.set(None);
1175 }
1176
1177 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1178 .turn_started(input.context.turn_id, input.context.input_message_id)
1179 .await?;
1180
1181 let atom = InputAtom::new(adapter.message_store());
1182 atom.execute(input).await
1183}
1184
1185pub(crate) struct UserPromptHookResult {
1192 pub(crate) decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1193 pub(crate) original_message: String,
1194}
1195
1196pub(crate) async fn run_user_prompt_submit_for_message<A: RuntimeHostAdapter>(
1197 adapter: &A,
1198 org_id: i64,
1199 input: &ReasonInput,
1200 message_text: String,
1201) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1202 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1203 adapter,
1204 org_id,
1205 input.context.session_id,
1206 input.harness_id,
1207 input.agent_id,
1208 )
1209 .await
1210 {
1211 Ok(pair) => pair,
1212 Err(error) => {
1213 warn!(
1214 session_id = %input.context.session_id,
1215 %error,
1216 "failed to collect user_prompt_submit hook specs; continuing without them"
1217 );
1218 return Ok(None);
1219 }
1220 };
1221 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1222 &specs,
1223 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1224 dispatcher,
1225 );
1226 if hooks.is_empty() {
1227 return Ok(None);
1228 }
1229
1230 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1231 session_id: input.context.session_id,
1232 turn_id: Some(input.context.turn_id),
1233 org_id: org_public_id_from_internal(org_id).parse().ok(),
1234 agent_id: input.agent_id.map(|a| a.to_string()),
1235 };
1236 let original_message = message_text.clone();
1237 let decision =
1238 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1239 .await;
1240 Ok(Some(UserPromptHookResult {
1241 decision,
1242 original_message,
1243 }))
1244}
1245
1246async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1247 adapter: &A,
1248 org_id: i64,
1249 input: &ReasonInput,
1250) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1251 let message_text = adapter
1252 .message_store()
1253 .get(input.context.session_id, input.context.input_message_id)
1254 .await
1255 .ok()
1256 .flatten()
1257 .map(|m| m.content_to_llm_string())
1258 .unwrap_or_default();
1259 run_user_prompt_submit_for_message(adapter, org_id, input, message_text).await
1260}
1261
1262pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1263 adapter: &A,
1264 org_id: i64,
1265 input: ReasonInput,
1266) -> everruns_provider::error::Result<ReasonResult> {
1267 let prompt_message_ids = (input.iteration <= 1)
1268 .then_some(input.context.input_message_id)
1269 .into_iter()
1270 .collect();
1271 execute_reason_activity_with_prompt_messages(adapter, org_id, input, prompt_message_ids).await
1272}
1273
1274pub async fn execute_reason_activity_with_prompt_messages<A: RuntimeHostAdapter>(
1279 adapter: &A,
1280 org_id: i64,
1281 input: ReasonInput,
1282 prompt_message_ids: Vec<MessageId>,
1283) -> everruns_provider::error::Result<ReasonResult> {
1284 if let Some(blocker) =
1285 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1286 {
1287 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1288 .dependency_blocked(
1289 input.context.turn_id,
1290 input.context.input_message_id,
1291 blocker,
1292 )
1293 .await?;
1294 return Ok(ReasonResult {
1295 native_counts: None,
1296 success: false,
1297 text: blocker.message().to_string(),
1298 tool_calls: vec![],
1299 has_tool_calls: false,
1300 tool_definitions: vec![],
1301 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1302 error: Some("dependency_unavailable".to_string()),
1303 user_facing_error: None,
1304 error_disclosure: None,
1305 usage: None,
1306 output_message_id: None,
1307 time_to_first_token_ms: None,
1308 response_id: None,
1309 finish_reason: None,
1310 locale: None,
1311 network_access: None,
1312 parallel_tool_calls: None,
1313 });
1314 }
1315
1316 let mut user_prompt_message_overrides = Vec::new();
1321 for message_id in prompt_message_ids {
1322 let mut hook_input = input.clone();
1323 hook_input.context.input_message_id = message_id;
1324 let Some(hook_result) =
1325 run_user_prompt_submit_for_turn(adapter, org_id, &hook_input).await?
1326 else {
1327 continue;
1328 };
1329 match hook_result.decision {
1330 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1331 reason,
1332 user_message,
1333 } => {
1334 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1335 .user_prompt_blocked(
1336 input.context.turn_id,
1337 input.context.input_message_id,
1338 &reason,
1339 user_message.as_deref(),
1340 )
1341 .await?;
1342 return Ok(ReasonResult {
1343 native_counts: None,
1344 success: false,
1345 text: user_message.unwrap_or_else(|| reason.clone()),
1346 tool_calls: vec![],
1347 has_tool_calls: false,
1348 tool_definitions: vec![],
1349 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1350 error: Some("blocked_by_user_prompt_hook".to_string()),
1351 user_facing_error: None,
1352 error_disclosure: None,
1353 usage: None,
1354 output_message_id: None,
1355 time_to_first_token_ms: None,
1356 response_id: None,
1357 finish_reason: None,
1358 locale: None,
1359 network_access: None,
1360 parallel_tool_calls: None,
1361 });
1362 }
1363 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1364 if message != hook_result.original_message {
1365 user_prompt_message_overrides.push((message_id, message));
1366 }
1367 }
1368 }
1369 }
1370
1371 let validation_session = adapter
1375 .session_store(org_id)
1376 .get_session(input.context.session_id)
1377 .await?
1378 .ok_or_else(|| {
1379 everruns_provider::error::AgentLoopError::session_not_found(input.context.session_id)
1380 })?;
1381 let validation_capabilities = load_execution_capabilities(
1382 adapter,
1383 org_id,
1384 input.context.session_id,
1385 input.harness_id,
1386 input.agent_id,
1387 validation_session.locale.clone(),
1388 validation_session.blueprint_id.as_deref(),
1389 )
1390 .await?;
1391 let query_history_allowed = validation_capabilities
1392 .tool_registry
1393 .get("query_history")
1394 .is_some();
1395 let validation_services = runtime_tool_context_services(
1396 adapter,
1397 org_id,
1398 input.context.session_id,
1399 input.agent_id,
1400 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1401 None,
1402 validation_capabilities.subagent_nesting_policy,
1403 );
1404 validation_capabilities
1405 .tool_registry
1406 .validate_context_services(&validation_services)?;
1407
1408 let mut turn_inputs = adapter
1409 .load_resolved_turn(org_id, input.context.session_id)
1410 .await?;
1411 if let Some(augmentor) = adapter.tool_augmentor() {
1412 augmentor
1413 .augment_reason_tools(
1414 input.context.session_id,
1415 adapter.session_store(org_id),
1416 adapter.session_task_registry(),
1417 &mut turn_inputs.mcp_tool_definitions,
1418 )
1419 .await?;
1420 }
1421
1422 let reason_capability_registry = {
1423 let mut registry = adapter.capability_registry();
1424 if !query_history_allowed {
1425 let query_history_owner = registry
1431 .list()
1432 .into_iter()
1433 .find(|capability| {
1434 capability
1435 .tool_definitions()
1436 .iter()
1437 .any(|tool| tool.name() == "query_history")
1438 })
1439 .map(Arc::clone);
1440 if let Some(capability) = query_history_owner {
1441 registry.register(MessageFilterOnlyCapability(capability));
1442 }
1443 }
1444 registry
1445 };
1446 let context_resolver = crate::runtime_context::StoreTurnContextResolver::new(
1447 adapter.harness_store(org_id),
1448 adapter.agent_store(org_id),
1449 adapter.session_store(org_id),
1450 adapter.message_store(),
1451 adapter.provider_store(org_id),
1452 reason_capability_registry.clone(),
1453 adapter.driver_registry(),
1454 )
1455 .with_file_store(adapter.file_store());
1456 let mut atom = ReasonAtom::new(
1457 context_resolver,
1458 adapter.message_store(),
1459 reason_capability_registry.clone(),
1460 adapter.event_emitter(),
1461 );
1462 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1463 atom = atom.with_image_resolver(image_resolver);
1464 }
1465 if let Some(file_resolver) = adapter.file_resolver(org_id) {
1466 atom = atom.with_file_resolver(file_resolver);
1467 }
1468 if let Some(hb) = adapter.stream_heartbeater() {
1469 atom = atom.with_stream_heartbeater(hb);
1470 }
1471 if let Some(timeout) = adapter.provider_stall_timeout() {
1472 atom = atom.with_provider_stall_timeout(timeout);
1473 }
1474 if let Some(config) = adapter.provider_retry_config() {
1475 atom = atom.with_provider_retry_config(config);
1476 }
1477 if let Some(store) = adapter.partial_stream_store() {
1478 atom = atom.with_partial_stream_store(store);
1479 }
1480 if let Some(store) = adapter.durable_tool_result_store() {
1481 atom = atom.with_durable_tool_result_store(store);
1482 }
1483 if let Some(store) = adapter.compaction_checkpoint_store() {
1484 atom = atom.with_compaction_checkpoint_store(store);
1485 }
1486 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1487 atom = atom.with_reasoning_effort_handle(handle);
1488 }
1489 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1490 atom = atom.with_utility_llm_service(utility_llm_service);
1491 }
1492 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1495 atom = atom.with_schedule_store(schedule_store);
1496 }
1497
1498 let mut assembled = crate::runtime_context::assemble_turn_context_from_snapshot(
1499 turn_inputs.snapshot,
1500 adapter.message_store().as_ref(),
1501 adapter.provider_store(org_id).as_ref(),
1502 &reason_capability_registry,
1503 &adapter.driver_registry(),
1504 &turn_inputs.mcp_tool_definitions,
1505 Some(adapter.file_store()),
1506 )
1507 .await?;
1508 let input = ReasonInput {
1509 mcp_tool_definitions: turn_inputs.mcp_tool_definitions,
1510 ..input
1511 };
1512
1513 if !user_prompt_message_overrides.is_empty() {
1514 for (message_id, message_override) in user_prompt_message_overrides {
1515 let message = assembled
1516 .messages
1517 .iter_mut()
1518 .find(|message| message.id == message_id)
1519 .ok_or_else(|| {
1520 everruns_provider::error::AgentLoopError::config(
1521 "user_prompt_submit mutation: input message not found in assembled context",
1522 )
1523 })?;
1524
1525 message
1528 .content
1529 .retain(|part| !matches!(part, ContentPart::Text(_)));
1530 message
1531 .content
1532 .insert(0, ContentPart::text(message_override));
1533 }
1534 }
1535 if input.iteration <= 1 {
1540 emit_model_change_if_switched(adapter, org_id, &input, &assembled).await;
1541 }
1542
1543 crate::native_async::execute_reason(adapter, org_id, input, assembled, atom).await
1544}
1545
1546async fn emit_model_change_if_switched<A: RuntimeHostAdapter>(
1551 adapter: &A,
1552 org_id: i64,
1553 input: &ReasonInput,
1554 assembled: &AssembledTurnContext,
1555) {
1556 let Some((previous_model_id, model_id)) = model_switch(&assembled.messages) else {
1557 return;
1558 };
1559 if assembled.resolved_model_id != Some(model_id) {
1560 return;
1563 }
1564
1565 let previous_model_name = adapter
1569 .provider_store(org_id)
1570 .get_model_spec(previous_model_id)
1571 .await
1572 .ok()
1573 .flatten()
1574 .map(|spec| spec.model);
1575
1576 let request = EventRequest::new(
1577 input.context.session_id,
1578 EventContext::turn(input.context.turn_id, input.context.input_message_id),
1579 SessionModelChangedData {
1580 previous_model_id: Some(previous_model_id),
1581 previous_model_name,
1582 model_id,
1583 model_name: assembled.model.model.clone(),
1584 },
1585 );
1586 if let Err(e) = adapter.event_emitter().emit(request).await {
1587 warn!(error = %e, "Failed to emit session.model.changed event");
1588 }
1589}
1590
1591fn model_switch(messages: &[Message]) -> Option<(ModelId, ModelId)> {
1598 let mut user_model_ids = messages
1599 .iter()
1600 .rev()
1601 .filter(|message| message.role == MessageRole::User)
1602 .map(|message| {
1603 message
1604 .controls
1605 .as_ref()
1606 .and_then(|controls| controls.model_id)
1607 });
1608
1609 let model_id = user_model_ids.next().flatten()?;
1613 let previous_model_id = user_model_ids.next().flatten()?;
1614 (previous_model_id != model_id).then_some((previous_model_id, model_id))
1615}
1616
1617pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1618 adapter: &A,
1619 input: ActInput,
1620) -> everruns_provider::error::Result<ActResult> {
1621 let org_id = input.org_id.ok_or_else(|| {
1622 everruns_provider::error::AgentLoopError::config(
1623 "ActInput.org_id must be set for runtime host execution",
1624 )
1625 })?;
1626
1627 if let Some(blocker) =
1628 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1629 {
1630 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1631 .dependency_blocked(
1632 input.context.turn_id,
1633 input.context.input_message_id,
1634 blocker,
1635 )
1636 .await?;
1637 return Ok(ActResult {
1638 results: vec![],
1639 completed: true,
1640 success_count: 0,
1641 error_count: 1,
1642 waiting_for_tool_results: false,
1643 waiting_for_url_elicitation: false,
1644 blocked: true,
1645 client_tool_calls: vec![],
1646 client_tool_definitions: vec![],
1647 });
1648 }
1649
1650 let execution_capabilities = load_execution_capabilities(
1651 adapter,
1652 org_id,
1653 input.context.session_id,
1654 input.harness_id,
1655 input.agent_id,
1656 input.locale.clone(),
1657 input.blueprint_id.as_deref(),
1658 )
1659 .await?;
1660 let mut tool_registry = execution_capabilities.tool_registry;
1661
1662 if let Some(augmentor) = adapter.tool_augmentor() {
1663 augmentor
1664 .augment_act_tools(
1665 input.context.session_id,
1666 adapter.session_store(org_id),
1667 adapter.session_task_registry(),
1668 adapter.file_store(),
1669 &input.tool_definitions,
1670 &mut tool_registry,
1671 )
1672 .await?;
1673 }
1674
1675 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1685 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1686 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1687 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1688 tool_registry.register_boxed(tool);
1689 }
1690 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1691 &input.tool_definitions,
1692 invoker,
1693 )));
1694 }
1695
1696 let builtin_tool_registry = Arc::new(tool_registry.clone());
1697 let context_services = runtime_tool_context_services(
1698 adapter,
1699 org_id,
1700 input.context.session_id,
1701 input.agent_id,
1702 Some(builtin_tool_registry),
1703 mcp_invoker,
1704 execution_capabilities.subagent_nesting_policy,
1705 );
1706 tool_registry.validate_context_services(&context_services)?;
1707 let executor: Arc<dyn everruns_core::tool_execution::ToolExecutor> = Arc::new(tool_registry);
1708
1709 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1710 .with_context_services(context_services)
1711 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1712 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1713 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1714
1715 #[cfg(feature = "builtins")]
1716 {
1717 atom = atom.with_final_post_tool_hook(Arc::new(everruns_builtins::PersistOutputHook));
1718 }
1719
1720 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1721 atom = atom.with_outbound_tool_rate_limiter(limiter);
1722 }
1723 if let Some(store) = adapter.durable_tool_result_store() {
1724 atom = atom.with_durable_tool_result_store(store);
1725 }
1726
1727 atom.execute(input).await
1728}