From 33b3e529dd6e23809ece9f7d4f1928d038c2e570 Mon Sep 17 00:00:00 2001 From: Taariq Lewis <701864+taariq@users.noreply.github.com> Date: Tue, 4 Aug 2026 17:19:28 +0200 Subject: [PATCH] fix(slack): acknowledge accepted inbound messages --- channels/slack/README.md | 8 +- channels/slack/src/main.rs | 32 ++- channels/slack/src/slack_api.rs | 64 +++++- channels/slack/tests/ingress_protocol.rs | 252 ++++++++++++++++++++++- 4 files changed, 343 insertions(+), 13 deletions(-) diff --git a/channels/slack/README.md b/channels/slack/README.md index 496b2f2..096d993 100644 --- a/channels/slack/README.md +++ b/channels/slack/README.md @@ -27,6 +27,7 @@ Behavior: - Events API ingress uses host-managed webhook callbacks - Socket Mode ingress supports both one-shot `poll_ingress` fetches and background `start_ingress` sessions - Socket Mode opens Slack's websocket via `apps.connections.open` and emits normalized inbound events back to Dispatch +- accepted inbound messages receive an immediate `:eyes:` reaction when bot-token delivery is configured - challenge and acknowledgement replies are returned through `callback_reply` - status frames render visible status messages into Slack conversations @@ -87,6 +88,7 @@ Bot-token mode: Recommended bot scopes: - `chat:write` - required for `deliver`, `push`, and `status` +- `reactions:write` - required for inbound `:eyes:` acknowledgements - `app_mentions:read` - required if you subscribe to `app_mention` - `channels:history` - required for `message.channels` - `im:history` - required for `message.im` @@ -105,9 +107,9 @@ Recommended event subscriptions: Common scope sets: -- mention-only Socket Mode: `chat:write`, `app_mentions:read` -- public-channel Events API or Socket Mode: `chat:write`, `channels:history` -- full plugin test setup: `chat:write`, `files:write`, `app_mentions:read`, `channels:history`, `groups:history`, `im:history`, `mpim:history` +- mention-only Socket Mode: `chat:write`, `reactions:write`, `app_mentions:read` +- public-channel Events API or Socket Mode: `chat:write`, `reactions:write`, `channels:history` +- full plugin test setup: `chat:write`, `reactions:write`, `files:write`, `app_mentions:read`, `channels:history`, `groups:history`, `im:history`, `mpim:history` Socket Mode setup: diff --git a/channels/slack/src/main.rs b/channels/slack/src/main.rs index ef10435..84187a2 100644 --- a/channels/slack/src/main.rs +++ b/channels/slack/src/main.rs @@ -60,6 +60,7 @@ const META_PATH: &str = "path"; const META_API_APP_ID: &str = "api_app_id"; const META_EVENT_CONTEXT: &str = "event_context"; const META_CHANNEL_TYPE: &str = "channel_type"; +const META_MESSAGE_TS: &str = "message_ts"; const META_STATUS_KIND: &str = "status_kind"; const META_REASON_CODE: &str = "reason_code"; @@ -68,6 +69,7 @@ const MODE_INCOMING_WEBHOOK: &str = "incoming_webhook"; const DELIVERY_MODE_CHAT_POST_MESSAGE: &str = "chat.postMessage"; const TRANSPORT_EVENTS_WEBHOOK: &str = "events_webhook"; const TRANSPORT_SOCKET_MODE: &str = "socket_mode"; +const THINKING_REACTION: &str = "eyes"; const MAX_SIGNATURE_AGE_SECS: i64 = 300; const ROUTE_CONVERSATION_ID: &str = "conversation_id"; @@ -773,13 +775,16 @@ fn build_inbound_event( if let Some(event_context) = &envelope.event_context { event_metadata.insert(META_EVENT_CONTEXT.to_string(), event_context.clone()); } + if let Some(message_ts) = &event.ts { + event_metadata.insert(META_MESSAGE_TS.to_string(), message_ts.clone()); + } let account_id = envelope.team_id.clone(); if !team_is_allowed(config, account_id.as_deref()) { return Ok(None); } - Ok(Some(InboundEventEnvelope { + let inbound_event = InboundEventEnvelope { event_id: envelope.event_id.clone().unwrap_or_else(|| { let ts = event .event_ts @@ -808,7 +813,30 @@ fn build_inbound_event( }, account_id, metadata: event_metadata, - })) + }; + + acknowledge_inbound_event(config, channel_id, event.ts.as_deref()); + + Ok(Some(inbound_event)) +} + +fn acknowledge_inbound_event(config: &ChannelConfig, channel_id: &str, message_ts: Option<&str>) { + if !has_optional_env(bot_token_env(config)) { + return; + } + + let result = message_ts + .ok_or_else(|| anyhow!("Slack event is missing its message timestamp")) + .and_then(|message_ts| { + SlackClient::from_env(bot_token_env(config)).map(|client| (client, message_ts)) + }) + .and_then(|(client, message_ts)| { + client.add_reaction(channel_id, message_ts, THINKING_REACTION) + }); + + if result.is_err() { + eprintln!("slack inbound acknowledgement failed"); + } } fn supports_inbound_event(event: &SlackEventPayload) -> bool { diff --git a/channels/slack/src/slack_api.rs b/channels/slack/src/slack_api.rs index 9f05ee5..b2cb508 100644 --- a/channels/slack/src/slack_api.rs +++ b/channels/slack/src/slack_api.rs @@ -10,6 +10,7 @@ use std::{ use tungstenite::{Message, WebSocket, stream::MaybeTlsStream}; const DEFAULT_API_BASE: &str = "https://slack.com/api"; +const REACTION_TIMEOUT: Duration = Duration::from_secs(2); #[derive(Debug)] pub struct SlackClient { @@ -133,6 +134,38 @@ impl SlackClient { }) } + pub fn add_reaction( + &self, + channel_id: &str, + message_ts: &str, + reaction_name: &str, + ) -> Result<()> { + let agent = ureq::Agent::config_builder() + .timeout_global(Some(REACTION_TIMEOUT)) + .build() + .new_agent(); + let body = self + .post_json_response( + Some(&agent), + "reactions.add", + json!({ + "channel": channel_id, + "timestamp": message_ts, + "name": reaction_name, + }), + "failed to add Slack reaction", + ) + .map_err(|_| anyhow!("failed to add Slack reaction"))?; + + match body.get("ok").and_then(Value::as_bool) { + Some(true) => Ok(()), + Some(false) if body.get("error").and_then(Value::as_str) == Some("already_reacted") => { + Ok(()) + } + _ => bail!("failed to add Slack reaction"), + } + } + fn upload_file( &self, channel_id: &str, @@ -200,13 +233,7 @@ impl SlackClient { } fn post_json(&self, method: &str, payload: Value, context: &str) -> Result { - let url = format!("{}/{}", self.base_url, method); - let mut response = ureq::post(&url) - .header("Authorization", &format!("Bearer {}", self.bot_token)) - .header("Content-Type", "application/json") - .send_json(payload) - .map_err(|error| anyhow!("{context}: {error}"))?; - let body = read_json_body(&mut response, context)?; + let body = self.post_json_response(None, method, payload, context)?; let ok = body .get("ok") .and_then(Value::as_bool) @@ -221,6 +248,29 @@ impl SlackClient { Ok(body) } + fn post_json_response( + &self, + agent: Option<&ureq::Agent>, + method: &str, + payload: Value, + context: &str, + ) -> Result { + let url = format!("{}/{}", self.base_url, method); + let mut response = match agent { + Some(agent) => agent + .post(&url) + .header("Authorization", &format!("Bearer {}", self.bot_token)) + .header("Content-Type", "application/json") + .send_json(payload), + None => ureq::post(&url) + .header("Authorization", &format!("Bearer {}", self.bot_token)) + .header("Content-Type", "application/json") + .send_json(payload), + } + .map_err(|error| anyhow!("{context}: {error}"))?; + read_json_body(&mut response, context) + } + fn upload_file_bytes(&self, upload_url: &str, upload: &SlackUpload) -> Result<()> { let mut response = ureq::post(upload_url) .header("Content-Type", &upload.mime_type) diff --git a/channels/slack/tests/ingress_protocol.rs b/channels/slack/tests/ingress_protocol.rs index 3ac9b49..b9dd260 100644 --- a/channels/slack/tests/ingress_protocol.rs +++ b/channels/slack/tests/ingress_protocol.rs @@ -2,7 +2,9 @@ use std::collections::BTreeMap; use std::io::{BufRead, BufReader, Read, Write}; use std::net::TcpListener; use std::process::{Command, Stdio}; +use std::sync::mpsc::{self, Receiver}; use std::thread; +use std::time::Duration; use serde_json::{Value, json}; use tungstenite::{ @@ -51,11 +53,21 @@ fn run_request(request: Value) -> Value { } fn run_request_with_env(request: Value, envs: BTreeMap) -> Value { + run_request_with_env_and_stderr(request, envs).0 +} + +fn run_request_with_env_and_stderr( + request: Value, + envs: BTreeMap, +) -> (Value, String) { let binary = std::env::var("CARGO_BIN_EXE_channel-slack").expect("channel-slack binary path"); let mut child = Command::new(binary) + .env_remove("SLACK_BOT_TOKEN") + .env_remove("SLACK_API_BASE_URL") .envs(envs) .stdin(Stdio::piped()) .stdout(Stdio::piped()) + .stderr(Stdio::piped()) .spawn() .expect("spawn channel-slack"); @@ -77,7 +89,90 @@ fn run_request_with_env(request: Value, envs: BTreeMap) -> Value .find(|line| !line.trim().is_empty()) .expect("response line"); let response: Value = serde_json::from_str(line).expect("parse response"); - response["result"].clone() + ( + response["result"].clone(), + String::from_utf8(output.stderr).expect("stderr utf-8"), + ) +} + +#[derive(Debug)] +struct CapturedApiRequest { + request_line: String, + authorization: Option, + body: Value, +} + +fn serve_slack_api( + responses: Vec, +) -> (String, Receiver, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind Slack API test listener"); + let address = listener.local_addr().expect("Slack API listener address"); + let (request_tx, request_rx) = mpsc::channel(); + let server = thread::spawn(move || { + for response in responses { + let (mut stream, _) = listener.accept().expect("accept Slack API request"); + let request = read_api_request(&mut stream); + request_tx.send(request).expect("capture Slack API request"); + + let body = response.to_string(); + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ) + .expect("write Slack API response"); + } + }); + + (format!("http://{address}/api"), request_rx, server) +} + +fn read_api_request(stream: &mut std::net::TcpStream) -> CapturedApiRequest { + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .expect("set Slack API request timeout"); + let mut buffer = Vec::new(); + let mut chunk = [0_u8; 4096]; + let header_end = loop { + let read = stream.read(&mut chunk).expect("read Slack API request"); + assert!(read > 0, "Slack API connection closed before headers"); + buffer.extend_from_slice(&chunk[..read]); + if let Some(position) = buffer.windows(4).position(|window| window == b"\r\n\r\n") { + break position; + } + }; + + let header_text = String::from_utf8(buffer[..header_end].to_vec()).expect("header utf-8"); + let mut lines = header_text.split("\r\n"); + let request_line = lines.next().expect("request line").to_string(); + let mut authorization = None; + let mut content_length = 0_usize; + for line in lines { + let Some((name, value)) = line.split_once(':') else { + continue; + }; + if name.eq_ignore_ascii_case("authorization") { + authorization = Some(value.trim().to_string()); + } else if name.eq_ignore_ascii_case("content-length") { + content_length = value.trim().parse().expect("content length"); + } + } + + let body_start = header_end + 4; + while buffer.len() < body_start + content_length { + let read = stream.read(&mut chunk).expect("read Slack API body"); + assert!(read > 0, "Slack API connection closed before request body"); + buffer.extend_from_slice(&chunk[..read]); + } + let body = serde_json::from_slice(&buffer[body_start..body_start + content_length]) + .expect("Slack API request JSON"); + + CapturedApiRequest { + request_line, + authorization, + body, + } } fn read_message(reader: &mut BufReader) -> Value { @@ -98,6 +193,8 @@ fn read_message(reader: &mut BufReader) -> Value { fn run_start_ingress_cycle(config: Value, envs: BTreeMap) -> (Value, Value) { let binary = std::env::var("CARGO_BIN_EXE_channel-slack").expect("channel-slack binary path"); let mut child = Command::new(binary) + .env_remove("SLACK_BOT_TOKEN") + .env_remove("SLACK_API_BASE_URL") .envs(envs) .stdin(Stdio::piped()) .stdout(Stdio::piped()) @@ -295,6 +392,159 @@ fn ingress_event_round_trips_slack_message_event() { assert_eq!(event["metadata"]["endpoint_id"], "slack-events"); } +#[test] +fn accepted_ingress_reacts_after_policy_checks_with_slack_timestamp() { + let bot_token_env = "SLACK_TEST_BOT_TOKEN_REACTIONS"; + let bot_token = "xoxb-test-reaction-token"; + let (base_url, request_rx, server) = serve_slack_api(vec![ + json!({ "ok": true }), + json!({ "ok": false, "error": "already_reacted" }), + ]); + let envs = BTreeMap::from([ + (bot_token_env.to_string(), bot_token.to_string()), + ("SLACK_API_BASE_URL".to_string(), base_url), + ]); + let event_body = |user: &str| { + json!({ + "type": "event_callback", + "team_id": "T123", + "event_id": "EvReaction123", + "event_time": 1712860000, + "event": { + "type": "app_mention", + "channel": "C123", + "channel_type": "channel", + "user": user, + "text": "hello from slack", + "client_msg_id": "client-generated-id", + "ts": "1712860000.100200", + "event_ts": "1712860000.100200" + } + }) + .to_string() + }; + let request = |body: String| { + json!({ + "protocol_version": 1, + "request": { + "kind": "ingress_event", + "config": { + "bot_token_env": bot_token_env, + "owner_id": "U123" + }, + "payload": { + "endpoint_id": "slack-events", + "method": "POST", + "path": "/slack/events", + "headers": {}, + "query": {}, + "body": body, + "trust_verified": true, + "received_at": "2026-04-11T18:00:00Z" + } + } + }) + }; + + let rejected = run_request_with_env(request(event_body("U999")), envs.clone()); + assert!( + rejected["events"] + .as_array() + .expect("events array") + .is_empty() + ); + assert!(request_rx.recv_timeout(Duration::from_millis(100)).is_err()); + + for _ in 0..2 { + let (accepted, stderr) = + run_request_with_env_and_stderr(request(event_body("U123")), envs.clone()); + let event = &accepted["events"][0]; + assert_eq!(event["message"]["id"], "client-generated-id"); + assert_eq!(event["metadata"]["message_ts"], "1712860000.100200"); + assert!(!stderr.contains("slack inbound acknowledgement failed")); + } + + for _ in 0..2 { + let api_request = request_rx + .recv_timeout(Duration::from_secs(2)) + .expect("reactions.add request"); + assert_eq!(api_request.request_line, "POST /api/reactions.add HTTP/1.1"); + assert_eq!( + api_request.authorization.as_deref(), + Some("Bearer xoxb-test-reaction-token") + ); + assert_eq!(api_request.body["channel"], "C123"); + assert_eq!(api_request.body["timestamp"], "1712860000.100200"); + assert_eq!(api_request.body["name"], "eyes"); + assert!(api_request.body.get("client_msg_id").is_none()); + } + server.join().expect("Slack API server"); +} + +#[test] +fn reaction_failure_is_fail_open_and_diagnostic_is_content_free() { + let bot_token_env = "SLACK_TEST_BOT_TOKEN_REACTION_FAILURE"; + let bot_token = "xoxb-sensitive-test-token"; + let (base_url, request_rx, server) = + serve_slack_api(vec![json!({ "ok": false, "error": "missing_scope" })]); + let body = json!({ + "type": "event_callback", + "team_id": "T123", + "event_id": "EvReactionFailure", + "event_time": 1712860000, + "event": { + "type": "app_mention", + "channel": "C-sensitive-test", + "channel_type": "channel", + "user": "U123", + "text": "sensitive message body", + "ts": "1712860000.100200" + } + }) + .to_string(); + let (response, stderr) = run_request_with_env_and_stderr( + json!({ + "protocol_version": 1, + "request": { + "kind": "ingress_event", + "config": { "bot_token_env": bot_token_env }, + "payload": { + "endpoint_id": "slack-events", + "method": "POST", + "path": "/slack/events", + "headers": {}, + "query": {}, + "body": body, + "trust_verified": true, + "received_at": "2026-04-11T18:00:00Z" + } + } + }), + BTreeMap::from([ + (bot_token_env.to_string(), bot_token.to_string()), + ("SLACK_API_BASE_URL".to_string(), base_url), + ]), + ); + + assert_eq!( + response["events"].as_array().expect("events array").len(), + 1 + ); + request_rx + .recv_timeout(Duration::from_secs(2)) + .expect("reactions.add request"); + assert_eq!(stderr.trim(), "slack inbound acknowledgement failed"); + for sensitive in [ + bot_token, + "C-sensitive-test", + "sensitive message body", + "missing_scope", + ] { + assert!(!stderr.contains(sensitive)); + } + server.join().expect("Slack API server"); +} + #[test] fn start_ingress_emits_slack_socket_mode_event() { let app_token_env = "SLACK_TEST_APP_TOKEN_SOCKET";