use std::collections::{HashMap, HashSet};
use std::fmt;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use everruns_core::InitialFile;
use everruns_host::{
AgentBuilder as RuntimeAgentBuilder, Environment, EnvironmentBindingError,
EnvironmentBindingStore, EventLogError, EventSink, HarnessBuilder, HostBackends,
InProcessRuntime, InProcessRuntimeBuilder, SessionBuilder, WorkspaceProvider,
WorkspaceProviderId,
};
use everruns_llmsim::{LlmSimConfig, LlmSimDriver};
use everruns_provider::model_spec::ModelSpec;
use everruns_provider::runtime_provider::Provider;
use everruns_provider::typed_id::SessionId;
use crate::tool::{FunctionTool, IntoTool, Tool, validate_tool_name, validate_tool_schema};
#[derive(Clone)]
pub struct Model {
id: String,
bundled_provider: Option<Provider>,
}
impl Model {
pub fn new(id: impl Into<String>, provider: impl Into<Provider>) -> Self {
Self::bundled(id, provider.into())
}
pub fn simulated(response: impl Into<String>) -> Self {
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(LlmSimConfig::fixed(response))),
)
}
pub fn simulated_with_config(config: LlmSimConfig) -> Self {
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(config)),
)
}
fn bundled(id: impl Into<String>, provider: Provider) -> Self {
Self {
id: id.into(),
bundled_provider: Some(provider),
}
}
}
impl From<&str> for Model {
fn from(id: &str) -> Self {
Self {
id: id.to_string(),
bundled_provider: None,
}
}
}
impl From<String> for Model {
fn from(id: String) -> Self {
Self {
id,
bundled_provider: None,
}
}
}
#[cfg(test)]
impl Model {
pub(crate) fn simulated_capturing(
response: impl Into<String>,
capture: std::sync::Arc<
std::sync::Mutex<Vec<Vec<everruns_provider::driver_registry::LlmMessage>>>,
>,
) -> Self {
let mut sim = LlmSimConfig::fixed(response);
sim.message_capture = Some(capture);
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(sim)),
)
}
pub(crate) fn simulated_delayed(
response: impl Into<String>,
delay: std::time::Duration,
) -> Self {
let sim = LlmSimConfig::fixed(response).with_response_delay(delay);
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(sim)),
)
}
pub(crate) fn simulated_error(message: impl Into<String>) -> Self {
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(LlmSimConfig::error(message))),
)
}
pub(crate) fn simulated_scripted(
response: impl Into<String>,
tool_call_sequence: Vec<Vec<everruns_provider::tool_types::ToolCall>>,
) -> Self {
let sim = LlmSimConfig::fixed(response).with_tool_call_sequence(tool_call_sequence);
Self::bundled(
"llmsim-model",
Provider::new("llmsim", LlmSimDriver::new(sim)),
)
}
}
impl fmt::Debug for Model {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Model")
.field("id", &self.id)
.field("bundled_provider", &self.bundled_provider)
.finish()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum BuildError {
BlankInstructions,
MissingModel,
InvalidToolName {
name: String,
reason: String,
},
InvalidToolSchema {
name: String,
reason: String,
},
DuplicateTool {
name: String,
},
InvalidCapability {
id: String,
reason: String,
},
DuplicateCapability {
id: String,
},
MissingProvider,
MultipleProviders {
registered: Vec<String>,
},
InvalidMcpServer {
reason: String,
},
DuplicateWorkspaceProvider {
id: String,
},
}
impl fmt::Display for BuildError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
BuildError::BlankInstructions => {
write!(f, "agent instructions must not be blank")
}
BuildError::MissingModel => write!(f, "agent requires a model"),
BuildError::InvalidToolName { name, reason } => {
write!(f, "invalid tool name {name:?}: {reason}")
}
BuildError::InvalidToolSchema { name, reason } => {
write!(f, "invalid JSON schema for tool {name:?}: {reason}")
}
BuildError::DuplicateTool { name } => {
write!(f, "duplicate tool name {name:?}")
}
BuildError::InvalidCapability { id, reason } => {
write!(f, "invalid capability {id:?}: {reason}")
}
BuildError::DuplicateCapability { id } => {
write!(f, "duplicate capability id {id:?}")
}
BuildError::MissingProvider => write!(f, "agent requires a provider for a live model"),
BuildError::MultipleProviders { registered } => write!(
f,
"agent accepts one provider; registered providers: [{}]",
registered.join(", ")
),
BuildError::InvalidMcpServer { reason } => {
write!(f, "invalid MCP server configuration: {reason}")
}
BuildError::DuplicateWorkspaceProvider { id } => {
write!(f, "duplicate workspace provider id {id:?}")
}
}
}
}
impl std::error::Error for BuildError {}
#[derive(Clone)]
pub struct Agent {
name: String,
instructions: String,
model: ModelSpec,
provider: Provider,
capabilities: Vec<everruns_capability::CapabilityRef>,
capability_implementations: Vec<CapabilityImplementation>,
initial_files: Vec<InitialFile>,
max_iterations: Option<usize>,
parallel_tool_calls: Option<bool>,
workspace_root: Option<PathBuf>,
default_workspace: crate::default_workspace::DefaultWorkspace,
workspace_policy: everruns_core::WorkspacePolicy,
mcp_servers: everruns_core::ScopedMcpServers,
plugin_warnings: Vec<String>,
#[cfg(feature = "local")]
local: Option<crate::LocalConfig>,
lifecycle_hooks: crate::hooks::LifecycleHooks,
state: Arc<AgentState>,
}
#[derive(Clone)]
enum CapabilityImplementation {
Function(FunctionTool),
#[cfg(feature = "capabilities")]
Definition(crate::capability::Definition),
}
impl CapabilityImplementation {
fn register(&self, builder: InProcessRuntimeBuilder) -> InProcessRuntimeBuilder {
match self {
Self::Function(tool) => builder.capability(tool.clone().into_capability()),
#[cfg(feature = "capabilities")]
Self::Definition(definition) => {
builder.capability(crate::capability::runtime_adapter(definition))
}
}
}
}
impl fmt::Debug for CapabilityImplementation {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Function(tool) => formatter
.debug_tuple("Function")
.field(&tool.name())
.finish(),
#[cfg(feature = "capabilities")]
Self::Definition(definition) => formatter
.debug_tuple("Definition")
.field(&definition.id())
.finish(),
}
}
}
struct AgentState {
workspace_providers: Mutex<HashMap<WorkspaceProviderId, Arc<dyn WorkspaceProvider>>>,
}
impl AgentState {
fn new(workspace_providers: HashMap<WorkspaceProviderId, Arc<dyn WorkspaceProvider>>) -> Self {
Self {
workspace_providers: Mutex::new(workspace_providers),
}
}
fn remember_provider(&self, provider: Arc<dyn WorkspaceProvider>) -> bool {
let mut providers = self
.workspace_providers
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let id = provider.id();
match providers.get(&id) {
Some(existing) => Arc::ptr_eq(existing, &provider),
None => {
providers.insert(id, provider);
true
}
}
}
fn provider(&self, id: &WorkspaceProviderId) -> Option<Arc<dyn WorkspaceProvider>> {
self.workspace_providers
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(id)
.cloned()
}
}
impl fmt::Debug for AgentState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("AgentState")
.field(
"workspace_provider_count",
&self
.workspace_providers
.lock()
.map_or(0, |providers| providers.len()),
)
.finish()
}
}
impl fmt::Debug for Agent {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let capability_ids = self
.capabilities
.iter()
.map(everruns_capability::CapabilityRef::capability_id)
.collect::<Vec<_>>();
f.debug_struct("Agent")
.field("name", &self.name)
.field("instructions", &self.instructions)
.field("model", &self.model)
.field("provider", &self.provider)
.field("capabilities", &capability_ids)
.field(
"capability_implementations",
&self.capability_implementations,
)
.field("initial_files", &self.initial_files)
.field("max_iterations", &self.max_iterations)
.field("parallel_tool_calls", &self.parallel_tool_calls)
.field("workspace_root", &self.workspace_root)
.field("mcp_servers", &self.mcp_servers)
.field("plugin_warnings", &self.plugin_warnings)
.field("state", &self.state)
.finish_non_exhaustive()
}
}
#[allow(dead_code)] pub(crate) enum BackendInitError {
Event(EventLogError),
Host(everruns_provider::error::AgentLoopError),
}
impl BackendInitError {
pub(crate) fn into_agent_loop(self) -> everruns_provider::error::AgentLoopError {
match self {
Self::Event(error) => {
everruns_provider::error::AgentLoopError::store(error.to_string())
}
Self::Host(error) => error,
}
}
pub(crate) fn history_error(&self) -> crate::HistoryError {
match self {
Self::Event(EventLogError::Corruption { .. }) => crate::HistoryError::Corrupt,
Self::Event(_) | Self::Host(_) => crate::HistoryError::Unavailable,
}
}
pub(crate) fn resume_error(&self) -> crate::ResumeError {
match self {
Self::Event(EventLogError::Corruption { .. }) => crate::ResumeError::Corrupt,
Self::Event(_) | Self::Host(_) => crate::ResumeError::Unavailable,
}
}
}
impl Agent {
pub fn builder() -> AgentBuilder {
AgentBuilder::default()
}
pub(crate) async fn bind_session_environment(
&self,
binding_store: &dyn EnvironmentBindingStore,
session_id: SessionId,
environment: &Environment,
) -> Result<(), crate::SessionEnvironmentError> {
let head = environment.workspace_head();
if !self.state.remember_provider(head.provider()) {
return Err(crate::SessionEnvironmentError::ProviderConflict);
}
binding_store
.bind(session_id, head.binding())
.await
.map_err(crate::SessionEnvironmentError::from)?;
Ok(())
}
pub(crate) async fn bind_default_session_environment(
&self,
binding_store: &dyn EnvironmentBindingStore,
session_id: SessionId,
) -> Result<Environment, crate::SessionEnvironmentError> {
let environment = self
.default_workspace
.environment(session_id)
.await
.map_err(crate::SessionEnvironmentError::Workspace)?;
self.bind_session_environment(binding_store, session_id, &environment)
.await?;
Ok(environment)
}
pub(crate) async fn reopen_session_environment(
&self,
binding_store: &dyn EnvironmentBindingStore,
session_id: SessionId,
) -> Result<Option<Environment>, crate::ResumeError> {
let Some(binding) = binding_store
.load(session_id)
.await
.map_err(|error| match error {
EnvironmentBindingError::Corrupt => crate::ResumeError::WorkspaceBindingCorrupt,
EnvironmentBindingError::Conflict | EnvironmentBindingError::Unavailable => {
crate::ResumeError::Unavailable
}
_ => crate::ResumeError::Unavailable,
})?
else {
return Ok(None);
};
let provider = self.state.provider(&binding.provider_id).ok_or_else(|| {
crate::ResumeError::WorkspaceProviderUnavailable {
provider_id: binding.provider_id.to_string(),
}
})?;
let descriptor = provider
.open_workspace_from_binding(&binding)
.await
.map_err(map_workspace_resume_error)?;
let workspace = everruns_host::Workspace::from_descriptor(provider, descriptor);
let head = workspace
.reopen(&binding)
.await
.map_err(map_workspace_resume_error)?;
Ok(Some(
Environment::builder()
.workspace(head)
.build()
.map_err(map_workspace_resume_error)?,
))
}
pub(crate) fn name(&self) -> &str {
&self.name
}
pub(crate) fn lifecycle_hooks(&self) -> crate::hooks::LifecycleHooks {
self.lifecycle_hooks.clone()
}
#[cfg(feature = "local")]
pub(crate) fn local_config(&self) -> Option<&crate::LocalConfig> {
self.local.as_ref()
}
pub(crate) async fn catalog_session(
&self,
backends: &HostBackends,
session_id: SessionId,
) -> Result<(), crate::HistoryError> {
if backends
.session_store
.get_session(session_id)
.await
.map_err(|_| crate::HistoryError::Unavailable)?
.is_some()
{
return Ok(());
}
let session = SessionBuilder::new(everruns_provider::typed_id::HarnessId::new())
.id(session_id)
.build();
backends
.session_store
.add_session(session)
.await
.map_err(|_| crate::HistoryError::Unavailable)
}
pub(crate) async fn build_runtime_with_event_sink(
&self,
backends: HostBackends,
session_id: SessionId,
environment: Option<Environment>,
event_sink: Arc<dyn EventSink>,
hook_state: Arc<crate::hooks::HookRunState>,
) -> Result<InProcessRuntime, everruns_provider::error::AgentLoopError> {
self.build_runtime_with_backends(
backends,
session_id,
environment,
Some(event_sink),
Some(hook_state),
)
.await
}
pub(crate) fn plugin_warnings(&self) -> Vec<String> {
self.plugin_warnings.clone()
}
async fn build_runtime_with_backends(
&self,
mut backends: HostBackends,
session_id: SessionId,
environment: Option<Environment>,
event_sink: Option<Arc<dyn EventSink>>,
hook_state: Option<Arc<crate::hooks::HookRunState>>,
) -> Result<InProcessRuntime, everruns_provider::error::AgentLoopError> {
let hook_capability = hook_state.and_then(|state| state.capability());
let mut capabilities = self.capabilities.clone();
if hook_capability.is_some() {
capabilities.retain(|capability| {
capability.capability_id() != crate::hooks::LIFECYCLE_HOOK_CAPABILITY_ID
});
capabilities.push(everruns_capability::CapabilityRef::new(
crate::hooks::LIFECYCLE_HOOK_CAPABILITY_ID,
));
}
let mut harness =
HarnessBuilder::new(&self.name, &self.instructions).capabilities(capabilities.clone());
if let Some(parallel) = self.parallel_tool_calls {
harness = harness.parallel_tool_calls(parallel);
}
for file in &self.initial_files {
harness = harness.initial_file(file.clone());
}
let harness_id = harness.harness_id();
let harness = harness.build();
let mut agent = RuntimeAgentBuilder::new(&self.name, &self.instructions)
.capabilities(capabilities.clone());
if let Some(parallel) = self.parallel_tool_calls {
agent = agent.parallel_tool_calls(parallel);
}
if let Some(max_iterations) = self.max_iterations {
agent = agent.max_iterations(max_iterations);
}
let agent_id = agent.agent_id();
let agent = agent.build();
let mut session = SessionBuilder::new(harness_id)
.id(session_id)
.agent(agent_id)
.capabilities(capabilities)
.mcp_servers(self.mcp_servers.clone());
if let Some(environment) = &environment {
session = session.workspace(environment.workspace_head().workspace_id());
}
if let Some(parallel) = self.parallel_tool_calls {
session = session.parallel_tool_calls(parallel);
}
for file in &self.initial_files {
session = session.initial_file(file.clone());
}
let session = session.build();
let mut builder = InProcessRuntimeBuilder::new()
.harness(harness)
.agent(agent)
.session(session)
.model_spec(self.model.clone())
.workspace_policy(self.workspace_policy.clone());
if let Some(event_sink) = event_sink {
backends = backends.with_event_sink(event_sink);
}
builder = builder.backends(backends);
let workspace_root = self.workspace_root.as_ref().or({
#[cfg(feature = "local")]
{
self.local.as_ref().map(|config| &config.workspace_root)
}
#[cfg(not(feature = "local"))]
{
None
}
});
if let Some(environment) = environment {
let registry = {
#[cfg(feature = "local")]
if self.local.is_some() {
everruns_local::local_capability_registry()
} else {
framework_capability_registry(false)
}
#[cfg(not(feature = "local"))]
{
framework_capability_registry(false)
}
};
let platform = everruns_host::HostComposition::builder()
.capability_registry(registry)
.driver_registry(everruns_provider::driver_registry::DriverRegistry::new())
.egress_service(everruns_host::runtime_egress_service())
.session_file_system_factory(Arc::new(
everruns_host::FixedSessionFileSystemFactory::new(
environment.workspace_head().file_system(),
),
))
.build();
builder = builder.host_composition(platform);
} else if let Some(root) = workspace_root {
let registry = {
#[cfg(feature = "local")]
if self.local.is_some() {
everruns_local::local_capability_registry()
} else {
framework_capability_registry(false)
}
#[cfg(not(feature = "local"))]
{
framework_capability_registry(false)
}
};
let platform = everruns_host::HostComposition::builder()
.capability_registry(registry)
.driver_registry(everruns_provider::DriverRegistry::new())
.egress_service(everruns_host::runtime_egress_service())
.session_file_system_factory(Arc::new(
everruns_host::RealDiskSessionFileSystemFactory::new(root),
))
.build();
builder = builder.host_composition(platform);
}
for implementation in &self.capability_implementations {
builder = implementation.register(builder);
}
if let Some(capability) = hook_capability {
builder = builder.capability(capability);
}
builder = builder.provider(self.provider.clone());
builder.build().await
}
}
#[derive(Clone, Default)]
pub struct AgentBuilder {
name: Option<String>,
instructions: Option<String>,
model: Option<Model>,
providers: Vec<Provider>,
capabilities: Vec<crate::CapabilitySpec>,
tools: Vec<Tool>,
initial_files: Vec<InitialFile>,
max_iterations: Option<usize>,
parallel_tool_calls: Option<bool>,
workspace_root: Option<PathBuf>,
workspace_policy: everruns_core::WorkspacePolicy,
workspace_providers: Vec<Arc<dyn WorkspaceProvider>>,
mcp_servers: Vec<crate::McpServer>,
plugin_warnings: Vec<String>,
#[cfg(feature = "local")]
local: Option<crate::LocalConfig>,
lifecycle_hooks: crate::hooks::LifecycleHooks,
}
impl AgentBuilder {
pub fn instructions(mut self, instructions: impl Into<String>) -> Self {
self.instructions = Some(instructions.into());
self
}
pub fn model(mut self, model: impl Into<Model>) -> Self {
self.model = Some(model.into());
self
}
pub fn provider(mut self, provider: impl Into<Provider>) -> Self {
self.providers.push(provider.into());
self
}
pub fn name(mut self, name: impl Into<String>) -> Self {
self.name = Some(name.into());
self
}
pub fn tool(mut self, tool: impl IntoTool) -> Self {
self.tools.push(tool.into_tool());
self
}
pub fn on_agent_start<F, Fut, O>(mut self, handler: F) -> Self
where
F: Fn(crate::AgentStartContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = O> + Send + 'static,
O: crate::IntoHookResult + 'static,
{
self.lifecycle_hooks.on_agent_start(handler);
self
}
pub fn on_turn_start<F, Fut, O>(mut self, handler: F) -> Self
where
F: Fn(crate::TurnStartContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = O> + Send + 'static,
O: crate::IntoHookResult + 'static,
{
self.lifecycle_hooks.on_turn_start(handler);
self
}
pub fn on_tool_start<F, Fut, O>(mut self, handler: F) -> Self
where
F: Fn(crate::ToolStartContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = O> + Send + 'static,
O: crate::IntoHookResult + 'static,
{
self.lifecycle_hooks.on_tool_start(handler);
self
}
pub fn on_tool_end<F, Fut, O>(mut self, handler: F) -> Self
where
F: Fn(crate::ToolEndContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = O> + Send + 'static,
O: crate::IntoHookResult + 'static,
{
self.lifecycle_hooks.on_tool_end(handler);
self
}
pub fn on_completion<F, Fut, O>(mut self, handler: F) -> Self
where
F: Fn(crate::CompletionContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = O> + Send + 'static,
O: crate::IntoHookResult + 'static,
{
self.lifecycle_hooks.on_completion(handler);
self
}
pub fn capability(mut self, capability: impl crate::IntoCapability) -> Self {
self.capabilities.push(capability.into_capability());
self
}
pub fn max_iterations(mut self, max_iterations: usize) -> Self {
self.max_iterations = Some(max_iterations);
self
}
pub fn parallel_tool_calls(mut self, parallel: bool) -> Self {
self.parallel_tool_calls = Some(parallel);
self
}
pub fn initial_file(mut self, file: InitialFile) -> Self {
self.initial_files.push(file);
self
}
pub fn file(mut self, path: impl Into<String>, content: impl Into<String>) -> Self {
self.initial_files.push(InitialFile {
path: path.into(),
content: content.into(),
encoding: "text".to_string(),
is_readonly: false,
});
self
}
pub fn readonly_file(mut self, path: impl Into<String>, content: impl Into<String>) -> Self {
self.initial_files.push(InitialFile {
path: path.into(),
content: content.into(),
encoding: "text".to_string(),
is_readonly: true,
});
self
}
pub fn workspace(mut self, root: impl Into<PathBuf>) -> Self {
self.workspace_root = Some(root.into());
if !self
.capabilities
.iter()
.any(|capability| capability.capability_ref().id() == "session_file_system")
{
self.capabilities
.push(crate::CapabilityRef::new("session_file_system").into());
}
self
}
pub fn workspace_policy(mut self, policy: crate::WorkspacePolicy) -> Self {
self.workspace_policy = policy;
self
}
pub fn workspace_provider(mut self, provider: Arc<dyn WorkspaceProvider>) -> Self {
self.workspace_providers.push(provider);
self
}
pub fn mcp_server(mut self, server: crate::McpServer) -> Self {
self.mcp_servers.push(server);
self
}
pub fn plugin(mut self, path: impl AsRef<Path>) -> Result<Self, crate::PluginError> {
let loaded = crate::plugin::load(path.as_ref())?;
self.capabilities.push(loaded.capability);
self.plugin_warnings.extend(loaded.warnings);
Ok(self)
}
#[cfg(feature = "local")]
pub fn local(mut self, config: crate::LocalConfig) -> Self {
self.local = Some(config);
if !self
.capabilities
.iter()
.any(|capability| capability.capability_ref().id() == "session_file_system")
{
self.capabilities
.push(crate::CapabilityRef::new("session_file_system").into());
}
self
}
pub fn build(self) -> Result<Agent, BuildError> {
let mut workspace_providers = HashMap::new();
for provider in &self.workspace_providers {
let id = provider.id();
if workspace_providers
.insert(id.clone(), provider.clone())
.is_some()
{
return Err(BuildError::DuplicateWorkspaceProvider { id: id.to_string() });
}
}
let default_workspace_root = self.workspace_root.clone().or_else(|| {
#[cfg(feature = "local")]
{
self.local
.as_ref()
.map(|configuration| configuration.workspace_root.clone())
}
#[cfg(not(feature = "local"))]
{
None
}
});
let default_workspace = default_workspace_root
.map(crate::default_workspace::DefaultWorkspace::directory)
.unwrap_or_else(crate::default_workspace::DefaultWorkspace::in_memory);
let default_provider = default_workspace.provider();
let default_provider_id = default_provider.id();
if workspace_providers
.insert(default_provider_id.clone(), default_provider)
.is_some()
{
return Err(BuildError::DuplicateWorkspaceProvider {
id: default_provider_id.to_string(),
});
}
let instructions = self.instructions.unwrap_or_default();
if instructions.trim().is_empty() {
return Err(BuildError::BlankInstructions);
}
let model = self.model.ok_or(BuildError::MissingModel)?;
let mut providers = self.providers;
if let Some(provider) = model.bundled_provider {
providers.push(provider);
}
if providers.is_empty() {
return Err(BuildError::MissingProvider);
}
if providers.len() > 1 {
let mut registered = providers
.iter()
.map(|provider| provider.id().to_string())
.collect::<Vec<_>>();
registered.sort();
return Err(BuildError::MultipleProviders { registered });
}
let provider = providers.pop().expect("provider count checked above");
let model = ModelSpec::on(provider.id().clone(), model.id);
let name = self.name.unwrap_or_else(|| "agent".to_string());
let mut function_tools = Vec::new();
let mut seen_tool_names: HashSet<String> = HashSet::new();
for tool in self.tools {
let tool_name = tool.name().to_string();
if !seen_tool_names.insert(tool_name.clone()) {
return Err(BuildError::DuplicateTool { name: tool_name });
}
let function_tool = tool.into_function();
validate_tool_name(function_tool.name()).map_err(|reason| {
BuildError::InvalidToolName {
name: tool_name.clone(),
reason,
}
})?;
validate_tool_schema(function_tool.schema()).map_err(|reason| {
BuildError::InvalidToolSchema {
name: tool_name,
reason,
}
})?;
function_tools.push(function_tool);
}
let mcp_servers = crate::mcp::into_scoped(self.mcp_servers)
.map_err(|reason| BuildError::InvalidMcpServer { reason })?;
let capability_registry = {
#[cfg(feature = "local")]
if self.local.is_some() {
everruns_local::local_capability_registry()
} else {
framework_capability_registry(false)
}
#[cfg(not(feature = "local"))]
{
framework_capability_registry(false)
}
};
let mut capabilities = Vec::new();
let mut capability_implementations = Vec::new();
let mut activated_capabilities = everruns_capability::ActivationSet::new();
for input in self.capabilities {
let parts = input.into_parts();
let input_id = parts.reference.id().to_string();
everruns_capability::validate_capability_id(&input_id).map_err(|error| {
BuildError::InvalidCapability {
id: input_id.clone(),
reason: error.reason(),
}
})?;
everruns_capability::validate_capability_config(
&input_id,
parts.reference.config_value(),
)
.map_err(|error| BuildError::InvalidCapability {
id: input_id.clone(),
reason: error.reason(),
})?;
let canonical_id = capability_registry
.canonical_id(&input_id)
.unwrap_or(&input_id)
.to_string();
if activated_capabilities.contains(&canonical_id) {
return Err(BuildError::DuplicateCapability { id: canonical_id });
}
#[cfg(feature = "capabilities")]
if let Some(definition) = parts.definition {
definition
.validate()
.map_err(|error| BuildError::InvalidCapability {
id: input_id.clone(),
reason: error.reason(),
})?;
if capability_registry.get(definition.id()).is_some() {
return Err(BuildError::DuplicateCapability {
id: definition.id().to_string(),
});
}
for tool in definition.tools() {
let spec = tool.spec();
validate_tool_name(spec.name()).map_err(|reason| {
BuildError::InvalidToolName {
name: spec.name().to_string(),
reason,
}
})?;
validate_tool_schema(spec.input_schema()).map_err(|reason| {
BuildError::InvalidToolSchema {
name: spec.name().to_string(),
reason,
}
})?;
if !seen_tool_names.insert(spec.name().to_string()) {
return Err(BuildError::DuplicateTool {
name: spec.name().to_string(),
});
}
}
capability_implementations.push(CapabilityImplementation::Definition(definition));
}
validate_registered_capability_config(
&capability_registry,
&canonical_id,
parts.reference.config_value(),
)?;
activated_capabilities
.activate(canonical_id.clone())
.map_err(|error| BuildError::DuplicateCapability {
id: error.id().to_string(),
})?;
capabilities.push(everruns_capability::CapabilityRef::with_config(
canonical_id,
parts.reference.config_value().clone(),
));
}
for function_tool in function_tools {
let id = function_tool.name().to_string();
everruns_capability::validate_capability_id(&id).map_err(|error| {
BuildError::InvalidCapability {
id: id.clone(),
reason: error.reason(),
}
})?;
if capability_registry.get(&id).is_some()
|| activated_capabilities.activate(id.clone()).is_err()
{
return Err(BuildError::DuplicateCapability { id });
}
capabilities.push(everruns_capability::CapabilityRef::new(id.as_str()));
capability_implementations.push(CapabilityImplementation::Function(function_tool));
}
Ok(Agent {
name,
instructions,
model,
provider,
capabilities,
capability_implementations,
initial_files: self.initial_files,
max_iterations: self.max_iterations,
parallel_tool_calls: self.parallel_tool_calls,
workspace_root: self.workspace_root,
default_workspace,
workspace_policy: self.workspace_policy,
mcp_servers,
plugin_warnings: self.plugin_warnings,
#[cfg(feature = "local")]
local: self.local,
lifecycle_hooks: self.lifecycle_hooks,
state: Arc::new(AgentState::new(workspace_providers)),
})
}
}
impl fmt::Debug for AgentBuilder {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("AgentBuilder")
.field("name", &self.name)
.field("instructions", &self.instructions)
.field("workspace_root", &self.workspace_root)
.field(
"workspace_providers",
&self
.workspace_providers
.iter()
.map(|provider| provider.id())
.collect::<Vec<_>>(),
)
.finish_non_exhaustive()
}
}
fn map_workspace_resume_error(error: everruns_host::WorkspaceError) -> crate::ResumeError {
match error {
everruns_host::WorkspaceError::BindingMismatch => crate::ResumeError::WorkspaceMismatch,
everruns_host::WorkspaceError::NotFound
| everruns_host::WorkspaceError::Archived
| everruns_host::WorkspaceError::ProviderUnavailable(_)
| everruns_host::WorkspaceError::Provider(_)
| everruns_host::WorkspaceError::Conflict
| everruns_host::WorkspaceError::InvalidRequest(_) => {
crate::ResumeError::WorkspaceUnavailable
}
_ => crate::ResumeError::WorkspaceUnavailable,
}
}
fn framework_capability_registry(hosted_base: bool) -> everruns_core::CapabilityRegistry {
#[cfg(not(feature = "builtins"))]
let _ = hosted_base;
let registry = everruns_core::CapabilityRegistry::new();
#[cfg(feature = "builtins")]
let registry = {
let mut registry = registry;
if hosted_base {
everruns_builtins::register_portable_capabilities(&mut registry)
.expect("portable built-in catalog must have unique capability IDs");
} else {
everruns_builtins::register_runtime_capabilities(&mut registry)
.expect("portable runtime catalog must have unique capability IDs");
}
registry
};
everruns_host::compose_runtime_capability_registry(registry)
}
fn validate_registered_capability_config(
registry: &everruns_core::CapabilityRegistry,
id: &str,
config: &serde_json::Value,
) -> Result<(), BuildError> {
if let Some(capability) = registry.get(id) {
capability
.validate_config(config)
.map_err(|reason| BuildError::InvalidCapability {
id: id.to_string(),
reason,
})?;
}
#[cfg(feature = "builtins")]
if id == everruns_builtins::AUTO_TOOL_SEARCH_CAPABILITY_ID {
validate_tool_search_config(config).map_err(|reason| BuildError::InvalidCapability {
id: id.to_string(),
reason,
})?;
}
if everruns_core::is_declarative_capability(id) || everruns_capability::is_plugin_capability(id)
{
let definition = serde_json::from_value::<everruns_core::DeclarativeCapabilityDefinition>(
config.clone(),
)
.map_err(|error| BuildError::InvalidCapability {
id: id.to_string(),
reason: format!("invalid declarative capability config: {error}"),
})?;
everruns_core::validate_declarative_capability_definition(&definition).map_err(
|reason| BuildError::InvalidCapability {
id: id.to_string(),
reason,
},
)?;
}
Ok(())
}
#[cfg(feature = "builtins")]
fn validate_tool_search_config(config: &serde_json::Value) -> Result<(), String> {
let object = config
.as_object()
.expect("capability config object validated before built-in schema");
for key in object.keys() {
if key != "threshold" && key != "never_defer" {
return Err(format!("unknown tool-search config field {key:?}"));
}
}
if let Some(threshold) = object.get("threshold") {
let value = threshold
.as_u64()
.ok_or_else(|| "threshold must be a non-negative integer".to_string())?;
usize::try_from(value)
.map_err(|_| "threshold is too large for this platform".to_string())?;
}
if let Some(never_defer) = object.get("never_defer") {
let names = never_defer
.as_array()
.ok_or_else(|| "never_defer must be an array of tool names".to_string())?;
for name in names {
let name = name
.as_str()
.ok_or_else(|| "never_defer must contain only tool names".to_string())?;
validate_tool_name(name)
.map_err(|reason| format!("invalid never_defer tool {name:?}: {reason}"))?;
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex};
use everruns_provider::tool_types::ToolCall;
use serde_json::{Value, json};
use super::*;
use crate::{FunctionTool, InMemoryEngine};
fn obj_schema() -> Value {
json!({ "type": "object", "properties": {}, "additionalProperties": false })
}
#[test]
fn build_rejects_invalid_tool_name() {
let err = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("ok"))
.tool(FunctionTool::new(
"bad name",
"desc",
obj_schema(),
|_: Value| async move { Ok::<_, String>(json!({})) },
))
.build()
.unwrap_err();
assert!(
matches!(err, BuildError::InvalidToolName { ref name, .. } if name == "bad name"),
"got {err:?}"
);
}
#[test]
fn build_rejects_invalid_tool_schema() {
let err = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("ok"))
.tool(FunctionTool::new(
"arr",
"desc",
json!({ "type": "array" }),
|_: Value| async move { Ok::<_, String>(json!({})) },
))
.build()
.unwrap_err();
assert!(
matches!(err, BuildError::InvalidToolSchema { ref name, .. } if name == "arr"),
"got {err:?}"
);
}
#[test]
fn build_rejects_duplicate_tool_names() {
let make = || {
FunctionTool::new("dup", "desc", obj_schema(), |_: Value| async move {
Ok::<_, String>(json!({}))
})
};
let err = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("ok"))
.tool(make())
.tool(make())
.build()
.unwrap_err();
assert_eq!(
err,
BuildError::DuplicateTool {
name: "dup".to_string()
}
);
}
#[tokio::test]
async fn function_tool_executes_end_to_end() {
let received: Arc<Mutex<Option<Value>>> = Arc::new(Mutex::new(None));
let sink = received.clone();
let tool = FunctionTool::new(
"greet",
"Greet a person by name.",
json!({
"type": "object",
"properties": { "name": { "type": "string" } },
"required": ["name"],
}),
move |args: Value| {
let sink = sink.clone();
async move {
*sink.lock().unwrap() = Some(args.clone());
let name = args["name"].as_str().unwrap_or("world");
Ok::<_, String>(json!({ "greeting": format!("Hello, {name}!") }))
}
},
);
let agent = Agent::builder()
.instructions("Call greet when asked to greet someone.")
.model(Model::simulated_scripted(
"All done.",
vec![
vec![ToolCall {
id: "call_greet_1".into(),
name: "greet".into(),
arguments: json!({ "name": "Ada" }),
}],
vec![],
],
))
.tool(tool)
.build()
.expect("valid agent");
let session = InMemoryEngine::new().create(agent.clone());
let turn = session.run("Please greet Ada.").await.expect("turn runs");
assert!(turn.success, "turn should succeed: {:?}", turn.error);
assert_eq!(turn.tool_calls, 1, "the function tool must have executed");
assert_eq!(turn.response, "All done.");
assert_eq!(
received.lock().unwrap().as_ref().expect("handler ran")["name"],
json!("Ada"),
"handler must receive the model's call arguments",
);
}
#[tokio::test]
async fn function_tool_handler_error_is_model_visible_not_a_panic() {
let tool = FunctionTool::new(
"always_fails",
"Always returns an error.",
obj_schema(),
|_: Value| async move { Err::<Value, String>("boom".to_string()) },
);
let agent = Agent::builder()
.instructions("Call the tool.")
.model(Model::simulated_scripted(
"Handled.",
vec![
vec![ToolCall {
id: "call_fail_1".into(),
name: "always_fails".into(),
arguments: json!({}),
}],
vec![],
],
))
.tool(tool)
.build()
.expect("valid agent");
let session = InMemoryEngine::new().create(agent.clone());
let turn = session.run("go").await.expect("turn runs");
assert!(turn.success, "turn should recover from a tool error");
assert_eq!(turn.tool_calls, 1);
}
#[test]
fn build_rejects_blank_instructions() {
let err = Agent::builder()
.instructions(" ")
.model(Model::simulated("hi"))
.build()
.unwrap_err();
assert_eq!(err, BuildError::BlankInstructions);
}
#[test]
fn build_rejects_missing_model() {
let err = Agent::builder()
.instructions("You are concise.")
.build()
.unwrap_err();
assert_eq!(err, BuildError::MissingModel);
}
#[test]
fn build_succeeds_with_simulator() {
let agent = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("Sure."))
.name("assistant")
.build()
.expect("valid agent");
assert_eq!(agent.name, "assistant");
}
#[test]
fn build_resolves_plain_model_id_against_provider() {
let agent = Agent::builder()
.instructions("You are concise.")
.provider(Provider::new(
"acme",
LlmSimDriver::new(LlmSimConfig::fixed("Sure.")),
))
.model("assistant-v1")
.build()
.expect("valid agent");
assert_eq!(agent.model, ModelSpec::on("acme", "assistant-v1"));
assert_eq!(agent.provider.id().as_str(), "acme");
}
#[test]
fn build_accepts_a_model_with_its_provider() {
let provider = Provider::new("acme", LlmSimDriver::new(LlmSimConfig::fixed("Sure.")));
let agent = Agent::builder()
.instructions("You are concise.")
.model(Model::new("assistant-v1", provider))
.build()
.expect("valid agent");
assert_eq!(agent.model, ModelSpec::on("acme", "assistant-v1"));
assert_eq!(agent.provider.id().as_str(), "acme");
}
#[test]
fn build_rejects_live_model_without_provider() {
let err = Agent::builder()
.instructions("You are concise.")
.model("assistant-v1")
.build()
.unwrap_err();
assert_eq!(err, BuildError::MissingProvider);
}
#[test]
fn build_rejects_multiple_providers() {
let provider = |id| Provider::new(id, LlmSimDriver::new(LlmSimConfig::fixed("Sure.")));
let err = Agent::builder()
.instructions("You are concise.")
.provider(provider("zeta"))
.provider(provider("alpha"))
.model("assistant-v1")
.build()
.unwrap_err();
assert_eq!(
err,
BuildError::MultipleProviders {
registered: vec!["alpha".to_string(), "zeta".to_string()],
}
);
}
#[test]
fn workspace_policy_defaults_to_read_only() {
let agent = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("Sure."))
.build()
.expect("valid agent");
assert!(agent.workspace_policy.permits_read("notes.txt"));
assert!(!agent.workspace_policy.permits_write("notes.txt"));
assert!(!agent.workspace_policy.permits_read(".env"));
}
#[tokio::test]
async fn custom_workspace_policy_reaches_the_high_level_runtime() {
let workspace = tempfile::tempdir().expect("workspace tempdir");
let policy = crate::WorkspacePolicy::builder()
.allow_read("public")
.build()
.expect("valid policy");
let agent = Agent::builder()
.instructions("Read the workspace.")
.model(Model::simulated("Sure."))
.workspace(workspace.path())
.workspace_policy(policy)
.file("public/note.txt", "visible")
.file("private/note.txt", "hidden")
.build()
.expect("valid agent");
let session_id = SessionId::new();
let runtime = agent
.build_runtime_with_backends(HostBackends::in_memory(), session_id, None, None, None)
.await
.expect("runtime builds");
let visible = runtime
.read_file(session_id, "public/note.txt")
.await
.expect("allowed read")
.expect("seeded file");
assert_eq!(visible.content.as_deref(), Some("visible"));
let denied = runtime.read_file(session_id, "private/note.txt").await;
assert!(denied.is_err());
}
#[cfg(feature = "openai")]
#[test]
fn openai_provider_converts_without_leaking_key() {
use crate::providers::openai::OpenAI;
let provider: Provider = OpenAI::new("sk-super-secret").into();
assert_eq!(provider.id().as_str(), "openai");
let rendered = format!("{provider:?}");
assert!(!rendered.contains("sk-super-secret"), "got {rendered}");
}
#[cfg(feature = "openai")]
#[tokio::test]
async fn openai_agent_builds_runtime_offline() {
use crate::providers::openai::OpenAI;
let agent = Agent::builder()
.instructions("You are concise.")
.provider(OpenAI::new("sk-test"))
.model("gpt-5-mini")
.build()
.expect("valid agent");
let runtime = agent
.build_runtime_with_backends(
HostBackends::in_memory(),
SessionId::new(),
None,
None,
None,
)
.await
.expect("openai runtime builds offline");
let _ = runtime;
}
#[tokio::test]
async fn build_runtime_seeds_the_requested_session_id() {
let agent = Agent::builder()
.instructions("You are concise.")
.model(Model::simulated("Sure."))
.build()
.expect("valid agent");
let session_id = SessionId::new();
let runtime = agent
.build_runtime_with_backends(HostBackends::in_memory(), session_id, None, None, None)
.await
.expect("runtime builds");
let result = runtime
.run_turn(session_id, everruns_core::InputMessage::user("hi"))
.await
.expect("turn runs");
assert!(result.success);
assert_eq!(result.response, "Sure.");
}
}