pub(crate) mod convert_conversation; mod convert_from; mod convert_to; mod r#impl; pub use ai::agent::convert::ConvertToAPITypeError; use ai::api_keys::ApiKeyManager; pub use convert_from::{ user_inputs_from_messages, ConversionParams, ConvertAPIMessageToClientOutputMessage, MaybeAIAgentOutputMessage, MessageToAIAgentOutputMessageError, }; pub use r#impl::generate_multi_agent_output; use futures_lite::Stream; use serde::Serialize; use std::path::Path; use std::pin::Pin; use std::sync::Arc; use warp_core::channel::ChannelState; use warp_core::execution_mode::AppExecutionMode; use warp_core::features::FeatureFlag; use crate::ai::agent::conversation::AIConversationId; use crate::ai::ambient_agents::AmbientAgentTaskId; use crate::{ ai::{blocklist::SessionContext, llms::LLMId}, server::server_api::AIApiError, }; use super::{AIAgentInput, MCPContext, MCPServer, RequestMetadata, Suggestions}; use crate::ai::blocklist::{BlocklistAIPermissions, RequestInput}; use crate::ai::mcp::templatable_manager::TemplatableMCPServerInfo; use crate::ai::mcp::TemplatableMCPServerManager; use crate::settings::AISettings; use crate::terminal::safe_mode_settings::get_secret_obfuscation_mode; use crate::workspaces::user_workspaces::UserWorkspaces; use warp_core::user_preferences::GetUserPreferences; use warpui::{AppContext, EntityId, SingletonEntity as _}; /// Unique, server-generated conversation-scoped token to be roundtripped to the API when sending /// requests that follow-up within a given conversation. #[derive(Serialize, Debug, Clone, PartialEq, Eq, Hash)] pub struct ServerConversationToken(String); impl ServerConversationToken { pub fn new(id: String) -> Self { Self(id) } pub fn as_str(&self) -> &str { &self.0 } pub fn debug_link(&self) -> String { format!( "{}/debug/maa/{}", ChannelState::server_root_url(), self.as_str() ) } pub fn conversation_link(&self) -> String { format!( "{}/conversation/{}", ChannelState::server_root_url(), self.as_str() ) } } impl From for String { fn from(value: ServerConversationToken) -> Self { value.0 } } // Conversions between AI ServerConversationToken and protocol ServerConversationToken impl From for ServerConversationToken { fn from(token: session_sharing_protocol::common::ServerConversationToken) -> Self { Self(token.to_string()) } } impl TryFrom for session_sharing_protocol::common::ServerConversationToken { type Error = uuid::Error; fn try_from(token: ServerConversationToken) -> Result { token.as_str().parse() } } #[derive(Debug, Clone)] pub struct RequestParams { pub input: Vec, pub conversation_token: Option, pub forked_from_conversation_token: Option, pub ambient_agent_task_id: Option, pub tasks: Vec, pub existing_suggestions: Option, pub metadata: Option, pub session_context: SessionContext, pub model: LLMId, #[allow(unused)] pub coding_model: LLMId, pub cli_agent_model: LLMId, pub computer_use_model: LLMId, pub is_memory_enabled: bool, pub warp_drive_context_enabled: bool, pub mcp_context: Option, pub planning_enabled: bool, should_redact_secrets: bool, /// User-provided API keys for AI providers (BYO API Key). pub api_keys: Option, pub allow_use_of_warp_credits_with_byok: bool, pub autonomy_level: warp_multi_agent_api::AutonomyLevel, pub isolation_level: warp_multi_agent_api::IsolationLevel, pub web_search_enabled: bool, pub computer_use_enabled: bool, pub ask_user_question_enabled: bool, pub research_agent_enabled: bool, pub orchestration_enabled: bool, pub supported_tools_override: Option>, /// The conversation ID of the parent agent that spawned this child agent, if any. pub parent_agent_id: Option, /// The display name for this agent (e.g. "Agent 1"), assigned by the orchestrator. pub agent_name: Option, } pub type Event = Result>; #[cfg(not(target_family = "wasm"))] pub type ResponseStream = Pin + Send + 'static>>; // The WASM version of this type has no bound on `Send`, which is an unnecessary bound when // targeting wasm because the browser is single-threaded (and we don't leverage WebWorkers for async // execution in WoW). #[cfg(target_family = "wasm")] pub type ResponseStream = Pin>>; #[derive(Debug, Clone)] pub struct ConversationData { pub id: AIConversationId, pub tasks: Vec, pub server_conversation_token: Option, pub forked_from_conversation_token: Option, pub ambient_agent_task_id: Option, pub existing_suggestions: Option, } impl RequestParams { pub fn new( terminal_view_id: Option, session_context: SessionContext, request_input: &RequestInput, conversation: ConversationData, metadata: Option, app: &AppContext, ) -> Self { let ai_settings = AISettings::as_ref(app); let is_memory_enabled = ai_settings.is_memory_enabled(app); let warp_drive_context_enabled = ai_settings.is_warp_drive_context_enabled(app); // Build MCP context - either grouped by server or flat lists based on feature flag let mcp_context = if FeatureFlag::MCPGroupedServerContext.is_enabled() { // Group MCP tools and resources by server let templatable_manager = TemplatableMCPServerManager::as_ref(app); let mut active_servers: Vec<&TemplatableMCPServerInfo> = templatable_manager .get_active_templatable_servers() .values() .copied() .collect(); // If file-based MCP servers are enabled, add active servers in scope of // the user's current working directory if let Some(cwd) = session_context.current_working_directory() { active_servers.extend( templatable_manager .get_active_file_based_servers(Path::new(cwd), app) .values(), ); } // Include any ephemeral MCP servers started via the Oz CLI. active_servers.extend( templatable_manager .get_active_cli_spawned_servers() .values(), ); let servers: Vec = active_servers .into_iter() .map(|server| MCPServer { name: server.name().to_string(), description: server.description().unwrap_or_default().to_string(), id: server.installation_id().to_string(), resources: server.resources().to_vec(), tools: server.tools().to_vec(), }) .collect(); if servers.is_empty() { None } else { #[allow(deprecated)] Some(MCPContext { resources: vec![], tools: vec![], servers, }) } } else { // Flat lists of resources and tools let templatable_mcp_manager = TemplatableMCPServerManager::as_ref(app); let resources = templatable_mcp_manager .resources() .cloned() .collect::>(); let tools = templatable_mcp_manager.tools().cloned().collect::>(); #[allow(deprecated)] (!resources.is_empty() || !tools.is_empty()).then_some(MCPContext { resources, tools, servers: vec![], }) }; let should_redact_secrets = get_secret_obfuscation_mode(app).should_redact_secret(); let user_workspaces = UserWorkspaces::as_ref(app); let api_keys = ApiKeyManager::as_ref(app).api_keys_for_request( user_workspaces.is_byo_api_key_enabled(), user_workspaces.is_aws_bedrock_credentials_enabled(app), ); let allow_use_of_warp_credits_with_byok = *AISettings::as_ref(app).can_use_warp_credits_with_byok; let app_execution_mode = AppExecutionMode::as_ref(app); let autonomy_level = if app_execution_mode.is_autonomous() { warp_multi_agent_api::AutonomyLevel::Unsupervised } else { warp_multi_agent_api::AutonomyLevel::Supervised }; let isolation_level = if app_execution_mode.is_sandboxed() { warp_multi_agent_api::IsolationLevel::Sandbox } else { warp_multi_agent_api::IsolationLevel::None }; let web_search_enabled = BlocklistAIPermissions::as_ref(app).get_web_search_enabled(app, terminal_view_id); let research_agent_enabled = app .private_user_preferences() .read_value("ResearchAgentEnabled") .ok() .flatten() .and_then(|s| s.parse().ok()) .unwrap_or_default(); let is_ambient_agent = conversation.ambient_agent_task_id.is_some(); let computer_use_enabled = FeatureFlag::AgentModeComputerUse.is_enabled() && BlocklistAIPermissions::as_ref(app) .get_computer_use_setting(app, terminal_view_id) .is_enabled() && computer_use::is_supported_on_current_platform() && (FeatureFlag::LocalComputerUse.is_enabled() || is_ambient_agent); let ask_user_question_enabled = BlocklistAIPermissions::as_ref(app) .get_ask_user_question_setting(app, terminal_view_id) != crate::ai::execution_profiles::AskUserQuestionPermission::Never; let orchestration_enabled = ai_settings.is_orchestration_enabled(app) && session_context .session_type() .as_ref() .is_none_or(|t| matches!(t, crate::terminal::model::session::SessionType::Local)); Self { input: request_input.all_inputs().cloned().collect(), conversation_token: conversation.server_conversation_token, forked_from_conversation_token: conversation.forked_from_conversation_token, ambient_agent_task_id: conversation.ambient_agent_task_id, tasks: conversation.tasks, existing_suggestions: conversation.existing_suggestions, metadata, session_context, model: request_input.model_id.clone(), coding_model: request_input.coding_model_id.clone(), cli_agent_model: request_input.cli_agent_model_id.clone(), computer_use_model: request_input.computer_use_model_id.clone(), is_memory_enabled, warp_drive_context_enabled, mcp_context, planning_enabled: true, should_redact_secrets, api_keys, allow_use_of_warp_credits_with_byok, autonomy_level, isolation_level, web_search_enabled, computer_use_enabled, ask_user_question_enabled, research_agent_enabled, orchestration_enabled, supported_tools_override: request_input.supported_tools_override.clone(), parent_agent_id: None, agent_name: None, } } }