use crate::plugin::rich_output::{transcript_text as plugin_transcript_text, view as plugin_view};
use crate::{
capability::{CapabilityBackend, CapabilityDescriptor, CapabilityId},
config::{Config, ConfigStore},
connection::{ClientChannel, IronConnection},
context::handoff::{HandoffExporter, HandoffImporter},
durable::{
ContentBlock, DurableSession, DurableToolRecord, SessionId, StructuredMessage,
TimelineEntry,
},
error::RuntimeError,
profile::{
default_identity_prompt, normalize_profile_name, AgentApproval, AgentProfile,
AgentProfileEntry, AgentProfileId, AgentProfileProvider, ProfileLoadDiagnostic,
ProfileLoadIssue, ProfileLoadReport, SkillFilter, ToolFilter, PROFILE_SCHEMA_VERSION,
},
provider_credential::domain::{ProviderAuthStatus, ProviderPromptContext, ProviderSlug},
provider_credential::store::DynCredentialStore,
runtime::{ConnectionId, IronRuntime},
stored_prompt::{
load_prompts, PromptLoadReport, StoredPrompt, StoredPromptEntry, StoredPromptRegistry,
},
tool::Tool,
};
use agent_client_protocol::schema::v1 as acp;
use futures::Stream;
use iron_providers::Provider;
use parking_lot::{Mutex, RwLock};
use std::{
cell::RefCell,
collections::HashMap,
pin::Pin,
rc::Rc,
sync::Arc,
task::{Context, Poll},
};
const PERMISSION_ALLOW_ONCE: &str = "allow_once";
const PERMISSION_REJECT_ONCE: &str = "reject_once";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PromptOutcome {
EndTurn,
Cancelled,
MaxTurnRequests,
Unknown,
}
impl From<acp::StopReason> for PromptOutcome {
fn from(reason: acp::StopReason) -> Self {
match reason {
acp::StopReason::EndTurn => PromptOutcome::EndTurn,
acp::StopReason::Cancelled => PromptOutcome::Cancelled,
acp::StopReason::MaxTurnRequests => PromptOutcome::MaxTurnRequests,
acp::StopReason::MaxTokens | acp::StopReason::Refusal | _ => PromptOutcome::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PermissionVerdict {
AllowOnce,
Deny,
Cancel,
}
#[derive(Debug, Clone)]
pub struct PermissionRequest {
pub call_id: String,
pub tool_name: String,
pub arguments: serde_json::Value,
}
#[derive(Debug, Clone)]
pub enum PromptEvent {
Status {
message: String,
},
Output {
text: String,
},
ToolCall {
call_id: String,
tool_name: String,
arguments: serde_json::Value,
},
ApprovalRequest {
call_id: String,
tool_name: String,
arguments: serde_json::Value,
},
ToolResult {
call_id: String,
tool_name: String,
status: ToolResultStatus,
result: Option<serde_json::Value>,
transcript_text: Option<String>,
view: Option<serde_json::Value>,
},
ScriptActivity {
script_id: String,
parent_call_id: String,
activity_type: ScriptActivityType,
status: ScriptActivityStatus,
detail: Option<serde_json::Value>,
},
AuthStateChange {
auth_id: String,
previous_state: crate::plugin::auth::AuthState,
new_state: crate::plugin::auth::AuthState,
},
ModelSwitched {
from_model: String,
to_model: String,
adapted: bool,
capability_diff: crate::context::CapabilityDiff,
},
ModelSwitchPending {
target_model: String,
target_provider: Option<String>,
},
CompactionStarted {
compaction_id: String,
method: String,
},
CompactionFinished {
compaction_id: String,
tokens_before: Option<u32>,
tokens_after: Option<u32>,
method: String,
},
CompactionFailed {
compaction_id: String,
reason: String,
},
Complete {
outcome: PromptOutcome,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ScriptActivityType {
ScriptStarted,
ScriptPhase,
ScriptCompleted,
ChildToolCallStarted,
ChildToolCallCompleted,
ChildToolCallFailed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ScriptActivityStatus {
Running,
Completed,
CompletedWithFailures,
Failed,
Cancelled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ToolResultStatus {
Completed,
Failed,
Denied,
Cancelled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PromptStatus {
Pending,
Running,
Completed,
Cancelled,
}
pub struct PromptHandle {
approval_resolvers:
Rc<RefCell<HashMap<String, tokio::sync::oneshot::Sender<PermissionVerdict>>>>,
session: AgentSession,
status: Rc<RefCell<PromptStatus>>,
}
impl PromptHandle {
pub fn approve(&self, call_id: &str) -> Result<(), String> {
let mut resolvers = self.approval_resolvers.borrow_mut();
match resolvers.remove(call_id) {
Some(tx) => {
let _ = tx.send(PermissionVerdict::AllowOnce);
Ok(())
}
None => Err(format!("no pending approval for call_id: {}", call_id)),
}
}
pub fn deny(&self, call_id: &str) -> Result<(), String> {
let mut resolvers = self.approval_resolvers.borrow_mut();
match resolvers.remove(call_id) {
Some(tx) => {
let _ = tx.send(PermissionVerdict::Deny);
Ok(())
}
None => Err(format!("no pending approval for call_id: {}", call_id)),
}
}
pub async fn cancel(&self) {
{
let mut resolvers = self.approval_resolvers.borrow_mut();
for (_, tx) in resolvers.drain() {
let _ = tx.send(PermissionVerdict::Cancel);
}
}
*self.status.borrow_mut() = PromptStatus::Cancelled;
self.session.cancel().await;
}
pub fn status(&self) -> PromptStatus {
*self.status.borrow()
}
}
impl std::fmt::Debug for PromptHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PromptHandle")
.field("status", &*self.status.borrow())
.finish()
}
}
pub struct PromptEvents {
rx: tokio::sync::mpsc::UnboundedReceiver<PromptEvent>,
}
impl PromptEvents {
pub async fn next(&mut self) -> Option<PromptEvent> {
self.rx.recv().await
}
pub fn try_next(&mut self) -> Option<PromptEvent> {
self.rx.try_recv().ok()
}
pub fn into_stream(self) -> Pin<Box<dyn Stream<Item = PromptEvent>>> {
Box::pin(self)
}
}
impl Stream for PromptEvents {
type Item = PromptEvent;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.rx.poll_recv(cx)
}
}
impl std::fmt::Debug for PromptEvents {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PromptEvents").finish()
}
}
struct StreamPromptState {
event_tx: tokio::sync::mpsc::UnboundedSender<PromptEvent>,
approval_resolvers:
Rc<RefCell<HashMap<String, tokio::sync::oneshot::Sender<PermissionVerdict>>>>,
tool_name_index: Rc<RefCell<HashMap<String, String>>>,
}
type ProfileRegistry = HashMap<AgentProfileId, AgentProfile>;
fn default_profile_registry() -> ProfileRegistry {
let mut registry = HashMap::new();
let default = AgentProfileEntry {
id: AgentProfileId::from("default"),
profile: AgentProfile {
name: "default".to_string(),
provider: AgentProfileProvider::RuntimeDefault,
tools: ToolFilter::Inherit,
skills: SkillFilter::Inherit,
approval: AgentApproval::PerTool,
identity_prompt: Some(default_identity_prompt().to_string()),
},
};
registry.insert(default.id, default.profile);
registry
}
pub struct IronAgent {
runtime: IronRuntime,
profile_registry: Arc<RwLock<ProfileRegistry>>,
stored_prompt_registry: Arc<RwLock<StoredPromptRegistry>>,
}
impl IronAgent {
pub fn new<P: Provider + 'static>(config: Config, provider: P) -> Self {
let runtime = IronRuntime::new(config, provider);
let profile_registry = Arc::new(RwLock::new(default_profile_registry()));
let stored_prompt_registry = Arc::new(RwLock::new(StoredPromptRegistry::new()));
runtime.set_profile_registry(profile_registry.clone());
runtime.set_stored_prompt_registry(stored_prompt_registry.clone());
Self {
runtime,
profile_registry,
stored_prompt_registry,
}
}
pub fn with_tokio_handle<P: Provider + 'static>(
config: Config,
provider: P,
handle: tokio::runtime::Handle,
) -> Self {
let runtime = IronRuntime::from_handle(config, provider, handle);
let profile_registry = Arc::new(RwLock::new(default_profile_registry()));
let stored_prompt_registry = Arc::new(RwLock::new(StoredPromptRegistry::new()));
runtime.set_profile_registry(profile_registry.clone());
runtime.set_stored_prompt_registry(stored_prompt_registry.clone());
Self {
runtime,
profile_registry,
stored_prompt_registry,
}
}
pub fn new_with_credential_store<P: Provider + 'static>(
config: Config,
provider: P,
credential_store: DynCredentialStore,
) -> Self {
let runtime = IronRuntime::new_with_credential_store(config, provider, credential_store);
let profile_registry = Arc::new(RwLock::new(default_profile_registry()));
let stored_prompt_registry = Arc::new(RwLock::new(StoredPromptRegistry::new()));
runtime.set_profile_registry(profile_registry.clone());
runtime.set_stored_prompt_registry(stored_prompt_registry.clone());
Self {
runtime,
profile_registry,
stored_prompt_registry,
}
}
pub fn runtime(&self) -> &IronRuntime {
&self.runtime
}
pub fn set_debug_sink(&self, sink: Option<Arc<dyn crate::debug::DebugSink>>) {
self.runtime.set_debug_sink(sink);
}
pub fn register_tool<T: Tool + 'static>(&self, tool: T) {
self.runtime.register_tool(tool);
}
pub fn register_builtin_tools(&self, config: &crate::builtin::BuiltinToolConfig) {
self.runtime.register_builtin_tools(config);
}
#[cfg(feature = "embedded-python")]
pub fn register_python_exec_tool(&self) {
self.runtime
.register_tool(crate::embedded_python::PythonExecTool::new());
}
pub fn register_activate_skill_tool(&self) {
self.runtime.register_activate_skill_tool();
}
pub fn tokio_handle(&self) -> &tokio::runtime::Handle {
self.runtime.tokio_handle()
}
pub fn register_capability(&self, descriptor: CapabilityDescriptor) {
self.runtime.register_capability(descriptor);
}
pub fn set_capability_backend(&self, id: CapabilityId, backend: CapabilityBackend) {
self.runtime.set_capability_backend(id, backend);
}
pub fn register_mcp_server(&self, config: crate::mcp::McpServerConfig) {
self.runtime.register_mcp_server(config);
}
pub fn mcp_registry(&self) -> parking_lot::RwLockReadGuard<'_, crate::mcp::McpServerRegistry> {
self.runtime.mcp_registry()
}
pub fn register_plugin(&self, config: crate::plugin::config::PluginConfig) {
self.runtime.register_plugin(config);
}
pub fn plugin_registry(
&self,
) -> parking_lot::RwLockReadGuard<'_, crate::plugin::registry::PluginRegistry> {
self.runtime.plugin_registry()
}
pub fn get_effective_tools(&self, session_id: SessionId) -> Vec<crate::tool::ToolDefinition> {
self.runtime.get_effective_tool_definitions(session_id)
}
pub fn get_plugin_inventory(&self) -> Vec<crate::plugin::status::PluginInfo> {
self.runtime.get_plugin_inventory()
}
pub fn get_auth_prompts(&self) -> Vec<crate::plugin::auth::AuthPrompt> {
self.runtime.get_auth_prompts()
}
pub fn get_plugin_status(
&self,
plugin_id: &str,
) -> Option<crate::plugin::status::PluginStatus> {
self.runtime.get_plugin_status(plugin_id)
}
pub fn set_plugin_credentials(
&self,
plugin_id: &str,
credentials: crate::plugin::auth::CredentialBinding,
) {
self.runtime.set_plugin_credentials(plugin_id, credentials);
}
pub fn clear_plugin_credentials(&self, plugin_id: &str) {
self.runtime.clear_plugin_credentials(plugin_id);
}
pub fn start_auth_flow(
&self,
plugin_id: &str,
) -> Result<crate::plugin::auth::AuthInteractionRequest, String> {
self.runtime.begin_plugin_auth_flow(plugin_id)
}
pub fn complete_auth_flow(
&self,
plugin_id: &str,
response: crate::plugin::auth::AuthInteractionResponse,
) -> Result<crate::plugin::auth::AuthStatusTransition, String> {
self.runtime.complete_plugin_auth_flow(plugin_id, response)
}
pub fn get_plugin_availability(
&self,
plugin_id: &str,
) -> Option<crate::plugin::registry::PluginAvailabilitySummary> {
self.runtime.get_plugin_availability(plugin_id)
}
pub fn get_session_tool_diagnostics(
&self,
session_id: SessionId,
) -> Option<Vec<crate::mcp::session_catalog::ToolDiagnostic>> {
self.runtime.get_session_tool_diagnostics(session_id)
}
pub fn default_profile(&self) -> AgentProfileEntry {
AgentProfileEntry {
id: AgentProfileId::from("default"),
profile: AgentProfile {
name: "default".to_string(),
provider: AgentProfileProvider::RuntimeDefault,
tools: ToolFilter::Inherit,
skills: SkillFilter::Inherit,
approval: AgentApproval::PerTool,
identity_prompt: Some(default_identity_prompt().to_string()),
},
}
}
pub async fn seed_default_profiles(
&self,
store: &ConfigStore,
policy: crate::profile::DefaultProfileSeedPolicy,
) -> Result<crate::profile::DefaultProfileSeedReport, crate::config::ConfigError> {
crate::profile::seed_default_profiles(store, policy).await
}
pub async fn load_provider_profiles(
&self,
store: &ConfigStore,
) -> Result<(), crate::config::ConfigError> {
self.runtime.load_provider_profiles(store).await
}
pub fn register_profile(
&self,
id: AgentProfileId,
profile: AgentProfile,
) -> Result<(), String> {
if id.as_str().eq_ignore_ascii_case("default") {
return Err("profile ID 'default' is reserved".to_string());
}
if !crate::profile::is_valid_profile_id(id.as_str()) {
return Err(format!("invalid profile ID: '{}'", id.as_str()));
}
let name = normalize_profile_name(&profile.name)
.ok_or_else(|| format!("invalid profile name: '{}'", profile.name))?;
if name.eq_ignore_ascii_case("default") {
return Err("profile name 'default' is reserved".to_string());
}
if !matches!(
profile.approval,
AgentApproval::PerTool | AgentApproval::AutoApprove
) {
return Err(format!(
"invalid profile approval value: {:?}. Only PerTool and AutoApprove are supported.",
profile.approval
));
}
let mut registry = self.profile_registry.write();
let name_conflict = registry.iter().any(|(existing_id, existing_profile)| {
existing_id.as_str() != id.as_str() && existing_profile.name == name
});
if name_conflict {
return Err(format!(
"profile name '{}' is already used by another profile",
name
));
}
let mut profile = profile;
profile.name = name;
registry.insert(id, profile);
Ok(())
}
pub fn unregister_profile(&self, id: &AgentProfileId) -> bool {
if id.as_str().eq_ignore_ascii_case("default") {
return false;
}
self.profile_registry.write().remove(id).is_some()
}
pub fn list_profiles(&self) -> Vec<AgentProfileEntry> {
let registry = self.profile_registry.read();
let mut ids: Vec<&AgentProfileId> = registry.keys().collect();
ids.sort_by(|a, b| a.as_str().cmp(b.as_str()));
ids.into_iter()
.map(|id| AgentProfileEntry {
id: id.clone(),
profile: registry[id].clone(),
})
.collect()
}
pub async fn load_profiles(
&self,
store: &ConfigStore,
) -> Result<ProfileLoadReport, crate::config::ConfigError> {
let record_ids = store.list_profile_ids().await?;
let mut record_ids = record_ids;
record_ids.sort();
let mut report = ProfileLoadReport::empty();
for record_id in record_ids {
let record = match store.get_profile(&record_id).await? {
Some(record) => record,
None => {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: AgentProfileId::from(record_id),
name: None,
issue: ProfileLoadIssue::MissingRecord,
});
continue;
}
};
if record.schema_version != PROFILE_SCHEMA_VERSION {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: AgentProfileId::from(record.id.clone()),
name: None,
issue: ProfileLoadIssue::UnsupportedSchemaVersion {
version: record.schema_version,
},
});
continue;
}
let profile: AgentProfile = match serde_json::from_value(record.payload.clone()) {
Ok(profile) => profile,
Err(_) => {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: AgentProfileId::from(record.id.clone()),
name: None,
issue: ProfileLoadIssue::InvalidPayload,
});
continue;
}
};
let profile_id = AgentProfileId::from(record.id.clone());
if profile_id.as_str().eq_ignore_ascii_case("default")
|| profile.name.trim().eq_ignore_ascii_case("default")
{
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: profile_id.clone(),
name: Some(profile.name.clone()),
issue: ProfileLoadIssue::ReservedDefault,
});
continue;
}
if !crate::profile::is_valid_profile_id(profile_id.as_str()) {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: profile_id.clone(),
name: Some(profile.name.clone()),
issue: ProfileLoadIssue::InvalidProfileId,
});
continue;
}
if normalize_profile_name(&profile.name).is_none() {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: profile_id.clone(),
name: Some(profile.name.clone()),
issue: ProfileLoadIssue::InvalidName,
});
continue;
}
let mut registry = self.profile_registry.write();
let normalized_name =
normalize_profile_name(&profile.name).unwrap_or_else(|| profile.name.clone());
let conflict = registry.iter().any(|(existing_id, existing_profile)| {
existing_id.as_str() != profile_id.as_str()
&& existing_profile.name == normalized_name
});
if conflict {
report.diagnostics.push(ProfileLoadDiagnostic {
profile_id: profile_id.clone(),
name: Some(profile.name.clone()),
issue: ProfileLoadIssue::DuplicateName,
});
continue;
}
let mut profile = profile;
profile.name = normalized_name;
registry.insert(profile_id.clone(), profile.clone());
report.loaded.push(AgentProfileEntry {
id: profile_id,
profile,
});
}
Ok(report)
}
pub fn get_profile(&self, id: &AgentProfileId) -> Option<AgentProfileEntry> {
self.profile_registry
.read()
.get(id)
.map(|profile| AgentProfileEntry {
id: id.clone(),
profile: profile.clone(),
})
}
pub fn register_stored_prompt(
&self,
id: impl Into<String>,
prompt: StoredPrompt,
) -> Result<(), String> {
self.stored_prompt_registry
.write()
.register(id.into(), prompt)
}
pub fn unregister_stored_prompt(&self, id: &str) -> bool {
self.stored_prompt_registry.write().unregister(id.trim())
}
pub async fn delete_stored_prompt(
&self,
store: &ConfigStore,
id: &str,
) -> Result<(), crate::config::ConfigError> {
let normalized_id = id.trim();
store.delete_prompt(normalized_id).await?;
self.stored_prompt_registry
.write()
.unregister(normalized_id);
Ok(())
}
pub fn list_stored_prompts(&self) -> Vec<StoredPromptEntry> {
self.stored_prompt_registry.read().list()
}
pub fn get_stored_prompt(&self, id: &str) -> Option<StoredPromptEntry> {
let normalized_id = id.trim();
self.stored_prompt_registry
.read()
.get(normalized_id)
.cloned()
.map(|prompt| StoredPromptEntry {
id: normalized_id.to_string(),
prompt,
identity_state: crate::stored_prompt::IdentityState::Ready,
})
}
pub async fn load_stored_prompts(
&self,
store: &ConfigStore,
) -> Result<PromptLoadReport, crate::config::ConfigError> {
let report = load_prompts(store).await?;
let mut registry = self.stored_prompt_registry.write();
for entry in &report.loaded {
registry
.register(entry.id.clone(), entry.prompt.clone())
.map_err(crate::config::ConfigError::Validation)?;
}
Ok(report)
}
pub fn connect(&self) -> AgentConnection {
AgentConnection::new(self.runtime.clone(), self.profile_registry.clone())
}
pub fn shutdown(&self) {
self.runtime.shutdown();
}
}
impl Clone for IronAgent {
fn clone(&self) -> Self {
Self {
runtime: self.runtime.clone(),
profile_registry: self.profile_registry.clone(),
stored_prompt_registry: self.stored_prompt_registry.clone(),
}
}
}
type AsyncPermissionHandler =
Box<dyn Fn(PermissionRequest) -> Pin<Box<dyn std::future::Future<Output = PermissionVerdict>>>>;
type SyncPermissionHandler = Rc<RefCell<Option<Box<dyn Fn(&str) -> PermissionVerdict>>>>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HiddenSessionInfo {
pub session_id: SessionId,
pub connection_id: ConnectionId,
pub parent_session_id: Option<SessionId>,
}
pub struct AgentConnection {
inner: Rc<IronConnection>,
permission_handler: SyncPermissionHandler,
async_permission_handler: Rc<RefCell<Option<AsyncPermissionHandler>>>,
active_streams: Rc<RefCell<HashMap<String, StreamPromptState>>>,
}
impl AgentConnection {
fn emit_auth_transition_to_all_streams(
&self,
auth_id: &str,
previous_state: crate::plugin::auth::AuthState,
new_state: crate::plugin::auth::AuthState,
) {
let streams = self.active_streams.borrow();
for (_, state) in streams.iter() {
let _ = state.event_tx.send(PromptEvent::AuthStateChange {
auth_id: auth_id.to_string(),
previous_state,
new_state,
});
}
}
fn new(runtime: IronRuntime, profile_registry: Arc<RwLock<ProfileRegistry>>) -> Self {
let inner = Rc::new(IronConnection::new_with_profile_registry(
runtime,
profile_registry,
));
let permission_handler = Rc::new(RefCell::new(None));
let async_permission_handler = Rc::new(RefCell::new(None));
let active_streams = Rc::new(RefCell::new(HashMap::new()));
let client: Rc<dyn ClientChannel> = Rc::new(FacadeClientChannel {
permission_handler: permission_handler.clone(),
async_permission_handler: async_permission_handler.clone(),
active_streams: active_streams.clone(),
});
inner.set_client(client);
Self {
inner,
permission_handler,
async_permission_handler,
active_streams,
}
}
pub fn id(&self) -> ConnectionId {
self.inner.id()
}
pub fn on_permission(&self, handler: impl Fn(&str) -> PermissionVerdict + 'static) {
*self.permission_handler.borrow_mut() = Some(Box::new(handler));
}
pub fn on_permission_async(
&self,
handler: impl Fn(PermissionRequest) -> Pin<Box<dyn std::future::Future<Output = PermissionVerdict>>>
+ 'static,
) {
*self.async_permission_handler.borrow_mut() = Some(Box::new(handler));
}
pub fn create_session(&self) -> Result<AgentSession, RuntimeError> {
let connection_id = self.inner.id();
let (session_id, durable) = self.inner.runtime().create_session(connection_id)?;
Ok(AgentSession {
id: session_id,
durable,
connection: self.inner.clone(),
active_streams: self.active_streams.clone(),
profile_id: None,
})
}
pub fn create_session_with_profile(
&self,
profile_id: AgentProfileId,
) -> Result<AgentSession, RuntimeError> {
let selected_profile = {
let registry = self.inner.profile_registry().read();
registry
.get(&profile_id)
.cloned()
.ok_or_else(|| RuntimeError::Session {
message: format!("Profile '{}' not found in registry", profile_id.as_str()),
})?
};
let connection_id = self.inner.id();
let (session_id, durable) = self.inner.runtime().create_session(connection_id)?;
{
let mut session = durable.lock();
session.profile_id = Some(profile_id.clone());
session.effective_tool_filter = Some(selected_profile.tools.clone());
session.effective_approval = Some(selected_profile.approval);
if let Some(ref identity) = selected_profile.identity_prompt {
if !identity.trim().is_empty() {
session.set_profile_identity(identity.clone());
}
}
}
Ok(AgentSession {
id: session_id,
durable,
connection: self.inner.clone(),
active_streams: self.active_streams.clone(),
profile_id: Some(profile_id),
})
}
pub fn close_session(&self, session: &AgentSession) -> Result<(), RuntimeError> {
let owner = self.inner.runtime().get_session_connection(session.id);
if owner != Some(self.inner.id()) {
return Err(RuntimeError::Connection(
"session not owned by this connection".into(),
));
}
self.inner.runtime().close_session(session.id);
Ok(())
}
pub fn active_sessions(&self) -> Vec<SessionId> {
self.inner
.runtime()
.sessions_for_connection(self.inner.id(), false)
}
pub fn active_sessions_include_hidden(&self) -> Vec<SessionId> {
self.inner
.runtime()
.sessions_for_connection(self.inner.id(), true)
}
pub fn hidden_sessions(&self) -> Vec<HiddenSessionInfo> {
let runtime = self.inner.runtime();
let connection_id = self.inner.id();
runtime
.list_hidden_sessions()
.into_iter()
.filter(|(_, child_connection_id, parent_session_id)| {
*child_connection_id == connection_id
|| parent_session_id.and_then(|parent| runtime.get_session_connection(parent))
== Some(connection_id)
})
.map(
|(session_id, connection_id, parent_session_id)| HiddenSessionInfo {
session_id,
connection_id,
parent_session_id,
},
)
.collect()
}
pub fn child_sessions(&self, parent: &AgentSession) -> Result<Vec<SessionId>, RuntimeError> {
let owner = self.inner.runtime().get_session_connection(parent.id);
if owner != Some(self.inner.id()) {
return Err(RuntimeError::Connection(
"session not owned by this connection".into(),
));
}
Ok(self.inner.runtime().list_child_sessions(parent.id))
}
pub fn start_auth_flow(
&self,
plugin_id: &str,
) -> Result<crate::plugin::auth::AuthInteractionRequest, String> {
let previous_state = self
.inner
.runtime()
.get_plugin_status(plugin_id)
.map(|status| status.auth.state)
.unwrap_or(crate::plugin::auth::AuthState::Unauthenticated);
let request = self.inner.runtime().begin_plugin_auth_flow(plugin_id)?;
self.emit_auth_transition_to_all_streams(
plugin_id,
previous_state,
crate::plugin::auth::AuthState::Authenticating,
);
Ok(request)
}
pub fn complete_auth_flow(
&self,
plugin_id: &str,
response: crate::plugin::auth::AuthInteractionResponse,
) -> Result<crate::plugin::auth::AuthStatusTransition, String> {
let transition = self
.inner
.runtime()
.complete_plugin_auth_flow(plugin_id, response)?;
self.emit_auth_transition_to_all_streams(
&transition.auth_id,
transition.previous_state,
transition.new_state,
);
Ok(transition)
}
pub fn get_auth_prompts(&self) -> Vec<crate::plugin::auth::AuthPrompt> {
self.inner.runtime().get_auth_prompts()
}
pub fn create_session_from_handoff(
&self,
bundle: crate::context::HandoffBundle,
) -> Result<AgentSession, RuntimeError> {
let connection_id = self.inner.id();
let durable = crate::context::HandoffImporter::hydrate_into_new(bundle);
let session_id = durable.id;
self.inner
.runtime()
.insert_session(session_id, durable, connection_id)?;
let durable = self
.inner
.runtime()
.get_session(session_id)
.ok_or_else(|| RuntimeError::Connection("Failed to retrieve created session".into()))?;
Ok(AgentSession {
id: session_id,
durable,
connection: self.inner.clone(),
active_streams: self.active_streams.clone(),
profile_id: None,
})
}
}
struct FacadeClientChannel {
permission_handler: SyncPermissionHandler,
async_permission_handler: Rc<RefCell<Option<AsyncPermissionHandler>>>,
active_streams: Rc<RefCell<HashMap<String, StreamPromptState>>>,
}
impl FacadeClientChannel {
fn emit_stream_event(&self, event: PromptEvent) {
let streams = self.active_streams.borrow();
for (_, state) in streams.iter() {
let _ = state.event_tx.send(event.clone());
}
}
}
impl ClientChannel for FacadeClientChannel {
fn send_notification(
&self,
notification: acp::SessionNotification,
) -> Pin<Box<dyn std::future::Future<Output = agent_client_protocol::Result<()>>>> {
let session_key = notification.session_id.to_string();
let streams = self.active_streams.borrow();
if let Some(state) = streams.get(&session_key) {
if let Some(prompt_event) = convert_notification_to_prompt_event_with_index(
¬ification,
&state.tool_name_index,
) {
let _ = state.event_tx.send(prompt_event);
}
}
drop(streams);
Box::pin(async { Ok(()) })
}
fn emit_script_activity(
&self,
script_id: &str,
parent_call_id: &str,
activity_type: &str,
status: &str,
detail: Option<serde_json::Value>,
) -> Pin<Box<dyn std::future::Future<Output = ()>>> {
let activity = match activity_type {
"script_started" => ScriptActivityType::ScriptStarted,
"script_phase" => ScriptActivityType::ScriptPhase,
"script_completed" => ScriptActivityType::ScriptCompleted,
"child_tool_call_started" => ScriptActivityType::ChildToolCallStarted,
"child_tool_call_completed" => ScriptActivityType::ChildToolCallCompleted,
"child_tool_call_failed" => ScriptActivityType::ChildToolCallFailed,
_ => return Box::pin(async {}),
};
let act_status = match status {
"running" => ScriptActivityStatus::Running,
"completed" => ScriptActivityStatus::Completed,
"completed_with_failures" => ScriptActivityStatus::CompletedWithFailures,
"failed" => ScriptActivityStatus::Failed,
"cancelled" => ScriptActivityStatus::Cancelled,
_ => return Box::pin(async {}),
};
self.emit_stream_event(PromptEvent::ScriptActivity {
script_id: script_id.to_string(),
parent_call_id: parent_call_id.to_string(),
activity_type: activity,
status: act_status,
detail,
});
Box::pin(async {})
}
fn emit_compaction_event(
&self,
event_type: &str,
tokens_before: Option<u32>,
tokens_after: Option<u32>,
method: &str,
reason: Option<&str>,
compaction_id: &str,
) -> Pin<Box<dyn std::future::Future<Output = ()>>> {
let event = match event_type {
"started" => Some(PromptEvent::CompactionStarted {
compaction_id: compaction_id.to_string(),
method: method.to_string(),
}),
"finished" => Some(PromptEvent::CompactionFinished {
compaction_id: compaction_id.to_string(),
tokens_before,
tokens_after,
method: method.to_string(),
}),
"failed" => reason.map(|r| PromptEvent::CompactionFailed {
compaction_id: compaction_id.to_string(),
reason: r.to_string(),
}),
_ => None,
};
if let Some(evt) = event {
self.emit_stream_event(evt);
}
Box::pin(async {})
}
fn request_permission(
&self,
request: acp::RequestPermissionRequest,
) -> Pin<
Box<
dyn std::future::Future<
Output = agent_client_protocol::Result<acp::RequestPermissionResponse>,
>,
>,
> {
let call_id = request.tool_call.tool_call_id.to_string();
let tool_name = request.tool_call.fields.title.clone().unwrap_or_default();
let arguments = request
.tool_call
.fields
.raw_input
.clone()
.unwrap_or_default();
let session_key = request.session_id.to_string();
let stream_state = self.active_streams.borrow();
if let Some(state) = stream_state.get(&session_key) {
let (tx, rx) = tokio::sync::oneshot::channel();
state
.approval_resolvers
.borrow_mut()
.insert(call_id.clone(), tx);
let event_tx = state.event_tx.clone();
drop(stream_state);
return Box::pin(async move {
let _ = event_tx.send(PromptEvent::ApprovalRequest {
call_id,
tool_name,
arguments,
});
match rx.await {
Ok(verdict) => Ok(acp::RequestPermissionResponse::new(verdict_to_outcome(
verdict,
))),
Err(_) => Ok(acp::RequestPermissionResponse::new(verdict_to_outcome(
PermissionVerdict::Deny,
))),
}
});
}
drop(stream_state);
let async_handler = self.async_permission_handler.borrow();
if let Some(ref handler) = *async_handler {
let perm_req = PermissionRequest {
call_id,
tool_name,
arguments,
};
let future = handler(perm_req);
drop(async_handler);
return Box::pin(async move {
let verdict = future.await;
Ok(acp::RequestPermissionResponse::new(verdict_to_outcome(
verdict,
)))
});
}
drop(async_handler);
let handler = self.permission_handler.borrow();
let verdict = handler
.as_ref()
.map(|h| h(&call_id))
.unwrap_or(PermissionVerdict::AllowOnce);
drop(handler);
Box::pin(async move {
Ok(acp::RequestPermissionResponse::new(verdict_to_outcome(
verdict,
)))
})
}
}
fn verdict_to_outcome(verdict: PermissionVerdict) -> acp::RequestPermissionOutcome {
match verdict {
PermissionVerdict::AllowOnce => {
acp::RequestPermissionOutcome::Selected(acp::SelectedPermissionOutcome::new(
acp::PermissionOptionId::new(PERMISSION_ALLOW_ONCE),
))
}
PermissionVerdict::Deny => {
acp::RequestPermissionOutcome::Selected(acp::SelectedPermissionOutcome::new(
acp::PermissionOptionId::new(PERMISSION_REJECT_ONCE),
))
}
PermissionVerdict::Cancel => acp::RequestPermissionOutcome::Cancelled,
}
}
fn convert_notification_to_prompt_event_with_index(
notification: &acp::SessionNotification,
tool_name_index: &Rc<RefCell<HashMap<String, String>>>,
) -> Option<PromptEvent> {
match ¬ification.update {
acp::SessionUpdate::AgentMessageChunk(chunk) => match &chunk.content {
acp::ContentBlock::Text(tc) => Some(PromptEvent::Output {
text: tc.text.clone(),
}),
_ => None,
},
acp::SessionUpdate::ToolCall(tc) => {
let call_id = tc.tool_call_id.to_string();
let tool_name = tc.title.clone();
let arguments = tc.raw_input.clone().unwrap_or_default();
tool_name_index
.borrow_mut()
.insert(call_id.clone(), tool_name.clone());
Some(PromptEvent::ToolCall {
call_id,
tool_name,
arguments,
})
}
acp::SessionUpdate::ToolCallUpdate(update) => {
let call_id = update.tool_call_id.to_string();
let tool_name = update
.fields
.title
.clone()
.or_else(|| tool_name_index.borrow().get(&call_id).cloned())
.unwrap_or_default();
let result = update.fields.raw_output.clone();
let transcript_text = result
.as_ref()
.and_then(plugin_transcript_text)
.map(str::to_string);
let view = result.as_ref().and_then(plugin_view).cloned();
let status = match update.fields.status {
Some(acp::ToolCallStatus::Completed) => ToolResultStatus::Completed,
Some(acp::ToolCallStatus::Failed) => {
let is_denied = result.as_ref().is_some_and(|r| {
r.get("error")
.and_then(|v| v.as_str())
.is_some_and(|s| s.contains("denied"))
});
if is_denied {
ToolResultStatus::Denied
} else {
ToolResultStatus::Failed
}
}
Some(acp::ToolCallStatus::Pending) => {
return None;
}
Some(acp::ToolCallStatus::InProgress) => {
return None;
}
_ => return None,
};
Some(PromptEvent::ToolResult {
call_id,
tool_name,
status,
result,
transcript_text,
view,
})
}
_ => None,
}
}
pub struct AgentSession {
id: SessionId,
durable: Arc<Mutex<DurableSession>>,
connection: Rc<IronConnection>,
active_streams: Rc<RefCell<HashMap<String, StreamPromptState>>>,
profile_id: Option<AgentProfileId>,
}
impl AgentSession {
fn emit_auth_transition_to_stream(
&self,
auth_id: &str,
previous_state: crate::plugin::auth::AuthState,
new_state: crate::plugin::auth::AuthState,
) {
let session_key = self.id.to_string();
let streams = self.active_streams.borrow();
if let Some(state) = streams.get(&session_key) {
let _ = state.event_tx.send(PromptEvent::AuthStateChange {
auth_id: auth_id.to_string(),
previous_state,
new_state,
});
}
}
fn emit_model_switch_event(&self, event: PromptEvent) {
let session_key = self.id.to_string();
let streams = self.active_streams.borrow();
if let Some(state) = streams.get(&session_key) {
let _ = state.event_tx.send(event);
}
}
pub fn id(&self) -> SessionId {
self.id
}
pub async fn prompt(&self, text: &str) -> PromptOutcome {
self.try_prompt(text)
.await
.unwrap_or(PromptOutcome::EndTurn)
}
pub async fn try_prompt(&self, text: &str) -> Result<PromptOutcome, String> {
let acp_session_id = acp::SessionId::new(self.id.to_string());
let request = acp::PromptRequest::new(
acp_session_id,
vec![acp::ContentBlock::Text(acp::TextContent::new(text))],
);
match self
.connection
.handle_prompt(request, self.profile_id.as_ref())
.await
{
Ok(response) => Ok(response.stop_reason.into()),
Err(e) => Err(format!("prompt failed: {}", e)),
}
}
pub fn profile_unavailable(&self) -> Option<String> {
self.durable.lock().profile_unavailable.clone()
}
pub async fn prompt_with_blocks(&self, blocks: &[ContentBlock]) -> PromptOutcome {
let acp_session_id = acp::SessionId::new(self.id.to_string());
let acp_blocks: Vec<_> = blocks.iter().map(to_acp_content_block).collect();
let request = acp::PromptRequest::new(acp_session_id, acp_blocks);
match self
.connection
.handle_prompt(request, self.profile_id.as_ref())
.await
{
Ok(response) => response.stop_reason.into(),
Err(_) => PromptOutcome::EndTurn,
}
}
pub async fn prompt_managed(
&self,
text: &str,
provider_context: ProviderPromptContext,
) -> PromptOutcome {
let acp_session_id = acp::SessionId::new(self.id.to_string());
let request = acp::PromptRequest::new(
acp_session_id,
vec![acp::ContentBlock::Text(acp::TextContent::new(text))],
);
match self
.connection
.handle_prompt_managed(request, provider_context, self.profile_id.as_ref())
.await
{
Ok(response) => response.stop_reason.into(),
Err(_) => PromptOutcome::EndTurn,
}
}
pub async fn prompt_with_blocks_managed(
&self,
blocks: &[ContentBlock],
provider_context: ProviderPromptContext,
) -> PromptOutcome {
let acp_session_id = acp::SessionId::new(self.id.to_string());
let acp_blocks: Vec<_> = blocks.iter().map(to_acp_content_block).collect();
let request = acp::PromptRequest::new(acp_session_id, acp_blocks);
match self
.connection
.handle_prompt_managed(request, provider_context, self.profile_id.as_ref())
.await
{
Ok(response) => response.stop_reason.into(),
Err(_) => PromptOutcome::EndTurn,
}
}
pub fn prompt_stream(&self, text: &str) -> (PromptHandle, PromptEvents) {
let acp_blocks = vec![acp::ContentBlock::Text(acp::TextContent::new(text))];
self.prompt_stream_with_acp_blocks(acp_blocks)
}
pub fn prompt_stream_with_blocks(
&self,
blocks: &[ContentBlock],
) -> (PromptHandle, PromptEvents) {
let acp_blocks: Vec<_> = blocks.iter().map(to_acp_content_block).collect();
self.prompt_stream_with_acp_blocks(acp_blocks)
}
fn prompt_stream_with_acp_blocks(
&self,
acp_blocks: Vec<acp::ContentBlock>,
) -> (PromptHandle, PromptEvents) {
let (event_tx, event_rx) = tokio::sync::mpsc::unbounded_channel();
let approval_resolvers: Rc<
RefCell<HashMap<String, tokio::sync::oneshot::Sender<PermissionVerdict>>>,
> = Rc::new(RefCell::new(HashMap::new()));
let prompt_key = self.id.to_string();
self.active_streams.borrow_mut().insert(
prompt_key.clone(),
StreamPromptState {
event_tx: event_tx.clone(),
approval_resolvers: approval_resolvers.clone(),
tool_name_index: Rc::new(RefCell::new(HashMap::new())),
},
);
let status = Rc::new(RefCell::new(PromptStatus::Running));
let handle = PromptHandle {
approval_resolvers,
session: self.clone(),
status: status.clone(),
};
let events = PromptEvents { rx: event_rx };
let acp_session_id = acp::SessionId::new(self.id.to_string());
let request = acp::PromptRequest::new(acp_session_id, acp_blocks);
let connection = self.connection.clone();
let active_streams = self.active_streams.clone();
let status_cell = status.clone();
let profile_id = self.profile_id.clone();
tokio::task::spawn_local(async move {
let outcome = match connection.handle_prompt(request, profile_id.as_ref()).await {
Ok(response) => response.stop_reason.into(),
Err(_) => PromptOutcome::EndTurn,
};
{
let mut streams = active_streams.borrow_mut();
if let Some(state) = streams.remove(&prompt_key) {
let _ = state.event_tx.send(PromptEvent::Complete { outcome });
}
}
*status_cell.borrow_mut() = match outcome {
PromptOutcome::Cancelled => PromptStatus::Cancelled,
_ => PromptStatus::Completed,
};
});
(handle, events)
}
pub fn prompt_stream_managed(
&self,
text: &str,
provider_context: ProviderPromptContext,
) -> (PromptHandle, PromptEvents) {
let acp_blocks = vec![acp::ContentBlock::Text(acp::TextContent::new(text))];
self.prompt_stream_with_acp_blocks_managed(acp_blocks, provider_context)
}
pub fn prompt_stream_with_blocks_managed(
&self,
blocks: &[ContentBlock],
provider_context: ProviderPromptContext,
) -> (PromptHandle, PromptEvents) {
let acp_blocks: Vec<_> = blocks.iter().map(to_acp_content_block).collect();
self.prompt_stream_with_acp_blocks_managed(acp_blocks, provider_context)
}
fn prompt_stream_with_acp_blocks_managed(
&self,
acp_blocks: Vec<acp::ContentBlock>,
provider_context: ProviderPromptContext,
) -> (PromptHandle, PromptEvents) {
let (event_tx, event_rx) = tokio::sync::mpsc::unbounded_channel();
let approval_resolvers: Rc<
RefCell<HashMap<String, tokio::sync::oneshot::Sender<PermissionVerdict>>>,
> = Rc::new(RefCell::new(HashMap::new()));
let prompt_key = self.id.to_string();
self.active_streams.borrow_mut().insert(
prompt_key.clone(),
StreamPromptState {
event_tx: event_tx.clone(),
approval_resolvers: approval_resolvers.clone(),
tool_name_index: Rc::new(RefCell::new(HashMap::new())),
},
);
let status = Rc::new(RefCell::new(PromptStatus::Running));
let handle = PromptHandle {
approval_resolvers,
session: self.clone(),
status: status.clone(),
};
let events = PromptEvents { rx: event_rx };
let acp_session_id = acp::SessionId::new(self.id.to_string());
let request = acp::PromptRequest::new(acp_session_id, acp_blocks);
let connection = self.connection.clone();
let active_streams = self.active_streams.clone();
let status_cell = status.clone();
let profile_id = self.profile_id.clone();
tokio::task::spawn_local(async move {
let outcome = match connection
.handle_prompt_managed(request, provider_context, profile_id.as_ref())
.await
{
Ok(response) => response.stop_reason.into(),
Err(_) => PromptOutcome::EndTurn,
};
{
let mut streams = active_streams.borrow_mut();
if let Some(state) = streams.remove(&prompt_key) {
let _ = state.event_tx.send(PromptEvent::Complete { outcome });
}
}
*status_cell.borrow_mut() = match outcome {
PromptOutcome::Cancelled => PromptStatus::Cancelled,
_ => PromptStatus::Completed,
};
});
(handle, events)
}
pub async fn cancel(&self) {
let acp_session_id = acp::SessionId::new(self.id.to_string());
let notification = acp::CancelNotification::new(acp_session_id);
let _ = self.connection.handle_cancel(notification).await;
}
pub fn timeline(&self) -> Vec<TimelineEntry> {
self.durable.lock().timeline.clone()
}
pub fn messages(&self) -> Vec<StructuredMessage> {
self.durable.lock().messages.clone()
}
pub fn tool_records(&self) -> Vec<DurableToolRecord> {
self.durable.lock().tool_records.clone()
}
pub fn is_empty(&self) -> bool {
self.durable.lock().is_empty()
}
pub fn set_instructions(&self, instructions: impl Into<String>) {
self.durable.lock().set_instructions(instructions);
}
pub fn set_workspace_roots(
&self,
roots: Vec<std::path::PathBuf>,
) -> Result<bool, RuntimeError> {
self.connection
.runtime()
.set_session_workspace_roots(self.id, roots)
}
pub async fn provider_auth_status(&self, slug: &ProviderSlug) -> Option<ProviderAuthStatus> {
self.connection.runtime().provider_auth_status(slug).await
}
pub async fn provider_auth_status_with_api_key(
&self,
slug: &ProviderSlug,
api_key: Option<&str>,
) -> Option<ProviderAuthStatus> {
self.connection
.runtime()
.provider_auth_status_with_api_key(slug, api_key)
.await
}
pub async fn provider_auth_status_for_context(
&self,
context: &ProviderPromptContext,
) -> Option<ProviderAuthStatus> {
self.connection
.runtime()
.provider_auth_status_for_context(context)
.await
}
pub async fn disconnect_provider_oauth(&self, slug: &ProviderSlug) {
self.connection
.runtime()
.disconnect_provider_oauth(slug)
.await;
}
pub fn refresh_skill_catalog(&self) -> Vec<crate::skill::SkillDiagnostic> {
let diagnostics = self.connection.runtime().refresh_skill_catalog();
let available_skills = self.connection.runtime().available_skill_snapshot();
self.durable.lock().set_available_skills(available_skills);
diagnostics
}
pub fn list_available_skills(&self) -> Vec<crate::skill::SkillMetadata> {
self.durable
.lock()
.list_available_skills()
.iter()
.map(|skill| skill.metadata.clone())
.collect()
}
pub fn activate_skill(&self, name: &str) -> Result<(), String> {
let skill = self
.durable
.lock()
.load_available_skill(name)
.ok_or_else(|| format!("Skill '{}' not found", name))?;
self.durable.lock().activate_skill(
&skill.metadata.id,
&skill.body,
skill.resources.clone(),
);
Ok(())
}
pub fn deactivate_skill(&self, name: &str) {
self.durable.lock().deactivate_skill(name);
}
pub fn list_active_skills(&self) -> Vec<String> {
self.durable
.lock()
.list_active_skills()
.into_iter()
.map(|s| s.to_string())
.collect()
}
pub fn active_context(
&self,
tool_registry: &crate::tool::ToolRegistry,
current_prompt: Option<&str>,
context_window_hint: Option<usize>,
) -> crate::context::ActiveContextSnapshot {
let session = self.durable.lock();
let tail = session.to_transcript();
let mut snapshot = crate::context::ContextTelemetry::for_session(
session.instruction_text_for_estimate().as_deref(),
&session.compressed_blocks,
&tail.messages,
tool_registry,
current_prompt,
context_window_hint,
crate::context::SessionModelInfo {
current_model: session.current_model.as_deref(),
model_switch_count: session.model_switch_history.len(),
},
Some(&session.token_tracker),
);
let context_config = &self.connection.runtime().config().context_management;
if context_config.enabled {
snapshot.compact_threshold_tokens = Some(context_config.maintenance_threshold);
}
snapshot
}
pub fn is_idle(&self) -> bool {
let durable_idle = self.durable.lock().is_idle();
let has_active_prompt = self.connection.runtime().has_active_prompt(self.id);
durable_idle && !has_active_prompt
}
pub fn uncompacted_tokens(&self) -> usize {
self.durable.lock().uncompacted_tokens
}
pub fn compressed_blocks(&self) -> Vec<crate::context::models::CompressedBlock> {
self.durable.lock().compressed_blocks.clone()
}
pub async fn checkpoint(&self) -> Result<(), String> {
if !self.is_idle() {
return Err("Cannot checkpoint: session is not idle".into());
}
let config = self.connection.runtime().config().clone();
if !config.context_management.enabled {
return Err("Context management is not enabled".into());
}
self.try_prompt("/compact").await.map(|_| ())
}
pub async fn export_handoff(
&self,
model: &str,
provider_name: Option<&str>,
) -> Result<crate::context::HandoffBundle, String> {
if !self.is_idle() {
return Err("Cannot export handoff: session is not idle".into());
}
let config = self.connection.runtime().config().clone();
let (compressed_blocks, tail) = {
let session = self.durable.lock();
let tail = session.messages.clone();
(session.compressed_blocks.clone(), tail)
};
let session = self.durable.lock();
HandoffExporter::export(
&session,
model,
&compressed_blocks,
tail,
&config.context_management,
provider_name,
)
}
pub fn set_mcp_server_enabled(&self, server_id: impl Into<String>, enabled: bool) {
self.durable
.lock()
.set_mcp_server_enabled(server_id, enabled);
}
pub fn is_mcp_server_enabled(&self, server_id: &str) -> Option<bool> {
self.durable.lock().is_mcp_server_enabled(server_id)
}
pub fn list_enabled_mcp_servers(&self) -> Vec<String> {
self.durable.lock().list_enabled_mcp_servers()
}
pub fn set_plugin_enabled(&self, plugin_id: impl Into<String>, enabled: bool) {
self.durable.lock().set_plugin_enabled(plugin_id, enabled);
}
pub fn is_plugin_enabled(&self, plugin_id: &str) -> Option<bool> {
self.durable.lock().is_plugin_enabled(plugin_id)
}
pub fn list_enabled_plugins(&self) -> Vec<String> {
self.durable.lock().list_enabled_plugins()
}
pub fn switch_model(
&self,
request: crate::context::model_switch::ModelSwitchRequest,
) -> Result<(), String> {
let runtime = self.connection.runtime();
if self.is_idle() {
let from_model = self
.durable
.lock()
.current_model
.clone()
.unwrap_or_else(|| "unknown".to_string());
let (capability_diff, compaction_info) = runtime
.apply_model_switch(self.id, request.clone())
.map_err(|e| e.to_string())?;
let (to_model, _to_provider) = match &request {
crate::context::model_switch::ModelSwitchRequest::Managed {
provider_slug,
model,
..
} => (model.clone(), Some(provider_slug.clone())),
crate::context::model_switch::ModelSwitchRequest::Unmanaged {
model,
provider_name,
} => (model.clone(), Some(provider_name.clone())),
};
if let Some(info) = compaction_info {
let compaction_id = uuid::Uuid::new_v4().to_string();
self.emit_model_switch_event(PromptEvent::CompactionStarted {
compaction_id: compaction_id.clone(),
method: info.method.clone(),
});
self.emit_model_switch_event(PromptEvent::CompactionFinished {
compaction_id,
tokens_before: Some(info.tokens_before),
tokens_after: Some(info.tokens_after),
method: info.method,
});
}
self.emit_model_switch_event(PromptEvent::ModelSwitched {
from_model,
to_model,
adapted: capability_diff.window_shrink.is_some()
|| !capability_diff.hidden_tools.is_empty()
|| !capability_diff.unsupported_modalities.is_empty(),
capability_diff,
});
Ok(())
} else {
runtime
.queue_model_switch(self.id, request.clone())
.map_err(|e| e.to_string())?;
let (target_model, target_provider) = match &request {
crate::context::model_switch::ModelSwitchRequest::Managed {
provider_slug,
model,
..
} => (model.clone(), Some(provider_slug.clone())),
crate::context::model_switch::ModelSwitchRequest::Unmanaged {
model,
provider_name,
} => (model.clone(), Some(provider_name.clone())),
};
self.emit_model_switch_event(PromptEvent::ModelSwitchPending {
target_model,
target_provider,
});
Ok(())
}
}
pub fn start_auth_flow(
&self,
plugin_id: &str,
) -> Result<crate::plugin::auth::AuthInteractionRequest, String> {
let previous_state = self
.connection
.runtime()
.get_plugin_status(plugin_id)
.map(|status| status.auth.state)
.unwrap_or(crate::plugin::auth::AuthState::Unauthenticated);
let request = self
.connection
.runtime()
.begin_plugin_auth_flow(plugin_id)?;
self.emit_auth_transition_to_stream(
plugin_id,
previous_state,
crate::plugin::auth::AuthState::Authenticating,
);
Ok(request)
}
pub fn complete_auth_flow(
&self,
plugin_id: &str,
response: crate::plugin::auth::AuthInteractionResponse,
) -> Result<crate::plugin::auth::AuthStatusTransition, String> {
let transition = self
.connection
.runtime()
.complete_plugin_auth_flow(plugin_id, response)?;
self.emit_auth_transition_to_stream(
&transition.auth_id,
transition.previous_state,
transition.new_state,
);
Ok(transition)
}
pub fn get_auth_prompts(&self) -> Vec<crate::plugin::auth::AuthPrompt> {
self.connection.runtime().get_auth_prompts()
}
pub fn get_plugin_tool_summary(
&self,
) -> crate::plugin::effective_tools::SessionPluginToolSummary {
self.connection
.runtime()
.get_session_plugin_summary(self.id)
.unwrap_or_default()
}
pub fn get_tool_diagnostics(&self) -> Vec<crate::mcp::session_catalog::ToolDiagnostic> {
self.connection
.runtime()
.get_session_tool_diagnostics(self.id)
.unwrap_or_default()
}
pub fn import_handoff(&self, bundle: crate::context::HandoffBundle) -> Result<(), String> {
if !self.is_idle() {
return Err("Cannot import handoff: session is not idle".into());
}
if !self.is_empty() {
return Err(
"Cannot import handoff: session must be empty. Use create_session_from_handoff instead.".into(),
);
}
let mut session = self.durable.lock();
HandoffImporter::hydrate(&mut session, bundle)
}
}
fn to_acp_content_block(block: &ContentBlock) -> acp::ContentBlock {
match block {
ContentBlock::Text { text } => acp::ContentBlock::Text(acp::TextContent::new(text)),
ContentBlock::Image { data, mime_type } => {
acp::ContentBlock::Image(acp::ImageContent::new(data, mime_type))
}
ContentBlock::Resource { uri, name } => acp::ContentBlock::ResourceLink(
acp::ResourceLink::new(name.as_deref().unwrap_or("resource"), uri),
),
}
}
impl Clone for AgentSession {
fn clone(&self) -> Self {
Self {
id: self.id,
durable: self.durable.clone(),
connection: self.connection.clone(),
active_streams: self.active_streams.clone(),
profile_id: self.profile_id.clone(),
}
}
}
impl std::fmt::Debug for IronAgent {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IronAgent").finish()
}
}
impl std::fmt::Debug for AgentConnection {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AgentConnection")
.field("id", &self.inner.id())
.finish()
}
}
impl std::fmt::Debug for AgentSession {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AgentSession")
.field("id", &self.id)
.finish()
}
}