use std::collections::BTreeMap;
use std::sync::Arc;
use crate::agent_definition::AgentDefinition;
use crate::capabilities::{CapabilityRegistry, SystemPromptContext, resolve_capability_configs};
use crate::compaction_policy::CompactionPolicy;
use crate::config_layer::AgentConfigOverlay;
use crate::driver_registry::ChatDriver;
use crate::error::Result;
use crate::events::TokenUsage;
use crate::harness_definition::HarnessDefinition;
use crate::message::{RuntimeMessage, RuntimeMessageRole};
use crate::provider::DriverId;
use crate::runtime_agent::{RuntimeAgent, RuntimeAgentBuilder};
use crate::session::ExecutionSession;
use crate::session_files::SessionFileSystem;
use crate::tool_types::ToolDefinition;
use crate::typed_id::{ModelId, SessionId};
use crate::{AgentCapabilityConfig, ResolvedExecutionSnapshot};
#[async_trait::async_trait]
pub trait TurnContextResolver: Send + Sync {
async fn resolve_turn_context(
&self,
request: TurnContextRequest,
) -> Result<AssembledTurnContext>;
}
#[derive(Debug, Clone)]
pub struct TurnContextRequest {
pub session_id: SessionId,
pub harness_id: crate::HarnessId,
pub agent_id: Option<crate::AgentId>,
pub mcp_tool_definitions: Vec<ToolDefinition>,
pub allow_provider_managed_reduction: bool,
}
#[derive(Clone)]
pub struct ResolvedModelExecution {
pub model: String,
pub provider: crate::ProviderKey,
pub provider_type: DriverId,
pub driver: Arc<dyn ChatDriver>,
pub provider_managed_reduction_option: Option<(String, serde_json::Value)>,
}
impl std::fmt::Debug for ResolvedModelExecution {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ResolvedModelExecution")
.field("model", &self.model)
.field("provider", &self.provider)
.field("provider_type", &self.provider_type)
.field(
"provider_managed_reduction",
&self.provider_managed_reduction_option.is_some(),
)
.field("driver", &"<opaque>")
.finish()
}
}
#[derive(Debug, Clone)]
pub struct ResolvedTurnContextInput {
pub snapshot: ResolvedExecutionSnapshot,
pub messages: Vec<RuntimeMessage>,
pub message_source_sequence: Option<i64>,
pub model: ResolvedModelExecution,
pub resolved_model_id: Option<ModelId>,
pub mcp_tool_definitions: Vec<ToolDefinition>,
}
#[derive(Debug, Clone)]
pub struct AssembledTurnContext {
pub snapshot: ResolvedExecutionSnapshot,
pub resolved_capability_configs: Vec<AgentCapabilityConfig>,
pub messages: Vec<RuntimeMessage>,
pub message_source_sequence: Option<i64>,
pub runtime_agent: RuntimeAgent,
pub model: ResolvedModelExecution,
pub resolved_model_id: Option<ModelId>,
pub resolved_locale: Option<String>,
pub compaction_policy: Option<Arc<dyn CompactionPolicy>>,
pub embedder_metadata: BTreeMap<String, String>,
}
impl AssembledTurnContext {
pub fn session_id(&self) -> SessionId {
self.snapshot.session_id
}
pub fn cumulative_usage(&self) -> Option<TokenUsage> {
self.snapshot.cumulative_usage.clone()
}
}
#[derive(Debug, Clone)]
pub struct ResolvedRuntimeCapabilities {
pub effective_overlay: AgentConfigOverlay,
pub resolved_capability_configs: Vec<AgentCapabilityConfig>,
}
pub async fn assemble_resolved_turn_context(
input: ResolvedTurnContextInput,
capability_registry: &CapabilityRegistry,
file_store: Option<Arc<dyn SessionFileSystem>>,
session_storage: Option<Arc<dyn crate::session_services::SessionStorageStore>>,
) -> Result<AssembledTurnContext> {
let ResolvedTurnContextInput {
snapshot,
messages,
message_source_sequence,
model,
resolved_model_id,
mcp_tool_definitions,
} = input;
let ResolvedRuntimeCapabilities {
effective_overlay,
resolved_capability_configs,
} = resolve_snapshot_capabilities(&snapshot, capability_registry);
let resolved_locale = extract_locale_override(&messages).or_else(|| snapshot.locale.clone());
let file_store =
file_store.map(|fs| crate::mount_fs::scoped_prompt_file_store(fs, snapshot.workspace_id));
let prompt_ctx = SystemPromptContext {
session_id: snapshot.session_id,
locale: resolved_locale.clone(),
file_store,
model: Some(model.model.clone()),
session_storage,
};
let compaction_policy = effective_overlay.capabilities.iter().find_map(|config| {
capability_registry
.get(config.capability_id())?
.compaction_policy(config.config_value())
});
let mut runtime_agent = build_runtime_agent(
&snapshot,
effective_overlay,
capability_registry,
&prompt_ctx,
&mcp_tool_definitions,
&model.model,
)
.await?;
if let Some((key, value)) = &model.provider_managed_reduction_option {
runtime_agent
.driver_options
.insert(key.clone(), value.clone());
} else if resolved_capability_configs
.iter()
.any(|config| config.capability_id() == "infinity_context")
{
runtime_agent.driver_options.insert(
"everruns/provider_managed_reduction_fallback".to_string(),
serde_json::json!({"reason":"provider_or_model_ineligible"}),
);
}
let embedder_metadata = snapshot.embedder_metadata.clone();
Ok(AssembledTurnContext {
snapshot,
resolved_capability_configs,
messages,
message_source_sequence,
runtime_agent,
model,
resolved_model_id,
resolved_locale,
compaction_policy,
embedder_metadata,
})
}
pub fn resolve_snapshot_capabilities(
snapshot: &ResolvedExecutionSnapshot,
capability_registry: &CapabilityRegistry,
) -> ResolvedRuntimeCapabilities {
let effective_overlay = AgentConfigOverlay {
system_prompt: snapshot.instructions.clone(),
capabilities: snapshot.capabilities.clone(),
initial_files: snapshot.initial_files.clone(),
network_access: snapshot.network_access.clone(),
default_model_id: snapshot.default_model_id,
tools: snapshot.tools.clone(),
max_iterations: snapshot.max_iterations,
parallel_tool_calls: snapshot.parallel_tool_calls,
mcp_servers: Default::default(),
};
let resolved_capability_configs =
resolve_capability_configs(&effective_overlay.capabilities, capability_registry)
.unwrap_or_else(|error| {
tracing::warn!(
error = ?error,
"failed to resolve capability configs; falling back to snapshot capabilities"
);
effective_overlay.capabilities.clone()
});
ResolvedRuntimeCapabilities {
effective_overlay,
resolved_capability_configs,
}
}
pub fn resolve_runtime_capabilities(
harness: &HarnessDefinition,
agent: Option<&AgentDefinition>,
session: &ExecutionSession,
capability_registry: &CapabilityRegistry,
) -> ResolvedRuntimeCapabilities {
let mut effective_overlay = AgentConfigOverlay::fold(
[AgentConfigOverlay::from(harness)]
.into_iter()
.chain(agent.into_iter().map(AgentConfigOverlay::from))
.chain([AgentConfigOverlay::from(session)]),
);
imply_mcp_connect_tool(&mut effective_overlay, capability_registry);
let resolved_capability_configs =
resolve_capability_configs(&effective_overlay.capabilities, capability_registry)
.unwrap_or_else(|error| {
tracing::warn!(error = ?error, "failed to resolve capability configs");
effective_overlay.capabilities.clone()
});
ResolvedRuntimeCapabilities {
effective_overlay,
resolved_capability_configs,
}
}
fn imply_mcp_connect_tool(overlay: &mut AgentConfigOverlay, registry: &CapabilityRegistry) {
use crate::mcp_server::{USER_MCP_CAPABILITY_ID, USER_MCP_CONNECT_SETTING};
if registry.get(USER_MCP_CAPABILITY_ID).is_none() {
return;
}
let signs_in_as_person = overlay
.mcp_servers
.values()
.chain(
crate::capabilities::collect_capability_mcp_servers(&overlay.capabilities, registry)
.values(),
)
.any(|server| server.acts_as.uses_user_grant());
if !signs_in_as_person {
return;
}
match overlay
.capabilities
.iter_mut()
.find(|config| config.capability_id() == USER_MCP_CAPABILITY_ID)
{
Some(config) => {
let value = config.config_mut();
if !value.is_object() {
*value = serde_json::json!({});
}
if let Some(object) = value.as_object_mut() {
object.insert(
USER_MCP_CONNECT_SETTING.into(),
serde_json::Value::Bool(true),
);
}
}
None => overlay
.capabilities
.push(AgentCapabilityConfig::with_config(
USER_MCP_CAPABILITY_ID,
serde_json::json!({ "use": false, USER_MCP_CONNECT_SETTING: true }),
)),
}
}
async fn build_runtime_agent(
snapshot: &ResolvedExecutionSnapshot,
mut effective_overlay: AgentConfigOverlay,
capability_registry: &CapabilityRegistry,
prompt_ctx: &SystemPromptContext,
mcp_tool_definitions: &[ToolDefinition],
model: &str,
) -> Result<RuntimeAgent> {
let mut runtime_agent = if let Some(ref blueprint_id) = snapshot.blueprint_id {
let blueprint = capability_registry.blueprint(blueprint_id).ok_or_else(|| {
anyhow::anyhow!(
"Unknown blueprint: \"{blueprint_id}\". Snapshot references a blueprint absent from the registry."
)
})?;
let blueprint_model = match &blueprint.model {
crate::capabilities::BlueprintModel::Fixed(model) => model.clone(),
crate::capabilities::BlueprintModel::Default(default) => snapshot
.blueprint_config
.as_ref()
.and_then(|config| config.get("model"))
.and_then(|value| value.as_str())
.map(str::to_owned)
.unwrap_or_else(|| default.clone()),
crate::capabilities::BlueprintModel::Inherit => model.to_owned(),
};
let mut prompt = blueprint.system_prompt.to_string();
if let Some(ref config) = snapshot.blueprint_config {
prompt.push_str(&format!("\n\n<config>\n{}\n</config>", config));
}
RuntimeAgentBuilder::new()
.system_prompt(&prompt)
.tools(blueprint.tool_definitions())
.model(&blueprint_model)
.max_iterations(blueprint.max_turns.unwrap_or(20))
.network_access(effective_overlay.network_access.clone())
.with_locale(prompt_ctx.locale.as_deref())
.build()
} else {
let overlay_tools = std::mem::take(&mut effective_overlay.tools);
RuntimeAgentBuilder::from_overlay(effective_overlay, capability_registry, prompt_ctx)
.await
.with_locale(prompt_ctx.locale.as_deref())
.tools(
mcp_tool_definitions
.iter()
.cloned()
.map(crate::mcp_deferred::normalize_deferred_mcp_server_definition),
)
.tools(overlay_tools)
.model(model)
.build()
};
if crate::channel_messaging::session_uses_channel_tools(&snapshot.tags) {
runtime_agent = crate::channel_messaging::apply_channel_message_mode(runtime_agent);
}
Ok(runtime_agent)
}
fn extract_locale_override(messages: &[RuntimeMessage]) -> Option<String> {
messages
.iter()
.rev()
.find(|message| message.role == RuntimeMessageRole::User)
.and_then(|message| message.controls.as_ref())
.and_then(|controls| controls.locale.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_owned)
}
#[cfg(test)]
mod tests {
use super::*;
struct CredentialCapturingDriver {
_secret: String,
}
#[async_trait::async_trait]
impl ChatDriver for CredentialCapturingDriver {
async fn chat_completion_stream(
&self,
_endpoint: &everruns_contracts::ProviderEndpoint,
_messages: Vec<crate::driver_registry::Message>,
_config: &crate::LlmCallConfig,
) -> crate::Result<crate::LlmResponseStream> {
unreachable!("debug-surface test never invokes the driver")
}
}
#[test]
fn resolved_model_execution_debug_is_credential_safe() {
let secret = "credential-marker-that-must-not-leak";
let resolved = ResolvedModelExecution {
model: "model-name".into(),
provider: crate::ProviderKey::new("provider-account"),
provider_type: DriverId::OpenAI,
driver: Arc::new(CredentialCapturingDriver {
_secret: secret.into(),
}),
provider_managed_reduction_option: None,
};
let debug = format!("{resolved:?}");
assert!(debug.contains("provider-account"));
assert!(debug.contains("<opaque>"));
assert!(!debug.contains(secret));
}
struct UserMcpStub;
impl crate::capabilities::Capability for UserMcpStub {
fn id(&self) -> &str {
crate::mcp_server::USER_MCP_CAPABILITY_ID
}
fn name(&self) -> &str {
"User MCP"
}
fn description(&self) -> &str {
"stub"
}
}
fn overlay_with_server(acts_as: &str) -> AgentConfigOverlay {
let server: crate::ScopedMcpServer = serde_json::from_value(
serde_json::json!({ "url": "https://mcp.example.com/mcp", "actsAs": acts_as }),
)
.unwrap();
AgentConfigOverlay {
mcp_servers: [("example".to_string(), server)].into_iter().collect(),
..Default::default()
}
}
fn user_mcp_config(overlay: &AgentConfigOverlay) -> Option<serde_json::Value> {
overlay
.capabilities
.iter()
.find(|config| config.capability_id() == crate::mcp_server::USER_MCP_CAPABILITY_ID)
.map(|config| config.config_value().clone())
}
#[test]
fn servers_signing_in_as_the_person_bring_the_connect_tool() {
let registry = crate::capabilities::CapabilityRegistryBuilder::new()
.capability(UserMcpStub)
.build();
for acts_as in ["user", "user_or_service"] {
let mut overlay = overlay_with_server(acts_as);
imply_mcp_connect_tool(&mut overlay, ®istry);
assert_eq!(
user_mcp_config(&overlay),
Some(serde_json::json!({ "use": false, "connect": true })),
"{acts_as}"
);
}
for acts_as in ["service", "none"] {
let mut overlay = overlay_with_server(acts_as);
imply_mcp_connect_tool(&mut overlay, ®istry);
assert_eq!(user_mcp_config(&overlay), None, "{acts_as}");
}
}
#[test]
fn an_existing_user_mcp_config_keeps_its_settings_and_gains_connect() {
let registry = crate::capabilities::CapabilityRegistryBuilder::new()
.capability(UserMcpStub)
.build();
let mut overlay = overlay_with_server("user");
overlay
.capabilities
.push(AgentCapabilityConfig::with_config(
crate::mcp_server::USER_MCP_CAPABILITY_ID,
serde_json::json!({ "manage": true }),
));
imply_mcp_connect_tool(&mut overlay, ®istry);
assert_eq!(
user_mcp_config(&overlay),
Some(serde_json::json!({ "manage": true, "connect": true }))
);
assert_eq!(overlay.capabilities.len(), 1);
}
#[test]
fn hosts_without_user_mcp_are_left_unchanged() {
let registry = crate::capabilities::CapabilityRegistryBuilder::new().build();
let mut overlay = overlay_with_server("user");
imply_mcp_connect_tool(&mut overlay, ®istry);
assert!(overlay.capabilities.is_empty());
}
}