use crate::HostComposition;
use crate::InProcessExecution;
use crate::SessionFileSystemFactoryContext;
use crate::SessionMutator;
use crate::backends::{
HostBackends, RuntimeAgentStore, RuntimeHarnessStore, RuntimeProviderStore, RuntimeSessionStore,
};
use crate::builders::SingleSessionBuilder;
use crate::events::{EventHistory, EventLog, EventReadLimit, EventReadRequest, HostEventEmitter};
use crate::host::{
ResolvedTurnInputs, RuntimeHostAdapter, execute_act_activity, execute_input_activity,
execute_reason_activity_with_prompt_messages, run_user_prompt_submit_for_message,
};
use crate::in_memory::{InMemorySessionFileStore, InMemorySessionFileSystemFactory};
use async_trait::async_trait;
use chrono::Utc;
use everruns_capability::plugin_capability_id;
use everruns_core::ExecutionContext;
use everruns_core::agent_definition::AgentDefinition;
#[cfg(feature = "mcp")]
use everruns_core::capabilities::collect_capability_mcp_servers;
use everruns_core::capabilities::{
Capability, CapabilityRegistry, CapabilityStatus, resolve_capability_configs,
};
use everruns_core::config_layer::AgentConfigOverlay;
use everruns_core::events::{
Event, EventContext, EventRequest, InputMessageData, SessionStartedData,
};
use everruns_core::harness_definition::HarnessDefinition;
use everruns_core::lifecycle_hooks::UserPromptDecision;
use everruns_core::message::{ContentPart, Message};
use everruns_core::plugins::{PluginFileSet, compile_plugin};
#[cfg(feature = "mcp")]
use everruns_core::resolve_runtime_capabilities;
use everruns_core::runtime_context::AssembledTurnContext;
use everruns_core::session::{ExecutionSession, SessionExecutionState};
use everruns_core::session_file::{InitialFile, SessionFile};
use everruns_core::turn::TurnStopReason;
use everruns_core::{
InputMessage, MessageRetriever, ResolvedExecutionSnapshot, session_files::SessionFileSystem,
};
use everruns_core::{
connection_services::UserConnectionResolver, event_emitter::EventEmitter,
execution_loading::AgentStore, execution_loading::HarnessStore,
execution_loading::SessionStore, provider_resolution::ProviderStore,
session_services::SessionStorageStore,
};
use everruns_engine::{
ActOutcome, ActivityOutcome, Execution, HostFacts, TurnPlan, TurnState, reason_schedules_act,
};
use everruns_engine::{InputAtomInput, ReasonInput};
use everruns_provider::driver_registry::DriverRegistry;
use everruns_provider::error::{AgentLoopError, Result};
use everruns_provider::typed_id::{AgentId, MessageId, OrgId, SessionId, TurnId};
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_provider::typed_id::TurnId,
}
#[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)
}
}
#[derive(Clone, Debug)]
pub struct TurnSteering {
state: Arc<Mutex<TurnSteeringState>>,
}
#[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_provider::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 {
host_composition: HostComposition,
default_provider: Option<(everruns_provider::runtime_provider::Provider, String)>,
replacement_providers: Vec<everruns_provider::runtime_provider::Provider>,
providers: Vec<everruns_provider::runtime_provider::Provider>,
model_spec: Option<everruns_provider::model_spec::ModelSpec>,
backends: Option<HostBackends>,
workspace_policy: Option<everruns_core::WorkspacePolicy>,
session_file_system_factory_context: SessionFileSystemFactoryContext,
harnesses: Vec<crate::builders::SeededHarness>,
agents: Vec<AgentDefinition>,
sessions: Vec<ExecutionSession>,
default_session_id: Option<SessionId>,
seeded_files: Vec<(SessionId, InitialFile)>,
#[cfg(feature = "mcp")]
mcp_auth_provider: Option<Arc<dyn everruns_mcp::McpAuthProvider>>,
provider_retry_config: Option<everruns_provider::llm_retry::LlmRetryConfig>,
provider_stall_timeout: Option<std::time::Duration>,
plugin_capability_configs: Vec<everruns_capability::CapabilityRef>,
plugin_warnings: Vec<String>,
}
impl Default for InProcessRuntimeBuilder {
fn default() -> Self {
Self::new()
}
}
impl InProcessRuntimeBuilder {
pub fn new() -> Self {
Self {
host_composition: HostComposition::builder()
.capability_registry(crate::runtime_capability_registry())
.driver_registry(DriverRegistry::new())
.egress_service(crate::runtime_egress_service())
.session_file_system_factory(Arc::new(InMemorySessionFileSystemFactory))
.build(),
default_provider: None,
replacement_providers: Vec::new(),
providers: Vec::new(),
model_spec: 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(),
#[cfg(feature = "mcp")]
mcp_auth_provider: None,
provider_retry_config: None,
provider_stall_timeout: None,
plugin_capability_configs: Vec::new(),
plugin_warnings: Vec::new(),
}
}
#[cfg(feature = "mcp")]
pub fn mcp_auth_provider(mut self, provider: Arc<dyn everruns_mcp::McpAuthProvider>) -> Self {
self.mcp_auth_provider = Some(provider);
self
}
pub fn host_composition(mut self, host_composition: HostComposition) -> Self {
self.host_composition = host_composition;
self
}
pub fn capability<C: Capability + 'static>(self, capability: C) -> Self {
self.host_composition
.register_capability_overriding(Arc::new(capability));
self
}
pub fn driver_registry(mut self, driver_registry: DriverRegistry) -> Self {
*self.host_composition.driver_registry_mut() = driver_registry;
self
}
pub fn provider(mut self, provider: everruns_provider::runtime_provider::Provider) -> Self {
self.providers.push(provider);
self
}
pub fn replace_provider(
mut self,
provider: everruns_provider::runtime_provider::Provider,
) -> Self {
self.replacement_providers.push(provider);
self
}
pub fn provider_with_default_model(
mut self,
provider: everruns_provider::runtime_provider::Provider,
model_id: impl Into<String>,
) -> Self {
self.default_provider = Some((provider, model_id.into()));
self
}
pub fn default_model(mut self, model: everruns_provider::model_spec::ModelSpec) -> Self {
self.model_spec = Some(model);
self
}
pub fn model_spec(mut self, model: everruns_provider::model_spec::ModelSpec) -> Self {
self.model_spec = Some(model);
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_tool_context_extensions_factory(
mut self,
factory: crate::ToolContextExtensionsFactory,
) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_tool_context_extensions_factory(factory));
self
}
pub fn with_subagent_delegate_factory(
mut self,
factory: crate::SubagentDelegateFactory,
) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_subagent_delegate_factory(factory));
self
}
pub fn with_tool_augmentor(mut self, augmentor: Arc<dyn crate::HostToolAugmentor>) -> Self {
let backends = self.backends.take().unwrap_or_else(HostBackends::in_memory);
self.backends = Some(backends.with_tool_augmentor(augmentor));
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_provider::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: crate::builders::SeededHarness) -> Self {
self.harnesses.push(harness);
self
}
pub fn agent(mut self, agent: AgentDefinition) -> Self {
self.agents.push(agent);
self
}
pub fn session(mut self, session: ExecutionSession) -> 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(everruns_capability::CapabilityRef::with_config(
cap_id,
hydrated_config,
));
Ok(self)
}
pub fn plugin_capability(&self, name: &str) -> Option<everruns_capability::CapabilityRef> {
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.host_composition,
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((provider, model_id)) = self.default_provider.take() {
let provider_key = provider.id().clone();
self.host_composition
.driver_registry_mut()
.replace_provider(provider);
if self.model_spec.is_none() {
self.model_spec = Some(everruns_provider::model_spec::ModelSpec::on(
provider_key,
model_id,
));
}
}
for provider in self.replacement_providers {
self.host_composition
.driver_registry_mut()
.replace_provider(provider);
}
for provider in self.providers {
self.host_composition
.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::model_spec(...) or \
InProcessRuntimeBuilder::provider_with_default_model(...) \
(e.g. the everruns-llmsim `.llm_sim_as_default(...)` extension)",
)
})?;
backends
.provider_store
.set_default_model_spec(model_spec)
.await?;
for harness in &mut self.harnesses {
hydrate_plugin_refs(
&mut harness.definition.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.id, harness.definition.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 {
host_composition: Arc::new(self.host_composition),
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,
native_async_store: backends.native_async_store,
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,
tool_context_extensions_factory: backends.tool_context_extensions_factory,
subagent_delegate_factory: backends.subagent_delegate_factory,
tool_augmentor: backends.tool_augmentor,
#[cfg(feature = "mcp")]
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,
#[cfg(feature = "mcp")]
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(
host_composition: &HostComposition,
file_system_factory_context: SessionFileSystemFactoryContext,
) -> Result<Arc<dyn SessionFileSystem>> {
let file_system_factory = host_composition.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 {
host_composition: Arc<HostComposition>,
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>,
native_async_store: Option<Arc<dyn everruns_core::native_async_store::NativeAsyncStore>>,
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>,
tool_context_extensions_factory: Option<crate::ToolContextExtensionsFactory>,
subagent_delegate_factory: Option<crate::SubagentDelegateFactory>,
tool_augmentor: Option<Arc<dyn crate::HostToolAugmentor>>,
#[cfg(feature = "mcp")]
mcp_auth_provider: Arc<dyn everruns_mcp::McpAuthProvider>,
provider_retry_config: Option<everruns_provider::llm_retry::LlmRetryConfig>,
provider_stall_timeout: Option<std::time::Duration>,
#[cfg(feature = "mcp")]
mcp_discovery_cache: Arc<crate::mcp_cache::McpDiscoveryCache>,
plugin_warnings: Vec<String>,
}
impl InProcessRuntime {
async fn ensure_session_started(&self, session: &ExecutionSession) -> 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()
}
#[cfg(feature = "mcp")]
fn mcp_client(&self) -> Arc<everruns_mcp::McpClient> {
Arc::new(everruns_mcp::McpClient::with_url_elicitation(
self.host_composition.egress_service(),
self.mcp_auth_provider.clone(),
Arc::new(everruns_mcp::RelayUrlElicitations),
))
}
#[cfg(feature = "mcp")]
async fn session_mcp_servers(
&self,
session: &ExecutionSession,
agent: Option<&AgentDefinition>,
) -> everruns_core::ScopedMcpServers {
let harness = self
.harness_store
.get_harness(session.harness_id)
.await
.ok()
.flatten()
.unwrap_or_default();
let resolved = resolve_runtime_capabilities(
&harness,
agent,
session,
&self.host_composition.capability_registry(),
);
let contributed = collect_capability_mcp_servers(
&resolved.resolved_capability_configs,
&self.host_composition.capability_registry(),
);
let explicit = crate::mcp::merge_session_scoped_servers(&harness, 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 fn register_capability(&self, capability: Arc<dyn Capability>) -> Result<()> {
self.host_composition
.register_capability(capability)
.map_err(|error| AgentLoopError::config(format!("cannot register capability: {error}")))
}
pub fn is_capability_registered(&self, capability_id: &str) -> bool {
self.host_composition
.is_capability_registered(capability_id)
}
pub async fn activate_capability(
&self,
session_id: SessionId,
capability: impl Into<everruns_capability::CapabilityRef>,
) -> Result<CapabilityDelta> {
let mut capability = capability.into();
let registry = self.host_composition.capability_registry();
let registered = registry.get(capability.capability_id()).ok_or_else(|| {
AgentLoopError::config(format!(
"unknown capability: {}",
capability.capability_id()
))
})?;
if !registered.status().is_active() {
let reason = if registered.status() == CapabilityStatus::Retired {
"capability has been removed"
} else {
"capability is not available"
};
return Err(AgentLoopError::config(format!(
"{reason}: {}",
capability.capability_id()
)));
}
registered
.validate_config(capability.config_value())
.map_err(|error| {
AgentLoopError::config(format!("invalid capability config: {error}"))
})?;
let canonical_id = registered.id().to_string();
capability.set_id(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.snapshot.capabilities;
candidate.push(capability.clone());
resolve_capability_configs(&candidate, ®istry)
.map_err(|error| AgentLoopError::config(error.to_string()))?;
self.session_store
.upsert_session_capability(session_id, capability)
.await?;
#[cfg(feature = "mcp")]
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.host_composition.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 = self
.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::session_not_found(session_id))?;
let session_capability_id = 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?;
#[cfg(feature = "mcp")]
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
}
pub async fn run_steerable_turn(
&self,
session_id: SessionId,
input: AcceptedTurnInput,
turn_id: TurnId,
steering: TurnSteering,
) -> Result<TurnResult> {
let snapshot = self.resolved_execution_snapshot(session_id).await?;
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(&snapshot.organization_id);
let state = TurnState {
org_id,
session_id,
harness_id: snapshot.harness_id,
agent_id: snapshot.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 = ExecutionContext::new(session_id, turn_id, input_message.id)
.with_workspace_id(snapshot.workspace_id);
if exec { context.next_exec() } else { context }
};
execute_input_activity(
self,
org_id,
InputAtomInput {
context: base_context(false),
},
)
.await?;
let mut execution = InProcessExecution::new(state);
let transition = execution.advance(
ActivityOutcome::ProcessInput {
turn_id: Some(turn_id),
},
0,
Utc::now(),
HostFacts::default(),
);
crate::turn_strategy::perform_effects(self, org_id, session_id, transition.effects).await?;
let mut plan = transition.plan;
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 {
TurnPlan::ScheduleReason(_) => {
let state = execution.state().clone();
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: snapshot.harness_id,
agent_id: snapshot.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 transition = execution.advance(
ActivityOutcome::Reason(Box::new(reason_result)),
pending_user_message_count,
Utc::now(),
HostFacts {
act_scheduling,
..HostFacts::default()
},
);
crate::turn_strategy::perform_effects(
self,
org_id,
session_id,
transition.effects,
)
.await?;
plan = transition.plan;
}
TurnPlan::ScheduleAct(act_plan) => {
tool_calls_count += act_plan.input.tool_calls.len();
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,
waiting_for_url_elicitation: act_result.waiting_for_url_elicitation,
};
let hints = crate::turn_strategy::resolve_pause_hints(
self, org_id, session_id, outcome,
)
.await;
let transition = execution.advance(
ActivityOutcome::Act(outcome),
0,
Utc::now(),
HostFacts {
setup_connection_hint_enabled: hints.setup_connection,
url_elicitation_hint_enabled: hints.url_elicitation,
..HostFacts::default()
},
);
crate::turn_strategy::perform_effects(
self,
org_id,
session_id,
transition.effects,
)
.await?;
plan = transition.plan;
}
TurnPlan::Complete { stop_reason, error } => {
steering.close();
return Ok(finish_turn(
turn_id,
stop_reason,
error,
last_response,
iterations,
tool_calls_count,
));
}
TurnPlan::WaitForToolResults { .. } => {
steering.close();
self.append_accepted_inputs(session_id, turn_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)
}
pub async fn append_accepted_inputs(
&self,
session_id: SessionId,
turn_id: TurnId,
inputs: Vec<AcceptedTurnInput>,
) -> Result<()> {
if inputs.is_empty() {
return Ok(());
}
let snapshot = self.resolved_execution_snapshot(session_id).await?;
let org_id = in_process_internal_org_id(&snapshot.organization_id);
for input in inputs {
let mut message = input.into_message();
let hook_input = ReasonInput {
context: ExecutionContext::new(session_id, turn_id, message.id)
.with_workspace_id(snapshot.workspace_id),
harness_id: snapshot.harness_id,
agent_id: snapshot.agent_id,
org_id,
mcp_tool_definitions: vec![],
previous_response_id: None,
iteration: 1,
};
if let Some(result) = run_user_prompt_submit_for_message(
self,
org_id,
&hook_input,
message.content_to_llm_string(),
)
.await?
{
let enforced = match result.decision {
UserPromptDecision::Continue { message }
if message != result.original_message =>
{
Some((message, false))
}
UserPromptDecision::Block {
reason,
user_message,
} => Some((user_message.unwrap_or(reason), true)),
UserPromptDecision::Continue { .. } => None,
};
if let Some((safe_text, blocked)) = enforced {
if blocked {
message.content.clear();
} else {
message
.content
.retain(|part| !matches!(part, ContentPart::Text(_)));
}
message.content.insert(0, ContentPart::text(safe_text));
}
}
self.event_emitter
.emit(EventRequest::new(
session_id,
EventContext::empty(),
InputMessageData::new(message),
))
.await?;
}
Ok(())
}
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}")))?;
#[cfg(feature = "mcp")]
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await?,
None => None,
};
#[cfg(feature = "mcp")]
let scoped_servers = self.session_mcp_servers(&session, agent.as_ref()).await;
#[cfg(feature = "mcp")]
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
};
#[cfg(not(feature = "mcp"))]
let mcp_tool_definitions = vec![];
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.host_composition.capability_registry();
let host = crate::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.host_composition.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.host_composition.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 attached_session(&self, session_id: SessionId) -> Result<ExecutionSession> {
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;
Ok(session)
}
async fn resolved_execution_snapshot(
&self,
session_id: SessionId,
) -> Result<ResolvedExecutionSnapshot> {
let session = self.attached_session(session_id).await?;
crate::load_execution_snapshot_for_session(
self.harness_store.as_ref(),
self.agent_store.as_ref(),
&session,
)
.await
}
async fn inspect_context_with_ids(
&self,
session_id: SessionId,
harness_id: everruns_provider::typed_id::HarnessId,
agent_id: Option<AgentId>,
mcp_tool_definitions: &[everruns_provider::tool_types::ToolDefinition],
) -> Result<AssembledTurnContext> {
crate::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.host_composition.capability_registry(),
self.host_composition.driver_registry(),
session_id,
harness_id,
agent_id,
mcp_tool_definitions,
Some(self.file_store.clone()),
)
.await
}
}
#[async_trait]
impl RuntimeHostAdapter for InProcessRuntime {
async fn set_session_status(
&self,
_org_id: i64,
session_id: SessionId,
_status: SessionExecutionState,
) -> Result<()> {
self.session_store
.get_session(session_id)
.await?
.ok_or_else(|| AgentLoopError::store(format!("session not found: {session_id}")))?;
Ok(())
}
async fn load_resolved_turn(
&self,
_org_id: i64,
session_id: SessionId,
) -> Result<ResolvedTurnInputs> {
let session = self.attached_session(session_id).await?;
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await?,
None => None,
};
let harness = self
.harness_store
.get_harness(session.harness_id)
.await?
.ok_or_else(|| AgentLoopError::harness_not_found(session.harness_id))?;
let snapshot = ResolvedExecutionSnapshot::project(&harness, agent.as_ref(), &session)?;
let messages = self.event_history.load(session_id).await?;
#[cfg(feature = "mcp")]
let scoped_servers = self.session_mcp_servers(&session, agent.as_ref()).await;
#[cfg(feature = "mcp")]
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
};
#[cfg(not(feature = "mcp"))]
let mcp_tool_definitions = vec![];
Ok(ResolvedTurnInputs {
snapshot,
messages,
mcp_tool_definitions,
})
}
#[cfg(feature = "mcp")]
async fn mcp_executor(
&self,
_org_id: i64,
session_id: SessionId,
) -> Option<Arc<dyn everruns_core::McpToolInvoker>> {
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)
.map(|executor| executor as Arc<dyn everruns_core::McpToolInvoker>)
}
fn capability_registry(&self) -> CapabilityRegistry {
(*self.host_composition.capability_registry()).clone()
}
fn driver_registry(&self) -> DriverRegistry {
self.host_composition.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 native_async_store(
&self,
) -> Option<Arc<dyn everruns_core::native_async_store::NativeAsyncStore>> {
self.native_async_store.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::session_services::SessionScheduleStore>> {
self.schedule_store_factory
.as_ref()
.map(|factory| factory(org_id))
}
fn tool_context_extensions(
&self,
org_id: i64,
session_id: SessionId,
) -> everruns_core::tool_context::ToolContextExtensions {
self.tool_context_extensions_factory
.as_ref()
.map(|factory| factory(org_id, session_id))
.unwrap_or_default()
}
fn subagent_delegate(
&self,
org_id: i64,
session_id: SessionId,
) -> Option<Arc<dyn everruns_core::subagent_delegation::SubagentSessionDelegate>> {
self.subagent_delegate_factory
.as_ref()
.map(|factory| factory(org_id, session_id))
}
fn tool_augmentor(&self) -> Option<Arc<dyn crate::HostToolAugmentor>> {
self.tool_augmentor.clone()
}
fn utility_llm_service(&self) -> Option<Arc<dyn everruns_core::UtilityLlmService>> {
Some(self.host_composition.utility_llm_service())
}
fn egress_service(&self) -> Option<Arc<dyn everruns_core::EgressService>> {
Some(self.host_composition.egress_service())
}
fn provider_retry_config(&self) -> Option<everruns_provider::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: &HarnessDefinition,
agent: Option<&AgentDefinition>,
session: &ExecutionSession,
) -> AgentConfigOverlay {
let agent_layers = agent.into_iter().map(AgentConfigOverlay::from);
AgentConfigOverlay::fold(
[AgentConfigOverlay::from(harness)]
.into_iter()
.chain(agent_layers)
.chain([AgentConfigOverlay::from(session)]),
)
}
fn hydrate_plugin_refs(
capabilities: &mut [everruns_capability::CapabilityRef],
plugin_configs: &[everruns_capability::CapabilityRef],
) {
for cap in capabilities.iter_mut() {
let cap_id = cap.capability_id();
if !everruns_capability::is_plugin_capability(cap_id) {
continue;
}
let is_bare = cap.config_value().is_null()
|| cap
.config_value()
.clone()
.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.set_config(hydrated.config_value().clone());
}
}
}
async fn seed_runtime_initial_files(
harness_store: &dyn RuntimeHarnessStore,
agent_store: &dyn RuntimeAgentStore,
file_store: &dyn SessionFileSystem,
session: &ExecutionSession,
) -> Result<()> {
let harness = harness_store
.get_harness(session.harness_id)
.await?
.ok_or_else(|| {
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, 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,
phase_source: 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);
}
}