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 session_storage: None,
584 };
585 let collected = collect_capabilities_with_configs(
586 &resolved.resolved_capability_configs,
587 &capability_registry,
588 &prompt_ctx,
589 )
590 .await;
591
592 let mut registry = ToolRegistry::with_defaults();
593 #[cfg(feature = "builtins")]
594 everruns_builtins::register_default_tools(&mut registry);
595 for tool in collected.tools {
596 registry.register_boxed(tool);
597 }
598
599 let mut post_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PostToolExecHook>> = resolved
604 .resolved_capability_configs
605 .iter()
606 .flat_map(|config| {
607 capability_registry
608 .get(config.capability_id())
609 .filter(|capability| capability.status().is_active())
610 .map(|capability| {
611 capability.post_tool_exec_hooks_with_config(config.config_value())
612 })
613 .unwrap_or_default()
614 })
615 .collect();
616 post_tool_hooks.sort_by_key(|hook| hook.priority());
619
620 let tool_augmentor = adapter.tool_augmentor();
627 let user_hook_specs = finalize_specs_from_configs(
628 &resolved.resolved_capability_configs,
629 &capability_registry,
630 tool_augmentor.as_deref(),
631 );
632 if user_hook_specs
637 .iter()
638 .any(|spec| spec.event == everruns_core::user_hook_types::HookEvent::UserPromptSubmit)
639 {
640 registry.unregister("query_history");
641 }
642 let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>> = resolved
645 .resolved_capability_configs
646 .iter()
647 .flat_map(|config| {
648 capability_registry
649 .get(config.capability_id())
650 .filter(|capability| capability.status().is_active())
651 .map(|capability| capability.pre_tool_use_hooks_with_config(config.config_value()))
652 .unwrap_or_default()
653 })
654 .collect();
655 if !user_hook_specs.is_empty() {
656 let dispatcher = bash_hook_dispatcher(adapter.file_store());
657 post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
658 &user_hook_specs,
659 dispatcher.clone(),
660 ));
661 pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
662 &user_hook_specs,
663 dispatcher,
664 ));
665 }
666
667 let tool_call_hooks = collected.tool_call_hooks;
678
679 Ok(RuntimeExecutionCapabilities {
680 tool_registry: registry,
681 post_tool_hooks,
682 pre_tool_hooks,
683 tool_call_hooks,
684 subagent_nesting_policy: subagent_nesting_policy_from_configs(
685 &resolved.resolved_capability_configs,
686 ),
687 })
688}
689
690fn runtime_tool_context_services<A: RuntimeHostAdapter>(
691 adapter: &A,
692 org_id: i64,
693 session_id: SessionId,
694 agent_id: Option<AgentId>,
695 tool_registry: Option<Arc<ToolRegistry>>,
696 mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
697 subagent_nesting_policy: everruns_core::delegation_services::SubagentNestingPolicy,
698) -> ToolContextServices {
699 let extensions = {
700 let mut extensions = adapter.tool_context_extensions(org_id, session_id);
701 extensions.insert(Arc::new(SessionMutatorExt(adapter.session_mutator(org_id))));
702 extensions
703 };
704 ToolContextServices {
705 file_store: Some(adapter.file_store()),
706 storage_store: adapter.storage_store(),
707 image_store: adapter.image_artifact_store(org_id),
708 provider_credential_store: adapter.provider_credential_store(org_id),
709 utility_llm_service: adapter.utility_llm_service(),
710 mcp_invoker,
711 egress_service: adapter.egress_service(),
712 message_retriever: Some(adapter.message_store()),
713 session_store: Some(adapter.session_store(org_id)),
714 agent_store: Some(adapter.agent_store(org_id)),
715 connection_resolver: adapter.connection_resolver(),
716 schedule_store: adapter.schedule_store(org_id),
717 subagent_delegate: adapter.subagent_delegate(org_id, session_id),
718 extensions,
719 leased_resource_store: adapter.leased_resource_store(),
720 session_resource_registry: adapter.session_resource_registry(),
721 session_task_registry: adapter.session_task_registry(),
722 event_emitter: Some(adapter.event_emitter()),
723 capability_registry: Some(adapter.capability_registry()),
724 tool_registry,
725 org_id: Some(
726 org_public_id_from_internal(org_id)
727 .parse()
728 .expect("internal org id converts to valid public org id"),
729 ),
730 network_access: None,
731 budget_checker: adapter.budget_checker(org_id, agent_id),
732 payment_authority: adapter.payment_authority(org_id, agent_id),
733 session_creation_authority: adapter.session_creation_authority(org_id, session_id),
734 subagent_spawn_store: adapter.subagent_spawn_store(),
735 subagent_nesting_policy,
736 reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
737 }
738}
739
740#[derive(Debug, Default)]
743struct TurnAgentIdentity {
744 id: Option<AgentId>,
745 name: Option<String>,
746 description: Option<String>,
747}
748
749pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
750 adapter: A,
751 org_id: i64,
752 session_id: SessionId,
753}
754
755impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
756 pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
757 Self {
758 adapter,
759 org_id,
760 session_id,
761 }
762 }
763
764 async fn set_session_status(
765 &self,
766 status: SessionExecutionState,
767 _action: &'static str,
768 ) -> everruns_provider::error::Result<()> {
769 self.adapter
770 .set_session_status(self.org_id, self.session_id, status)
771 .await
772 }
773
774 async fn emit_event(&self, request: EventRequest) -> everruns_provider::error::Result<()> {
775 self.adapter.event_emitter().emit(request).await.map(|_| ())
776 }
777
778 pub async fn turn_started(
779 &self,
780 turn_id: TurnId,
781 input_message_id: MessageId,
782 ) -> everruns_provider::error::Result<()> {
783 let input_content = self
784 .adapter
785 .message_store()
786 .get(self.session_id, input_message_id)
787 .await
788 .ok()
789 .flatten()
790 .map(|message| message.content_to_llm_string());
791
792 self.set_session_status(SessionExecutionState::Active, "turn_started")
793 .await?;
794
795 self.emit_event(EventRequest::new(
796 self.session_id,
797 EventContext::turn(turn_id, input_message_id),
798 SessionActivatedData {
799 turn_id,
800 input_message_id,
801 },
802 ))
803 .await?;
804
805 let agent = self.agent_identity().await;
806 self.emit_event(EventRequest::new(
807 self.session_id,
808 EventContext::turn(turn_id, input_message_id),
809 TurnStartedData {
810 turn_id,
811 input_message_id,
812 input_content,
813 agent_id: agent.id,
814 agent_name: agent.name,
815 agent_description: agent.description,
816 },
817 ))
818 .await?;
819 Ok(())
820 }
821
822 async fn agent_identity(&self) -> TurnAgentIdentity {
826 let agent_id = self
827 .adapter
828 .session_store(self.org_id)
829 .get_session(self.session_id)
830 .await
831 .ok()
832 .flatten()
833 .and_then(|session| session.agent_id);
834 let Some(agent_id) = agent_id else {
835 return TurnAgentIdentity::default();
836 };
837 let agent = self
838 .adapter
839 .agent_store(self.org_id)
840 .get_agent(agent_id)
841 .await
842 .ok()
843 .flatten();
844 TurnAgentIdentity {
845 id: Some(agent_id),
846 name: agent
849 .as_ref()
850 .map(|a| a.display_name.clone().unwrap_or_else(|| a.name.clone())),
851 description: agent.and_then(|a| a.description),
852 }
853 }
854
855 pub async fn emit_turn_completed(
856 &self,
857 input_message_id: MessageId,
858 data: TurnCompletedData,
859 ) -> everruns_provider::error::Result<()> {
860 let turn_id = data.turn_id;
861 self.emit_event(EventRequest::new(
862 self.session_id,
863 EventContext::turn(turn_id, input_message_id),
864 data,
865 ))
866 .await
867 }
868
869 pub async fn emit_session_idled(
870 &self,
871 turn_id: TurnId,
872 input_message_id: MessageId,
873 iterations: Option<u32>,
874 usage: Option<TokenUsage>,
875 ) -> everruns_provider::error::Result<()> {
876 self.set_session_status(SessionExecutionState::Idle, "emit_session_idled")
877 .await?;
878
879 self.emit_event(EventRequest::new(
880 self.session_id,
881 EventContext::turn(turn_id, input_message_id),
882 SessionIdledData {
883 turn_id,
884 iterations,
885 usage,
886 },
887 ))
888 .await
889 }
890
891 pub async fn turn_completed(
892 &self,
893 turn_id: TurnId,
894 input_message_id: MessageId,
895 iterations: u32,
896 usage: Option<TokenUsage>,
897 input_content: Option<String>,
898 ) -> everruns_provider::error::Result<()> {
899 self.emit_turn_completed(
900 input_message_id,
901 TurnCompletedData {
902 turn_id,
903 iterations,
904 duration_ms: None,
905 usage: usage.clone(),
906 input_content,
907 final_message_id: None,
908 final_answer_preview: None,
909 time_to_first_token_ms: None,
910 tool_call_count: None,
911 llm_call_count: None,
912 status: Some("completed".to_string()),
913 },
914 )
915 .await?;
916 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
917 .await
918 }
919
920 pub async fn turn_sealed(
927 &self,
928 turn_id: TurnId,
929 input_message_id: MessageId,
930 reason: &str,
931 iterations: u32,
932 usage: Option<TokenUsage>,
933 ) -> everruns_provider::error::Result<()> {
934 let context = EventContext::turn(turn_id, input_message_id);
935
936 self.emit_event(EventRequest::new(
937 self.session_id,
938 context.clone(),
939 everruns_core::events::TurnSealedData {
940 turn_id,
941 reason: reason.to_string(),
942 detail: None,
943 iterations: Some(iterations),
944 usage: usage.clone(),
945 },
946 ))
947 .await?;
948
949 self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
950 .await
951 }
952
953 pub async fn fire_turn_end_hooks(
957 &self,
958 harness_id: HarnessId,
959 agent_id: Option<AgentId>,
960 turn_id: TurnId,
961 success: bool,
962 ) {
963 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
964 &self.adapter,
965 self.org_id,
966 self.session_id,
967 harness_id,
968 agent_id,
969 )
970 .await
971 {
972 Ok(pair) => pair,
973 Err(error) => {
974 warn!(
975 session_id = %self.session_id,
976 %error,
977 "failed to collect turn_end hook specs; skipping"
978 );
979 return;
980 }
981 };
982 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
983 &specs,
984 everruns_core::user_hook_types::HookEvent::TurnEnd,
985 dispatcher,
986 );
987 if hooks.is_empty() {
988 return;
989 }
990 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
991 session_id: self.session_id,
992 turn_id: Some(turn_id),
993 org_id: org_public_id_from_internal(self.org_id).parse().ok(),
994 agent_id: agent_id.map(|a| a.to_string()),
995 };
996 everruns_core::lifecycle_hooks::run_turn_end_hooks(
997 &hooks,
998 &ctx,
999 serde_json::json!({ "success": success }),
1000 )
1001 .await;
1002 }
1003
1004 pub async fn user_prompt_blocked(
1009 &self,
1010 turn_id: TurnId,
1011 input_message_id: MessageId,
1012 reason: &str,
1013 user_message: Option<&str>,
1014 ) -> everruns_provider::error::Result<()> {
1015 let user_error =
1016 UserFacingError::new(everruns_provider::user_facing_error::codes::BLOCKED_BY_HOOK);
1017 let shown = user_message.unwrap_or(reason);
1018 let mut error_message = Message::assistant(shown);
1019 let mut metadata = std::collections::HashMap::new();
1020 user_error.apply_to_message_metadata(&mut metadata);
1021 error_message.metadata = Some(metadata);
1022
1023 self.emit_event(EventRequest::new(
1024 self.session_id,
1025 EventContext::turn(turn_id, input_message_id),
1026 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1027 ))
1028 .await?;
1029
1030 self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
1031 .await
1032 }
1033
1034 pub async fn turn_failed(
1035 &self,
1036 turn_id: TurnId,
1037 input_message_id: MessageId,
1038 error: &str,
1039 user_error: Option<&UserFacingError>,
1040 ) -> everruns_provider::error::Result<()> {
1041 self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
1042 .await
1043 }
1044
1045 pub async fn turn_failed_with_disclosure(
1049 &self,
1050 turn_id: TurnId,
1051 input_message_id: MessageId,
1052 error: &str,
1053 user_error: Option<&UserFacingError>,
1054 disclosure: Option<ErrorDisclosure>,
1055 ) -> everruns_provider::error::Result<()> {
1056 self.set_session_status(SessionExecutionState::Idle, "turn_failed")
1057 .await?;
1058
1059 self.emit_event(EventRequest::new(
1060 self.session_id,
1061 EventContext::turn(turn_id, input_message_id),
1062 {
1063 let mut data = TurnFailedData {
1064 turn_id,
1065 error: error.to_string(),
1066 error_code: None,
1067 error_fields: None,
1068 error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
1069 };
1070 if let Some(user_error) = user_error {
1071 user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
1072 }
1073 data
1074 },
1075 ))
1076 .await?;
1077
1078 self.emit_event(EventRequest::new(
1079 self.session_id,
1080 EventContext::turn(turn_id, input_message_id),
1081 SessionIdledData {
1082 turn_id,
1083 iterations: None,
1084 usage: None,
1085 },
1086 ))
1087 .await
1088 }
1089
1090 pub async fn waiting_for_tool_results(&self) -> everruns_provider::error::Result<()> {
1091 self.set_session_status(
1092 SessionExecutionState::WaitingForToolResults,
1093 "waiting_for_tool_results",
1094 )
1095 .await
1096 }
1097
1098 pub async fn dependency_blocked(
1099 &self,
1100 turn_id: TurnId,
1101 input_message_id: MessageId,
1102 blocker: DependencyBlocker,
1103 ) -> everruns_provider::error::Result<()> {
1104 let user_error = UserFacingError::new(blocker.error_code())
1105 .with_field(
1106 "dependency",
1107 match blocker {
1108 DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
1109 "harness"
1110 }
1111 DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
1112 },
1113 )
1114 .with_field(
1115 "state",
1116 match blocker {
1117 DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
1118 "archived"
1119 }
1120 DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
1121 "deleted"
1122 }
1123 },
1124 );
1125 let mut error_message = Message::assistant(blocker.message());
1126 let mut metadata = std::collections::HashMap::new();
1127 user_error.apply_to_message_metadata(&mut metadata);
1128 error_message.metadata = Some(metadata);
1129
1130 self.emit_event(EventRequest::new(
1131 self.session_id,
1132 EventContext::turn(turn_id, input_message_id),
1133 OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1134 ))
1135 .await?;
1136
1137 self.turn_failed(
1138 turn_id,
1139 input_message_id,
1140 blocker.message(),
1141 Some(&user_error),
1142 )
1143 .await
1144 }
1145}
1146
1147pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1148 adapter: &A,
1149 org_id: i64,
1150 harness_id: HarnessId,
1151 agent_id: Option<AgentId>,
1152) -> everruns_provider::error::Result<Option<DependencyBlocker>> {
1153 let harness_store = adapter.harness_store(org_id);
1154 let agent_store = adapter.agent_store(org_id);
1155 if let Some(blocker) = harness_store.get_harness_blocker(harness_id).await? {
1156 return Ok(Some(blocker));
1157 }
1158 if let Some(agent_id) = agent_id
1159 && let Some(blocker) = agent_store.get_agent_blocker(agent_id).await?
1160 {
1161 return Ok(Some(blocker));
1162 }
1163 Ok(None)
1164}
1165
1166pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1167 adapter: &A,
1168 org_id: i64,
1169 input: InputAtomInput,
1170) -> everruns_provider::error::Result<InputAtomResult> {
1171 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1175 handle.set(None);
1176 }
1177
1178 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1179 .turn_started(input.context.turn_id, input.context.input_message_id)
1180 .await?;
1181
1182 let atom = InputAtom::new(adapter.message_store());
1183 atom.execute(input).await
1184}
1185
1186pub(crate) struct UserPromptHookResult {
1193 pub(crate) decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1194 pub(crate) original_message: String,
1195}
1196
1197pub(crate) async fn run_user_prompt_submit_for_message<A: RuntimeHostAdapter>(
1198 adapter: &A,
1199 org_id: i64,
1200 input: &ReasonInput,
1201 message_text: String,
1202) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1203 let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1204 adapter,
1205 org_id,
1206 input.context.session_id,
1207 input.harness_id,
1208 input.agent_id,
1209 )
1210 .await
1211 {
1212 Ok(pair) => pair,
1213 Err(error) => {
1214 warn!(
1215 session_id = %input.context.session_id,
1216 %error,
1217 "failed to collect user_prompt_submit hook specs; continuing without them"
1218 );
1219 return Ok(None);
1220 }
1221 };
1222 let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1223 &specs,
1224 everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1225 dispatcher,
1226 );
1227 if hooks.is_empty() {
1228 return Ok(None);
1229 }
1230
1231 let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1232 session_id: input.context.session_id,
1233 turn_id: Some(input.context.turn_id),
1234 org_id: org_public_id_from_internal(org_id).parse().ok(),
1235 agent_id: input.agent_id.map(|a| a.to_string()),
1236 };
1237 let original_message = message_text.clone();
1238 let decision =
1239 everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1240 .await;
1241 Ok(Some(UserPromptHookResult {
1242 decision,
1243 original_message,
1244 }))
1245}
1246
1247async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1248 adapter: &A,
1249 org_id: i64,
1250 input: &ReasonInput,
1251) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1252 let message_text = adapter
1253 .message_store()
1254 .get(input.context.session_id, input.context.input_message_id)
1255 .await
1256 .ok()
1257 .flatten()
1258 .map(|m| m.content_to_llm_string())
1259 .unwrap_or_default();
1260 run_user_prompt_submit_for_message(adapter, org_id, input, message_text).await
1261}
1262
1263pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1264 adapter: &A,
1265 org_id: i64,
1266 input: ReasonInput,
1267) -> everruns_provider::error::Result<ReasonResult> {
1268 let prompt_message_ids = (input.iteration <= 1)
1269 .then_some(input.context.input_message_id)
1270 .into_iter()
1271 .collect();
1272 execute_reason_activity_with_prompt_messages(adapter, org_id, input, prompt_message_ids).await
1273}
1274
1275pub async fn execute_reason_activity_with_prompt_messages<A: RuntimeHostAdapter>(
1280 adapter: &A,
1281 org_id: i64,
1282 input: ReasonInput,
1283 prompt_message_ids: Vec<MessageId>,
1284) -> everruns_provider::error::Result<ReasonResult> {
1285 if let Some(blocker) =
1286 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1287 {
1288 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1289 .dependency_blocked(
1290 input.context.turn_id,
1291 input.context.input_message_id,
1292 blocker,
1293 )
1294 .await?;
1295 return Ok(ReasonResult {
1296 native_counts: None,
1297 success: false,
1298 text: blocker.message().to_string(),
1299 tool_calls: vec![],
1300 has_tool_calls: false,
1301 tool_definitions: vec![],
1302 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1303 error: Some("dependency_unavailable".to_string()),
1304 user_facing_error: None,
1305 error_disclosure: None,
1306 usage: None,
1307 output_message_id: None,
1308 time_to_first_token_ms: None,
1309 response_id: None,
1310 finish_reason: None,
1311 locale: None,
1312 network_access: None,
1313 parallel_tool_calls: None,
1314 });
1315 }
1316
1317 let mut user_prompt_message_overrides = Vec::new();
1322 for message_id in prompt_message_ids {
1323 let mut hook_input = input.clone();
1324 hook_input.context.input_message_id = message_id;
1325 let Some(hook_result) =
1326 run_user_prompt_submit_for_turn(adapter, org_id, &hook_input).await?
1327 else {
1328 continue;
1329 };
1330 match hook_result.decision {
1331 everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1332 reason,
1333 user_message,
1334 } => {
1335 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1336 .user_prompt_blocked(
1337 input.context.turn_id,
1338 input.context.input_message_id,
1339 &reason,
1340 user_message.as_deref(),
1341 )
1342 .await?;
1343 return Ok(ReasonResult {
1344 native_counts: None,
1345 success: false,
1346 text: user_message.unwrap_or_else(|| reason.clone()),
1347 tool_calls: vec![],
1348 has_tool_calls: false,
1349 tool_definitions: vec![],
1350 max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1351 error: Some("blocked_by_user_prompt_hook".to_string()),
1352 user_facing_error: None,
1353 error_disclosure: None,
1354 usage: None,
1355 output_message_id: None,
1356 time_to_first_token_ms: None,
1357 response_id: None,
1358 finish_reason: None,
1359 locale: None,
1360 network_access: None,
1361 parallel_tool_calls: None,
1362 });
1363 }
1364 everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1365 if message != hook_result.original_message {
1366 user_prompt_message_overrides.push((message_id, message));
1367 }
1368 }
1369 }
1370 }
1371
1372 let validation_session = adapter
1376 .session_store(org_id)
1377 .get_session(input.context.session_id)
1378 .await?
1379 .ok_or_else(|| {
1380 everruns_provider::error::AgentLoopError::session_not_found(input.context.session_id)
1381 })?;
1382 let validation_capabilities = load_execution_capabilities(
1383 adapter,
1384 org_id,
1385 input.context.session_id,
1386 input.harness_id,
1387 input.agent_id,
1388 validation_session.locale.clone(),
1389 validation_session.blueprint_id.as_deref(),
1390 )
1391 .await?;
1392 let query_history_allowed = validation_capabilities
1393 .tool_registry
1394 .get("query_history")
1395 .is_some();
1396 let validation_services = runtime_tool_context_services(
1397 adapter,
1398 org_id,
1399 input.context.session_id,
1400 input.agent_id,
1401 Some(Arc::new(validation_capabilities.tool_registry.clone())),
1402 None,
1403 validation_capabilities.subagent_nesting_policy,
1404 );
1405 validation_capabilities
1406 .tool_registry
1407 .validate_context_services(&validation_services)?;
1408
1409 let mut turn_inputs = adapter
1410 .load_resolved_turn(org_id, input.context.session_id)
1411 .await?;
1412 if let Some(augmentor) = adapter.tool_augmentor() {
1413 augmentor
1414 .augment_reason_tools(
1415 input.context.session_id,
1416 adapter.session_store(org_id),
1417 adapter.session_task_registry(),
1418 &mut turn_inputs.mcp_tool_definitions,
1419 )
1420 .await?;
1421 }
1422
1423 let reason_capability_registry = {
1424 let mut registry = adapter.capability_registry();
1425 if !query_history_allowed {
1426 let query_history_owner = registry
1432 .list()
1433 .into_iter()
1434 .find(|capability| {
1435 capability
1436 .tool_definitions()
1437 .iter()
1438 .any(|tool| tool.name() == "query_history")
1439 })
1440 .map(Arc::clone);
1441 if let Some(capability) = query_history_owner {
1442 registry.register(MessageFilterOnlyCapability(capability));
1443 }
1444 }
1445 registry
1446 };
1447 let context_resolver = crate::runtime_context::StoreTurnContextResolver::new(
1448 adapter.harness_store(org_id),
1449 adapter.agent_store(org_id),
1450 adapter.session_store(org_id),
1451 adapter.message_store(),
1452 adapter.provider_store(org_id),
1453 reason_capability_registry.clone(),
1454 adapter.driver_registry(),
1455 )
1456 .with_file_store(adapter.file_store());
1457 let context_resolver = match adapter.storage_store() {
1458 Some(store) => context_resolver.with_session_storage(store),
1459 None => context_resolver,
1460 };
1461 let mut atom = ReasonAtom::new(
1462 context_resolver,
1463 adapter.message_store(),
1464 reason_capability_registry.clone(),
1465 adapter.event_emitter(),
1466 );
1467 if let Some(image_resolver) = adapter.image_resolver(org_id) {
1468 atom = atom.with_image_resolver(image_resolver);
1469 }
1470 if let Some(file_resolver) = adapter.file_resolver(org_id) {
1471 atom = atom.with_file_resolver(file_resolver);
1472 }
1473 if let Some(hb) = adapter.stream_heartbeater() {
1474 atom = atom.with_stream_heartbeater(hb);
1475 }
1476 if let Some(timeout) = adapter.provider_stall_timeout() {
1477 atom = atom.with_provider_stall_timeout(timeout);
1478 }
1479 if let Some(config) = adapter.provider_retry_config() {
1480 atom = atom.with_provider_retry_config(config);
1481 }
1482 if let Some(store) = adapter.partial_stream_store() {
1483 atom = atom.with_partial_stream_store(store);
1484 }
1485 if let Some(store) = adapter.durable_tool_result_store() {
1486 atom = atom.with_durable_tool_result_store(store);
1487 }
1488 if let Some(store) = adapter.compaction_checkpoint_store() {
1489 atom = atom.with_compaction_checkpoint_store(store);
1490 }
1491 if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1492 atom = atom.with_reasoning_effort_handle(handle);
1493 }
1494 if let Some(utility_llm_service) = adapter.utility_llm_service() {
1495 atom = atom.with_utility_llm_service(utility_llm_service);
1496 }
1497 if let Some(schedule_store) = adapter.schedule_store(org_id) {
1500 atom = atom.with_schedule_store(schedule_store);
1501 }
1502
1503 let mut assembled = crate::runtime_context::assemble_turn_context_from_snapshot(
1504 turn_inputs.snapshot,
1505 adapter.message_store().as_ref(),
1506 adapter.provider_store(org_id).as_ref(),
1507 &reason_capability_registry,
1508 &adapter.driver_registry(),
1509 &turn_inputs.mcp_tool_definitions,
1510 Some(adapter.file_store()),
1511 adapter.storage_store(),
1515 )
1516 .await?;
1517 let input = ReasonInput {
1518 mcp_tool_definitions: turn_inputs.mcp_tool_definitions,
1519 ..input
1520 };
1521
1522 if !user_prompt_message_overrides.is_empty() {
1523 for (message_id, message_override) in user_prompt_message_overrides {
1524 let message = assembled
1525 .messages
1526 .iter_mut()
1527 .find(|message| message.id == message_id)
1528 .ok_or_else(|| {
1529 everruns_provider::error::AgentLoopError::config(
1530 "user_prompt_submit mutation: input message not found in assembled context",
1531 )
1532 })?;
1533
1534 message
1537 .content
1538 .retain(|part| !matches!(part, ContentPart::Text(_)));
1539 message
1540 .content
1541 .insert(0, ContentPart::text(message_override));
1542 }
1543 }
1544 if input.iteration <= 1 {
1549 emit_model_change_if_switched(adapter, org_id, &input, &assembled).await;
1550 }
1551
1552 crate::native_async::execute_reason(adapter, org_id, input, assembled, atom).await
1553}
1554
1555async fn emit_model_change_if_switched<A: RuntimeHostAdapter>(
1560 adapter: &A,
1561 org_id: i64,
1562 input: &ReasonInput,
1563 assembled: &AssembledTurnContext,
1564) {
1565 let Some((previous_model_id, model_id)) = model_switch(&assembled.messages) else {
1566 return;
1567 };
1568 if assembled.resolved_model_id != Some(model_id) {
1569 return;
1572 }
1573
1574 let previous_model_name = adapter
1578 .provider_store(org_id)
1579 .get_model_spec(previous_model_id)
1580 .await
1581 .ok()
1582 .flatten()
1583 .map(|spec| spec.model);
1584
1585 let request = EventRequest::new(
1586 input.context.session_id,
1587 EventContext::turn(input.context.turn_id, input.context.input_message_id),
1588 SessionModelChangedData {
1589 previous_model_id: Some(previous_model_id),
1590 previous_model_name,
1591 model_id,
1592 model_name: assembled.model.model.clone(),
1593 },
1594 );
1595 if let Err(e) = adapter.event_emitter().emit(request).await {
1596 warn!(error = %e, "Failed to emit session.model.changed event");
1597 }
1598}
1599
1600fn model_switch(messages: &[Message]) -> Option<(ModelId, ModelId)> {
1607 let mut user_model_ids = messages
1608 .iter()
1609 .rev()
1610 .filter(|message| message.role == MessageRole::User)
1611 .map(|message| {
1612 message
1613 .controls
1614 .as_ref()
1615 .and_then(|controls| controls.model_id)
1616 });
1617
1618 let model_id = user_model_ids.next().flatten()?;
1622 let previous_model_id = user_model_ids.next().flatten()?;
1623 (previous_model_id != model_id).then_some((previous_model_id, model_id))
1624}
1625
1626pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1627 adapter: &A,
1628 input: ActInput,
1629) -> everruns_provider::error::Result<ActResult> {
1630 let org_id = input.org_id.ok_or_else(|| {
1631 everruns_provider::error::AgentLoopError::config(
1632 "ActInput.org_id must be set for runtime host execution",
1633 )
1634 })?;
1635
1636 if let Some(blocker) =
1637 detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1638 {
1639 RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1640 .dependency_blocked(
1641 input.context.turn_id,
1642 input.context.input_message_id,
1643 blocker,
1644 )
1645 .await?;
1646 return Ok(ActResult {
1647 results: vec![],
1648 completed: true,
1649 success_count: 0,
1650 error_count: 1,
1651 waiting_for_tool_results: false,
1652 waiting_for_url_elicitation: false,
1653 blocked: true,
1654 client_tool_calls: vec![],
1655 client_tool_definitions: vec![],
1656 });
1657 }
1658
1659 let execution_capabilities = load_execution_capabilities(
1660 adapter,
1661 org_id,
1662 input.context.session_id,
1663 input.harness_id,
1664 input.agent_id,
1665 input.locale.clone(),
1666 input.blueprint_id.as_deref(),
1667 )
1668 .await?;
1669 let mut tool_registry = execution_capabilities.tool_registry;
1670
1671 if let Some(augmentor) = adapter.tool_augmentor() {
1672 augmentor
1673 .augment_act_tools(
1674 input.context.session_id,
1675 adapter.session_store(org_id),
1676 adapter.session_task_registry(),
1677 adapter.file_store(),
1678 &input.tool_definitions,
1679 &mut tool_registry,
1680 )
1681 .await?;
1682 }
1683
1684 let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1694 if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1695 let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1696 for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1697 tool_registry.register_boxed(tool);
1698 }
1699 mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1700 &input.tool_definitions,
1701 invoker,
1702 )));
1703 }
1704
1705 let builtin_tool_registry = Arc::new(tool_registry.clone());
1706 let context_services = runtime_tool_context_services(
1707 adapter,
1708 org_id,
1709 input.context.session_id,
1710 input.agent_id,
1711 Some(builtin_tool_registry),
1712 mcp_invoker,
1713 execution_capabilities.subagent_nesting_policy,
1714 );
1715 tool_registry.validate_context_services(&context_services)?;
1716 let executor: Arc<dyn everruns_core::tool_execution::ToolExecutor> = Arc::new(tool_registry);
1717
1718 let mut atom = ActAtom::new(executor, adapter.event_emitter())
1719 .with_context_services(context_services)
1720 .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1721 .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1722 .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1723
1724 #[cfg(feature = "builtins")]
1725 {
1726 atom = atom.with_final_post_tool_hook(Arc::new(everruns_builtins::PersistOutputHook));
1727 }
1728
1729 if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1730 atom = atom.with_outbound_tool_rate_limiter(limiter);
1731 }
1732 if let Some(store) = adapter.durable_tool_result_store() {
1733 atom = atom.with_durable_tool_result_store(store);
1734 }
1735
1736 atom.execute(input).await
1737}