Skip to main content

everruns_runtime/
host.rs

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