use crate::backends::{
HostBackends, RuntimeAgentStore, RuntimeHarnessStore, RuntimeProviderStore, RuntimeSessionStore,
};
use crate::builders::SingleSessionBuilder;
use crate::events::{EventHistory, EventLog, EventReadLimit, EventReadRequest, HostEventEmitter};
use crate::host::{
RuntimeHostAdapter, RuntimeHostTurnContext, execute_act_activity, execute_input_activity,
execute_reason_activity_with_prompt_messages,
};
use crate::in_memory::{InMemorySessionFileStore, InMemorySessionFileSystemFactory};
use crate::turn_strategy::{RuntimeTurnPlan, RuntimeTurnState};
use async_trait::async_trait;
use chrono::Utc;
use everruns_core::agent::Agent;
use everruns_core::atoms::{AtomContext, InputAtomInput, ReasonInput};
use everruns_core::capabilities::{
Capability, CapabilityRegistry, CapabilityStatus, collect_capability_mcp_servers,
resolve_capability_configs,
};
use everruns_core::config_layer::AgentConfigOverlay;
use everruns_core::driver_registry::{DriverId, DriverRegistry};
use everruns_core::error::{AgentLoopError, Result};
use everruns_core::events::{
Event, EventContext, EventRequest, InputMessageData, SessionStartedData,
};
use everruns_core::harness::Harness;
use everruns_core::llmsim_driver::{LlmSimConfig, LlmSimDriver};
use everruns_core::message::Message;
use everruns_core::platform_definition::PlatformDefinition;
use everruns_core::plugins::{PluginFileSet, compile_plugin};
use everruns_core::runtime_context::{AssembledTurnContext, inspect_turn_context};
use everruns_core::session::{Session, SessionStatus};
use everruns_core::session_file::{InitialFile, SessionFile};
use everruns_core::traits::{
AgentStore, EventEmitter, HarnessStore, ProviderStore, ResolvedModel, SessionMutator,
SessionStorageStore, SessionStore, UserConnectionResolver,
};
use everruns_core::turn::TurnStopReason;
use everruns_core::typed_id::{AgentId, HarnessId, MessageId, OrgId, SessionId, TurnId};
use everruns_core::{
AgentCapabilityConfig, CapabilityId, InputMessage, MessageRetriever, SessionFileSystem,
SessionFileSystemFactoryContext, plugin_capability_id, resolve_runtime_capabilities,
};
use everruns_engine::{
ActOutcome, plan_after_act, plan_after_process_input, plan_after_reason, reason_schedules_act,
};
use sha2::{Digest, Sha256};
use std::collections::VecDeque;
use std::path::Path;
use std::sync::{Arc, Mutex};
const HASH_INPUT_CAP_BYTES: usize = 128;
pub fn in_process_internal_org_id(public_org_id: &str) -> i64 {
if public_org_id == everruns_core::DEFAULT_ORG_PUBLIC_ID {
return everruns_core::DEFAULT_ORG_ID;
}
let Ok(parsed) = public_org_id.parse::<OrgId>() else {
return hash_public_org_id(public_org_id);
};
let raw: u128 = parsed.uuid().as_u128();
if raw == 0 {
return hash_public_org_id(public_org_id);
}
if raw <= i64::MAX as u128 {
return raw as i64;
}
hash_public_org_id(public_org_id)
}
fn hash_public_org_id(public_org_id: &str) -> i64 {
let bytes = public_org_id.as_bytes();
let bounded = &bytes[..bytes.len().min(HASH_INPUT_CAP_BYTES)];
let digest = Sha256::digest(bounded);
let mut buf = [0u8; 8];
buf.copy_from_slice(&digest[..8]);
let raw = u64::from_be_bytes(buf);
((raw % ((i64::MAX - 1) as u64)) as i64) + 2
}
#[derive(Debug, Clone)]
pub struct TurnResult {
pub response: String,
pub iterations: usize,
pub tool_calls_count: usize,
pub success: bool,
pub error: Option<String>,
pub stop_reason: TurnStopReason,
pub turn_id: everruns_core::typed_id::TurnId,
}
#[doc(hidden)]
#[derive(Clone, Debug)]
pub struct AcceptedTurnInput {
message_id: MessageId,
input: InputMessage,
}
impl AcceptedTurnInput {
pub fn new(input: impl Into<InputMessage>) -> Self {
Self {
message_id: MessageId::new(),
input: input.into(),
}
}
pub fn message_id(&self) -> MessageId {
self.message_id
}
pub fn input(&self) -> &InputMessage {
&self.input
}
fn into_message(self) -> Message {
message_from_input_with_id(self.message_id, self.input)
}
}
#[doc(hidden)]
#[derive(Clone, Debug)]
pub struct TurnSteering {
state: Arc<Mutex<TurnSteeringState>>,
}
#[doc(hidden)]
#[derive(Debug)]
pub enum TurnSteeringPushError {
Closed(Box<AcceptedTurnInput>),
Full(Box<AcceptedTurnInput>),
}
const TURN_STEERING_CAPACITY: usize = 256;
#[derive(Debug, Default)]
struct TurnSteeringState {
open: bool,
inputs: VecDeque<AcceptedTurnInput>,
}
impl TurnSteering {
pub fn new() -> Self {
Self {
state: Arc::new(Mutex::new(TurnSteeringState {
open: true,
inputs: VecDeque::new(),
})),
}
}
pub fn try_push(
&self,
input: AcceptedTurnInput,
) -> std::result::Result<(), TurnSteeringPushError> {
let mut state = self.state.lock().expect("turn steering lock poisoned");
if !state.open {
return Err(TurnSteeringPushError::Closed(Box::new(input)));
}
if state.inputs.len() >= TURN_STEERING_CAPACITY {
return Err(TurnSteeringPushError::Full(Box::new(input)));
}
state.inputs.push_back(input);
Ok(())
}
fn drain(&self) -> Vec<AcceptedTurnInput> {
let mut state = self.state.lock().expect("turn steering lock poisoned");
state.inputs.drain(..).collect()
}
fn drain_or_close(&self) -> Vec<AcceptedTurnInput> {
let mut state = self.state.lock().expect("turn steering lock poisoned");
if state.inputs.is_empty() {
state.open = false;
return vec![];
}
state.inputs.drain(..).collect()
}
pub fn close(&self) {
self.state.lock().expect("turn steering lock poisoned").open = false;
}
pub fn close_and_drain(&self) -> Vec<AcceptedTurnInput> {
let mut state = self.state.lock().expect("turn steering lock poisoned");
state.open = false;
state.inputs.drain(..).collect()
}
}
impl Default for TurnSteering {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CapabilityDelta {
pub capability_id: String,
pub changed: bool,
pub active: bool,
pub surfaces_dirty: bool,
}
fn finish_turn(
turn_id: everruns_core::typed_id::TurnId,
stop_reason: TurnStopReason,
error: Option<String>,
response: String,
iterations: usize,
tool_calls_count: usize,
) -> TurnResult {
if error.is_some() {
return TurnResult {
response: String::new(),
iterations,
tool_calls_count: 0,
success: false,
error,
stop_reason,
turn_id,
};
}
TurnResult {
response,
iterations,
tool_calls_count,
success: true,
error: None,
stop_reason,
turn_id,
}
}
pub struct InProcessRuntimeBuilder {
platform_definition: PlatformDefinition,
llm_sim_config: Option<LlmSimConfig>,
providers: Vec<everruns_core::Provider>,
model_spec: Option<everruns_core::ModelSpec>,
legacy_provider_config: Option<everruns_core::ProviderConfig>,
backends: Option<HostBackends>,
workspace_policy: Option<everruns_core::WorkspacePolicy>,
session_file_system_factory_context: SessionFileSystemFactoryContext,
harnesses: Vec<Harness>,
agents: Vec<Agent>,
sessions: Vec<Session>,
default_session_id: Option<SessionId>,
seeded_files: Vec<(SessionId, InitialFile)>,
mcp_auth_provider: Option<Arc<dyn everruns_mcp::McpAuthProvider>>,
provider_retry_config: Option<everruns_core::llm_retry::LlmRetryConfig>,
provider_stall_timeout: Option<std::time::Duration>,
plugin_capability_configs: Vec<AgentCapabilityConfig>,
plugin_warnings: Vec<String>,
}
impl Default for InProcessRuntimeBuilder {
fn default() -> Self {
Self::new()
}
}
impl InProcessRuntimeBuilder {
pub fn new() -> Self {
Self {
platform_definition: PlatformDefinition::builder()
.capability_registry(CapabilityRegistry::runtime_builtins())
.driver_registry(DriverRegistry::new())
.session_file_system_factory(Arc::new(InMemorySessionFileSystemFactory))
.build(),
llm_sim_config: None,
providers: Vec::new(),
model_spec: None,
legacy_provider_config: None,
backends: None,
workspace_policy: None,
session_file_system_factory_context: SessionFileSystemFactoryContext::new(),
harnesses: Vec::new(),
agents: Vec::new(),
sessions: Vec::new(),
default_session_id: None,
seeded_files: Vec::new(),
mcp_auth_provider: None,
provider_retry_config: None,
provider_stall_timeout: None,
plugin_capability_configs: Vec::new(),
plugin_warnings: Vec::new(),
}
}
pub fn mcp_auth_provider(mut self, provider: Arc<dyn everruns_mcp::McpAuthProvider>) -> Self {
self.mcp_auth_provider = Some(provider);
self
}
pub fn platform_definition(mut self, platform_definition: PlatformDefinition) -> Self {
self.platform_definition = platform_definition;
self
}
pub fn capability<C: Capability + 'static>(mut self, capability: C) -> Self {
self.platform_definition
.capability_registry_mut()
.register(capability);
self
}
pub fn driver_registry(mut self, driver_registry: DriverRegistry) -> Self {
*self.platform_definition.driver_registry_mut() = driver_registry;
self
}
pub fn provider(mut self, provider: everruns_core::Provider) -> Self {
self.providers.push(provider);
self
}
pub fn llm_sim(mut self, config: LlmSimConfig) -> Self {
self.llm_sim_config = Some(config);
self
}
pub fn default_model(mut self, model: ResolvedModel) -> Self {
let (spec, provider_config) = model.canonical_parts();
self.model_spec = Some(spec);
self.legacy_provider_config = Some(provider_config);
self
}
pub fn model_spec(mut self, model: everruns_core::ModelSpec) -> Self {
self.model_spec = Some(model);
self.legacy_provider_config = None;
self
}
pub fn backends(mut self, backends: HostBackends) -> Self {
self.backends = Some(backends);
self
}
pub fn workspace_policy(mut self, policy: everruns_core::WorkspacePolicy) -> Self {
self.workspace_policy = Some(policy);
self
}
pub fn with_session_task_registry(
mut self,
registry: Arc<dyn everruns_core::session_task::SessionTaskRegistry>,
) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_session_task_registry(registry));
self
}
pub fn with_schedule_store_factory(
mut self,
factory: crate::backends::ScheduleStoreFactory,
) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_schedule_store_factory(factory));
self
}
pub fn with_platform_store_factory(
mut self,
factory: crate::backends::PlatformStoreFactory,
) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_platform_store_factory(factory));
self
}
pub fn session_file_system_factory_context(
mut self,
context: SessionFileSystemFactoryContext,
) -> Self {
self.session_file_system_factory_context = context;
self
}
pub fn provider_retry_config(
mut self,
config: everruns_core::llm_retry::LlmRetryConfig,
) -> Self {
self.provider_retry_config = Some(config);
self
}
pub fn provider_stall_timeout(mut self, timeout: std::time::Duration) -> Self {
self.provider_stall_timeout = Some(timeout);
self
}
pub fn harness(mut self, harness: Harness) -> Self {
self.harnesses.push(harness);
self
}
pub fn agent(mut self, agent: Agent) -> Self {
self.agents.push(agent);
self
}
pub fn session(mut self, session: Session) -> Self {
self.sessions.push(session);
self
}
pub fn single_session<F>(mut self, configure: F) -> Self
where
F: FnOnce(SingleSessionBuilder) -> SingleSessionBuilder,
{
let (harness, agent, session, session_id) =
configure(SingleSessionBuilder::default()).build();
self.harnesses.push(harness);
self.agents.push(agent);
self.sessions.push(session);
self.default_session_id = Some(session_id);
self
}
pub fn seed_text_file(
mut self,
session_id: SessionId,
path: impl Into<String>,
content: impl Into<String>,
) -> Self {
self.seeded_files.push((
session_id,
InitialFile {
path: path.into(),
content: content.into(),
encoding: "text".to_string(),
is_readonly: false,
},
));
self
}
pub fn with_plugin_dir(mut self, path: &Path) -> Result<Self> {
let file_set = PluginFileSet::from_dir(path)
.map_err(|e| AgentLoopError::config(format!("plugin directory load failed: {e}")))?;
let compiled = compile_plugin(&file_set)
.map_err(|e| AgentLoopError::config(format!("plugin compilation failed: {e}")))?;
for warning in &compiled.warnings {
tracing::warn!(plugin = %compiled.definition.name, warning = %warning, "plugin compile warning");
}
self.plugin_warnings.extend(compiled.warnings);
let cap_id = plugin_capability_id(&compiled.definition.name);
let hydrated_config = serde_json::to_value(&compiled.definition)
.unwrap_or(serde_json::Value::Object(serde_json::Map::new()));
self.plugin_capability_configs
.push(AgentCapabilityConfig::with_config(cap_id, hydrated_config));
Ok(self)
}
pub fn plugin_capability(&self, name: &str) -> Option<AgentCapabilityConfig> {
let cap_id = plugin_capability_id(name);
self.plugin_capability_configs
.iter()
.find(|c| c.capability_id() == cap_id)
.cloned()
}
pub async fn build(mut self) -> Result<InProcessRuntime> {
let backends = match self.backends.take() {
Some(backends) => backends,
None => HostBackends::in_memory(),
};
let mut file_store = resolve_session_file_system(
&self.platform_definition,
self.session_file_system_factory_context.clone(),
)
.await?;
if let Some(policy) = self.workspace_policy.take() {
file_store = Arc::new(crate::PolicyFileStore::new(file_store, policy));
}
if let Some(config) = self.llm_sim_config.take() {
let driver = LlmSimDriver::new(config);
self.platform_definition
.driver_registry_mut()
.replace_provider(everruns_core::Provider::new("llmsim", driver));
if self.model_spec.is_none() {
self.model_spec = Some(everruns_core::ModelSpec::on("llmsim", "llmsim-model"));
}
}
for provider in self.providers {
self.platform_definition
.driver_registry_mut()
.register_provider(provider)?;
}
let model_spec = self.model_spec.ok_or_else(|| {
AgentLoopError::config(
"in-process runtime requires a default model; call \
InProcessRuntimeBuilder::default_model(...) or \
InProcessRuntimeBuilder::llm_sim(...)",
)
})?;
let legacy_config = self.legacy_provider_config.take();
let provider_type = legacy_config
.as_ref()
.map(|config| config.provider_type.clone())
.unwrap_or_else(|| DriverId::external(model_spec.provider.as_str()));
if let Some(config) = legacy_config
&& self
.platform_definition
.driver_registry()
.provider(&model_spec.provider)
.is_none()
{
let driver = self
.platform_definition
.driver_registry()
.create_chat_driver(&config)?;
self.platform_definition
.driver_registry_mut()
.register_provider(everruns_core::Provider::new(
model_spec.provider.clone(),
driver,
))?;
}
let provider_metadata = if provider_type.as_str() == model_spec.provider.as_str() {
None
} else {
Some(everruns_core::ProviderMetadata {
extra: Some(serde_json::json!({
"provider_id": model_spec.provider.as_str(),
})),
..Default::default()
})
};
let default_model = ResolvedModel {
model: model_spec.model,
provider_type,
api_key: None,
base_url: None,
provider_metadata,
};
backends
.provider_store
.set_default_model(default_model)
.await?;
for harness in &mut self.harnesses {
hydrate_plugin_refs(&mut harness.capabilities, &self.plugin_capability_configs);
}
for agent in &mut self.agents {
hydrate_plugin_refs(&mut agent.capabilities, &self.plugin_capability_configs);
}
for session in &mut self.sessions {
hydrate_plugin_refs(&mut session.capabilities, &self.plugin_capability_configs);
}
for harness in &self.harnesses {
backends.harness_store.add_harness(harness.clone()).await?;
}
for agent in &self.agents {
backends.agent_store.add_agent(agent.clone()).await?;
}
for session in &self.sessions {
backends.session_store.add_session(session.clone()).await?;
}
for session in &self.sessions {
seed_runtime_initial_files(
backends.harness_store.as_ref(),
backends.agent_store.as_ref(),
file_store.as_ref(),
session,
)
.await?;
}
for (session_id, file) in &self.seeded_files {
file_store.seed_initial_file(*session_id, file).await?;
}
let event_log = backends.event_log.clone();
let event_history = Arc::new(EventHistory::new(event_log.clone()));
let event_emitter = Arc::new(HostEventEmitter::new(
event_log.clone(),
backends.event_sink.clone(),
));
let (session_task_registry, session_wake_queue) = match backends.session_task_registry {
Some(inner) => {
let wake_queue = Arc::new(everruns_core::SessionWakeQueue::new());
let observing = everruns_core::ObservingTaskRegistry::new(inner)
.with_observer(wake_queue.clone());
let wrapped: Arc<dyn everruns_core::session_task::SessionTaskRegistry> =
Arc::new(observing);
(Some(wrapped), Some(wake_queue))
}
None => (None, None),
};
let seeded_session_ids = self.sessions.iter().map(|session| session.id).collect();
let runtime = InProcessRuntime {
platform_definition: Arc::new(self.platform_definition),
harness_store: backends.harness_store,
agent_store: backends.agent_store,
session_store: backends.session_store,
default_session_id: self.default_session_id,
seeded_session_ids,
event_log,
event_history,
compaction_checkpoint_store: backends.compaction_checkpoint_store,
provider_store: backends.provider_store,
event_emitter,
file_store,
storage_store: backends.storage_store,
connection_resolver: backends.connection_resolver,
session_task_registry,
session_wake_queue,
schedule_store_factory: backends.schedule_store_factory,
platform_store_factory: backends.platform_store_factory,
mcp_auth_provider: self
.mcp_auth_provider
.unwrap_or_else(|| Arc::new(everruns_mcp::NoAuthProvider)),
provider_retry_config: self.provider_retry_config,
provider_stall_timeout: self.provider_stall_timeout,
mcp_discovery_cache: Arc::new(crate::mcp_cache::McpDiscoveryCache::new()),
plugin_warnings: self.plugin_warnings,
};
for session in &self.sessions {
runtime.ensure_session_started(session).await?;
}
Ok(runtime)
}
}
async fn resolve_session_file_system(
platform_definition: &PlatformDefinition,
file_system_factory_context: SessionFileSystemFactoryContext,
) -> Result<Arc<dyn SessionFileSystem>> {
let file_system_factory = platform_definition.session_file_system_factory();
if file_system_factory.is_disabled() {
Ok(Arc::new(InMemorySessionFileStore::new()))
} else {
Ok(file_system_factory
.create_session_file_system(file_system_factory_context)
.await?)
}
}
#[derive(Clone)]
pub struct InProcessRuntime {
platform_definition: Arc<PlatformDefinition>,
harness_store: Arc<dyn RuntimeHarnessStore>,
agent_store: Arc<dyn RuntimeAgentStore>,
session_store: Arc<dyn RuntimeSessionStore>,
default_session_id: Option<SessionId>,
seeded_session_ids: Vec<SessionId>,
event_log: Arc<dyn EventLog>,
event_history: Arc<EventHistory>,
compaction_checkpoint_store: Arc<dyn everruns_core::CompactionCheckpointStore>,
provider_store: Arc<dyn RuntimeProviderStore>,
event_emitter: Arc<HostEventEmitter>,
file_store: Arc<dyn SessionFileSystem>,
storage_store: Arc<dyn SessionStorageStore>,
connection_resolver: Option<Arc<dyn UserConnectionResolver>>,
session_task_registry: Option<Arc<dyn everruns_core::session_task::SessionTaskRegistry>>,
session_wake_queue: Option<Arc<everruns_core::SessionWakeQueue>>,
schedule_store_factory: Option<crate::backends::ScheduleStoreFactory>,
platform_store_factory: Option<crate::backends::PlatformStoreFactory>,
mcp_auth_provider: Arc<dyn everruns_mcp::McpAuthProvider>,
provider_retry_config: Option<everruns_core::llm_retry::LlmRetryConfig>,
provider_stall_timeout: Option<std::time::Duration>,
mcp_discovery_cache: Arc<crate::mcp_cache::McpDiscoveryCache>,
plugin_warnings: Vec<String>,
}
impl InProcessRuntime {
async fn ensure_session_started(&self, session: &Session) -> Result<()> {
if self
.event_history
.contains_event_type(session.id, everruns_core::events::SESSION_STARTED)
.await
.map_err(|error| AgentLoopError::store(error.to_string()))?
{
return Ok(());
}
self.event_emitter
.emit(EventRequest::new(
session.id,
EventContext::empty(),
SessionStartedData {
harness_id: session.harness_id,
agent_id: session.agent_id,
model_id: session.model_id,
},
))
.await?;
Ok(())
}
pub fn host_event_emitter(&self) -> Arc<HostEventEmitter> {
self.event_emitter.clone()
}
pub fn event_log(&self) -> Arc<dyn EventLog> {
self.event_log.clone()
}
fn mcp_client(&self) -> Arc<everruns_mcp::McpClient> {
Arc::new(everruns_mcp::McpClient::new(
self.platform_definition.egress_service(),
self.mcp_auth_provider.clone(),
))
}
async fn session_mcp_servers(
&self,
session: &Session,
agent: Option<&Agent>,
) -> everruns_core::ScopedMcpServers {
let harness_chain = self
.harness_store
.get_harness_chain(session.harness_id)
.await
.unwrap_or_default();
let resolved = resolve_runtime_capabilities(
&harness_chain,
agent,
session,
self.platform_definition.capability_registry(),
);
let contributed = collect_capability_mcp_servers(
&resolved.resolved_capability_configs,
self.platform_definition.capability_registry(),
);
let explicit = crate::mcp::merge_session_scoped_servers(&harness_chain, agent, session);
everruns_core::merge_scoped_mcp_servers(&contributed, &explicit)
}
pub fn builder() -> InProcessRuntimeBuilder {
InProcessRuntimeBuilder::new()
}
pub fn default_session_id(&self) -> Option<SessionId> {
self.default_session_id
}
pub fn plugin_warnings(&self) -> &[String] {
&self.plugin_warnings
}
pub async fn activate_capability(
&self,
session_id: SessionId,
capability: impl Into<AgentCapabilityConfig>,
) -> Result<CapabilityDelta> {
let mut capability = capability.into();
let registry = self.platform_definition.capability_registry();
let registered = registry.get(capability.capability_id()).ok_or_else(|| {
AgentLoopError::config(format!(
"unknown capability: {}",
capability.capability_id()
))
})?;
if registered.status() != CapabilityStatus::Available {
return Err(AgentLoopError::config(format!(
"capability is not available: {}",
capability.capability_id()
)));
}
registered
.validate_config(&capability.config)
.map_err(|error| {
AgentLoopError::config(format!("invalid capability config: {error}"))
})?;
let canonical_id = registered.id().to_string();
capability.capability_ref = CapabilityId::new(canonical_id.clone());
let context = self.load_context(session_id).await?;
if context
.resolved_capability_configs
.iter()
.any(|config| config.capability_id() == canonical_id)
{
return Ok(CapabilityDelta {
capability_id: canonical_id,
changed: false,
active: true,
surfaces_dirty: false,
});
}
let mut candidate = context.effective_overlay.capabilities;
candidate.push(capability.clone());
resolve_capability_configs(&candidate, registry)
.map_err(|error| AgentLoopError::config(error.to_string()))?;
self.session_store
.upsert_session_capability(session_id, capability)
.await?;
self.mcp_discovery_cache
.invalidate_session(session_id.uuid());
Ok(CapabilityDelta {
capability_id: canonical_id,
changed: true,
active: true,
surfaces_dirty: true,
})
}
pub async fn deactivate_capability(
&self,
session_id: SessionId,
capability_id: &str,
) -> Result<CapabilityDelta> {
let registry = self.platform_definition.capability_registry();
let registered = registry.get(capability_id).ok_or_else(|| {
AgentLoopError::config(format!("unknown capability: {capability_id}"))
})?;
let canonical_id = registered.id().to_string();
let context = self.load_context(session_id).await?;
if !context
.resolved_capability_configs
.iter()
.any(|config| config.capability_id() == canonical_id)
{
return Ok(CapabilityDelta {
capability_id: canonical_id,
changed: false,
active: false,
surfaces_dirty: false,
});
}
let session_capability_id = context
.session
.capabilities
.iter()
.find(|config| {
registry
.get(config.capability_id())
.is_some_and(|capability| capability.id() == canonical_id)
})
.map(|config| config.capability_id().to_string())
.ok_or_else(|| {
AgentLoopError::config(format!(
"capability {canonical_id} is inherited and cannot be deactivated at the session layer"
))
})?;
self.session_store
.remove_session_capability(session_id, &session_capability_id)
.await?;
self.mcp_discovery_cache
.invalidate_session(session_id.uuid());
Ok(CapabilityDelta {
capability_id: canonical_id,
changed: true,
active: false,
surfaces_dirty: true,
})
}
pub async fn run_turn(
&self,
session_id: SessionId,
input: impl Into<InputMessage>,
) -> Result<TurnResult> {
self.run_steerable_turn(
session_id,
AcceptedTurnInput::new(input),
TurnId::new(),
TurnSteering::new(),
)
.await
}
#[doc(hidden)]
pub async fn run_steerable_turn(
&self,
session_id: SessionId,
input: AcceptedTurnInput,
turn_id: TurnId,
steering: TurnSteering,
) -> Result<TurnResult> {
let session = self
.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("session not found: {session_id}")))?;
let input_message = input.into_message();
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::empty(),
InputMessageData::new(input_message.clone()),
))
.await?;
let org_id = in_process_internal_org_id(&session.organization_id);
let mut state = RuntimeTurnState {
org_id,
session_id,
harness_id: session.harness_id,
agent_id: session.agent_id,
input_message_id: input_message.id,
turn_id: None,
previous_response_id: None,
iteration: 1,
request_id: None,
started_at: None,
cumulative_usage: None,
tool_call_count: 0,
llm_call_count: 0,
time_to_first_token_ms: None,
final_message_id: None,
final_answer_preview: None,
};
let base_context = |exec: bool| {
let context = AtomContext::new(session_id, turn_id, input_message.id)
.with_workspace_id(session.workspace_id);
if exec { context.next_exec() } else { context }
};
execute_input_activity(
self,
org_id,
InputAtomInput {
context: base_context(false),
},
)
.await?;
let mut plan = plan_after_process_input(&state, Some(turn_id), Utc::now());
let mut iterations: usize = 0;
let mut tool_calls_count: usize = 0;
let mut last_response = String::new();
let mut pending_prompt_message_ids = Vec::new();
loop {
match plan {
RuntimeTurnPlan::ScheduleReason(next_state) => {
state = next_state;
let mut prompt_message_ids = std::mem::take(&mut pending_prompt_message_ids);
prompt_message_ids.extend(
self.inject_steering_inputs(session_id, steering.drain())
.await?,
);
prompt_message_ids.extend(self.drain_and_inject_wakes(session_id).await?);
if state.iteration == 1 {
prompt_message_ids.insert(0, input_message.id);
}
let reason_result = execute_reason_activity_with_prompt_messages(
self,
org_id,
ReasonInput {
context: base_context(true),
harness_id: session.harness_id,
agent_id: session.agent_id,
org_id,
mcp_tool_definitions: vec![],
previous_response_id: state.previous_response_id.clone(),
iteration: state.iteration,
},
prompt_message_ids,
)
.await?;
iterations += 1;
if !reason_result.text.is_empty() {
last_response = reason_result.text.clone();
}
let pending_wake_count = usize::from(
self.session_wake_queue
.as_ref()
.is_some_and(|q| q.has_pending(session_id)),
);
let can_continue = reason_result.success
&& state.iteration < reason_result.max_iterations as u32;
let pending_steering = if !reason_schedules_act(&state, &reason_result)
&& can_continue
&& pending_wake_count == 0
{
steering.drain_or_close()
} else {
vec![]
};
let pending_steering_count = pending_steering.len();
if !pending_steering.is_empty() {
pending_prompt_message_ids.extend(
self.inject_steering_inputs(session_id, pending_steering)
.await?,
);
}
if !can_continue {
steering.close();
}
let pending_user_message_count = pending_wake_count + pending_steering_count;
let act_scheduling =
crate::turn_strategy::resolve_act_scheduling(self, &state, &reason_result)
.await?;
let (next, effects) = plan_after_reason(
&state,
reason_result,
pending_user_message_count,
Utc::now(),
act_scheduling,
);
crate::turn_strategy::perform_effects(self, org_id, session_id, effects).await;
plan = next;
}
RuntimeTurnPlan::ScheduleAct(act_plan) => {
tool_calls_count += act_plan.input.tool_calls.len();
state = *act_plan.resume_state;
state.previous_response_id = act_plan.previous_response_id;
state.iteration = act_plan.iteration;
state.request_id = act_plan.request_id;
let act_result = execute_act_activity(self, act_plan.input).await?;
let outcome = ActOutcome {
blocked: act_result.blocked,
waiting_for_tool_results: act_result.waiting_for_tool_results,
};
let setup_connection_hint_enabled =
crate::turn_strategy::resolve_setup_connection_hint(
self, org_id, session_id, outcome,
)
.await;
let (next, effects) =
plan_after_act(&state, outcome, setup_connection_hint_enabled);
crate::turn_strategy::perform_effects(self, org_id, session_id, effects).await;
plan = next;
}
RuntimeTurnPlan::Complete { stop_reason, error } => {
steering.close();
return Ok(finish_turn(
turn_id,
stop_reason,
error,
last_response,
iterations,
tool_calls_count,
));
}
RuntimeTurnPlan::WaitForToolResults { .. } => {
steering.close();
self.inject_steering_inputs(session_id, steering.drain())
.await?;
return Ok(finish_turn(
turn_id,
TurnStopReason::EndTurn,
None,
last_response,
iterations,
tool_calls_count,
));
}
}
}
}
pub async fn run_text_turn(
&self,
session_id: SessionId,
text: impl Into<String>,
) -> Result<TurnResult> {
self.run_turn(session_id, InputMessage::user(text)).await
}
async fn inject_steering_inputs(
&self,
session_id: SessionId,
inputs: Vec<AcceptedTurnInput>,
) -> Result<Vec<MessageId>> {
let mut message_ids = Vec::with_capacity(inputs.len());
for input in inputs {
let message = input.into_message();
message_ids.push(message.id);
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::empty(),
InputMessageData::new(message),
))
.await?;
}
Ok(message_ids)
}
#[doc(hidden)]
pub async fn append_accepted_inputs(
&self,
session_id: SessionId,
inputs: Vec<AcceptedTurnInput>,
) -> Result<()> {
self.inject_steering_inputs(session_id, inputs)
.await
.map(|_| ())
}
async fn drain_and_inject_wakes(&self, session_id: SessionId) -> Result<Vec<MessageId>> {
let Some(queue) = &self.session_wake_queue else {
return Ok(vec![]);
};
let wakes = queue.drain(session_id);
if wakes.is_empty() {
return Ok(vec![]);
}
let mut message_ids = Vec::with_capacity(wakes.len());
for wake in wakes {
let message = message_from_input(InputMessage::user(wake.text));
message_ids.push(message.id);
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::empty(),
InputMessageData::new(message),
))
.await?;
}
Ok(message_ids)
}
pub async fn messages(&self, session_id: SessionId) -> Result<Vec<Message>> {
self.event_history.load(session_id).await
}
pub async fn read_file(
&self,
session_id: SessionId,
path: &str,
) -> Result<Option<SessionFile>> {
self.file_store.read_file(session_id, path).await
}
pub async fn load_context(&self, session_id: SessionId) -> Result<AssembledTurnContext> {
let session = self
.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("session not found: {session_id}")))?;
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await?,
None => None,
};
let scoped_servers = self.session_mcp_servers(&session, agent.as_ref()).await;
let mcp_tool_definitions = if scoped_servers.is_empty() {
vec![]
} else {
crate::mcp::discover_tool_definitions(
&self.mcp_discovery_cache,
self.mcp_client(),
session_id.uuid(),
&scoped_servers,
)
.await
};
self.inspect_context_with_ids(
session_id,
session.harness_id,
session.agent_id,
&mcp_tool_definitions,
)
.await
}
pub async fn events(&self) -> Result<Vec<Event>> {
let limit = EventReadLimit::default();
let session_id = self
.default_session_id
.or_else(|| (self.seeded_session_ids.len() == 1).then(|| self.seeded_session_ids[0]))
.ok_or_else(|| {
AgentLoopError::config(
"events() requires exactly one seeded session; use event_log() for bounded per-session replay",
)
})?;
let mut request = EventReadRequest::new(session_id, limit);
let mut events = Vec::new();
loop {
let page = self
.event_log
.read_page(request)
.await
.map_err(|error| AgentLoopError::store(error.to_string()))?;
if events.len().saturating_add(page.events.len())
> crate::events::MAX_EVENT_HISTORY_REPLAY
{
return Err(AgentLoopError::store("event replay bound exceeded"));
}
events.extend(page.events);
let Some(cursor) = page.next_cursor else {
return Ok(events);
};
request = EventReadRequest::from_cursor(cursor, limit);
}
}
pub async fn execute_command(
&self,
session_id: SessionId,
request: everruns_core::command::ExecuteCommandRequest,
) -> Result<everruns_core::command::CommandResult> {
let ctx = self.load_context(session_id).await?;
let registry = self.platform_definition.capability_registry();
let host = everruns_core::command_host::StoreCommandHost::new(
session_id,
self.harness_store.clone(),
self.agent_store.clone(),
self.session_store.clone(),
self.event_history.clone(),
self.provider_store.clone(),
registry.clone(),
self.platform_definition.driver_registry().clone(),
)
.with_file_store(self.file_store.clone())
.with_assembled_context(ctx.clone());
let exec_ctx =
everruns_core::command::CommandExecutionContext::new(session_id, Arc::new(host));
for config in &ctx.resolved_capability_configs {
let Some(capability) = registry.get(config.capability_id()) else {
continue;
};
if capability.commands().iter().any(|c| c.name == request.name) {
return capability.execute_command(&request, &exec_ctx).await;
}
}
Err(AgentLoopError::config(format!(
"no capability declares command /{}",
request.name
)))
}
pub async fn list_commands(
&self,
session_id: SessionId,
) -> Result<Vec<everruns_core::command::CommandDescriptor>> {
let ctx = self.load_context(session_id).await?;
let registry = self.platform_definition.capability_registry();
let mut seen = std::collections::HashSet::new();
let mut commands = Vec::new();
for config in &ctx.resolved_capability_configs {
let Some(capability) = registry.get(config.capability_id()) else {
continue;
};
for command in capability.commands() {
if seen.insert(command.name.clone()) {
commands.push(command);
}
}
}
Ok(commands)
}
async fn inspect_context_with_ids(
&self,
session_id: SessionId,
harness_id: everruns_core::HarnessId,
agent_id: Option<AgentId>,
mcp_tool_definitions: &[everruns_core::ToolDefinition],
) -> Result<AssembledTurnContext> {
inspect_turn_context(
self.harness_store.as_ref(),
self.agent_store.as_ref(),
self.session_store.as_ref(),
self.event_history.as_ref(),
self.provider_store.as_ref(),
self.platform_definition.capability_registry(),
session_id,
harness_id,
agent_id,
mcp_tool_definitions,
Some(self.file_store.clone()),
)
.await
}
}
#[async_trait]
impl RuntimeHostAdapter for InProcessRuntime {
async fn get_agent(&self, _org_id: i64, agent_id: AgentId) -> Result<Option<Agent>> {
self.agent_store.get_agent(agent_id).await
}
async fn get_harness(&self, _org_id: i64, harness_id: HarnessId) -> Result<Option<Harness>> {
let chain = self.harness_store.get_harness_chain(harness_id).await?;
Ok(chain.into_iter().last())
}
async fn set_session_status(
&self,
_org_id: i64,
session_id: SessionId,
_status: SessionStatus,
) -> Result<Session> {
self.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("session not found: {session_id}")))
}
async fn load_turn_context(
&self,
_org_id: i64,
session_id: SessionId,
) -> Result<RuntimeHostTurnContext> {
let mut session = self
.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("session not found: {session_id}")))?;
everruns_core::ard_attachment::apply_session_attachments(
self.storage_store.as_ref(),
&mut session,
)
.await;
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await?,
None => None,
};
let messages = self.event_history.load(session_id).await?;
let model = self.provider_store.get_default_model().await?;
let scoped_servers = self.session_mcp_servers(&session, agent.as_ref()).await;
let mcp_tool_definitions = if scoped_servers.is_empty() {
vec![]
} else {
crate::mcp::discover_tool_definitions(
&self.mcp_discovery_cache,
self.mcp_client(),
session_id.uuid(),
&scoped_servers,
)
.await
};
Ok(RuntimeHostTurnContext {
agent,
session,
messages,
model,
mcp_tool_definitions,
})
}
async fn mcp_executor(
&self,
_org_id: i64,
session_id: SessionId,
) -> Option<Arc<everruns_mcp::McpExecutor>> {
let session = self.session_store.get_session(session_id).await.ok()??;
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await.ok().flatten(),
None => None,
};
let scoped_servers = self.session_mcp_servers(&session, agent.as_ref()).await;
crate::mcp::build_executor(self.mcp_client(), &scoped_servers)
}
fn capability_registry(&self) -> CapabilityRegistry {
self.platform_definition.capability_registry().clone()
}
fn driver_registry(&self) -> DriverRegistry {
self.platform_definition.driver_registry().clone()
}
fn harness_store(&self, _org_id: i64) -> Arc<dyn HarnessStore> {
self.harness_store.clone()
}
fn agent_store(&self, _org_id: i64) -> Arc<dyn AgentStore> {
self.agent_store.clone()
}
fn session_store(&self, _org_id: i64) -> Arc<dyn SessionStore> {
self.session_store.clone()
}
fn session_mutator(&self, _org_id: i64) -> Arc<dyn SessionMutator> {
self.session_store.clone()
}
fn provider_store(&self, _org_id: i64) -> Arc<dyn ProviderStore> {
self.provider_store.clone()
}
fn message_store(&self) -> Arc<dyn MessageRetriever> {
self.event_history.clone()
}
fn compaction_checkpoint_store(
&self,
) -> Option<Arc<dyn everruns_core::CompactionCheckpointStore>> {
Some(self.compaction_checkpoint_store.clone())
}
fn event_emitter(&self) -> Arc<dyn EventEmitter> {
self.event_emitter.clone()
}
fn file_store(&self) -> Arc<dyn SessionFileSystem> {
self.file_store.clone()
}
fn storage_store(&self) -> Option<Arc<dyn SessionStorageStore>> {
Some(self.storage_store.clone())
}
fn connection_resolver(&self) -> Option<Arc<dyn UserConnectionResolver>> {
self.connection_resolver.clone()
}
fn session_task_registry(
&self,
) -> Option<Arc<dyn everruns_core::session_task::SessionTaskRegistry>> {
self.session_task_registry.clone()
}
fn schedule_store(
&self,
org_id: i64,
) -> Option<Arc<dyn everruns_core::traits::SessionScheduleStore>> {
self.schedule_store_factory
.as_ref()
.map(|factory| factory(org_id))
}
fn platform_store(
&self,
org_id: i64,
session_id: SessionId,
) -> Option<Arc<dyn everruns_platform::PlatformStore>> {
self.platform_store_factory
.as_ref()
.map(|factory| factory(org_id, session_id))
}
fn utility_llm_service(&self) -> Option<Arc<dyn everruns_core::UtilityLlmService>> {
Some(self.platform_definition.utility_llm_service())
}
fn egress_service(&self) -> Option<Arc<dyn everruns_core::EgressService>> {
Some(self.platform_definition.egress_service())
}
fn provider_retry_config(&self) -> Option<everruns_core::llm_retry::LlmRetryConfig> {
self.provider_retry_config.clone()
}
fn provider_stall_timeout(&self) -> Option<std::time::Duration> {
self.provider_stall_timeout
}
}
fn effective_overlay(
harness_chain: &[Harness],
agent: Option<&Agent>,
session: &Session,
) -> AgentConfigOverlay {
let harness_layers = harness_chain.iter().map(AgentConfigOverlay::from);
let agent_layers = agent.into_iter().map(AgentConfigOverlay::from);
AgentConfigOverlay::fold(
harness_layers
.chain(agent_layers)
.chain([AgentConfigOverlay::from(session)]),
)
}
fn hydrate_plugin_refs(
capabilities: &mut [AgentCapabilityConfig],
plugin_configs: &[AgentCapabilityConfig],
) {
for cap in capabilities.iter_mut() {
let cap_id = cap.capability_id();
if !everruns_core::is_plugin_capability(cap_id) {
continue;
}
let is_bare = cap.config.is_null()
|| cap
.config
.as_object()
.map(|o| o.is_empty())
.unwrap_or(false);
if !is_bare {
continue;
}
if let Some(hydrated) = plugin_configs.iter().find(|c| c.capability_id() == cap_id) {
cap.config = hydrated.config.clone();
}
}
}
async fn seed_runtime_initial_files(
harness_store: &dyn RuntimeHarnessStore,
agent_store: &dyn RuntimeAgentStore,
file_store: &dyn SessionFileSystem,
session: &Session,
) -> Result<()> {
let harness_chain = harness_store.get_harness_chain(session.harness_id).await?;
if harness_chain.is_empty() {
return Err(AgentLoopError::store(format!(
"harness not found while seeding files: {}",
session.harness_id
)));
}
let agent = match session.agent_id {
Some(agent_id) => Some(
agent_store
.get_agent(agent_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("agent not found: {agent_id}")))?,
),
None => None,
};
let overlay = effective_overlay(&harness_chain, agent.as_ref(), session);
let seed_key = SessionId::from_uuid(session.workspace_id.uuid());
for file in &overlay.initial_files {
file_store.seed_initial_file(seed_key, file).await?;
}
Ok(())
}
fn message_from_input(input: InputMessage) -> Message {
message_from_input_with_id(MessageId::new(), input)
}
fn message_from_input_with_id(message_id: MessageId, input: InputMessage) -> Message {
Message {
id: message_id,
role: input.role,
content: input.content,
phase: None,
thinking: None,
thinking_signature: None,
controls: input.controls,
metadata: input.metadata,
external_actor: None,
created_at: Utc::now(),
}
}
#[cfg(test)]
mod org_id_mapping_tests {
use super::*;
use everruns_core::{DEFAULT_ORG_ID, DEFAULT_ORG_PUBLIC_ID, org_public_id_from_internal};
#[test]
fn default_public_id_maps_to_default_org() {
assert_eq!(
in_process_internal_org_id(DEFAULT_ORG_PUBLIC_ID),
DEFAULT_ORG_ID
);
}
#[test]
fn invalid_public_id_does_not_fall_back_to_default() {
for invalid in [
"",
"not-an-org",
"org_short",
"org_ZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZ",
"ORG_00000000000000000000000000000001",
] {
let mapped = in_process_internal_org_id(invalid);
assert_ne!(mapped, everruns_core::DEFAULT_ORG_ID);
assert!(
mapped >= 2,
"invalid input {invalid:?} should not map to default"
);
}
}
#[test]
fn zero_public_id_does_not_fall_back_to_default() {
let mapped = in_process_internal_org_id("org_00000000000000000000000000000000");
assert_ne!(mapped, everruns_core::DEFAULT_ORG_ID);
assert!(mapped >= 2, "all-zero id should not map to default");
}
#[test]
fn synthetic_public_id_round_trips_with_internal_helper() {
for internal in [1_i64, 2, 42, 1_000_000, i64::MAX - 1, i64::MAX] {
let public = org_public_id_from_internal(internal);
assert_eq!(
in_process_internal_org_id(&public),
internal,
"round-trip failed for internal={internal}"
);
}
}
#[test]
fn distinct_synthetic_ids_map_to_distinct_internal_ids() {
let a = org_public_id_from_internal(7);
let b = org_public_id_from_internal(8);
assert_ne!(a, b);
assert_ne!(
in_process_internal_org_id(&a),
in_process_internal_org_id(&b)
);
}
#[test]
fn high_entropy_uuid_style_id_hashes_into_reserved_range() {
let high = "org_80000000000000000000000000000000";
let mapped = in_process_internal_org_id(high);
assert!(mapped >= 2, "mapped id {mapped} must be >= 2");
assert_ne!(mapped, DEFAULT_ORG_ID);
assert_eq!(mapped, in_process_internal_org_id(high));
}
#[test]
fn high_entropy_ids_are_isolated_from_each_other() {
let a = in_process_internal_org_id("org_80000000000000000000000000000001");
let b = in_process_internal_org_id("org_80000000000000000000000000000002");
assert_ne!(a, b);
assert_ne!(a, DEFAULT_ORG_ID);
assert_ne!(b, DEFAULT_ORG_ID);
}
#[test]
fn hash_uses_stable_sha256_truncation() {
let mapped = in_process_internal_org_id("org_80000000000000000000000000000000");
let expected = {
let digest = sha2::Sha256::digest(b"org_80000000000000000000000000000000");
let mut buf = [0u8; 8];
buf.copy_from_slice(&digest[..8]);
let raw = u64::from_be_bytes(buf);
((raw % ((i64::MAX - 1) as u64)) as i64) + 2
};
assert_eq!(mapped, expected);
}
#[test]
fn oversize_input_is_bounded_and_does_not_collide_silently() {
let oversize = "x".repeat(super::HASH_INPUT_CAP_BYTES * 4);
let mapped = in_process_internal_org_id(&oversize);
assert!(mapped >= 2);
assert_ne!(mapped, DEFAULT_ORG_ID);
}
}