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