Skip to main content

everruns_runtime/
host.rs

1// Shared host orchestration for embedded and durable execution hosts.
2// Decision: everruns-runtime owns worker-facing turn phase execution so
3// durable/server-backed hosts reuse the same input/reason/act wiring.
4
5use async_trait::async_trait;
6use everruns_core::atoms::{
7    ActAtom, ActInput, ActResult, Atom, InputAtom, InputAtomInput, InputAtomResult, ReasonAtom,
8    ReasonInput, ReasonResult,
9};
10use everruns_core::capabilities::{
11    SystemPromptContext, collect_capabilities_with_configs, report_result_tool_for_child_session,
12    report_task_progress_tool_for_child_session,
13};
14use everruns_core::events::{
15    EventContext, EventRequest, OutputMessageCompletedData, SessionActivatedData, SessionIdledData,
16    TurnCompletedData, TurnFailedData, TurnStartedData,
17};
18use everruns_core::message::{ContentPart, Message};
19use everruns_core::message_retriever::MessageRetriever;
20use everruns_core::platform_store::PlatformStore;
21use everruns_core::session::SessionStatus;
22use everruns_core::tools::Tool;
23use everruns_core::traits::{
24    AgentStore, BudgetChecker, EventEmitter, HarnessStore, ImageArtifactStore, ImageResolver,
25    LeasedResourceStore, PaymentAuthority, ProviderCredentialStore, ProviderStore, ResolvedModel,
26    SessionCreationAuthority, SessionFileSystem, SessionMutator, SessionResourceRegistry,
27    SessionScheduleStore, SessionSqlDbStoreRef, SessionStorageStore, SessionStore,
28    ToolContextServices, UserConnectionResolver,
29};
30use everruns_core::typed_id::{AgentId, HarnessId, MessageId, SessionId, TurnId};
31use everruns_core::vector_store::KnowledgeIndexSearch;
32use everruns_core::{
33    Agent, CapabilityRegistry, CapabilityStatus, DependencyBlocker, DriverRegistry, EgressService,
34    ErrorDisclosure, Harness, Session, TokenUsage, ToolDefinition, ToolRegistry, UserFacingError,
35    UtilityLlmService, assemble_turn_context, org_public_id_from_internal,
36    resolve_runtime_capabilities,
37};
38use std::sync::Arc;
39use tracing::warn;
40
41/// Turn context loaded in one batched call for runtime host execution.
42#[derive(Debug, Clone)]
43pub struct RuntimeHostTurnContext {
44    pub agent: Option<Agent>,
45    pub session: Session,
46    pub messages: Vec<Message>,
47    pub model: Option<ResolvedModel>,
48    pub mcp_tool_definitions: Vec<ToolDefinition>,
49}
50
51/// Public adapter contract for server-backed or durable runtime hosts.
52///
53/// `everruns-runtime` owns shared host orchestration for both embedded and
54/// durable execution. That includes phase execution (`input -> reason -> act`),
55/// lifecycle emission, and the generic turn-strategy decisions used by durable
56/// or custom hosts.
57///
58/// Host crates implement this trait to provide persistence, session-lifecycle
59/// plumbing, event delivery, and their own orchestration backend. The durable
60/// engine itself remains outside this crate.
61#[async_trait]
62pub trait RuntimeHostAdapter: Send + Sync + Clone + 'static {
63    async fn get_agent(
64        &self,
65        org_id: i64,
66        agent_id: AgentId,
67    ) -> everruns_core::error::Result<Option<Agent>>;
68
69    async fn get_harness(
70        &self,
71        org_id: i64,
72        harness_id: HarnessId,
73    ) -> everruns_core::error::Result<Option<Harness>>;
74
75    async fn set_session_status(
76        &self,
77        org_id: i64,
78        session_id: SessionId,
79        status: SessionStatus,
80    ) -> everruns_core::error::Result<Session>;
81
82    async fn load_turn_context(
83        &self,
84        org_id: i64,
85        session_id: SessionId,
86    ) -> everruns_core::error::Result<RuntimeHostTurnContext>;
87
88    fn capability_registry(&self) -> CapabilityRegistry;
89
90    fn driver_registry(&self) -> DriverRegistry;
91
92    fn harness_store(&self, org_id: i64) -> Arc<dyn HarnessStore>;
93
94    fn agent_store(&self, org_id: i64) -> Arc<dyn AgentStore>;
95
96    fn session_store(&self, org_id: i64) -> Arc<dyn SessionStore>;
97
98    fn session_mutator(&self, org_id: i64) -> Arc<dyn SessionMutator>;
99
100    fn provider_store(&self, org_id: i64) -> Arc<dyn ProviderStore>;
101
102    fn message_store(&self) -> Arc<dyn MessageRetriever>;
103
104    fn compaction_checkpoint_store(
105        &self,
106    ) -> Option<Arc<dyn everruns_core::CompactionCheckpointStore>> {
107        None
108    }
109
110    fn event_emitter(&self) -> Arc<dyn EventEmitter>;
111
112    fn file_store(&self) -> Arc<dyn SessionFileSystem>;
113
114    fn image_resolver(&self, _org_id: i64) -> Option<Arc<dyn ImageResolver>> {
115        None
116    }
117
118    fn image_artifact_store(&self, _org_id: i64) -> Option<Arc<dyn ImageArtifactStore>> {
119        None
120    }
121
122    fn provider_credential_store(&self, _org_id: i64) -> Option<Arc<dyn ProviderCredentialStore>> {
123        None
124    }
125
126    fn utility_llm_service(&self) -> Option<Arc<dyn UtilityLlmService>> {
127        None
128    }
129
130    fn egress_service(&self) -> Option<Arc<dyn EgressService>> {
131        None
132    }
133
134    fn storage_store(&self) -> Option<Arc<dyn SessionStorageStore>> {
135        None
136    }
137
138    /// Knowledge store backing the `search_knowledge` tool. Default: none.
139    fn knowledge_store(&self) -> Option<Arc<dyn everruns_core::traits::KnowledgeStore>> {
140        None
141    }
142
143    fn connection_resolver(&self) -> Option<Arc<dyn UserConnectionResolver>> {
144        None
145    }
146
147    fn sqldb_store(&self) -> Option<SessionSqlDbStoreRef> {
148        None
149    }
150
151    fn leased_resource_store(&self) -> Option<Arc<dyn LeasedResourceStore>> {
152        None
153    }
154
155    fn session_resource_registry(&self) -> Option<Arc<dyn SessionResourceRegistry>> {
156        None
157    }
158
159    fn session_task_registry(
160        &self,
161    ) -> Option<Arc<dyn everruns_core::session_task::SessionTaskRegistry>> {
162        None
163    }
164
165    fn schedule_store(&self, _org_id: i64) -> Option<Arc<dyn SessionScheduleStore>> {
166        None
167    }
168
169    fn platform_store(
170        &self,
171        _org_id: i64,
172        _session_id: SessionId,
173    ) -> Option<Arc<dyn PlatformStore>> {
174        None
175    }
176
177    /// Get the Knowledge Index search service for the `search_index` tool.
178    /// Org-scoped; returns None when retrieval is not available (e.g. gRPC
179    /// workers without a search RPC, or in-memory test backends).
180    fn knowledge_index_search(&self, _org_id: i64) -> Option<Arc<dyn KnowledgeIndexSearch>> {
181        None
182    }
183
184    fn budget_checker(
185        &self,
186        _org_id: i64,
187        _agent_id: Option<AgentId>,
188    ) -> Option<Arc<dyn BudgetChecker>> {
189        None
190    }
191
192    fn payment_authority(
193        &self,
194        _org_id: i64,
195        _agent_id: Option<AgentId>,
196    ) -> Option<Arc<dyn PaymentAuthority>> {
197        None
198    }
199
200    fn session_creation_authority(
201        &self,
202        _org_id: i64,
203        _session_id: SessionId,
204    ) -> Option<Arc<dyn SessionCreationAuthority>> {
205        None
206    }
207
208    /// Per-org outbound tool-call rate limiter (TM-TOOL-009).
209    /// Default: `None` (no rate limiting — suitable for in-process / test environments).
210    fn outbound_tool_rate_limiter(
211        &self,
212        _org_id: i64,
213    ) -> Option<Arc<dyn everruns_core::OutboundToolRateLimiter>> {
214        None
215    }
216
217    /// Per-turn durable tool result store for act-activity idempotency (EVE-530).
218    /// Default: `None` (no durable claim/settle — every execution runs tools fresh).
219    fn durable_tool_result_store(&self) -> Option<Arc<dyn everruns_core::DurableToolResultStore>> {
220        None
221    }
222
223    /// Durable subagent spawn handle store for reattach on reclaim (EVE-535).
224    /// Default: `None` (no spawn dedup — dev/test mode or hosts without durable execution).
225    fn subagent_spawn_store(&self) -> Option<Arc<dyn everruns_core::SubagentSpawnStore>> {
226        None
227    }
228
229    /// Stream-liveness heartbeater for the Reason activity (EVE-531).
230    /// Default: `None` (no heartbeats sent — durable workers supply one).
231    fn stream_heartbeater(&self) -> Option<Arc<dyn everruns_core::StreamHeartbeater>> {
232        None
233    }
234
235    /// Partial-stream store for ContinuePartial recovery (EVE-532).
236    /// Default: `None` (no recovery; in-memory and dev hosts use this default).
237    fn partial_stream_store(&self) -> Option<Arc<dyn everruns_core::PartialStreamStore>> {
238        None
239    }
240
241    /// Live, turn-scoped reasoning-effort handle for the given session (EVE-595).
242    ///
243    /// When a host returns a handle, the Reason activity re-reads it on every
244    /// LLM step and the Act activity hands the same instance to each tool's
245    /// `ToolContext`. A tool can then change effort mid-turn and have subsequent
246    /// LLM steps in the same turn observe it. Hosts MUST return the *same*
247    /// handle instance for a session across reason/act activities of one turn.
248    /// Default: `None` (effort is resolved solely from message controls).
249    fn reasoning_effort_handle(
250        &self,
251        _session_id: SessionId,
252    ) -> Option<everruns_core::ReasoningEffortHandle> {
253        None
254    }
255
256    /// Provider stall timeout for the Reason activity (EVE-531).
257    /// Default: `None` (use built-in 120s default).
258    fn provider_stall_timeout(&self) -> Option<std::time::Duration> {
259        None
260    }
261
262    /// MCP executor routing `mcp_*` tool calls for this session, if the host
263    /// configures MCP (specs/runtime-mcp.md D4). Default: `None`, so hosts
264    /// without scoped MCP servers keep the plain tool registry unchanged.
265    async fn mcp_executor(
266        &self,
267        _org_id: i64,
268        _session_id: SessionId,
269    ) -> Option<Arc<everruns_mcp::McpExecutor>> {
270        None
271    }
272}
273
274struct RuntimeExecutionCapabilities {
275    tool_registry: ToolRegistry,
276    post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>>,
277    pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>>,
278    tool_call_hooks: Vec<Arc<dyn everruns_core::ToolCallHook>>,
279    subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
280}
281
282fn subagent_nesting_policy_from_configs(
283    resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
284) -> everruns_core::SubagentNestingPolicy {
285    let subagents_config = resolved_capability_configs.iter().find(|config| {
286        config.capability_id() == everruns_core::capabilities::SUBAGENTS_CAPABILITY_ID
287    });
288
289    let configured_depth = subagents_config
290        .and_then(|config| {
291            config
292                .config
293                .get("max_subagent_depth")
294                .or_else(|| config.config.get("max_depth"))
295        })
296        .and_then(|value| value.as_u64())
297        .and_then(|value| u32::try_from(value).ok());
298    let configured_max_active = subagents_config
299        .and_then(|config| {
300            config
301                .config
302                .get("max_active_descendant_tasks")
303                .or_else(|| config.config.get("max_concurrent_descendant_tasks"))
304        })
305        .and_then(|value| value.as_u64())
306        .and_then(|value| u32::try_from(value).ok());
307    let configured_max_total = subagents_config
308        .and_then(|config| config.config.get("max_total_descendant_tasks"))
309        .and_then(|value| value.as_u64())
310        .and_then(|value| u32::try_from(value).ok());
311    let configured_max_active_detached = subagents_config
312        .and_then(|config| config.config.get("max_active_detached_tasks"))
313        .and_then(|value| value.as_u64())
314        .and_then(|value| u32::try_from(value).ok());
315    let configured_max_total_detached = subagents_config
316        .and_then(|config| config.config.get("max_total_detached_tasks"))
317        .and_then(|value| value.as_u64())
318        .and_then(|value| u32::try_from(value).ok());
319
320    everruns_core::SubagentNestingPolicy::default()
321        .with_agent_override(configured_depth)
322        .with_agent_task_caps_override(configured_max_active, configured_max_total)
323        .with_agent_detached_task_caps_override(
324            configured_max_active_detached,
325            configured_max_total_detached,
326        )
327}
328
329/// Collect and finalize user-hook specs for a session from its resolved
330/// capability configs, plus the shared bash dispatcher used to run them.
331///
332/// This is the single place hook specs are gathered so every firing point —
333/// the act path (`load_execution_capabilities`) and the lifecycle firing
334/// points (`execute_reason_activity` for `user_prompt_submit`, turn completion
335/// for `turn_end`, and the server session paths) — applies identical
336/// `finalize_hook_specs` semantics: `{capability_id}:` namespace stamping,
337/// stable default ids, and `disabled_contributions` muting (TM-HOOK-004).
338fn finalize_specs_from_configs(
339    resolved_capability_configs: &[everruns_core::capability_types::AgentCapabilityConfig],
340    capability_registry: &CapabilityRegistry,
341) -> Vec<everruns_core::user_hook_types::UserHookSpec> {
342    let mut hook_contributions: Vec<(String, Vec<everruns_core::user_hook_types::UserHookSpec>)> =
343        Vec::new();
344    let mut disabled_contributions: Vec<String> = Vec::new();
345    for config in resolved_capability_configs {
346        let Some(capability) = capability_registry.get(config.capability_id()) else {
347            continue;
348        };
349        let specs = capability.user_hooks_with_config(&config.config);
350        if !specs.is_empty() {
351            hook_contributions.push((config.capability_id().to_string(), specs));
352        }
353        if config.capability_id() == "user_hooks" {
354            disabled_contributions.extend(
355                everruns_core::capabilities::user_hooks::disabled_contributions(&config.config),
356            );
357        }
358    }
359    everruns_core::hook_adapter::finalize_hook_specs(hook_contributions, &disabled_contributions)
360}
361
362/// Resolve a session's capability configs and collect finalized hook specs.
363/// Used by the lifecycle firing points, which need specs outside the act path.
364/// Returns `(specs, dispatcher)`; `specs` is empty when the session has no
365/// hook-contributing capabilities.
366async fn collect_lifecycle_hook_specs<A: RuntimeHostAdapter>(
367    adapter: &A,
368    org_id: i64,
369    session_id: SessionId,
370    harness_id: HarnessId,
371    agent_id: Option<AgentId>,
372) -> everruns_core::error::Result<(
373    Vec<everruns_core::user_hook_types::UserHookSpec>,
374    Arc<dyn everruns_core::hook_executor::BashHookDispatcher>,
375)> {
376    let capability_registry = adapter.capability_registry();
377    let harness_chain = adapter
378        .harness_store(org_id)
379        .get_harness_chain(harness_id)
380        .await?;
381    if harness_chain.is_empty() {
382        return Err(everruns_core::error::AgentLoopError::harness_not_found(
383            harness_id,
384        ));
385    }
386    let session = adapter
387        .session_store(org_id)
388        .get_session(session_id)
389        .await?
390        .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
391    let agent = match agent_id {
392        Some(agent_id) => adapter.agent_store(org_id).get_agent(agent_id).await?,
393        None => None,
394    };
395    let resolved = resolve_runtime_capabilities(
396        &harness_chain,
397        agent.as_ref(),
398        &session,
399        &capability_registry,
400    );
401    let specs =
402        finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
403    let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
404        everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
405    );
406    Ok((specs, dispatcher))
407}
408
409async fn load_execution_capabilities<A: RuntimeHostAdapter>(
410    adapter: &A,
411    org_id: i64,
412    session_id: SessionId,
413    harness_id: HarnessId,
414    agent_id: Option<AgentId>,
415    locale: Option<String>,
416    blueprint_id: Option<&str>,
417) -> everruns_core::error::Result<RuntimeExecutionCapabilities> {
418    let capability_registry = adapter.capability_registry();
419    if let Some(blueprint_id) = blueprint_id {
420        let mut registry = ToolRegistry::with_defaults();
421        let blueprint = capability_registry.blueprint(blueprint_id).ok_or_else(|| {
422            everruns_core::error::AgentLoopError::config(format!(
423                "Blueprint \"{blueprint_id}\" not found in registry"
424            ))
425        })?;
426        for tool in blueprint.tools {
427            registry.register_boxed(tool);
428        }
429        return Ok(RuntimeExecutionCapabilities {
430            tool_registry: registry,
431            post_tool_hooks: Vec::new(),
432            pre_tool_hooks: Vec::new(),
433            tool_call_hooks: Vec::new(),
434            subagent_nesting_policy: everruns_core::SubagentNestingPolicy::default(),
435        });
436    }
437
438    let harness_chain = adapter
439        .harness_store(org_id)
440        .get_harness_chain(harness_id)
441        .await?;
442    if harness_chain.is_empty() {
443        return Err(everruns_core::error::AgentLoopError::harness_not_found(
444            harness_id,
445        ));
446    }
447
448    let session = adapter
449        .session_store(org_id)
450        .get_session(session_id)
451        .await?
452        .ok_or_else(|| everruns_core::error::AgentLoopError::session_not_found(session_id))?;
453
454    let agent_store = adapter.agent_store(org_id);
455    let agent = match agent_id {
456        Some(agent_id) => Some(
457            agent_store
458                .get_agent(agent_id)
459                .await?
460                .ok_or_else(|| everruns_core::error::AgentLoopError::agent_not_found(agent_id))?,
461        ),
462        None => None,
463    };
464
465    let resolved = resolve_runtime_capabilities(
466        &harness_chain,
467        agent.as_ref(),
468        &session,
469        &capability_registry,
470    );
471    // Executor (act) path: this builds the worker-side tool registry, not the
472    // model-visible tool list. The model is left unset, so a model-adaptive
473    // capability like `auto_tool_search` resolves to its provider-agnostic
474    // client-side mechanism here. That registers the `tool_search` tool in the
475    // executor, which is a harmless superset: on native models the reason path
476    // never shows that tool to the model, so it is simply never called.
477    let prompt_ctx = SystemPromptContext {
478        session_id,
479        locale: locale.or(session.locale.clone()),
480        // Pin system-prompt file reads to the session's workspace (the default
481        // 1:1 case is a transparent pass-through), then resolve through the
482        // mount resolver (EVE-660): `/workspace` is a mount + cwd.
483        // `scoped_prompt_file_store` wraps with `wrap_if_needed` so a local
484        // embedder's backend-native display policy survives here too (it must
485        // match the reason path — see its doc); server stores stay on `/workspace`.
486        file_store: Some(everruns_core::scoped_prompt_file_store(
487            adapter.file_store(),
488            session.workspace_id,
489        )),
490        model: None,
491    };
492    let collected = collect_capabilities_with_configs(
493        &resolved.resolved_capability_configs,
494        &capability_registry,
495        &prompt_ctx,
496    )
497    .await;
498
499    let mut registry = ToolRegistry::with_defaults();
500    for tool in collected.tools {
501        registry.register_boxed(tool);
502    }
503
504    // Only `Available` capabilities contribute hooks, matching
505    // `collect_capabilities_with_configs` (which skips non-available
506    // capabilities). This keeps a `ComingSoon`/unavailable capability from
507    // affecting execution via any of its hook seams.
508    let mut post_tool_hooks: Vec<Arc<dyn everruns_core::PostToolExecHook>> = resolved
509        .resolved_capability_configs
510        .iter()
511        .flat_map(|config| {
512            capability_registry
513                .get(config.capability_id())
514                .filter(|capability| capability.status() == CapabilityStatus::Available)
515                .map(|capability| capability.post_tool_exec_hooks_with_config(&config.config))
516                .unwrap_or_default()
517        })
518        .collect();
519    // Tool-output guardrails must inspect the original result before other
520    // capability hooks can persist or compact it into secondary surfaces.
521    post_tool_hooks.sort_by_key(|hook| hook.priority());
522
523    // User-hook contributions (see `specs/user-hooks.md`). `finalize_specs_from_configs`
524    // gathers specs across every resolved capability — both the user-facing
525    // `user_hooks` capability and any capability that bundles hooks — and applies
526    // `finalize_hook_specs` (namespace stamping, stable ids, `disabled_contributions`
527    // muting; TM-HOOK-004). The same helper backs the lifecycle firing points so
528    // every event finalizes specs identically.
529    let user_hook_specs =
530        finalize_specs_from_configs(&resolved.resolved_capability_configs, &capability_registry);
531    // Persisted messages remain the immutable audit record, so they can contain
532    // text removed by a provider-bound user_prompt_submit hook. Until there is a
533    // durable provider-visible history view, fail closed rather than let
534    // query_history bypass that enforcement boundary.
535    if user_hook_specs
536        .iter()
537        .any(|spec| spec.event == everruns_core::user_hook_types::HookEvent::UserPromptSubmit)
538    {
539        registry.unregister("query_history");
540    }
541    // Capability-contributed pre-tool hooks run first (e.g. approval gating),
542    // then user-hook (`PreToolUse`) specs. The first hook to block wins.
543    let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::atoms::PreToolUseHook>> = resolved
544        .resolved_capability_configs
545        .iter()
546        .flat_map(|config| {
547            capability_registry
548                .get(config.capability_id())
549                .filter(|capability| capability.status() == CapabilityStatus::Available)
550                .map(|capability| capability.pre_tool_use_hooks_with_config(&config.config))
551                .unwrap_or_default()
552        })
553        .collect();
554    if !user_hook_specs.is_empty() {
555        let dispatcher: Arc<dyn everruns_core::hook_executor::BashHookDispatcher> = Arc::new(
556            everruns_core::hook_dispatch::BashkitShellHookDispatcher::new(adapter.file_store()),
557        );
558        post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
559            &user_hook_specs,
560            dispatcher.clone(),
561        ));
562        pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
563            &user_hook_specs,
564            dispatcher,
565        ));
566    }
567
568    // Use the hook list assembled by `collect_capabilities_with_configs` as the
569    // single source of truth. It already contains every explicit capability
570    // `tool_call_hooks()` followed by the generated `CapabilityNarrationHook`
571    // adapters — one per collected capability plus any auto-activated
572    // cross-cutting capability such as `background_execution`. Re-deriving only
573    // the explicit subset here dropped capability-owned narration, so tools fell
574    // back to generic `Ran {display_name}` lines (EVE-601). Explicit hooks stay
575    // first in this list, so model-authored narration (`human_intent`) keeps its
576    // precedence over default `Tool::narrate()`, and only available capabilities
577    // contributed because collection skips non-available ones.
578    let tool_call_hooks = collected.tool_call_hooks;
579
580    Ok(RuntimeExecutionCapabilities {
581        tool_registry: registry,
582        post_tool_hooks,
583        pre_tool_hooks,
584        tool_call_hooks,
585        subagent_nesting_policy: subagent_nesting_policy_from_configs(
586            &resolved.resolved_capability_configs,
587        ),
588    })
589}
590
591fn runtime_tool_context_services<A: RuntimeHostAdapter>(
592    adapter: &A,
593    org_id: i64,
594    session_id: SessionId,
595    agent_id: Option<AgentId>,
596    tool_registry: Option<Arc<ToolRegistry>>,
597    mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
598    subagent_nesting_policy: everruns_core::SubagentNestingPolicy,
599) -> ToolContextServices {
600    ToolContextServices {
601        file_store: Some(adapter.file_store()),
602        storage_store: adapter.storage_store(),
603        image_store: adapter.image_artifact_store(org_id),
604        provider_credential_store: adapter.provider_credential_store(org_id),
605        utility_llm_service: adapter.utility_llm_service(),
606        mcp_invoker,
607        egress_service: adapter.egress_service(),
608        sqldb_store: adapter.sqldb_store(),
609        message_retriever: Some(adapter.message_store()),
610        session_store: Some(adapter.session_store(org_id)),
611        session_mutator: Some(adapter.session_mutator(org_id)),
612        agent_store: Some(adapter.agent_store(org_id)),
613        connection_resolver: adapter.connection_resolver(),
614        schedule_store: adapter.schedule_store(org_id),
615        platform_store: adapter.platform_store(org_id, session_id),
616        knowledge_store: adapter.knowledge_store(),
617        knowledge_index_search: adapter.knowledge_index_search(org_id),
618        leased_resource_store: adapter.leased_resource_store(),
619        session_resource_registry: adapter.session_resource_registry(),
620        session_task_registry: adapter.session_task_registry(),
621        event_emitter: Some(adapter.event_emitter()),
622        capability_registry: Some(adapter.capability_registry()),
623        tool_registry,
624        org_id: Some(
625            org_public_id_from_internal(org_id)
626                .parse()
627                .expect("internal org id converts to valid public org id"),
628        ),
629        network_access: None,
630        budget_checker: adapter.budget_checker(org_id, agent_id),
631        payment_authority: adapter.payment_authority(org_id, agent_id),
632        session_creation_authority: adapter.session_creation_authority(org_id, session_id),
633        subagent_spawn_store: adapter.subagent_spawn_store(),
634        subagent_nesting_policy,
635        reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
636    }
637}
638
639/// Shared lifecycle helper for runtime-backed hosts.
640pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
641    adapter: A,
642    org_id: i64,
643    session_id: SessionId,
644}
645
646impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
647    pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
648        Self {
649            adapter,
650            org_id,
651            session_id,
652        }
653    }
654
655    async fn set_session_status(&self, status: SessionStatus, action: &'static str) {
656        if let Err(error) = self
657            .adapter
658            .set_session_status(self.org_id, self.session_id, status)
659            .await
660        {
661            warn!(
662                session_id = %self.session_id,
663                org_id = self.org_id,
664                action,
665                %error,
666                "runtime host lifecycle status update failed"
667            );
668        }
669    }
670
671    async fn emit_event(&self, request: EventRequest) {
672        let event_type = request.event_type.clone();
673        if let Err(error) = self.adapter.event_emitter().emit(request).await {
674            warn!(
675                session_id = %self.session_id,
676                org_id = self.org_id,
677                event_type,
678                %error,
679                "runtime host lifecycle event emission failed"
680            );
681        }
682    }
683
684    pub async fn turn_started(&self, turn_id: TurnId, input_message_id: MessageId) {
685        let input_content = self
686            .adapter
687            .message_store()
688            .get(self.session_id, input_message_id)
689            .await
690            .ok()
691            .flatten()
692            .map(|message| message.content_to_llm_string());
693
694        self.set_session_status(SessionStatus::Active, "turn_started")
695            .await;
696
697        self.emit_event(EventRequest::new(
698            self.session_id,
699            EventContext::turn(turn_id, input_message_id),
700            SessionActivatedData {
701                turn_id,
702                input_message_id,
703            },
704        ))
705        .await;
706
707        self.emit_event(EventRequest::new(
708            self.session_id,
709            EventContext::turn(turn_id, input_message_id),
710            TurnStartedData {
711                turn_id,
712                input_message_id,
713                input_content,
714            },
715        ))
716        .await;
717    }
718
719    pub async fn emit_turn_completed(&self, input_message_id: MessageId, data: TurnCompletedData) {
720        let turn_id = data.turn_id;
721        self.emit_event(EventRequest::new(
722            self.session_id,
723            EventContext::turn(turn_id, input_message_id),
724            data,
725        ))
726        .await;
727    }
728
729    pub async fn emit_session_idled(
730        &self,
731        turn_id: TurnId,
732        input_message_id: MessageId,
733        iterations: Option<u32>,
734        usage: Option<TokenUsage>,
735    ) {
736        self.set_session_status(SessionStatus::Idle, "emit_session_idled")
737            .await;
738
739        self.emit_event(EventRequest::new(
740            self.session_id,
741            EventContext::turn(turn_id, input_message_id),
742            SessionIdledData {
743                turn_id,
744                iterations,
745                usage,
746            },
747        ))
748        .await;
749    }
750
751    pub async fn turn_completed(
752        &self,
753        turn_id: TurnId,
754        input_message_id: MessageId,
755        iterations: u32,
756        usage: Option<TokenUsage>,
757        input_content: Option<String>,
758    ) {
759        self.emit_turn_completed(
760            input_message_id,
761            TurnCompletedData {
762                turn_id,
763                iterations,
764                duration_ms: None,
765                usage: usage.clone(),
766                input_content,
767                final_message_id: None,
768                final_answer_preview: None,
769                time_to_first_token_ms: None,
770                tool_call_count: None,
771                llm_call_count: None,
772                status: Some("completed".to_string()),
773            },
774        )
775        .await;
776        self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
777            .await;
778    }
779
780    /// Turn was deliberately sealed (EVE-534): emit `turn.sealed` + a
781    /// user-facing message + `session.idled`, and idle the session.
782    ///
783    /// Distinct from `turn_completed` (success) and `turn_failed` (error). The
784    /// session returns to `idle` so the UI unblocks; the Sealed state is
785    /// observable via the `turn.sealed` event and its `reason`.
786    pub async fn turn_sealed(
787        &self,
788        turn_id: TurnId,
789        input_message_id: MessageId,
790        reason: &str,
791        iterations: u32,
792        usage: Option<TokenUsage>,
793    ) {
794        let context = EventContext::turn(turn_id, input_message_id);
795
796        self.emit_event(EventRequest::new(
797            self.session_id,
798            context.clone(),
799            everruns_core::events::TurnSealedData {
800                turn_id,
801                reason: reason.to_string(),
802                detail: None,
803                iterations: Some(iterations),
804                usage: usage.clone(),
805            },
806        ))
807        .await;
808
809        self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
810            .await;
811    }
812
813    /// Fire `turn_end` lifecycle hooks (advisory). Collects the session's hook
814    /// specs and runs every `turn_end` hook; failures are logged, never fatal.
815    /// `harness_id`/`agent_id` are required to resolve the capability chain.
816    pub async fn fire_turn_end_hooks(
817        &self,
818        harness_id: HarnessId,
819        agent_id: Option<AgentId>,
820        turn_id: TurnId,
821        success: bool,
822    ) {
823        let (specs, dispatcher) = match collect_lifecycle_hook_specs(
824            &self.adapter,
825            self.org_id,
826            self.session_id,
827            harness_id,
828            agent_id,
829        )
830        .await
831        {
832            Ok(pair) => pair,
833            Err(error) => {
834                warn!(
835                    session_id = %self.session_id,
836                    %error,
837                    "failed to collect turn_end hook specs; skipping"
838                );
839                return;
840            }
841        };
842        let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
843            &specs,
844            everruns_core::user_hook_types::HookEvent::TurnEnd,
845            dispatcher,
846        );
847        if hooks.is_empty() {
848            return;
849        }
850        let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
851            session_id: self.session_id,
852            turn_id: Some(turn_id),
853            org_id: org_public_id_from_internal(self.org_id).parse().ok(),
854            agent_id: agent_id.map(|a| a.to_string()),
855        };
856        everruns_core::lifecycle_hooks::run_turn_end_hooks(
857            &hooks,
858            &ctx,
859            serde_json::json!({ "success": success }),
860        )
861        .await;
862    }
863
864    /// Abort a turn because a `user_prompt_submit` hook returned `Block`.
865    /// Reuses the dependency-blocked failure shape: emit a user-facing message
866    /// carrying the hook's `user_message` (or `reason`), then mark the turn
867    /// failed and idle the session.
868    pub async fn user_prompt_blocked(
869        &self,
870        turn_id: TurnId,
871        input_message_id: MessageId,
872        reason: &str,
873        user_message: Option<&str>,
874    ) {
875        let user_error =
876            UserFacingError::new(everruns_core::user_facing_error_codes::BLOCKED_BY_HOOK);
877        let shown = user_message.unwrap_or(reason);
878        let mut error_message = Message::assistant(shown);
879        let mut metadata = std::collections::HashMap::new();
880        user_error.apply_to_message_metadata(&mut metadata);
881        error_message.metadata = Some(metadata);
882
883        self.emit_event(EventRequest::new(
884            self.session_id,
885            EventContext::turn(turn_id, input_message_id),
886            OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
887        ))
888        .await;
889
890        self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
891            .await;
892    }
893
894    pub async fn turn_failed(
895        &self,
896        turn_id: TurnId,
897        input_message_id: MessageId,
898        error: &str,
899        user_error: Option<&UserFacingError>,
900    ) {
901        self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
902            .await;
903    }
904
905    /// `turn_failed` with the applied error-disclosure mode recorded on the
906    /// event. `user_error` (and the `error` text shown alongside it) must
907    /// already be disclosure-filtered by the caller.
908    pub async fn turn_failed_with_disclosure(
909        &self,
910        turn_id: TurnId,
911        input_message_id: MessageId,
912        error: &str,
913        user_error: Option<&UserFacingError>,
914        disclosure: Option<ErrorDisclosure>,
915    ) {
916        self.set_session_status(SessionStatus::Idle, "turn_failed")
917            .await;
918
919        self.emit_event(EventRequest::new(
920            self.session_id,
921            EventContext::turn(turn_id, input_message_id),
922            {
923                let mut data = TurnFailedData {
924                    turn_id,
925                    error: error.to_string(),
926                    error_code: None,
927                    error_fields: None,
928                    error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
929                };
930                if let Some(user_error) = user_error {
931                    user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
932                }
933                data
934            },
935        ))
936        .await;
937
938        self.emit_event(EventRequest::new(
939            self.session_id,
940            EventContext::turn(turn_id, input_message_id),
941            SessionIdledData {
942                turn_id,
943                iterations: None,
944                usage: None,
945            },
946        ))
947        .await;
948    }
949
950    pub async fn waiting_for_tool_results(&self) {
951        self.set_session_status(
952            SessionStatus::WaitingForToolResults,
953            "waiting_for_tool_results",
954        )
955        .await;
956    }
957
958    pub async fn dependency_blocked(
959        &self,
960        turn_id: TurnId,
961        input_message_id: MessageId,
962        blocker: DependencyBlocker,
963    ) {
964        let user_error = UserFacingError::new(blocker.error_code())
965            .with_field(
966                "dependency",
967                match blocker {
968                    DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
969                        "harness"
970                    }
971                    DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
972                },
973            )
974            .with_field(
975                "state",
976                match blocker {
977                    DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
978                        "archived"
979                    }
980                    DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
981                        "deleted"
982                    }
983                },
984            );
985        let mut error_message = Message::assistant(blocker.message());
986        let mut metadata = std::collections::HashMap::new();
987        user_error.apply_to_message_metadata(&mut metadata);
988        error_message.metadata = Some(metadata);
989
990        self.emit_event(EventRequest::new(
991            self.session_id,
992            EventContext::turn(turn_id, input_message_id),
993            OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
994        ))
995        .await;
996
997        self.turn_failed(
998            turn_id,
999            input_message_id,
1000            blocker.message(),
1001            Some(&user_error),
1002        )
1003        .await;
1004    }
1005}
1006
1007pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1008    adapter: &A,
1009    org_id: i64,
1010    harness_id: HarnessId,
1011    agent_id: Option<AgentId>,
1012) -> everruns_core::error::Result<Option<DependencyBlocker>> {
1013    let harness_store = adapter.harness_store(org_id);
1014    let agent_store = adapter.agent_store(org_id);
1015    everruns_core::detect_dependency_blocker(
1016        harness_store.as_ref(),
1017        agent_store.as_ref(),
1018        harness_id,
1019        agent_id,
1020    )
1021    .await
1022}
1023
1024pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1025    adapter: &A,
1026    org_id: i64,
1027    input: InputAtomInput,
1028) -> everruns_core::error::Result<InputAtomResult> {
1029    // The live effort override is turn-scoped. Clear any value left by the
1030    // previous turn before ReasonAtom can prefer it over this turn's message
1031    // controls.
1032    if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1033        handle.set(None);
1034    }
1035
1036    RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1037        .turn_started(input.context.turn_id, input.context.input_message_id)
1038        .await;
1039
1040    let atom = InputAtom::new(adapter.message_store());
1041    atom.execute(input).await
1042}
1043
1044/// Collect `user_prompt_submit` hooks for this turn and run them against the
1045/// inbound user message text. Returns `None` when the session has no such
1046/// hooks (the common case — no overhead beyond the spec collection, which is
1047/// skipped early). Errors loading specs are logged and treated as "no hooks"
1048/// so a hook-collection failure never blocks a turn that wasn't asking to be
1049/// hooked.
1050struct UserPromptHookResult {
1051    decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1052    original_message: String,
1053}
1054
1055async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1056    adapter: &A,
1057    org_id: i64,
1058    input: &ReasonInput,
1059) -> everruns_core::error::Result<Option<UserPromptHookResult>> {
1060    let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1061        adapter,
1062        org_id,
1063        input.context.session_id,
1064        input.harness_id,
1065        input.agent_id,
1066    )
1067    .await
1068    {
1069        Ok(pair) => pair,
1070        Err(error) => {
1071            warn!(
1072                session_id = %input.context.session_id,
1073                %error,
1074                "failed to collect user_prompt_submit hook specs; continuing without them"
1075            );
1076            return Ok(None);
1077        }
1078    };
1079    let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1080        &specs,
1081        everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1082        dispatcher,
1083    );
1084    if hooks.is_empty() {
1085        return Ok(None);
1086    }
1087
1088    let message_text = adapter
1089        .message_store()
1090        .get(input.context.session_id, input.context.input_message_id)
1091        .await
1092        .ok()
1093        .flatten()
1094        .map(|m| m.content_to_llm_string())
1095        .unwrap_or_default();
1096
1097    let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1098        session_id: input.context.session_id,
1099        turn_id: Some(input.context.turn_id),
1100        org_id: org_public_id_from_internal(org_id).parse().ok(),
1101        agent_id: input.agent_id.map(|a| a.to_string()),
1102    };
1103    let original_message = message_text.clone();
1104    let decision =
1105        everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1106            .await;
1107    Ok(Some(UserPromptHookResult {
1108        decision,
1109        original_message,
1110    }))
1111}
1112
1113pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1114    adapter: &A,
1115    org_id: i64,
1116    input: ReasonInput,
1117) -> everruns_core::error::Result<ReasonResult> {
1118    if let Some(blocker) =
1119        detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1120    {
1121        RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1122            .dependency_blocked(
1123                input.context.turn_id,
1124                input.context.input_message_id,
1125                blocker,
1126            )
1127            .await;
1128        return Ok(ReasonResult {
1129            success: false,
1130            text: blocker.message().to_string(),
1131            tool_calls: vec![],
1132            has_tool_calls: false,
1133            tool_definitions: vec![],
1134            max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1135            error: Some("dependency_unavailable".to_string()),
1136            user_facing_error: None,
1137            error_disclosure: None,
1138            usage: None,
1139            output_message_id: None,
1140            time_to_first_token_ms: None,
1141            response_id: None,
1142            finish_reason: None,
1143            locale: None,
1144            network_access: None,
1145            parallel_tool_calls: None,
1146        });
1147    }
1148
1149    // user_prompt_submit hook (see `specs/user-hooks.md`). Fires once per turn,
1150    // on the first reason iteration, before the LLM is consulted — the closest
1151    // choke point to "inbound user message accepted, before reason" that both
1152    // the in-process loop and the durable worker share. A `Block` aborts the
1153    // turn by reusing the same failure path as `dependency_blocked`: emit a
1154    // user-facing message + turn.failed, idle the session, and return a
1155    // non-success `ReasonResult` so no LLM/act work runs.
1156    let mut user_prompt_message_override = None;
1157    if input.iteration <= 1
1158        && let Some(hook_result) = run_user_prompt_submit_for_turn(adapter, org_id, &input).await?
1159    {
1160        match hook_result.decision {
1161            everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1162                reason,
1163                user_message,
1164            } => {
1165                RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1166                    .user_prompt_blocked(
1167                        input.context.turn_id,
1168                        input.context.input_message_id,
1169                        &reason,
1170                        user_message.as_deref(),
1171                    )
1172                    .await;
1173                return Ok(ReasonResult {
1174                    success: false,
1175                    text: user_message.unwrap_or_else(|| reason.clone()),
1176                    tool_calls: vec![],
1177                    has_tool_calls: false,
1178                    tool_definitions: vec![],
1179                    max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1180                    error: Some("blocked_by_user_prompt_hook".to_string()),
1181                    user_facing_error: None,
1182                    error_disclosure: None,
1183                    usage: None,
1184                    output_message_id: None,
1185                    time_to_first_token_ms: None,
1186                    response_id: None,
1187                    finish_reason: None,
1188                    locale: None,
1189                    network_access: None,
1190                    parallel_tool_calls: None,
1191                });
1192            }
1193            everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1194                if message != hook_result.original_message {
1195                    user_prompt_message_override = Some(message);
1196                }
1197            }
1198        }
1199    }
1200
1201    // Validate the executor-side registry before ReasonAtom exposes its tool
1202    // definitions to the model. This catches host wiring errors as
1203    // configuration failures instead of late tool-call failures.
1204    let validation_session = adapter
1205        .session_store(org_id)
1206        .get_session(input.context.session_id)
1207        .await?
1208        .ok_or_else(|| {
1209            everruns_core::error::AgentLoopError::session_not_found(input.context.session_id)
1210        })?;
1211    let validation_capabilities = load_execution_capabilities(
1212        adapter,
1213        org_id,
1214        input.context.session_id,
1215        input.harness_id,
1216        input.agent_id,
1217        validation_session.locale.clone(),
1218        validation_session.blueprint_id.as_deref(),
1219    )
1220    .await?;
1221    let query_history_allowed = validation_capabilities
1222        .tool_registry
1223        .get("query_history")
1224        .is_some();
1225    let validation_services = runtime_tool_context_services(
1226        adapter,
1227        org_id,
1228        input.context.session_id,
1229        input.agent_id,
1230        Some(Arc::new(validation_capabilities.tool_registry.clone())),
1231        None,
1232        validation_capabilities.subagent_nesting_policy,
1233    );
1234    validation_capabilities
1235        .tool_registry
1236        .validate_context_services(&validation_services)?;
1237
1238    let mut turn_context = adapter
1239        .load_turn_context(org_id, input.context.session_id)
1240        .await?;
1241    if let Some(registry) = adapter.session_task_registry() {
1242        let session_store = adapter.session_store(org_id);
1243        if let Some(tool) = report_result_tool_for_child_session(
1244            input.context.session_id,
1245            session_store.as_ref(),
1246            registry.as_ref(),
1247        )
1248        .await?
1249        {
1250            turn_context.mcp_tool_definitions.push(tool.to_definition());
1251        }
1252        if let Some(tool) = report_task_progress_tool_for_child_session(
1253            input.context.session_id,
1254            session_store.as_ref(),
1255            registry.as_ref(),
1256        )
1257        .await?
1258        {
1259            turn_context.mcp_tool_definitions.push(tool.to_definition());
1260        }
1261    }
1262
1263    let mut reason_capability_registry = adapter.capability_registry();
1264    if !query_history_allowed {
1265        // Keep Infinity Context's message filter active: persisted history may
1266        // contain raw text removed by earlier prompt hooks. Replace only its
1267        // model-visible prompt/tool contributions.
1268        reason_capability_registry
1269            .register(everruns_core::capabilities::InfinityContextFilterOnlyCapability);
1270    }
1271    let mut atom = ReasonAtom::new(
1272        adapter.harness_store(org_id),
1273        adapter.agent_store(org_id),
1274        adapter.session_store(org_id),
1275        adapter.message_store(),
1276        adapter.provider_store(org_id),
1277        reason_capability_registry.clone(),
1278        adapter.driver_registry(),
1279        adapter.event_emitter(),
1280    )
1281    .with_file_store(adapter.file_store());
1282    if let Some(image_resolver) = adapter.image_resolver(org_id) {
1283        atom = atom.with_image_resolver(image_resolver);
1284    }
1285    if let Some(hb) = adapter.stream_heartbeater() {
1286        atom = atom.with_stream_heartbeater(hb);
1287    }
1288    if let Some(timeout) = adapter.provider_stall_timeout() {
1289        atom = atom.with_provider_stall_timeout(timeout);
1290    }
1291    if let Some(store) = adapter.partial_stream_store() {
1292        atom = atom.with_partial_stream_store(store);
1293    }
1294    if let Some(store) = adapter.durable_tool_result_store() {
1295        atom = atom.with_durable_tool_result_store(store);
1296    }
1297    if let Some(store) = adapter.compaction_checkpoint_store() {
1298        atom = atom.with_compaction_checkpoint_store(store);
1299    }
1300    if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1301        atom = atom.with_reasoning_effort_handle(handle);
1302    }
1303    if let Some(utility_llm_service) = adapter.utility_llm_service() {
1304        atom = atom.with_utility_llm_service(utility_llm_service);
1305    }
1306    // Schedule store powers the `usage_limit_auto_continue` capability, which
1307    // schedules a continuation after a provider usage limit resets.
1308    if let Some(schedule_store) = adapter.schedule_store(org_id) {
1309        atom = atom.with_schedule_store(schedule_store);
1310    }
1311
1312    let input = ReasonInput {
1313        mcp_tool_definitions: turn_context.mcp_tool_definitions,
1314        ..input
1315    };
1316
1317    if let Some(message_override) = user_prompt_message_override {
1318        let mut assembled = assemble_turn_context(
1319            adapter.harness_store(org_id).as_ref(),
1320            adapter.agent_store(org_id).as_ref(),
1321            adapter.session_store(org_id).as_ref(),
1322            adapter.message_store().as_ref(),
1323            adapter.provider_store(org_id).as_ref(),
1324            &reason_capability_registry,
1325            input.context.session_id,
1326            input.harness_id,
1327            input.agent_id,
1328            &input.mcp_tool_definitions,
1329            Some(adapter.file_store()),
1330        )
1331        .await?;
1332
1333        let message = assembled
1334            .messages
1335            .iter_mut()
1336            .find(|message| message.id == input.context.input_message_id)
1337            .ok_or_else(|| {
1338                everruns_core::error::AgentLoopError::config(
1339                    "user_prompt_submit mutation: input message not found in assembled context",
1340                )
1341            })?;
1342
1343        // user_prompt_submit mutations are enforcement controls for the
1344        // provider-bound prompt. Apply them to the assembled context only
1345        // so persisted user history remains an audit record of the input.
1346        // Preserve non-text parts (images, files); replace only text parts.
1347        message
1348            .content
1349            .retain(|part| !matches!(part, ContentPart::Text(_)));
1350        message
1351            .content
1352            .insert(0, ContentPart::text(message_override));
1353
1354        return atom.execute_with_assembled_context(input, assembled).await;
1355    }
1356
1357    atom.execute(input).await
1358}
1359
1360pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1361    adapter: &A,
1362    input: ActInput,
1363) -> everruns_core::error::Result<ActResult> {
1364    let org_id = input.org_id.ok_or_else(|| {
1365        everruns_core::error::AgentLoopError::config(
1366            "ActInput.org_id must be set for runtime host execution",
1367        )
1368    })?;
1369
1370    if let Some(blocker) =
1371        detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1372    {
1373        RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1374            .dependency_blocked(
1375                input.context.turn_id,
1376                input.context.input_message_id,
1377                blocker,
1378            )
1379            .await;
1380        return Ok(ActResult {
1381            results: vec![],
1382            completed: true,
1383            success_count: 0,
1384            error_count: 1,
1385            waiting_for_tool_results: false,
1386            blocked: true,
1387            client_tool_calls: vec![],
1388            client_tool_definitions: vec![],
1389        });
1390    }
1391
1392    let execution_capabilities = load_execution_capabilities(
1393        adapter,
1394        org_id,
1395        input.context.session_id,
1396        input.harness_id,
1397        input.agent_id,
1398        input.locale.clone(),
1399        input.blueprint_id.as_deref(),
1400    )
1401    .await?;
1402    let mut tool_registry = execution_capabilities.tool_registry;
1403
1404    if input
1405        .tool_definitions
1406        .iter()
1407        .any(|definition| definition.name() == "report_result")
1408        && let Some(registry) = adapter.session_task_registry()
1409        && let Some(tool) = report_result_tool_for_child_session(
1410            input.context.session_id,
1411            adapter.session_store(org_id).as_ref(),
1412            registry.as_ref(),
1413        )
1414        .await?
1415    {
1416        tool_registry.register_boxed(Box::new(tool.with_file_store(adapter.file_store())));
1417    }
1418    if input
1419        .tool_definitions
1420        .iter()
1421        .any(|definition| definition.name() == "report_task_progress")
1422        && let Some(registry) = adapter.session_task_registry()
1423        && let Some(tool) = report_task_progress_tool_for_child_session(
1424            input.context.session_id,
1425            adapter.session_store(org_id).as_ref(),
1426            registry.as_ref(),
1427        )
1428        .await?
1429    {
1430        tool_registry.register_boxed(Box::new(tool));
1431    }
1432
1433    // Register the session's MCP tools as first-class registry tools, so they
1434    // execute through the regular `ToolExecutor` path and are visible to
1435    // everything that introspects the registry (spawn_background, tool_search,
1436    // openai_tool_search namespaces, ...). The turn's tool definitions already
1437    // include the discovered MCP tools, so no re-discovery is needed; the host's
1438    // MCP executor supplies execution (specs/runtime-mcp.md D5).
1439    // The MCP invoker is reused below for the guardrails `mcp` check, which
1440    // delegates a guardrail decision to an external endpoint over the same
1441    // scoped-MCP client/auth (specs/guardrails.md).
1442    let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1443    if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1444        let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1445        for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1446            tool_registry.register_boxed(tool);
1447        }
1448        mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1449            &input.tool_definitions,
1450            invoker,
1451        )));
1452    }
1453
1454    let builtin_tool_registry = Arc::new(tool_registry.clone());
1455    let context_services = runtime_tool_context_services(
1456        adapter,
1457        org_id,
1458        input.context.session_id,
1459        input.agent_id,
1460        Some(builtin_tool_registry),
1461        mcp_invoker,
1462        execution_capabilities.subagent_nesting_policy,
1463    );
1464    tool_registry.validate_context_services(&context_services)?;
1465    let executor: Arc<dyn everruns_core::traits::ToolExecutor> = Arc::new(tool_registry);
1466
1467    let mut atom = ActAtom::new(executor, adapter.event_emitter())
1468        .with_context_services(context_services)
1469        .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1470        .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1471        .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1472
1473    if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1474        atom = atom.with_outbound_tool_rate_limiter(limiter);
1475    }
1476    if let Some(store) = adapter.durable_tool_result_store() {
1477        atom = atom.with_durable_tool_result_store(store);
1478    }
1479
1480    atom.execute(input).await
1481}