Skip to main content

everruns_host/
host.rs

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