use crate::capabilities::{SystemPromptContext, collect_capabilities_with_configs};
use crate::engine::{
ActAtom, ActInput, ActResult, InputAtom, InputAtomInput, InputAtomResult, ReasonAtom,
ReasonInput, ReasonResult,
};
use crate::events::{
EventContext, EventRequest, OutputMessageCompletedData, SessionActivatedData, SessionIdledData,
SessionModelChangedData, TurnCompletedData, TurnFailedData, TurnStartedData,
};
use crate::host::SessionMutator;
use crate::host::turn_tool_context::{RuntimeToolCapabilityContext, runtime_tool_context_services};
use crate::message::{ContentPart, RuntimeMessage, RuntimeMessageRole};
use crate::message_retriever::MessageRetriever;
use crate::runtime_context::AssembledTurnContext;
use crate::session::SessionExecutionState;
use crate::{
CapabilityRegistry, DecisionsService, DependencyBlocker, EgressService,
ResolvedExecutionSnapshot, TokenUsage, ToolRegistry, UtilityLlmService,
org_public_id_from_internal, resolve_runtime_capabilities,
};
use crate::{
CompactionCheckpointStore, agents_api_store::AgentsApiStore,
connection_services::ProviderCredentialStore, connection_services::UserConnectionResolver,
delegation_services::SessionCreationAuthority, durability::DurableToolResultStore,
durability::PartialStreamStore, event_emitter::EventEmitter, execution_loading::AgentStore,
execution_loading::HarnessStore, execution_loading::SessionStore, file_services::FileResolver,
image_services::ImageArtifactStore, image_services::ImageResolver,
native_async_store::NativeAsyncStore, provider_resolution::ProviderStore,
session_files::SessionFileSystem, session_services::LeasedResourceStore,
session_services::SessionResourceRegistry, session_services::SessionScheduleStore,
session_services::SessionStorageStore, tool_execution::BudgetChecker,
tool_execution::PaymentAuthority,
};
use async_trait::async_trait;
use everruns_contracts::CapabilityRef;
use everruns_contracts::driver_registry::DriverRegistry;
use everruns_contracts::tool_types::ToolDefinition;
use everruns_contracts::typed_id::{AgentId, HarnessId, MessageId, ModelId, SessionId, TurnId};
use everruns_contracts::user_facing_error::{ErrorDisclosure, UserFacingError};
use std::sync::Arc;
use tracing::warn;
#[path = "host/message_filter_only.rs"]
mod message_filter_only;
use message_filter_only::MessageFilterOnlyCapability;
#[derive(Debug, Clone)]
pub struct ResolvedTurnInputs {
pub snapshot: ResolvedExecutionSnapshot,
pub messages: Vec<RuntimeMessage>,
pub mcp_tool_definitions: Vec<ToolDefinition>,
}
#[async_trait]
pub trait RuntimeHostAdapter: Send + Sync + Clone + 'static {
fn turn_cancellation(&self) -> Option<tokio::sync::watch::Receiver<bool>> {
None
}
fn turn_cancel_requested(&self) -> Option<tokio::sync::watch::Receiver<bool>> {
None
}
async fn set_session_status(
&self,
org_id: i64,
session_id: SessionId,
status: SessionExecutionState,
) -> everruns_contracts::error::Result<()>;
async fn load_resolved_turn(
&self,
org_id: i64,
session_id: SessionId,
) -> everruns_contracts::error::Result<ResolvedTurnInputs>;
async fn load_resolved_turn_for_execution(
&self,
org_id: i64,
session_id: SessionId,
_input_message_id: MessageId,
) -> everruns_contracts::error::Result<ResolvedTurnInputs> {
self.load_resolved_turn(org_id, session_id).await
}
fn capability_registry(&self) -> CapabilityRegistry;
fn driver_registry(&self) -> DriverRegistry;
fn harness_store(&self, org_id: i64) -> Arc<dyn HarnessStore>;
fn agent_store(&self, org_id: i64) -> Arc<dyn AgentStore>;
fn session_store(&self, org_id: i64) -> Arc<dyn SessionStore>;
fn session_mutator(&self, org_id: i64) -> Arc<dyn SessionMutator>;
fn provider_store(&self, org_id: i64) -> Arc<dyn ProviderStore>;
fn message_store(&self) -> Arc<dyn MessageRetriever>;
fn native_async_store(&self) -> Option<Arc<dyn NativeAsyncStore>> {
None
}
fn agents_api_store(&self) -> Option<Arc<dyn AgentsApiStore>> {
None
}
fn compaction_checkpoint_store(&self) -> Option<Arc<dyn CompactionCheckpointStore>> {
None
}
fn event_emitter(&self) -> Arc<dyn EventEmitter>;
fn file_store(&self, org_id: i64) -> Arc<dyn SessionFileSystem>;
fn bash_hook_dispatcher(
&self,
_org_id: i64,
) -> Arc<dyn crate::hook_executor::BashHookDispatcher> {
Arc::new(crate::host::DisabledBashHookDispatcher)
}
fn image_resolver(&self, _org_id: i64) -> Option<Arc<dyn ImageResolver>> {
None
}
fn file_resolver(&self, _org_id: i64) -> Option<Arc<dyn FileResolver>> {
None
}
fn image_artifact_store(&self, _org_id: i64) -> Option<Arc<dyn ImageArtifactStore>> {
None
}
fn provider_credential_store(&self, _org_id: i64) -> Option<Arc<dyn ProviderCredentialStore>> {
None
}
fn utility_llm_service(&self) -> Option<Arc<dyn UtilityLlmService>> {
None
}
fn decisions(&self) -> Option<Arc<dyn DecisionsService>> {
None
}
fn egress_service(&self) -> Option<Arc<dyn EgressService>> {
None
}
fn storage_store(&self, _org_id: i64) -> Option<Arc<dyn SessionStorageStore>> {
None
}
fn connection_resolver(&self) -> Option<Arc<dyn UserConnectionResolver>> {
None
}
fn tool_context_extensions(
&self,
_request: ToolContextRequest<'_>,
) -> crate::tool_context::ToolContextExtensions {
Default::default()
}
fn subagent_delegate(
&self,
_org_id: i64,
_session_id: SessionId,
) -> Option<Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>> {
None
}
fn tool_augmentor(&self) -> Option<Arc<dyn crate::host::HostToolAugmentor>> {
None
}
fn leased_resource_store(&self) -> Option<Arc<dyn LeasedResourceStore>> {
None
}
fn session_resource_registry(&self) -> Option<Arc<dyn SessionResourceRegistry>> {
None
}
fn session_task_registry(&self) -> Option<Arc<dyn crate::session_task::SessionTaskRegistry>> {
None
}
fn schedule_store(&self, _org_id: i64) -> Option<Arc<dyn SessionScheduleStore>> {
None
}
fn budget_checker(
&self,
_org_id: i64,
_agent_id: Option<AgentId>,
) -> Option<Arc<dyn BudgetChecker>> {
None
}
fn payment_authority(
&self,
_org_id: i64,
_agent_id: Option<AgentId>,
) -> Option<Arc<dyn PaymentAuthority>> {
None
}
fn session_creation_authority(
&self,
_org_id: i64,
_session_id: SessionId,
) -> Option<Arc<dyn SessionCreationAuthority>> {
None
}
fn outbound_tool_rate_limiter(
&self,
_org_id: i64,
) -> Option<Arc<dyn crate::tool_execution::OutboundToolRateLimiter>> {
None
}
fn durable_tool_result_store(&self) -> Option<Arc<dyn DurableToolResultStore>> {
None
}
fn subagent_spawn_store(
&self,
) -> Option<Arc<dyn crate::delegation_services::SubagentSpawnStore>> {
None
}
fn stream_heartbeater(&self) -> Option<Arc<dyn crate::durability::StreamHeartbeater>> {
None
}
fn partial_stream_store(&self) -> Option<Arc<dyn PartialStreamStore>> {
None
}
fn reasoning_effort_handle(
&self,
_session_id: SessionId,
) -> Option<crate::tool_context::ReasoningEffortHandle> {
None
}
fn provider_stall_timeout(&self) -> Option<std::time::Duration> {
None
}
fn provider_retry_config(&self) -> Option<everruns_contracts::llm_retry::LlmRetryConfig> {
None
}
async fn mcp_executor(
&self,
_org_id: i64,
_session_id: SessionId,
_agent_id: Option<AgentId>,
) -> Option<Arc<dyn crate::McpToolInvoker>> {
None
}
fn hosted_mcp_resolver(
&self,
_org_id: i64,
_session_id: SessionId,
_agent_id: Option<AgentId>,
_input_message_id: MessageId,
) -> Option<Arc<dyn everruns_contracts::hosted_mcp::HostedMcpResolver>> {
None
}
}
pub struct ToolContextRequest<'a> {
pub org_id: i64,
pub session_id: SessionId,
pub resolved_capabilities: &'a [CapabilityRef],
}
struct RuntimeExecutionCapabilities {
tool_registry: ToolRegistry,
post_tool_hooks: Vec<Arc<dyn crate::tool_hooks::PostToolExecHook>>,
pre_tool_hooks: Vec<Arc<dyn crate::tool_hooks::PreToolUseHook>>,
tool_call_hooks: Vec<Arc<dyn crate::ToolCallHook>>,
subagent_nesting_policy: crate::delegation_services::SubagentNestingPolicy,
resolved_capabilities: Vec<CapabilityRef>,
}
fn subagent_nesting_policy_from_configs(
resolved_capability_configs: &[everruns_contracts::CapabilityRef],
) -> crate::delegation_services::SubagentNestingPolicy {
let subagents_config = resolved_capability_configs
.iter()
.find(|config| config.capability_id() == "subagents");
let configured_depth = subagents_config
.and_then(|config| {
config
.config_value()
.get("max_subagent_depth")
.or_else(|| config.config_value().get("max_depth"))
})
.and_then(|value| value.as_u64())
.and_then(|value| u32::try_from(value).ok());
let configured_max_active = subagents_config
.and_then(|config| {
config
.config_value()
.get("max_active_descendant_tasks")
.or_else(|| config.config_value().get("max_concurrent_descendant_tasks"))
})
.and_then(|value| value.as_u64())
.and_then(|value| u32::try_from(value).ok());
let configured_max_total = subagents_config
.and_then(|config| config.config_value().get("max_total_descendant_tasks"))
.and_then(|value| value.as_u64())
.and_then(|value| u32::try_from(value).ok());
let configured_max_active_detached = subagents_config
.and_then(|config| config.config_value().get("max_active_detached_tasks"))
.and_then(|value| value.as_u64())
.and_then(|value| u32::try_from(value).ok());
let configured_max_total_detached = subagents_config
.and_then(|config| config.config_value().get("max_total_detached_tasks"))
.and_then(|value| value.as_u64())
.and_then(|value| u32::try_from(value).ok());
crate::delegation_services::SubagentNestingPolicy::default()
.with_agent_override(configured_depth)
.with_agent_task_caps_override(configured_max_active, configured_max_total)
.with_agent_detached_task_caps_override(
configured_max_active_detached,
configured_max_total_detached,
)
}
fn finalize_specs_from_configs(
resolved_capability_configs: &[everruns_contracts::CapabilityRef],
capability_registry: &CapabilityRegistry,
tool_augmentor: Option<&dyn crate::host::HostToolAugmentor>,
) -> Vec<crate::user_hook_types::UserHookSpec> {
let mut hook_contributions: Vec<(String, Vec<crate::user_hook_types::UserHookSpec>)> =
Vec::new();
let mut disabled_contributions: Vec<String> = Vec::new();
for config in resolved_capability_configs {
let Some(capability) = capability_registry.get(config.capability_id()) else {
continue;
};
let specs = capability.user_hooks_with_config(config.config_value());
if !specs.is_empty() {
hook_contributions.push((config.capability_id().to_string(), specs));
}
if let Some(augmentor) = tool_augmentor {
disabled_contributions.extend(
augmentor
.disabled_hook_contributions(config.capability_id(), config.config_value()),
);
}
}
crate::hook_adapter::finalize_hook_specs(hook_contributions, &disabled_contributions)
}
async fn collect_lifecycle_hook_specs<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
session_id: SessionId,
harness_id: HarnessId,
agent_id: Option<AgentId>,
) -> everruns_contracts::error::Result<(
Vec<crate::user_hook_types::UserHookSpec>,
Arc<dyn crate::hook_executor::BashHookDispatcher>,
)> {
let capability_registry = adapter.capability_registry();
let harness = adapter
.harness_store(org_id)
.get_harness(harness_id)
.await?
.ok_or_else(|| everruns_contracts::error::AgentLoopError::harness_not_found(harness_id))?;
let session = adapter
.session_store(org_id)
.get_session(session_id)
.await?
.ok_or_else(|| everruns_contracts::error::AgentLoopError::session_not_found(session_id))?;
let agent = match agent_id {
Some(agent_id) => adapter.agent_store(org_id).get_agent(agent_id).await?,
None => None,
};
let resolved =
resolve_runtime_capabilities(&harness, agent.as_ref(), &session, &capability_registry);
let tool_augmentor = adapter.tool_augmentor();
let specs = finalize_specs_from_configs(
&resolved.resolved_capability_configs,
&capability_registry,
tool_augmentor.as_deref(),
);
let dispatcher = adapter.bash_hook_dispatcher(org_id);
Ok((specs, dispatcher))
}
#[derive(Debug, Default)]
struct TurnAgentIdentity {
id: Option<AgentId>,
name: Option<String>,
description: Option<String>,
}
pub struct RuntimeSessionLifecycle<A: RuntimeHostAdapter> {
pub(crate) adapter: A,
pub(super) org_id: i64,
pub(crate) session_id: SessionId,
}
impl<A: RuntimeHostAdapter> RuntimeSessionLifecycle<A> {
pub fn new(adapter: A, org_id: i64, session_id: SessionId) -> Self {
Self {
adapter,
org_id,
session_id,
}
}
async fn set_session_status(
&self,
status: SessionExecutionState,
_action: &'static str,
) -> everruns_contracts::error::Result<()> {
self.adapter
.set_session_status(self.org_id, self.session_id, status)
.await
}
async fn emit_event(&self, request: EventRequest) -> everruns_contracts::error::Result<()> {
self.adapter.event_emitter().emit(request).await.map(|_| ())
}
pub async fn turn_started(
&self,
turn_id: TurnId,
input_message_id: MessageId,
) -> everruns_contracts::error::Result<()> {
let input_content = async {
self.adapter
.message_store()
.get(self.session_id, input_message_id)
.await
.ok()
.flatten()
.map(|message| message.content_to_llm_string())
};
let activated = async {
self.set_session_status(SessionExecutionState::Active, "turn_started")
.await?;
Ok::<_, everruns_contracts::error::AgentLoopError>(self.agent_identity().await)
};
let (input_content, agent) = futures::join!(input_content, activated);
let agent = agent?;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
SessionActivatedData {
turn_id,
input_message_id,
},
))
.await?;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
TurnStartedData {
turn_id,
input_message_id,
input_content,
agent_id: agent.id,
agent_name: agent.name,
agent_description: agent.description,
},
))
.await?;
Ok(())
}
async fn agent_identity(&self) -> TurnAgentIdentity {
let agent_id = self
.adapter
.session_store(self.org_id)
.get_session(self.session_id)
.await
.ok()
.flatten()
.and_then(|session| session.agent_id);
let Some(agent_id) = agent_id else {
return TurnAgentIdentity::default();
};
let agent = self
.adapter
.agent_store(self.org_id)
.get_agent(agent_id)
.await
.ok()
.flatten();
TurnAgentIdentity {
id: Some(agent_id),
name: agent
.as_ref()
.map(|a| a.display_name.clone().unwrap_or_else(|| a.name.clone())),
description: agent.and_then(|a| a.description),
}
}
pub async fn emit_turn_completed(
&self,
input_message_id: MessageId,
data: TurnCompletedData,
) -> everruns_contracts::error::Result<()> {
let turn_id = data.turn_id;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
data,
))
.await
}
pub async fn emit_session_idled(
&self,
turn_id: TurnId,
input_message_id: MessageId,
iterations: Option<u32>,
usage: Option<TokenUsage>,
) -> everruns_contracts::error::Result<()> {
self.set_session_status(SessionExecutionState::Idle, "emit_session_idled")
.await?;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
SessionIdledData {
turn_id,
iterations,
usage,
},
))
.await
}
pub async fn turn_completed(
&self,
turn_id: TurnId,
input_message_id: MessageId,
iterations: u32,
usage: Option<TokenUsage>,
input_content: Option<String>,
) -> everruns_contracts::error::Result<()> {
self.emit_turn_completed(
input_message_id,
TurnCompletedData {
turn_id,
iterations,
duration_ms: None,
usage: usage.clone(),
input_content,
status: Some("completed".to_string()),
..Default::default()
},
)
.await?;
self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
.await
}
pub async fn turn_sealed(
&self,
turn_id: TurnId,
input_message_id: MessageId,
reason: &str,
iterations: u32,
usage: Option<TokenUsage>,
) -> everruns_contracts::error::Result<()> {
let context = EventContext::turn(turn_id, input_message_id);
self.emit_event(EventRequest::new(
self.session_id,
context.clone(),
crate::events::TurnSealedData {
turn_id,
reason: reason.to_string(),
detail: None,
iterations: Some(iterations),
usage: usage.clone(),
},
))
.await?;
self.emit_session_idled(turn_id, input_message_id, Some(iterations), usage)
.await
}
pub async fn fire_turn_end_hooks(
&self,
harness_id: HarnessId,
agent_id: Option<AgentId>,
turn_id: TurnId,
success: bool,
) {
let (specs, dispatcher) = match collect_lifecycle_hook_specs(
&self.adapter,
self.org_id,
self.session_id,
harness_id,
agent_id,
)
.await
{
Ok(pair) => pair,
Err(error) => {
warn!(
session_id = %self.session_id,
%error,
"failed to collect turn_end hook specs; skipping"
);
return;
}
};
let hooks = crate::lifecycle_hooks::build_turn_lifecycle_hooks(
&specs,
crate::user_hook_types::HookEvent::TurnEnd,
dispatcher,
);
if hooks.is_empty() {
return;
}
let ctx = crate::lifecycle_hooks::TurnHookContext {
session_id: self.session_id,
turn_id: Some(turn_id),
org_id: org_public_id_from_internal(self.org_id).parse().ok(),
agent_id: agent_id.map(|a| a.to_string()),
};
crate::lifecycle_hooks::run_turn_end_hooks(
&hooks,
&ctx,
serde_json::json!({ "success": success }),
)
.await;
}
pub async fn user_prompt_blocked(
&self,
turn_id: TurnId,
input_message_id: MessageId,
reason: &str,
user_message: Option<&str>,
) -> everruns_contracts::error::Result<()> {
let user_error =
UserFacingError::new(everruns_contracts::user_facing_error::codes::BLOCKED_BY_HOOK);
let shown = user_message.unwrap_or(reason);
let mut error_message = RuntimeMessage::assistant(shown);
let mut metadata = std::collections::HashMap::new();
user_error.apply_to_message_metadata(&mut metadata);
error_message.metadata = Some(metadata);
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
))
.await?;
self.turn_failed(turn_id, input_message_id, reason, Some(&user_error))
.await
}
pub async fn turn_failed(
&self,
turn_id: TurnId,
input_message_id: MessageId,
error: &str,
user_error: Option<&UserFacingError>,
) -> everruns_contracts::error::Result<()> {
self.turn_failed_with_disclosure(turn_id, input_message_id, error, user_error, None)
.await
}
pub async fn turn_failed_with_disclosure(
&self,
turn_id: TurnId,
input_message_id: MessageId,
error: &str,
user_error: Option<&UserFacingError>,
disclosure: Option<ErrorDisclosure>,
) -> everruns_contracts::error::Result<()> {
self.set_session_status(SessionExecutionState::Idle, "turn_failed")
.await?;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
{
let mut data = TurnFailedData {
turn_id,
error: error.to_string(),
error_code: None,
error_fields: None,
error_disclosure: disclosure.map(|mode| mode.as_str().to_string()),
};
if let Some(user_error) = user_error {
user_error.apply_to_event_fields(&mut data.error_code, &mut data.error_fields);
}
data
},
))
.await?;
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
SessionIdledData {
turn_id,
iterations: None,
usage: None,
},
))
.await
}
pub async fn waiting_for_tool_results(&self) -> everruns_contracts::error::Result<()> {
self.set_session_status(
SessionExecutionState::WaitingForToolResults,
"waiting_for_tool_results",
)
.await
}
pub async fn dependency_blocked(
&self,
turn_id: TurnId,
input_message_id: MessageId,
blocker: DependencyBlocker,
) -> everruns_contracts::error::Result<()> {
let user_error = UserFacingError::new(blocker.error_code())
.with_field(
"dependency",
match blocker {
DependencyBlocker::HarnessArchived | DependencyBlocker::HarnessDeleted => {
"harness"
}
DependencyBlocker::AgentArchived | DependencyBlocker::AgentDeleted => "agent",
},
)
.with_field(
"state",
match blocker {
DependencyBlocker::HarnessArchived | DependencyBlocker::AgentArchived => {
"archived"
}
DependencyBlocker::HarnessDeleted | DependencyBlocker::AgentDeleted => {
"deleted"
}
},
);
let mut error_message = RuntimeMessage::assistant(blocker.message());
let mut metadata = std::collections::HashMap::new();
user_error.apply_to_message_metadata(&mut metadata);
error_message.metadata = Some(metadata);
self.emit_event(EventRequest::new(
self.session_id,
EventContext::turn(turn_id, input_message_id),
OutputMessageCompletedData::new(error_message).with_user_facing_error(&user_error),
))
.await?;
self.turn_failed(
turn_id,
input_message_id,
blocker.message(),
Some(&user_error),
)
.await
}
}
pub async fn detect_dependency_blocker<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
harness_id: HarnessId,
agent_id: Option<AgentId>,
) -> everruns_contracts::error::Result<Option<DependencyBlocker>> {
let harness_store = adapter.harness_store(org_id);
let agent_store = adapter.agent_store(org_id);
if let Some(blocker) = harness_store.get_harness_blocker(harness_id).await? {
return Ok(Some(blocker));
}
if let Some(agent_id) = agent_id
&& let Some(blocker) = agent_store.get_agent_blocker(agent_id).await?
{
return Ok(Some(blocker));
}
Ok(None)
}
pub async fn execute_input_activity<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: InputAtomInput,
) -> everruns_contracts::error::Result<InputAtomResult> {
if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
handle.set(None);
}
RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
.turn_started(input.context.turn_id, input.context.input_message_id)
.await?;
let atom = InputAtom::new(adapter.message_store());
atom.execute(input).await
}
pub(crate) struct UserPromptHookResult {
pub(crate) decision: crate::lifecycle_hooks::UserPromptDecision,
pub(crate) original_message: String,
}
pub(crate) async fn run_user_prompt_submit_for_message<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: &ReasonInput,
message_text: String,
) -> everruns_contracts::error::Result<Option<UserPromptHookResult>> {
let (specs, dispatcher) = match collect_lifecycle_hook_specs(
adapter,
org_id,
input.context.session_id,
input.harness_id,
input.agent_id,
)
.await
{
Ok(pair) => pair,
Err(error) => {
warn!(
session_id = %input.context.session_id,
%error,
"failed to collect user_prompt_submit hook specs; continuing without them"
);
return Ok(None);
}
};
let hooks = crate::lifecycle_hooks::build_turn_lifecycle_hooks(
&specs,
crate::user_hook_types::HookEvent::UserPromptSubmit,
dispatcher,
);
if hooks.is_empty() {
return Ok(None);
}
let ctx = crate::lifecycle_hooks::TurnHookContext {
session_id: input.context.session_id,
turn_id: Some(input.context.turn_id),
org_id: org_public_id_from_internal(org_id).parse().ok(),
agent_id: input.agent_id.map(|a| a.to_string()),
};
let original_message = message_text.clone();
let decision =
crate::lifecycle_hooks::run_user_prompt_submit_hooks(&hooks, &ctx, message_text).await;
Ok(Some(UserPromptHookResult {
decision,
original_message,
}))
}
async fn run_user_prompt_submit_for_turn<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: &ReasonInput,
) -> everruns_contracts::error::Result<Option<UserPromptHookResult>> {
let message_text = adapter
.message_store()
.get(input.context.session_id, input.context.input_message_id)
.await
.ok()
.flatten()
.map(|m| m.content_to_llm_string())
.unwrap_or_default();
run_user_prompt_submit_for_message(adapter, org_id, input, message_text).await
}
pub async fn execute_reason_activity<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: ReasonInput,
) -> everruns_contracts::error::Result<ReasonResult> {
let prompt_message_ids = (input.iteration <= 1)
.then_some(input.context.input_message_id)
.into_iter()
.collect();
execute_reason_activity_with_prompt_messages(adapter, org_id, input, prompt_message_ids).await
}
fn blocked_reason_result(text: String, error: &str) -> ReasonResult {
ReasonResult {
text,
max_iterations: crate::runtime_agent::default_max_iterations(),
error: Some(error.to_string()),
..Default::default()
}
}
pub async fn execute_reason_activity_with_prompt_messages<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: ReasonInput,
prompt_message_ids: Vec<MessageId>,
) -> everruns_contracts::error::Result<ReasonResult> {
if let Some(blocker) =
detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
{
RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
.dependency_blocked(
input.context.turn_id,
input.context.input_message_id,
blocker,
)
.await?;
return Ok(blocked_reason_result(
blocker.message().to_string(),
"dependency_unavailable",
));
}
let mut user_prompt_message_overrides = Vec::new();
for message_id in prompt_message_ids {
let mut hook_input = input.clone();
hook_input.context.input_message_id = message_id;
let Some(hook_result) =
run_user_prompt_submit_for_turn(adapter, org_id, &hook_input).await?
else {
continue;
};
match hook_result.decision {
crate::lifecycle_hooks::UserPromptDecision::Block {
reason,
user_message,
} => {
RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
.user_prompt_blocked(
input.context.turn_id,
input.context.input_message_id,
&reason,
user_message.as_deref(),
)
.await?;
return Ok(blocked_reason_result(
user_message.unwrap_or_else(|| reason.clone()),
"blocked_by_user_prompt_hook",
));
}
crate::lifecycle_hooks::UserPromptDecision::Continue { message } => {
if message != hook_result.original_message {
user_prompt_message_overrides.push((message_id, message));
}
}
}
}
let validation_session = adapter
.session_store(org_id)
.get_session(input.context.session_id)
.await?
.ok_or_else(|| {
everruns_contracts::error::AgentLoopError::session_not_found(input.context.session_id)
})?;
let validation_capabilities = load_execution_capabilities(
adapter,
org_id,
input.context.session_id,
input.harness_id,
input.agent_id,
validation_session.locale.clone(),
validation_session.blueprint_id.as_deref(),
)
.await?;
let query_history_allowed = validation_capabilities
.tool_registry
.get("query_history")
.is_some();
let validation_services = runtime_tool_context_services(
adapter,
org_id,
input.context.session_id,
input.agent_id,
Some(Arc::new(validation_capabilities.tool_registry.clone())),
None,
RuntimeToolCapabilityContext {
subagent_nesting_policy: validation_capabilities.subagent_nesting_policy,
resolved_capabilities: validation_capabilities.resolved_capabilities.clone(),
},
);
validation_capabilities
.tool_registry
.validate_context_services(&validation_services)?;
let mut turn_inputs = adapter
.load_resolved_turn_for_execution(
org_id,
input.context.session_id,
input.context.input_message_id,
)
.await?;
if let Some(augmentor) = adapter.tool_augmentor() {
augmentor
.augment_reason_tools(
input.context.session_id,
adapter.session_store(org_id),
adapter.session_task_registry(),
&mut turn_inputs.mcp_tool_definitions,
)
.await?;
}
let reason_capability_registry = {
let mut registry = adapter.capability_registry();
if !query_history_allowed {
let query_history_owner = registry
.list()
.into_iter()
.find(|capability| {
capability
.tool_definitions()
.iter()
.any(|tool| tool.name() == "query_history")
})
.map(Arc::clone);
if let Some(capability) = query_history_owner {
registry.register(MessageFilterOnlyCapability(capability));
}
}
registry
};
let context_resolver = crate::host::runtime_context::StoreTurnContextResolver::new(
adapter.harness_store(org_id),
adapter.agent_store(org_id),
adapter.session_store(org_id),
adapter.message_store(),
adapter.provider_store(org_id),
reason_capability_registry.clone(),
adapter.driver_registry(),
)
.with_file_store(adapter.file_store(org_id))
.with_hosted_mcp_resolver(adapter.hosted_mcp_resolver(
org_id,
input.context.session_id,
input.agent_id,
input.context.input_message_id,
));
let context_resolver = match adapter.storage_store(org_id) {
Some(store) => context_resolver.with_session_storage(store),
None => context_resolver,
};
let mut atom = ReasonAtom::new(
context_resolver,
adapter.message_store(),
reason_capability_registry.clone(),
adapter.event_emitter(),
);
if let Some(image_resolver) = adapter.image_resolver(org_id) {
atom = atom.with_image_resolver(image_resolver);
}
if let Some(file_resolver) = adapter.file_resolver(org_id) {
atom = atom.with_file_resolver(file_resolver);
}
if let Some(hb) = adapter.stream_heartbeater() {
atom = atom.with_stream_heartbeater(hb);
}
if let Some(timeout) = adapter.provider_stall_timeout() {
atom = atom.with_provider_stall_timeout(timeout);
}
if let Some(config) = adapter.provider_retry_config() {
atom = atom.with_provider_retry_config(config);
}
if let Some(store) = adapter.partial_stream_store() {
atom = atom.with_partial_stream_store(store);
}
if let Some(store) = adapter.durable_tool_result_store() {
atom = atom.with_durable_tool_result_store(store);
}
if let Some(store) = adapter.compaction_checkpoint_store() {
atom = atom.with_compaction_checkpoint_store(store);
}
if let Some(handle) = adapter.reasoning_effort_handle(input.context.session_id) {
atom = atom.with_reasoning_effort_handle(handle);
}
if let Some(utility_llm_service) = adapter.utility_llm_service() {
atom = atom.with_utility_llm_service(utility_llm_service);
}
if let Some(decisions) = adapter.decisions() {
atom = atom.with_decisions(decisions);
}
if let Some(schedule_store) = adapter.schedule_store(org_id) {
atom = atom.with_schedule_store(schedule_store);
}
let mut assembled = crate::host::runtime_context::assemble_turn_context_from_snapshot(
turn_inputs.snapshot,
adapter.message_store().as_ref(),
adapter.provider_store(org_id).as_ref(),
&reason_capability_registry,
&adapter.driver_registry(),
&turn_inputs.mcp_tool_definitions,
Some(adapter.file_store(org_id)),
adapter.storage_store(org_id),
)
.await?;
let input = ReasonInput {
mcp_tool_definitions: turn_inputs.mcp_tool_definitions,
..input
};
if !user_prompt_message_overrides.is_empty() {
for (message_id, message_override) in user_prompt_message_overrides {
let message = assembled
.messages
.iter_mut()
.find(|message| message.id == message_id)
.ok_or_else(|| {
everruns_contracts::error::AgentLoopError::config(
"user_prompt_submit mutation: input message not found in assembled context",
)
})?;
message
.content
.retain(|part| !matches!(part, ContentPart::Text(_)));
message
.content
.insert(0, ContentPart::text(message_override));
}
}
if input.iteration <= 1 {
emit_model_change_if_switched(adapter, org_id, &input, &assembled).await;
}
crate::host::reason_backend::execute_reason(adapter, org_id, input, assembled, atom).await
}
async fn emit_model_change_if_switched<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
input: &ReasonInput,
assembled: &AssembledTurnContext,
) {
let Some((previous_model_id, model_id)) = model_switch(&assembled.messages) else {
return;
};
if assembled.resolved_model_id != Some(model_id) {
return;
}
let previous_model_name = adapter
.provider_store(org_id)
.get_model_spec(previous_model_id)
.await
.ok()
.flatten()
.map(|spec| spec.model);
let request = EventRequest::new(
input.context.session_id,
EventContext::turn(input.context.turn_id, input.context.input_message_id),
SessionModelChangedData {
previous_model_id: Some(previous_model_id),
previous_model_name,
model_id,
model_name: assembled.model.model.clone(),
},
);
if let Err(e) = adapter.event_emitter().emit(request).await {
warn!(error = %e, "Failed to emit session.model.changed event");
}
}
fn model_switch(messages: &[RuntimeMessage]) -> Option<(ModelId, ModelId)> {
let mut user_model_ids = messages
.iter()
.rev()
.filter(|message| message.role == RuntimeMessageRole::User)
.map(|message| {
message
.controls
.as_ref()
.and_then(|controls| controls.model_id)
});
let model_id = user_model_ids.next().flatten()?;
let previous_model_id = user_model_ids.next().flatten()?;
(previous_model_id != model_id).then_some((previous_model_id, model_id))
}
pub async fn execute_act_activity<A: RuntimeHostAdapter>(
adapter: &A,
input: ActInput,
) -> everruns_contracts::error::Result<ActResult> {
let org_id = input.org_id.ok_or_else(|| {
everruns_contracts::error::AgentLoopError::config(
"ActInput.org_id must be set for runtime host execution",
)
})?;
if let Some(blocker) =
detect_dependency_blocker(adapter, org_id, input.harness_id, input.agent_id).await?
{
RuntimeSessionLifecycle::new(adapter.clone(), org_id, input.context.session_id)
.dependency_blocked(
input.context.turn_id,
input.context.input_message_id,
blocker,
)
.await?;
return Ok(ActResult {
results: vec![],
completed: true,
success_count: 0,
error_count: 1,
waiting_for_tool_results: false,
waiting_for_url_elicitation: false,
blocked: true,
client_tool_calls: vec![],
client_tool_definitions: vec![],
});
}
let execution_capabilities = load_execution_capabilities(
adapter,
org_id,
input.context.session_id,
input.harness_id,
input.agent_id,
input.locale.clone(),
input.blueprint_id.as_deref(),
)
.await?;
let mut tool_registry = execution_capabilities.tool_registry;
if let Some(augmentor) = adapter.tool_augmentor() {
augmentor
.augment_act_tools(
input.context.session_id,
adapter.session_store(org_id),
adapter.session_task_registry(),
adapter.file_store(org_id),
&input.tool_definitions,
&mut tool_registry,
)
.await?;
}
let mut mcp_invoker: Option<Arc<dyn crate::McpToolInvoker>> = None;
if let Some(mcp) = adapter
.mcp_executor(org_id, input.context.session_id, input.agent_id)
.await
{
let invoker: Arc<dyn crate::McpToolInvoker> = mcp;
for tool in crate::build_mcp_proxy_tools(&input.tool_definitions, invoker.clone()) {
tool_registry.register_boxed(tool);
}
mcp_invoker = Some(Arc::new(crate::ScopedMcpToolInvoker::new(
&input.tool_definitions,
invoker,
)));
}
let builtin_tool_registry = Arc::new(tool_registry.clone());
let context_services = runtime_tool_context_services(
adapter,
org_id,
input.context.session_id,
input.agent_id,
Some(builtin_tool_registry),
mcp_invoker,
RuntimeToolCapabilityContext {
subagent_nesting_policy: execution_capabilities.subagent_nesting_policy,
resolved_capabilities: execution_capabilities.resolved_capabilities.clone(),
},
);
tool_registry.validate_context_services(&context_services)?;
let executor: Arc<dyn crate::tool_execution::ToolExecutor> = Arc::new(tool_registry);
let mut atom = ActAtom::new(executor, adapter.event_emitter())
.with_context_services(context_services)
.with_post_tool_hooks(execution_capabilities.post_tool_hooks)
.with_pre_tool_hooks(execution_capabilities.pre_tool_hooks)
.with_tool_call_hooks(execution_capabilities.tool_call_hooks);
#[cfg(feature = "builtins")]
{
atom = atom.with_final_post_tool_hook(Arc::new(crate::builtins::PersistOutputHook));
}
if let Some(limiter) = adapter.outbound_tool_rate_limiter(org_id) {
atom = atom.with_outbound_tool_rate_limiter(limiter);
}
if let Some(store) = adapter.durable_tool_result_store() {
atom = atom.with_durable_tool_result_store(store);
}
atom.execute(input).await
}
#[path = "host/execution_capabilities.rs"]
mod execution_capabilities;
use execution_capabilities::load_execution_capabilities;