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    };
584    let collected = collect_capabilities_with_configs(
585        &resolved.resolved_capability_configs,
586        &capability_registry,
587        &prompt_ctx,
588    )
589    .await;
590
591    let mut registry = ToolRegistry::with_defaults();
592    #[cfg(feature = "builtins")]
593    everruns_builtins::register_default_tools(&mut registry);
594    for tool in collected.tools {
595        registry.register_boxed(tool);
596    }
597
598    // Only `Available` capabilities contribute hooks, matching
599    // `collect_capabilities_with_configs` (which skips non-available
600    // capabilities). This keeps a `ComingSoon`/unavailable capability from
601    // affecting execution via any of its hook seams.
602    let mut post_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PostToolExecHook>> = resolved
603        .resolved_capability_configs
604        .iter()
605        .flat_map(|config| {
606            capability_registry
607                .get(config.capability_id())
608                .filter(|capability| capability.status().is_active())
609                .map(|capability| {
610                    capability.post_tool_exec_hooks_with_config(config.config_value())
611                })
612                .unwrap_or_default()
613        })
614        .collect();
615    // Tool-output guardrails must inspect the original result before other
616    // capability hooks can persist or compact it into secondary surfaces.
617    post_tool_hooks.sort_by_key(|hook| hook.priority());
618
619    // User-hook contributions (see `knowledge/runtime-resources/user-hooks.md`). `finalize_specs_from_configs`
620    // gathers specs across every resolved capability — both the user-facing
621    // `user_hooks` capability and any capability that bundles hooks — and applies
622    // `finalize_hook_specs` (namespace stamping, stable ids, `disabled_contributions`
623    // muting; TM-HOOK-004). The same helper backs the lifecycle firing points so
624    // every event finalizes specs identically.
625    let tool_augmentor = adapter.tool_augmentor();
626    let user_hook_specs = finalize_specs_from_configs(
627        &resolved.resolved_capability_configs,
628        &capability_registry,
629        tool_augmentor.as_deref(),
630    );
631    // Persisted messages remain the immutable audit record, so they can contain
632    // text removed by a provider-bound user_prompt_submit hook. Until there is a
633    // durable provider-visible history view, fail closed rather than let
634    // query_history bypass that enforcement boundary.
635    if user_hook_specs
636        .iter()
637        .any(|spec| spec.event == everruns_core::user_hook_types::HookEvent::UserPromptSubmit)
638    {
639        registry.unregister("query_history");
640    }
641    // Capability-contributed pre-tool hooks run first (e.g. approval gating),
642    // then user-hook (`PreToolUse`) specs. The first hook to block wins.
643    let mut pre_tool_hooks: Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>> = resolved
644        .resolved_capability_configs
645        .iter()
646        .flat_map(|config| {
647            capability_registry
648                .get(config.capability_id())
649                .filter(|capability| capability.status().is_active())
650                .map(|capability| capability.pre_tool_use_hooks_with_config(config.config_value()))
651                .unwrap_or_default()
652        })
653        .collect();
654    if !user_hook_specs.is_empty() {
655        let dispatcher = bash_hook_dispatcher(adapter.file_store());
656        post_tool_hooks.extend(everruns_core::hook_adapter::build_post_tool_use_hooks(
657            &user_hook_specs,
658            dispatcher.clone(),
659        ));
660        pre_tool_hooks.extend(everruns_core::hook_adapter::build_pre_tool_use_hooks(
661            &user_hook_specs,
662            dispatcher,
663        ));
664    }
665
666    // Use the hook list assembled by `collect_capabilities_with_configs` as the
667    // single source of truth. It already contains every explicit capability
668    // `tool_call_hooks()` followed by the generated `CapabilityNarrationHook`
669    // adapters — one per collected capability plus any auto-activated
670    // cross-cutting capability such as `background_execution`. Re-deriving only
671    // the explicit subset here dropped capability-owned narration, so tools fell
672    // back to generic `Ran {display_name}` lines (EVE-601). Explicit hooks stay
673    // first in this list, so model-authored narration (`human_intent`) keeps its
674    // precedence over default `Tool::narrate()`, and only available capabilities
675    // contributed because collection skips non-available ones.
676    let tool_call_hooks = collected.tool_call_hooks;
677
678    Ok(RuntimeExecutionCapabilities {
679        tool_registry: registry,
680        post_tool_hooks,
681        pre_tool_hooks,
682        tool_call_hooks,
683        subagent_nesting_policy: subagent_nesting_policy_from_configs(
684            &resolved.resolved_capability_configs,
685        ),
686    })
687}
688
689fn runtime_tool_context_services<A: RuntimeHostAdapter>(
690    adapter: &A,
691    org_id: i64,
692    session_id: SessionId,
693    agent_id: Option<AgentId>,
694    tool_registry: Option<Arc<ToolRegistry>>,
695    mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>>,
696    subagent_nesting_policy: everruns_core::delegation_services::SubagentNestingPolicy,
697) -> ToolContextServices {
698    let extensions = {
699        let mut extensions = adapter.tool_context_extensions(org_id, session_id);
700        extensions.insert(Arc::new(SessionMutatorExt(adapter.session_mutator(org_id))));
701        extensions
702    };
703    ToolContextServices {
704        file_store: Some(adapter.file_store()),
705        storage_store: adapter.storage_store(),
706        image_store: adapter.image_artifact_store(org_id),
707        provider_credential_store: adapter.provider_credential_store(org_id),
708        utility_llm_service: adapter.utility_llm_service(),
709        mcp_invoker,
710        egress_service: adapter.egress_service(),
711        message_retriever: Some(adapter.message_store()),
712        session_store: Some(adapter.session_store(org_id)),
713        agent_store: Some(adapter.agent_store(org_id)),
714        connection_resolver: adapter.connection_resolver(),
715        schedule_store: adapter.schedule_store(org_id),
716        subagent_delegate: adapter.subagent_delegate(org_id, session_id),
717        extensions,
718        leased_resource_store: adapter.leased_resource_store(),
719        session_resource_registry: adapter.session_resource_registry(),
720        session_task_registry: adapter.session_task_registry(),
721        event_emitter: Some(adapter.event_emitter()),
722        capability_registry: Some(adapter.capability_registry()),
723        tool_registry,
724        org_id: Some(
725            org_public_id_from_internal(org_id)
726                .parse()
727                .expect("internal org id converts to valid public org id"),
728        ),
729        network_access: None,
730        budget_checker: adapter.budget_checker(org_id, agent_id),
731        payment_authority: adapter.payment_authority(org_id, agent_id),
732        session_creation_authority: adapter.session_creation_authority(org_id, session_id),
733        subagent_spawn_store: adapter.subagent_spawn_store(),
734        subagent_nesting_policy,
735        reasoning_effort_handle: adapter.reasoning_effort_handle(session_id),
736    }
737}
738
739/// Shared lifecycle helper for runtime-backed hosts.
740/// Agent identity snapshot carried on `turn.started`.
741#[derive(Debug, Default)]
742struct TurnAgentIdentity {
743    id: Option<AgentId>,
744    name: Option<String>,
745    description: Option<String>,
746}
747
748pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
749    adapter: A,
750    org_id: i64,
751    session_id: SessionId,
752}
753
754impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
755    pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
756        Self {
757            adapter,
758            org_id,
759            session_id,
760        }
761    }
762
763    async fn set_session_status(
764        &self,
765        status: SessionExecutionState,
766        _action: &'static str,
767    ) -> everruns_provider::error::Result<()> {
768        self.adapter
769            .set_session_status(self.org_id, self.session_id, status)
770            .await
771    }
772
773    async fn emit_event(&self, request: EventRequest) -> everruns_provider::error::Result<()> {
774        self.adapter.event_emitter().emit(request).await.map(|_| ())
775    }
776
777    pub async fn turn_started(
778        &self,
779        turn_id: TurnId,
780        input_message_id: MessageId,
781    ) -> everruns_provider::error::Result<()> {
782        let input_content = self
783            .adapter
784            .message_store()
785            .get(self.session_id, input_message_id)
786            .await
787            .ok()
788            .flatten()
789            .map(|message| message.content_to_llm_string());
790
791        self.set_session_status(SessionExecutionState::Active, "turn_started")
792            .await?;
793
794        self.emit_event(EventRequest::new(
795            self.session_id,
796            EventContext::turn(turn_id, input_message_id),
797            SessionActivatedData {
798                turn_id,
799                input_message_id,
800            },
801        ))
802        .await?;
803
804        let agent = self.agent_identity().await;
805        self.emit_event(EventRequest::new(
806            self.session_id,
807            EventContext::turn(turn_id, input_message_id),
808            TurnStartedData {
809                turn_id,
810                input_message_id,
811                input_content,
812                agent_id: agent.id,
813                agent_name: agent.name,
814                agent_description: agent.description,
815            },
816        ))
817        .await?;
818        Ok(())
819    }
820
821    /// Agent identity for the turn root event, resolved best-effort: a store
822    /// miss or failure yields `None` fields rather than failing the turn, since
823    /// the identity only labels traces.
824    async fn agent_identity(&self) -> TurnAgentIdentity {
825        let agent_id = self
826            .adapter
827            .session_store(self.org_id)
828            .get_session(self.session_id)
829            .await
830            .ok()
831            .flatten()
832            .and_then(|session| session.agent_id);
833        let Some(agent_id) = agent_id else {
834            return TurnAgentIdentity::default();
835        };
836        let agent = self
837            .adapter
838            .agent_store(self.org_id)
839            .get_agent(agent_id)
840            .await
841            .ok()
842            .flatten();
843        TurnAgentIdentity {
844            id: Some(agent_id),
845            // The conventions want the human-readable name; the slug is the
846            // fallback when no display name was set.
847            name: agent
848                .as_ref()
849                .map(|a| a.display_name.clone().unwrap_or_else(|| a.name.clone())),
850            description: agent.and_then(|a| a.description),
851        }
852    }
853
854    pub async fn emit_turn_completed(
855        &self,
856        input_message_id: MessageId,
857        data: TurnCompletedData,
858    ) -> everruns_provider::error::Result<()> {
859        let turn_id = data.turn_id;
860        self.emit_event(EventRequest::new(
861            self.session_id,
862            EventContext::turn(turn_id, input_message_id),
863            data,
864        ))
865        .await
866    }
867
868    pub async fn emit_session_idled(
869        &self,
870        turn_id: TurnId,
871        input_message_id: MessageId,
872        iterations: Option<u32>,
873        usage: Option<TokenUsage>,
874    ) -> everruns_provider::error::Result<()> {
875        self.set_session_status(SessionExecutionState::Idle, "emit_session_idled")
876            .await?;
877
878        self.emit_event(EventRequest::new(
879            self.session_id,
880            EventContext::turn(turn_id, input_message_id),
881            SessionIdledData {
882                turn_id,
883                iterations,
884                usage,
885            },
886        ))
887        .await
888    }
889
890    pub async fn turn_completed(
891        &self,
892        turn_id: TurnId,
893        input_message_id: MessageId,
894        iterations: u32,
895        usage: Option<TokenUsage>,
896        input_content: Option<String>,
897    ) -> everruns_provider::error::Result<()> {
898        self.emit_turn_completed(
899            input_message_id,
900            TurnCompletedData {
901                turn_id,
902                iterations,
903                duration_ms: None,
904                usage: usage.clone(),
905                input_content,
906                final_message_id: None,
907                final_answer_preview: None,
908                time_to_first_token_ms: None,
909                tool_call_count: None,
910                llm_call_count: None,
911                status: Some("completed".to_string()),
912            },
913        )
914        .await?;
915        self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
916            .await
917    }
918
919    /// Turn was deliberately sealed (EVE-534): emit `turn.sealed` + a
920    /// user-facing message + `session.idled`, and idle the session.
921    ///
922    /// Distinct from `turn_completed` (success) and `turn_failed` (error). The
923    /// session returns to `idle` so the UI unblocks; the Sealed state is
924    /// observable via the `turn.sealed` event and its `reason`.
925    pub async fn turn_sealed(
926        &self,
927        turn_id: TurnId,
928        input_message_id: MessageId,
929        reason: &str,
930        iterations: u32,
931        usage: Option<TokenUsage>,
932    ) -> everruns_provider::error::Result<()> {
933        let context = EventContext::turn(turn_id, input_message_id);
934
935        self.emit_event(EventRequest::new(
936            self.session_id,
937            context.clone(),
938            everruns_core::events::TurnSealedData {
939                turn_id,
940                reason: reason.to_string(),
941                detail: None,
942                iterations: Some(iterations),
943                usage: usage.clone(),
944            },
945        ))
946        .await?;
947
948        self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
949            .await
950    }
951
952    /// Fire `turn_end` lifecycle hooks (advisory). Collects the session's hook
953    /// specs and runs every `turn_end` hook; failures are logged, never fatal.
954    /// `harness_id`/`agent_id` are required to resolve the capability chain.
955    pub async fn fire_turn_end_hooks(
956        &self,
957        harness_id: HarnessId,
958        agent_id: Option<AgentId>,
959        turn_id: TurnId,
960        success: bool,
961    ) {
962        let (specs, dispatcher) = match collect_lifecycle_hook_specs(
963            &self.adapter,
964            self.org_id,
965            self.session_id,
966            harness_id,
967            agent_id,
968        )
969        .await
970        {
971            Ok(pair) => pair,
972            Err(error) => {
973                warn!(
974                    session_id = %self.session_id,
975                    %error,
976                    "failed to collect turn_end hook specs; skipping"
977                );
978                return;
979            }
980        };
981        let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
982            &specs,
983            everruns_core::user_hook_types::HookEvent::TurnEnd,
984            dispatcher,
985        );
986        if hooks.is_empty() {
987            return;
988        }
989        let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
990            session_id: self.session_id,
991            turn_id: Some(turn_id),
992            org_id: org_public_id_from_internal(self.org_id).parse().ok(),
993            agent_id: agent_id.map(|a| a.to_string()),
994        };
995        everruns_core::lifecycle_hooks::run_turn_end_hooks(
996            &hooks,
997            &ctx,
998            serde_json::json!({ "success": success }),
999        )
1000        .await;
1001    }
1002
1003    /// Abort a turn because a `user_prompt_submit` hook returned `Block`.
1004    /// Reuses the dependency-blocked failure shape: emit a user-facing message
1005    /// carrying the hook's `user_message` (or `reason`), then mark the turn
1006    /// failed and idle the session.
1007    pub async fn user_prompt_blocked(
1008        &self,
1009        turn_id: TurnId,
1010        input_message_id: MessageId,
1011        reason: &str,
1012        user_message: Option<&str>,
1013    ) -> everruns_provider::error::Result<()> {
1014        let user_error =
1015            UserFacingError::new(everruns_provider::user_facing_error::codes::BLOCKED_BY_HOOK);
1016        let shown = user_message.unwrap_or(reason);
1017        let mut error_message = Message::assistant(shown);
1018        let mut metadata = std::collections::HashMap::new();
1019        user_error.apply_to_message_metadata(&mut metadata);
1020        error_message.metadata = Some(metadata);
1021
1022        self.emit_event(EventRequest::new(
1023            self.session_id,
1024            EventContext::turn(turn_id, input_message_id),
1025            OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1026        ))
1027        .await?;
1028
1029        self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
1030            .await
1031    }
1032
1033    pub async fn turn_failed(
1034        &self,
1035        turn_id: TurnId,
1036        input_message_id: MessageId,
1037        error: &str,
1038        user_error: Option<&UserFacingError>,
1039    ) -> everruns_provider::error::Result<()> {
1040        self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
1041            .await
1042    }
1043
1044    /// `turn_failed` with the applied error-disclosure mode recorded on the
1045    /// event. `user_error` (and the `error` text shown alongside it) must
1046    /// already be disclosure-filtered by the caller.
1047    pub async fn turn_failed_with_disclosure(
1048        &self,
1049        turn_id: TurnId,
1050        input_message_id: MessageId,
1051        error: &str,
1052        user_error: Option<&UserFacingError>,
1053        disclosure: Option<ErrorDisclosure>,
1054    ) -> everruns_provider::error::Result<()> {
1055        self.set_session_status(SessionExecutionState::Idle, "turn_failed")
1056            .await?;
1057
1058        self.emit_event(EventRequest::new(
1059            self.session_id,
1060            EventContext::turn(turn_id, input_message_id),
1061            {
1062                let mut data = TurnFailedData {
1063                    turn_id,
1064                    error: error.to_string(),
1065                    error_code: None,
1066                    error_fields: None,
1067                    error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
1068                };
1069                if let Some(user_error) = user_error {
1070                    user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
1071                }
1072                data
1073            },
1074        ))
1075        .await?;
1076
1077        self.emit_event(EventRequest::new(
1078            self.session_id,
1079            EventContext::turn(turn_id, input_message_id),
1080            SessionIdledData {
1081                turn_id,
1082                iterations: None,
1083                usage: None,
1084            },
1085        ))
1086        .await
1087    }
1088
1089    pub async fn waiting_for_tool_results(&self) -> everruns_provider::error::Result<()> {
1090        self.set_session_status(
1091            SessionExecutionState::WaitingForToolResults,
1092            "waiting_for_tool_results",
1093        )
1094        .await
1095    }
1096
1097    pub async fn dependency_blocked(
1098        &self,
1099        turn_id: TurnId,
1100        input_message_id: MessageId,
1101        blocker: DependencyBlocker,
1102    ) -> everruns_provider::error::Result<()> {
1103        let user_error = UserFacingError::new(blocker.error_code())
1104            .with_field(
1105                "dependency",
1106                match blocker {
1107                    DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
1108                        "harness"
1109                    }
1110                    DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
1111                },
1112            )
1113            .with_field(
1114                "state",
1115                match blocker {
1116                    DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
1117                        "archived"
1118                    }
1119                    DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
1120                        "deleted"
1121                    }
1122                },
1123            );
1124        let mut error_message = Message::assistant(blocker.message());
1125        let mut metadata = std::collections::HashMap::new();
1126        user_error.apply_to_message_metadata(&mut metadata);
1127        error_message.metadata = Some(metadata);
1128
1129        self.emit_event(EventRequest::new(
1130            self.session_id,
1131            EventContext::turn(turn_id, input_message_id),
1132            OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
1133        ))
1134        .await?;
1135
1136        self.turn_failed(
1137            turn_id,
1138            input_message_id,
1139            blocker.message(),
1140            Some(&user_error),
1141        )
1142        .await
1143    }
1144}
1145
1146pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
1147    adapter: &A,
1148    org_id: i64,
1149    harness_id: HarnessId,
1150    agent_id: Option<AgentId>,
1151) -> everruns_provider::error::Result<Option<DependencyBlocker>> {
1152    let harness_store = adapter.harness_store(org_id);
1153    let agent_store = adapter.agent_store(org_id);
1154    if let Some(blocker) = harness_store.get_harness_blocker(harness_id).await? {
1155        return Ok(Some(blocker));
1156    }
1157    if let Some(agent_id) = agent_id
1158        && let Some(blocker) = agent_store.get_agent_blocker(agent_id).await?
1159    {
1160        return Ok(Some(blocker));
1161    }
1162    Ok(None)
1163}
1164
1165pub async fn execute_input_activity<A: RuntimeHostAdapter>(
1166    adapter: &A,
1167    org_id: i64,
1168    input: InputAtomInput,
1169) -> everruns_provider::error::Result<InputAtomResult> {
1170    // The live effort override is turn-scoped. Clear any value left by the
1171    // previous turn before ReasonAtom can prefer it over this turn's message
1172    // controls.
1173    if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1174        handle.set(None);
1175    }
1176
1177    RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1178        .turn_started(input.context.turn_id, input.context.input_message_id)
1179        .await?;
1180
1181    let atom = InputAtom::new(adapter.message_store());
1182    atom.execute(input).await
1183}
1184
1185/// Collect `user_prompt_submit` hooks for this turn and run them against the
1186/// inbound user message text. Returns `None` when the session has no such
1187/// hooks (the common case — no overhead beyond the spec collection, which is
1188/// skipped early). Errors loading specs are logged and treated as "no hooks"
1189/// so a hook-collection failure never blocks a turn that wasn't asking to be
1190/// hooked.
1191pub(crate) struct UserPromptHookResult {
1192    pub(crate) decision: everruns_core::lifecycle_hooks::UserPromptDecision,
1193    pub(crate) original_message: String,
1194}
1195
1196pub(crate) async fn run_user_prompt_submit_for_message<A: RuntimeHostAdapter>(
1197    adapter: &A,
1198    org_id: i64,
1199    input: &ReasonInput,
1200    message_text: String,
1201) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1202    let (specs, dispatcher) = match collect_lifecycle_hook_specs(
1203        adapter,
1204        org_id,
1205        input.context.session_id,
1206        input.harness_id,
1207        input.agent_id,
1208    )
1209    .await
1210    {
1211        Ok(pair) => pair,
1212        Err(error) => {
1213            warn!(
1214                session_id = %input.context.session_id,
1215                %error,
1216                "failed to collect user_prompt_submit hook specs; continuing without them"
1217            );
1218            return Ok(None);
1219        }
1220    };
1221    let hooks = everruns_core::lifecycle_hooks::build_turn_lifecycle_hooks(
1222        &specs,
1223        everruns_core::user_hook_types::HookEvent::UserPromptSubmit,
1224        dispatcher,
1225    );
1226    if hooks.is_empty() {
1227        return Ok(None);
1228    }
1229
1230    let ctx = everruns_core::lifecycle_hooks::TurnHookContext {
1231        session_id: input.context.session_id,
1232        turn_id: Some(input.context.turn_id),
1233        org_id: org_public_id_from_internal(org_id).parse().ok(),
1234        agent_id: input.agent_id.map(|a| a.to_string()),
1235    };
1236    let original_message = message_text.clone();
1237    let decision =
1238        everruns_core::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text)
1239            .await;
1240    Ok(Some(UserPromptHookResult {
1241        decision,
1242        original_message,
1243    }))
1244}
1245
1246async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
1247    adapter: &A,
1248    org_id: i64,
1249    input: &ReasonInput,
1250) -> everruns_provider::error::Result<Option<UserPromptHookResult>> {
1251    let message_text = adapter
1252        .message_store()
1253        .get(input.context.session_id, input.context.input_message_id)
1254        .await
1255        .ok()
1256        .flatten()
1257        .map(|m| m.content_to_llm_string())
1258        .unwrap_or_default();
1259    run_user_prompt_submit_for_message(adapter, org_id, input, message_text).await
1260}
1261
1262pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
1263    adapter: &A,
1264    org_id: i64,
1265    input: ReasonInput,
1266) -> everruns_provider::error::Result<ReasonResult> {
1267    let prompt_message_ids = (input.iteration <= 1)
1268        .then_some(input.context.input_message_id)
1269        .into_iter()
1270        .collect();
1271    execute_reason_activity_with_prompt_messages(adapter, org_id, input, prompt_message_ids).await
1272}
1273
1274/// Execute a reason activity while applying `user_prompt_submit` hooks to the
1275/// supplied messages. Hosts that inject synthetic user messages between reason
1276/// iterations must include their ids here so they cross the same policy
1277/// boundary as the turn's original input.
1278pub async fn execute_reason_activity_with_prompt_messages<A: RuntimeHostAdapter>(
1279    adapter: &A,
1280    org_id: i64,
1281    input: ReasonInput,
1282    prompt_message_ids: Vec<MessageId>,
1283) -> everruns_provider::error::Result<ReasonResult> {
1284    if let Some(blocker) =
1285        detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1286    {
1287        RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1288            .dependency_blocked(
1289                input.context.turn_id,
1290                input.context.input_message_id,
1291                blocker,
1292            )
1293            .await?;
1294        return Ok(ReasonResult {
1295            native_counts: None,
1296            success: false,
1297            text: blocker.message().to_string(),
1298            tool_calls: vec![],
1299            has_tool_calls: false,
1300            tool_definitions: vec![],
1301            max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1302            error: Some("dependency_unavailable".to_string()),
1303            user_facing_error: None,
1304            error_disclosure: None,
1305            usage: None,
1306            output_message_id: None,
1307            time_to_first_token_ms: None,
1308            response_id: None,
1309            finish_reason: None,
1310            locale: None,
1311            network_access: None,
1312            parallel_tool_calls: None,
1313        });
1314    }
1315
1316    // A `Block` aborts the turn by reusing the same failure path as
1317    // `dependency_blocked`. Hosts pass the original input on iteration one and
1318    // any synthetic user messages injected later, ensuring every provider-bound
1319    // user message crosses this policy boundary.
1320    let mut user_prompt_message_overrides = Vec::new();
1321    for message_id in prompt_message_ids {
1322        let mut hook_input = input.clone();
1323        hook_input.context.input_message_id = message_id;
1324        let Some(hook_result) =
1325            run_user_prompt_submit_for_turn(adapter, org_id, &hook_input).await?
1326        else {
1327            continue;
1328        };
1329        match hook_result.decision {
1330            everruns_core::lifecycle_hooks::UserPromptDecision::Block {
1331                reason,
1332                user_message,
1333            } => {
1334                RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1335                    .user_prompt_blocked(
1336                        input.context.turn_id,
1337                        input.context.input_message_id,
1338                        &reason,
1339                        user_message.as_deref(),
1340                    )
1341                    .await?;
1342                return Ok(ReasonResult {
1343                    native_counts: None,
1344                    success: false,
1345                    text: user_message.unwrap_or_else(|| reason.clone()),
1346                    tool_calls: vec![],
1347                    has_tool_calls: false,
1348                    tool_definitions: vec![],
1349                    max_iterations: everruns_core::runtime_agent::default_max_iterations(),
1350                    error: Some("blocked_by_user_prompt_hook".to_string()),
1351                    user_facing_error: None,
1352                    error_disclosure: None,
1353                    usage: None,
1354                    output_message_id: None,
1355                    time_to_first_token_ms: None,
1356                    response_id: None,
1357                    finish_reason: None,
1358                    locale: None,
1359                    network_access: None,
1360                    parallel_tool_calls: None,
1361                });
1362            }
1363            everruns_core::lifecycle_hooks::UserPromptDecision::Continue { message } => {
1364                if message != hook_result.original_message {
1365                    user_prompt_message_overrides.push((message_id, message));
1366                }
1367            }
1368        }
1369    }
1370
1371    // Validate the executor-side registry before ReasonAtom exposes its tool
1372    // definitions to the model. This catches host wiring errors as
1373    // configuration failures instead of late tool-call failures.
1374    let validation_session = adapter
1375        .session_store(org_id)
1376        .get_session(input.context.session_id)
1377        .await?
1378        .ok_or_else(|| {
1379            everruns_provider::error::AgentLoopError::session_not_found(input.context.session_id)
1380        })?;
1381    let validation_capabilities = load_execution_capabilities(
1382        adapter,
1383        org_id,
1384        input.context.session_id,
1385        input.harness_id,
1386        input.agent_id,
1387        validation_session.locale.clone(),
1388        validation_session.blueprint_id.as_deref(),
1389    )
1390    .await?;
1391    let query_history_allowed = validation_capabilities
1392        .tool_registry
1393        .get("query_history")
1394        .is_some();
1395    let validation_services = runtime_tool_context_services(
1396        adapter,
1397        org_id,
1398        input.context.session_id,
1399        input.agent_id,
1400        Some(Arc::new(validation_capabilities.tool_registry.clone())),
1401        None,
1402        validation_capabilities.subagent_nesting_policy,
1403    );
1404    validation_capabilities
1405        .tool_registry
1406        .validate_context_services(&validation_services)?;
1407
1408    let mut turn_inputs = adapter
1409        .load_resolved_turn(org_id, input.context.session_id)
1410        .await?;
1411    if let Some(augmentor) = adapter.tool_augmentor() {
1412        augmentor
1413            .augment_reason_tools(
1414                input.context.session_id,
1415                adapter.session_store(org_id),
1416                adapter.session_task_registry(),
1417                &mut turn_inputs.mcp_tool_definitions,
1418            )
1419            .await?;
1420    }
1421
1422    let reason_capability_registry = {
1423        let mut registry = adapter.capability_registry();
1424        if !query_history_allowed {
1425            // Persisted history may contain raw text removed by earlier prompt
1426            // hooks. Preserve the owning capability's message filter while
1427            // suppressing its prompt/tool contributions for this turn. Tool
1428            // ownership is discovered through the neutral capability contract;
1429            // host does not depend on a concrete implementation or capability ID.
1430            let query_history_owner = registry
1431                .list()
1432                .into_iter()
1433                .find(|capability| {
1434                    capability
1435                        .tool_definitions()
1436                        .iter()
1437                        .any(|tool| tool.name() == "query_history")
1438                })
1439                .map(Arc::clone);
1440            if let Some(capability) = query_history_owner {
1441                registry.register(MessageFilterOnlyCapability(capability));
1442            }
1443        }
1444        registry
1445    };
1446    let context_resolver = crate::runtime_context::StoreTurnContextResolver::new(
1447        adapter.harness_store(org_id),
1448        adapter.agent_store(org_id),
1449        adapter.session_store(org_id),
1450        adapter.message_store(),
1451        adapter.provider_store(org_id),
1452        reason_capability_registry.clone(),
1453        adapter.driver_registry(),
1454    )
1455    .with_file_store(adapter.file_store());
1456    let mut atom = ReasonAtom::new(
1457        context_resolver,
1458        adapter.message_store(),
1459        reason_capability_registry.clone(),
1460        adapter.event_emitter(),
1461    );
1462    if let Some(image_resolver) = adapter.image_resolver(org_id) {
1463        atom = atom.with_image_resolver(image_resolver);
1464    }
1465    if let Some(file_resolver) = adapter.file_resolver(org_id) {
1466        atom = atom.with_file_resolver(file_resolver);
1467    }
1468    if let Some(hb) = adapter.stream_heartbeater() {
1469        atom = atom.with_stream_heartbeater(hb);
1470    }
1471    if let Some(timeout) = adapter.provider_stall_timeout() {
1472        atom = atom.with_provider_stall_timeout(timeout);
1473    }
1474    if let Some(config) = adapter.provider_retry_config() {
1475        atom = atom.with_provider_retry_config(config);
1476    }
1477    if let Some(store) = adapter.partial_stream_store() {
1478        atom = atom.with_partial_stream_store(store);
1479    }
1480    if let Some(store) = adapter.durable_tool_result_store() {
1481        atom = atom.with_durable_tool_result_store(store);
1482    }
1483    if let Some(store) = adapter.compaction_checkpoint_store() {
1484        atom = atom.with_compaction_checkpoint_store(store);
1485    }
1486    if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
1487        atom = atom.with_reasoning_effort_handle(handle);
1488    }
1489    if let Some(utility_llm_service) = adapter.utility_llm_service() {
1490        atom = atom.with_utility_llm_service(utility_llm_service);
1491    }
1492    // Schedule store powers the `usage_limit_auto_continue` capability, which
1493    // schedules a continuation after a provider usage limit resets.
1494    if let Some(schedule_store) = adapter.schedule_store(org_id) {
1495        atom = atom.with_schedule_store(schedule_store);
1496    }
1497
1498    let mut assembled = crate::runtime_context::assemble_turn_context_from_snapshot(
1499        turn_inputs.snapshot,
1500        adapter.message_store().as_ref(),
1501        adapter.provider_store(org_id).as_ref(),
1502        &reason_capability_registry,
1503        &adapter.driver_registry(),
1504        &turn_inputs.mcp_tool_definitions,
1505        Some(adapter.file_store()),
1506    )
1507    .await?;
1508    let input = ReasonInput {
1509        mcp_tool_definitions: turn_inputs.mcp_tool_definitions,
1510        ..input
1511    };
1512
1513    if !user_prompt_message_overrides.is_empty() {
1514        for (message_id, message_override) in user_prompt_message_overrides {
1515            let message = assembled
1516                .messages
1517                .iter_mut()
1518                .find(|message| message.id == message_id)
1519                .ok_or_else(|| {
1520                    everruns_provider::error::AgentLoopError::config(
1521                        "user_prompt_submit mutation: input message not found in assembled context",
1522                    )
1523                })?;
1524
1525            // Apply enforcement mutations to provider context only, retaining
1526            // persisted history as an audit record of the original content.
1527            message
1528                .content
1529                .retain(|part| !matches!(part, ContentPart::Text(_)));
1530            message
1531                .content
1532                .insert(0, ContentPart::text(message_override));
1533        }
1534    }
1535    // A model switch is otherwise visible only in `llm.generation`, which is
1536    // diagnostic and off the reading path. Mark it here, where every host —
1537    // the durable worker and the in-process framework runtime alike — resolves
1538    // the turn's model, rather than at one API's message-create path.
1539    if input.iteration <= 1 {
1540        emit_model_change_if_switched(adapter, org_id, &input, &assembled).await;
1541    }
1542
1543    crate::native_async::execute_reason(adapter, org_id, input, assembled, atom).await
1544}
1545
1546/// Emit `session.model.changed` when this turn's input selects a model
1547/// different from the previous turn's.
1548///
1549/// Best effort: a missing marker must not fail the turn.
1550async fn emit_model_change_if_switched<A: RuntimeHostAdapter>(
1551    adapter: &A,
1552    org_id: i64,
1553    input: &ReasonInput,
1554    assembled: &AssembledTurnContext,
1555) {
1556    let Some((previous_model_id, model_id)) = model_switch(&assembled.messages) else {
1557        return;
1558    };
1559    if assembled.resolved_model_id != Some(model_id) {
1560        // The requested model did not survive resolution (unknown or removed),
1561        // so the turn runs on a fallback this marker would misname.
1562        return;
1563    }
1564
1565    // The previous model is named by looking it up; the current one is already
1566    // resolved for this turn. Names are captured now so the transcript stays
1567    // readable after a model is removed from the org.
1568    let previous_model_name = adapter
1569        .provider_store(org_id)
1570        .get_model_spec(previous_model_id)
1571        .await
1572        .ok()
1573        .flatten()
1574        .map(|spec| spec.model);
1575
1576    let request = EventRequest::new(
1577        input.context.session_id,
1578        EventContext::turn(input.context.turn_id, input.context.input_message_id),
1579        SessionModelChangedData {
1580            previous_model_id: Some(previous_model_id),
1581            previous_model_name,
1582            model_id,
1583            model_name: assembled.model.model.clone(),
1584        },
1585    );
1586    if let Err(e) = adapter.event_emitter().emit(request).await {
1587        warn!(error = %e, "Failed to emit session.model.changed event");
1588    }
1589}
1590
1591/// `(previous, current)` model overrides when the turn's input switched models.
1592///
1593/// Only an explicit override replacing a different explicit override counts. A
1594/// turn without an override runs on an inherited default, and history visible
1595/// here is capability-filtered: treating a missing override as "the default"
1596/// would report a switch whenever an older message was filtered out.
1597fn model_switch(messages: &[Message]) -> Option<(ModelId, ModelId)> {
1598    let mut user_model_ids = messages
1599        .iter()
1600        .rev()
1601        .filter(|message| message.role == MessageRole::User)
1602        .map(|message| {
1603            message
1604                .controls
1605                .as_ref()
1606                .and_then(|controls| controls.model_id)
1607        });
1608
1609    // `latest_model_override` in `runtime_context` resolves the turn's model
1610    // from the last user message alone, so the comparison is between the last
1611    // two user messages — not the last two overrides anywhere in history.
1612    let model_id = user_model_ids.next().flatten()?;
1613    let previous_model_id = user_model_ids.next().flatten()?;
1614    (previous_model_id != model_id).then_some((previous_model_id, model_id))
1615}
1616
1617pub async fn execute_act_activity<A: RuntimeHostAdapter>(
1618    adapter: &A,
1619    input: ActInput,
1620) -> everruns_provider::error::Result<ActResult> {
1621    let org_id = input.org_id.ok_or_else(|| {
1622        everruns_provider::error::AgentLoopError::config(
1623            "ActInput.org_id must be set for runtime host execution",
1624        )
1625    })?;
1626
1627    if let Some(blocker) =
1628        detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
1629    {
1630        RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
1631            .dependency_blocked(
1632                input.context.turn_id,
1633                input.context.input_message_id,
1634                blocker,
1635            )
1636            .await?;
1637        return Ok(ActResult {
1638            results: vec![],
1639            completed: true,
1640            success_count: 0,
1641            error_count: 1,
1642            waiting_for_tool_results: false,
1643            waiting_for_url_elicitation: false,
1644            blocked: true,
1645            client_tool_calls: vec![],
1646            client_tool_definitions: vec![],
1647        });
1648    }
1649
1650    let execution_capabilities = load_execution_capabilities(
1651        adapter,
1652        org_id,
1653        input.context.session_id,
1654        input.harness_id,
1655        input.agent_id,
1656        input.locale.clone(),
1657        input.blueprint_id.as_deref(),
1658    )
1659    .await?;
1660    let mut tool_registry = execution_capabilities.tool_registry;
1661
1662    if let Some(augmentor) = adapter.tool_augmentor() {
1663        augmentor
1664            .augment_act_tools(
1665                input.context.session_id,
1666                adapter.session_store(org_id),
1667                adapter.session_task_registry(),
1668                adapter.file_store(),
1669                &input.tool_definitions,
1670                &mut tool_registry,
1671            )
1672            .await?;
1673    }
1674
1675    // Register the session's MCP tools as first-class registry tools, so they
1676    // execute through the regular `ToolExecutor` path and are visible to
1677    // everything that introspects the registry (spawn_background, tool_search,
1678    // openai_tool_search namespaces, ...). The turn's tool definitions already
1679    // include the discovered MCP tools, so no re-discovery is needed; the host's
1680    // MCP executor supplies execution (knowledge/integrations/runtime-mcp.md D5).
1681    // The MCP invoker is reused below for the guardrails `mcp` check, which
1682    // delegates a guardrail decision to an external endpoint over the same
1683    // scoped-MCP client/auth (knowledge/execution/guardrails.md).
1684    let mut mcp_invoker: Option<Arc<dyn everruns_core::McpToolInvoker>> = None;
1685    if let Some(mcp) = adapter.mcp_executor(org_id, input.context.session_id).await {
1686        let invoker: Arc<dyn everruns_core::McpToolInvoker> = mcp;
1687        for tool in everruns_core::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
1688            tool_registry.register_boxed(tool);
1689        }
1690        mcp_invoker = Some(Arc::new(everruns_core::ScopedMcpToolInvoker::new(
1691            &input.tool_definitions,
1692            invoker,
1693        )));
1694    }
1695
1696    let builtin_tool_registry = Arc::new(tool_registry.clone());
1697    let context_services = runtime_tool_context_services(
1698        adapter,
1699        org_id,
1700        input.context.session_id,
1701        input.agent_id,
1702        Some(builtin_tool_registry),
1703        mcp_invoker,
1704        execution_capabilities.subagent_nesting_policy,
1705    );
1706    tool_registry.validate_context_services(&context_services)?;
1707    let executor: Arc<dyn everruns_core::tool_execution::ToolExecutor> = Arc::new(tool_registry);
1708
1709    let mut atom = ActAtom::new(executor, adapter.event_emitter())
1710        .with_context_services(context_services)
1711        .with_post_tool_hooks(execution_capabilities.post_tool_hooks)
1712        .with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
1713        .with_tool_call_hooks(execution_capabilities.tool_call_hooks);
1714
1715    #[cfg(feature = "builtins")]
1716    {
1717        atom = atom.with_final_post_tool_hook(Arc::new(everruns_builtins::PersistOutputHook));
1718    }
1719
1720    if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
1721        atom = atom.with_outbound_tool_rate_limiter(limiter);
1722    }
1723    if let Some(store) = adapter.durable_tool_result_store() {
1724        atom = atom.with_durable_tool_result_store(store);
1725    }
1726
1727    atom.execute(input).await
1728}