diff --git a/src/app.rs b/src/app.rs index d3f15eb0db9..931d4af003a 100644 --- a/src/app.rs +++ b/src/app.rs @@ -695,6 +695,7 @@ impl AppBuilder { tools: &Arc, hooks: &Arc, settings_store_override: Option>, + ownership_cache: Arc, ) -> Result< ( Arc, @@ -989,6 +990,13 @@ impl AppBuilder { if let Some(ref ss) = settings_store_override { em = em.with_settings_store(Arc::clone(ss)); } + if let Some(ref db) = self.db { + let ps = Arc::new(crate::pairing::PairingStore::new( + Arc::clone(db), + Arc::clone(&ownership_cache), + )); + em = em.with_pairing_store(ps); + } let manager = Arc::new(em); tools.register_extension_tools(Arc::clone(&manager)); @@ -1181,6 +1189,7 @@ impl AppBuilder { _ => (None, None), }; + let ownership_cache = Arc::new(crate::ownership::OwnershipCache::new()); let ( mcp_session_manager, mcp_process_manager, @@ -1189,7 +1198,12 @@ impl AppBuilder { catalog_entries, dev_loaded_tool_names, ) = self - .init_extensions(&tools, &hooks, settings_store.clone()) + .init_extensions( + &tools, + &hooks, + settings_store.clone(), + Arc::clone(&ownership_cache), + ) .await?; // Load bootstrap-completed flag from settings so that existing users @@ -1333,7 +1347,7 @@ impl AppBuilder { catalog_entries, dev_loaded_tool_names, builder, - ownership_cache: Arc::new(crate::ownership::OwnershipCache::new()), + ownership_cache, }) } } diff --git a/src/channels/relay/channel.rs b/src/channels/relay/channel.rs index 192e48d00da..25827cb4b87 100644 --- a/src/channels/relay/channel.rs +++ b/src/channels/relay/channel.rs @@ -6,6 +6,7 @@ //! proxy API (Slack). use std::collections::HashMap; +use std::sync::Arc; use async_trait::async_trait; use tokio::sync::mpsc; @@ -15,6 +16,7 @@ use crate::channels::{ Channel, ChatApprovalPrompt, IncomingMessage, MessageStream, OutgoingResponse, StatusUpdate, }; use crate::error::ChannelError; +use crate::pairing::PairingStore; /// Default channel name for the Slack relay integration. pub const DEFAULT_RELAY_NAME: &str = "slack-relay"; @@ -51,6 +53,8 @@ pub struct RelayChannel { event_tx: mpsc::Sender, /// Receiver side — taken once by `start()`. event_rx: tokio::sync::Mutex>>, + /// Resolves Slack sender_id → internal UserId for multi-tenant support. + pairing_store: Option>, } impl RelayChannel { @@ -88,9 +92,16 @@ impl RelayChannel { instance_id, event_tx, event_rx: tokio::sync::Mutex::new(Some(event_rx)), + pairing_store: None, } } + /// Set the pairing store for multi-tenant identity resolution. + pub fn with_pairing_store(mut self, store: Arc) -> Self { + self.pairing_store = Some(store); + self + } + /// Get a clone of the event sender for wiring into the webhook endpoint. pub fn event_sender(&self) -> mpsc::Sender { self.event_tx.clone() @@ -123,9 +134,10 @@ impl RelayChannel { team_id: &str, method: &str, body: serde_json::Value, + slack_user_id: Option<&str>, ) -> Result { self.client - .proxy_provider(self.provider.as_str(), team_id, method, body) + .proxy_provider_with_user(self.provider.as_str(), team_id, method, body, slack_user_id) .await } @@ -205,6 +217,8 @@ impl Channel for RelayChannel { let (tx, rx) = mpsc::channel(64); let provider_str = self.provider.as_str().to_string(); let relay_name = channel_name.clone(); + let pairing_store = self.pairing_store.clone(); + let pairing_client = self.client.clone(); // Spawn a task that reads events from the webhook handler and converts to IncomingMessage tokio::spawn(async move { @@ -240,7 +254,80 @@ impl Channel for RelayChannel { "Relay: received message from {}", provider_str ); - let mut msg = IncomingMessage::new(&relay_name, &event.sender_id, event.text()) + // Resolve sender_id → internal UserId via PairingStore. + // External ID is scoped to workspace: "team_id:sender_id". + let scoped_external_id = format!("{}:{}", event.provider_scope, event.sender_id); + let resolved_user_id: String = if let Some(ref store) = pairing_store { + match store + .resolve_identity(&relay_name, &scoped_external_id) + .await + { + Ok(Some(uid)) => { + let user_str = uid.as_str().to_string(); + tracing::debug!( + sender_id = %event.sender_id, + resolved_user = %user_str, + "Relay: resolved sender to internal user" + ); + user_str + } + Ok(None) => { + tracing::info!( + sender_id = %event.sender_id, + "Relay: sender not paired, sending pairing code" + ); + let meta = serde_json::json!({ + "sender_name": event.display_name(), + "channel_id": event.channel_id, + }); + match store + .upsert_request(&relay_name, &scoped_external_id, Some(meta)) + .await + { + Ok(record) => { + let instructions = format!( + "Enter this code in IronClaw to pair your Slack account: `{}`", + record.code + ); + let team_id = event.team_id().to_string(); + let body = serde_json::json!({ + "channel": event.channel_id, + "text": instructions, + "thread_ts": event.thread_id.as_deref().unwrap_or(&event.id), + }); + if let Err(e) = pairing_client + .proxy_provider( + &provider_str, + &team_id, + "chat.postMessage", + body, + ) + .await + { + tracing::warn!(error = %e, "Relay: failed to send pairing code reply"); + } + } + Err(e) => { + tracing::warn!(error = %e, "Relay: failed to create pairing request"); + } + } + continue; + } + Err(e) => { + tracing::warn!( + sender_id = %event.sender_id, + error = %e, + "Relay: pairing resolution failed, dropping message" + ); + continue; + } + } + } else { + event.sender_id.clone() + }; + + let mut msg = IncomingMessage::new(&relay_name, &resolved_user_id, event.text()) + .with_sender_id(event.sender_id.clone()) .with_user_name(event.display_name()) .with_metadata(serde_json::json!({ "team_id": event.team_id(), @@ -323,8 +410,9 @@ impl Channel for RelayChannel { .filter(|s| !s.is_empty()); let (method, body) = self.build_send_body(channel_id, &response.content, thread_id); + let sender_id = metadata.get("sender_id").and_then(|v| v.as_str()); - self.proxy_send(team_id, &method, body) + self.proxy_send(team_id, &method, body, sender_id) .await .map_err(|e| ChannelError::SendFailed { name: channel_name, @@ -384,7 +472,7 @@ impl Channel for RelayChannel { })?; let body = self.build_approval_body(channel_id, thread_id, &prompt, &approval_token); - self.proxy_send(team_id, "chat.postMessage", body) + self.proxy_send(team_id, "chat.postMessage", body, None) .await .map_err(|e| ChannelError::SendFailed { name: self.name().to_string(), @@ -410,7 +498,7 @@ impl Channel for RelayChannel { let (method, body) = self.build_send_body(target, &response.content, thread_id); - self.proxy_send(&self.team_id, &method, body) + self.proxy_send(&self.team_id, &method, body, None) .await .map_err(|e| ChannelError::SendFailed { name: channel_name, diff --git a/src/channels/relay/client.rs b/src/channels/relay/client.rs index 16f40f66474..f09aeab1b13 100644 --- a/src/channels/relay/client.rs +++ b/src/channels/relay/client.rs @@ -81,9 +81,13 @@ impl ChannelEvent { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Connection { pub provider: String, + #[serde(alias = "provider_scope")] pub team_id: String, + #[serde(alias = "provider_scope_name")] pub team_name: Option, + #[serde(default)] pub connected: bool, + pub authed_user_id: Option, } /// HTTP client for the channel-relay service. @@ -237,6 +241,18 @@ impl RelayClient { team_id: &str, method: &str, body: serde_json::Value, + ) -> Result { + self.proxy_provider_with_user(provider, team_id, method, body, None) + .await + } + + pub async fn proxy_provider_with_user( + &self, + provider: &str, + team_id: &str, + method: &str, + body: serde_json::Value, + slack_user_id: Option<&str>, ) -> Result { let url = format!("{}/proxy/{}/{}", self.base_url, provider, method); tracing::trace!( @@ -245,7 +261,10 @@ impl RelayClient { method = %method, "RelayClient::proxy_provider: sending request" ); - let query: Vec<(&str, &str)> = vec![("team_id", team_id)]; + let mut query: Vec<(&str, &str)> = vec![("team_id", team_id)]; + if let Some(uid) = slack_user_id { + query.push(("slack_user_id", uid)); + } let resp = self .http .post(&url) diff --git a/src/channels/web/features/oauth/mod.rs b/src/channels/web/features/oauth/mod.rs index 8d3ef00ece6..8a1731ed39c 100644 --- a/src/channels/web/features/oauth/mod.rs +++ b/src/channels/web/features/oauth/mod.rs @@ -736,7 +736,8 @@ pub(crate) async fn slack_relay_oauth_callback_handler( format!("Failed to persist relay team_id: {e}") })?; - // Activate the relay channel + // Activate the relay channel first — this creates the relay client and + // verifies the connection is usable. tracing::info!( relay = %relay_extension_name, owner_id = %state.owner_id, @@ -747,6 +748,84 @@ pub(crate) async fn slack_relay_oauth_callback_handler( .await .map_err(|e| format!("Failed to activate relay channel: {}", e))?; + // Create channel identity pairing: Slack authed_user_id → IronClaw user. + // Fetch authed_user_id from the relay's connections API (server-side, + // not from the redirect URL which could be tampered). + if let Some(pairing_store) = ext_mgr.pairing_store() { + let relay_config = ext_mgr + .relay_config() + .map_err(|e| format!("Relay config not available: {e}"))?; + let effective_url = ext_mgr + .effective_relay_url(&relay_extension_name) + .await + .unwrap_or_else(|| relay_config.url.clone()); + let client = crate::channels::relay::RelayClient::new( + effective_url, + relay_config.api_key.clone(), + relay_config.request_timeout_secs, + ) + .map_err(|e| format!("Failed to create relay client: {e}"))?; + + let connections = client + .list_connections("") + .await + .map_err(|e| format!("Failed to fetch relay connections: {e}"))?; + let authed_user_id = connections + .iter() + .find(|c| c.team_id == team_id) + .and_then(|c| c.authed_user_id.clone()) + .ok_or_else(|| { + "No connection with authed_user_id found for this team".to_string() + })?; + + let user_key = format!("relay:{}:oauth_user", relay_extension_name); + let oauth_user = ext_mgr + .secrets() + .get_decrypted(&state.owner_id, &user_key) + .await + .ok() + .map(|s| s.expose().to_string()) + .unwrap_or_else(|| state.owner_id.clone()); + let _ = ext_mgr.secrets().delete(&state.owner_id, &user_key).await; + + let user_record = if let Some(ref db) = state.store { + db.get_user(&oauth_user).await.ok().flatten() + } else { + None + }; + let Some(ref record) = user_record else { + return Err(format!( + "OAuth user '{oauth_user}' not found — cannot create relay identity" + )); + }; + if record.status != "active" { + return Err(format!( + "OAuth user '{oauth_user}' is not active (status: {})", + record.status + )); + } + let role = match record.role.as_str() { + "owner" => crate::ownership::UserRole::Owner, + "admin" => crate::ownership::UserRole::Admin, + _ => crate::ownership::UserRole::Regular, + }; + let Ok(user_id) = crate::ownership::UserId::new(&oauth_user, role) else { + return Err(format!( + "OAuth user '{oauth_user}' has invalid user_id format" + )); + }; + // Scope external_id to workspace: "team_id:slack_user_id" + let scoped_external_id = format!("{}:{}", team_id, authed_user_id); + pairing_store + .create_identity( + crate::channels::relay::channel::DEFAULT_RELAY_NAME, + &scoped_external_id, + &user_id, + ) + .await + .map_err(|e| format!("Failed to create relay identity: {e}"))?; + } + Ok(()) } .await; diff --git a/src/db/libsql/pairing.rs b/src/db/libsql/pairing.rs index 7b114e8d58c..bf942b7dd02 100644 --- a/src/db/libsql/pairing.rs +++ b/src/db/libsql/pairing.rs @@ -462,6 +462,27 @@ impl ChannelPairingStore for LibSqlBackend { .map_err(|e| DatabaseError::Query(e.to_string()))?; Ok(()) } + + async fn create_channel_identity( + &self, + channel: &str, + external_id: &str, + owner_id: &str, + ) -> Result<(), DatabaseError> { + let channel = crate::pairing::normalize_channel_name(channel); + let id = uuid::Uuid::new_v4().to_string(); + let conn = self.connect().await?; + conn.execute( + "INSERT INTO channel_identities (id, owner_id, channel, external_id) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT (channel, external_id) + DO UPDATE SET owner_id = ?2", + params![id, owner_id, channel, external_id], + ) + .await + .map_err(|e| DatabaseError::Query(e.to_string()))?; + Ok(()) + } } #[cfg(test)] diff --git a/src/db/mod.rs b/src/db/mod.rs index 7b605773ace..272245df731 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -1192,6 +1192,15 @@ pub trait ChannelPairingStore: Send + Sync { channel: &str, external_id: &str, ) -> Result<(), DatabaseError>; + + /// Create or update a channel identity directly (trusted path, e.g. OAuth). + /// Inserts into channel_identities without requiring a pairing code. + async fn create_channel_identity( + &self, + channel: &str, + external_id: &str, + owner_id: &str, + ) -> Result<(), DatabaseError>; } /// Generates an 8-character pairing code from an unambiguous alphabet. diff --git a/src/db/postgres.rs b/src/db/postgres.rs index fda565eeb3c..c7621bd8d62 100644 --- a/src/db/postgres.rs +++ b/src/db/postgres.rs @@ -1467,6 +1467,31 @@ impl ChannelPairingStore for PgBackend { .map_err(|e| DatabaseError::Query(e.to_string()))?; Ok(()) } + + async fn create_channel_identity( + &self, + channel: &str, + external_id: &str, + owner_id: &str, + ) -> Result<(), DatabaseError> { + let channel = crate::pairing::normalize_channel_name(channel); + let client = self + .pool() + .get() + .await + .map_err(|e| DatabaseError::Pool(e.to_string()))?; + client + .execute( + "INSERT INTO channel_identities (owner_id, channel, external_id) + VALUES ($1, $2, $3) + ON CONFLICT (channel, external_id) + DO UPDATE SET owner_id = $1", + &[&owner_id, &channel, &external_id], + ) + .await + .map_err(|e| DatabaseError::Query(e.to_string()))?; + Ok(()) + } } // ==================== IdentityStore ==================== diff --git a/src/extensions/manager.rs b/src/extensions/manager.rs index 65e0a0334aa..88ce7a9a125 100644 --- a/src/extensions/manager.rs +++ b/src/extensions/manager.rs @@ -441,6 +441,8 @@ pub struct ExtensionManager { /// Stored here so the web gateway can verify incoming callbacks without /// any env var or shared secret. relay_signing_secret_cache: Arc>>>, + /// PairingStore for multi-tenant relay identity resolution. + pairing_store: Option>, /// When `true`, OAuth flows always return an auth URL to the caller /// instead of opening a browser on the server via `open::that()`. /// Set by the web gateway at startup via `enable_gateway_mode()`. @@ -670,6 +672,7 @@ impl ExtensionManager { relay_config: crate::config::RelayConfig::from_env(), relay_event_tx: Arc::new(tokio::sync::Mutex::new(None)), relay_signing_secret_cache: Arc::new(std::sync::Mutex::new(None)), + pairing_store: None, gateway_mode: std::sync::atomic::AtomicBool::new(false), gateway_base_url: RwLock::new(None), channel_activation_locks: RwLock::new(HashMap::new()), @@ -750,7 +753,7 @@ impl ExtensionManager { } /// Get the relay config stored at startup. - fn relay_config(&self) -> Result<&crate::config::RelayConfig, ExtensionError> { + pub(crate) fn relay_config(&self) -> Result<&crate::config::RelayConfig, ExtensionError> { self.relay_config.as_ref().ok_or_else(|| { ExtensionError::Config( "CHANNEL_RELAY_URL and CHANNEL_RELAY_API_KEY must be set".to_string(), @@ -771,7 +774,7 @@ impl ExtensionManager { /// and the URL must not contain userinfo (embedded credentials). This /// prevents a malicious override from exfiltrating the instance-wide relay /// API key to an attacker-controlled host. - async fn effective_relay_url(&self, name: &str) -> Option { + pub(crate) async fn effective_relay_url(&self, name: &str) -> Option { if let Some(ref store) = self.store { let key = format!("extensions.{name}.relay_url"); if let Ok(Some(v)) = store.get_setting(&self.user_id, &key).await { @@ -1156,6 +1159,10 @@ impl ExtensionManager { &self.secrets } + pub fn pairing_store(&self) -> Option<&Arc> { + self.pairing_store.as_ref() + } + /// Expose the per-user MCP client store. Tool wrappers registered in /// the global `ToolRegistry` hold an `Arc` and resolve /// the caller's client at dispatch time via @@ -1412,6 +1419,11 @@ impl ExtensionManager { self } + pub fn with_pairing_store(mut self, store: Arc) -> Self { + self.pairing_store = Some(store); + self + } + async fn clear_pending_extension_auth(&self, name: &str, user_id: &str) { { let mut pending = self.pending_auth.write().await; @@ -6462,6 +6474,15 @@ impl ExtensionManager { ExtensionError::AuthFailed(format!("Failed to store OAuth state: {e}")) })?; + // Store the initiating user_id so the OAuth callback knows which IronClaw + // user to pair with the Slack authed_user_id. + let user_key = format!("relay:{}:oauth_user", name); + let _ = self.secrets.delete(&self.user_id, &user_key).await; + self.secrets + .create(&self.user_id, CreateSecretParams::new(&user_key, user_id)) + .await + .map_err(|e| ExtensionError::AuthFailed(format!("Failed to store OAuth user: {e}")))?; + // Channel-relay derives all URLs from trusted instance_url in chat-api. // We only pass the nonce for CSRF validation on the callback. tracing::trace!( @@ -6612,7 +6633,7 @@ impl ExtensionManager { // Create the event channel for webhook callbacks let (event_tx, event_rx) = tokio::sync::mpsc::channel(64); - let channel = crate::channels::relay::RelayChannel::new_with_provider( + let mut channel = crate::channels::relay::RelayChannel::new_with_provider( client.clone(), crate::channels::relay::channel::RelayProvider::Slack, team_id.clone(), @@ -6620,6 +6641,9 @@ impl ExtensionManager { event_tx.clone(), event_rx, ); + if let Some(ref ps) = self.pairing_store { + channel = channel.with_pairing_store(Arc::clone(ps)); + } // Hot-add to channel manager let cm_guard = self.relay_channel_manager.read().await; diff --git a/src/pairing/store.rs b/src/pairing/store.rs index 77777f4335d..932fc129c4f 100644 --- a/src/pairing/store.rs +++ b/src/pairing/store.rs @@ -197,4 +197,22 @@ impl PairingStore { self.cache.evict(&channel, external_id); Ok(()) } + + /// Create a channel identity directly (trusted path, e.g. OAuth completion). + /// Inserts into channel_identities and populates the cache without requiring + /// a pairing code flow. + pub async fn create_identity( + &self, + channel: &str, + external_id: &str, + owner_id: &UserId, + ) -> Result<(), DatabaseError> { + let channel = crate::pairing::normalize_channel_name(channel); + if let Some(ref db) = self.db { + db.create_channel_identity(&channel, external_id, owner_id.as_str()) + .await?; + } + self.cache.insert(&channel, external_id, owner_id.clone()); + Ok(()) + } }