From ac9def34fab047196b94dc44dec86d7ca9fe460b Mon Sep 17 00:00:00 2001 From: Oded Falik Date: Sat, 1 Aug 2026 19:47:37 -0400 Subject: [PATCH] Remove the passive Slack listener; outbound only The listener spawned a full `claude -p` for every unowned message in a watched channel. That is expensive, answers people who were not talking to us, and is superseded by QM, which handles inbound Slack with per-person scopes and its own turn detection. Removed: Socket Mode connection, leader election and its lock file, the JSONL event bus, the event-bus tailer, cold replies, and channel notifications -- along with the claude/channel capability, since nothing pushes any more. SLACK_APP_TOKEN is no longer read at all. Kept: every outbound tool (reply, read_history, get_thread, fetch_file, add_reaction, remove_reaction, list_channels, debug) and per-conversation thread tracking. 1234 -> 887 lines. Verified against the venv actually in use (mcp 1.28.1): module imports, all eight tools present, no listener symbols, and auth_test still succeeds as @golem. Note: a fresh resolve picks mcp 2.0.0, where the lowlevel Server API has no list_tools and the module fails to import. That pre-dates this branch -- pristine main fails identically -- but pyproject's `mcp>=1.26.0` should probably be capped. --- README.md | 4 + slack_channel/server.py | 379 ++-------------------------------------- 2 files changed, 20 insertions(+), 363 deletions(-) diff --git a/README.md b/README.md index c45eb2c..2531f5e 100644 --- a/README.md +++ b/README.md @@ -1,5 +1,9 @@ # slack-channel-plugin +> **Outbound only.** The passive Socket Mode listener was removed — this server posts +> to Slack and reads history, but never pushes inbound messages to an agent. Inbound +> Slack is handled by QM. + Slack channel plugin for [Claude Code](https://claude.ai/code) — real-time messaging via Socket Mode with per-conversation thread routing. ## What it does diff --git a/slack_channel/server.py b/slack_channel/server.py index ac83b60..f4fad0c 100644 --- a/slack_channel/server.py +++ b/slack_channel/server.py @@ -1,14 +1,13 @@ -"""Slack channel plugin — MCP server over stdio with Socket Mode listener. +"""Slack channel plugin — outbound MCP server over stdio. -Connects to Slack via Socket Mode, declares claude/channel capability, -emits notifications/claude/channel on incoming messages, and exposes -reply + read tools. +Exposes reply + read tools so an agent can post to Slack, read history and react. +It does not listen: there is no Socket Mode connection, no event bus and no +cold-reply path. -Architecture: - - One Socket Mode connection across all instances (leader election via lock file) - - Leader writes events to a shared JSONL file (event bus) - - ALL instances tail the event bus and send channel notifications for their threads - - This ensures every session gets real-time push for threads it created +The passive listener was removed deliberately. It spawned a `claude -p` for every +unowned message in a watched channel, which was expensive, answered people who +were not talking to us, and is superseded by QM — which handles inbound Slack +properly, with per-person scopes and its own turn detection. """ from __future__ import annotations @@ -32,9 +31,6 @@ from mcp.server.models import InitializationOptions from mcp.server.session import ServerSession from mcp.server.stdio import stdio_server -from mcp.shared.message import SessionMessage -from slack_bolt.adapter.socket_mode.async_handler import AsyncSocketModeHandler -from slack_bolt.async_app import AsyncApp # Load env vars. Priority: real env vars > plugin channels dir > local .env # The channels dir (~/.claude/channels/slack-channel/.env) is the standard @@ -54,7 +50,6 @@ _slack_client: Any = None # AsyncWebClient (bot token) — used for writes + notifications _read_client: Any = None # AsyncWebClient (user xoxp token if set, else falls back to bot) — used for reads _bot_user_id: str | None = None -_pending_eyes: dict[str, str] = {} # thread_ts → message_ts (messages with eyes to remove) _user_name_cache: dict[str, str] = {} CHANNEL_ID = os.environ.get("SLACK_CHANNEL_ID", "") @@ -66,8 +61,6 @@ ) if os.environ.get("SLACK_CHANNEL_STATE_DIR") else ( Path.home() / ".config" / "slack-channel" ) -_EVENTS_FILE = _STATE_DIR / "events.jsonl" -_LOCK_FILE = _STATE_DIR / "listener.lock" # Thread tracking — persisted per conversation ID. # The conversation ID is stable across --resume (it's what you pass to --resume). @@ -204,36 +197,9 @@ def _save_threads() -> None: json.dump(other + mine, f, indent=2) f.write("\n") -# Leader election state -_is_leader = False -_lock_fd: Any = None # file descriptor held for lifetime - -def _try_become_leader() -> bool: - """Try to acquire the listener lock. Non-blocking.""" - global _lock_fd - _STATE_DIR.mkdir(parents=True, exist_ok=True) - _lock_fd = open(_LOCK_FILE, "w") - try: - fcntl.flock(_lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) - _lock_fd.write(str(os.getpid())) - _lock_fd.flush() - return True - except OSError: - _lock_fd.close() - _lock_fd = None - return False - - -def _append_event(event: dict) -> None: - """Append an event to the shared event bus (JSONL).""" - with open(_EVENTS_FILE, "a") as f: - fcntl.flock(f, fcntl.LOCK_EX) - f.write(json.dumps(event) + "\n") - - # --------------------------------------------------------------------------- # Channel/DM name resolution # --------------------------------------------------------------------------- @@ -564,17 +530,6 @@ async def _handle_reply(args: dict) -> list[types.TextContent]: _save_threads() logger.info("Tracking thread %s in %s", owned_ts, channel) - # Auto-remove eyes reaction from the message we're responding to - effective_thread = thread_ts or ts - if effective_thread in _pending_eyes: - eyes_ts = _pending_eyes.pop(effective_thread) - try: - await _slack_client.reactions_remove( - channel=channel, timestamp=eyes_ts, name="eyes", - ) - except Exception: - pass - return [types.TextContent(type="text", text=f"Sent (ts={ts}, channel={channel})")] @@ -674,7 +629,6 @@ async def _handle_debug(args: dict) -> list[types.TextContent]: "pid": os.getpid(), "ppid": os.getppid(), "conversation_id": _conversation_id, - "is_leader": _is_leader, "bot_user_id": _bot_user_id, "owned_threads": _owned_threads, "channel_filter": CHANNEL_ID, @@ -837,47 +791,6 @@ async def _resolve_user(user_id: str) -> str: # Channel notifications → MCP client # --------------------------------------------------------------------------- -async def _send_channel_notification(event: dict) -> None: - """Push a notifications/claude/channel message over the MCP session. - - Claude Code expects params = { content: str, meta: dict }. - It renders as: content - """ - if _session is None: - return - - sender = event.get("_sender_name", event.get("user", "unknown")) - text = event.get("text", "") - - meta: dict[str, str] = { - "channel": event.get("channel", ""), - "sender": sender, - "sender_id": event.get("user", ""), - "ts": event.get("ts", ""), - } - thread_ts = event.get("thread_ts") - if thread_ts: - meta["thread_ts"] = thread_ts - - try: - # Build the raw JSON-RPC notification and write directly to the - # session's write stream. We bypass session.send_message() because - # the MCP SDK validates against ServerNotificationType (a closed union - # of standard methods) and rejects our custom method. - notification = types.JSONRPCNotification( - jsonrpc="2.0", - method="notifications/claude/channel", - params={ - "content": f"{sender}: {text}", - "meta": meta, - }, - ) - raw_msg = SessionMessage(message=types.JSONRPCMessage(notification)) - await _session._write_stream.send(raw_msg) - except Exception: - logger.debug("Failed to send channel notification", exc_info=True) - - # --------------------------------------------------------------------------- # Cold replies (leader only) — for messages no active session owns # --------------------------------------------------------------------------- @@ -886,254 +799,20 @@ async def _send_channel_notification(event: dict) -> None: REACT_PREFIX = "REACT:" -async def _fetch_thread_context(channel: str, thread_ts: str) -> str: - """Fetch all messages in a thread and format as context.""" - try: - result = await _slack_client.conversations_replies( - channel=channel, ts=thread_ts, limit=50, - ) - except Exception: - return "(failed to fetch thread)" - - lines = [] - for msg in result.get("messages", []): - user = await _resolve_user(msg.get("user", "")) - text = msg.get("text", "") + _file_annotations(msg) - lines.append(f"{user}: {text}") - return "\n".join(lines) - - -async def _cold_reply(event: dict, channel: str, thread_ts: str, msg_ts: str) -> None: - """Spawn claude -p to respond to an unowned message.""" - import subprocess as sp - - sender = event.get("_sender_name", "someone") - logger.info("Cold reply: thread=%s sender=%s", thread_ts, sender) - - # Add eyes reaction - try: - await _slack_client.reactions_add(channel=channel, timestamp=msg_ts, name="eyes") - except Exception: - pass - - # Fetch thread context - context = await _fetch_thread_context(channel, thread_ts) - - system_prompt = ( - "You are a helpful AI assistant replying in a Slack thread. " - "Keep responses concise. " - "You have three response modes:\n" - f"1. If no response is needed — respond with exactly: {NO_RESPONSE_SENTINEL}\n" - f"2. If a simple emoji reaction is better — respond with: {REACT_PREFIX} " - f"(e.g. REACT:thumbsup, REACT:white_check_mark)\n" - "3. Otherwise, respond normally with text.\n" - "Use your judgment — a thumbs-up is often better than a wordy acknowledgment." - ) - - full_prompt = f"{system_prompt}\n\n## Thread context\n{context}" - - try: - result = sp.run( - [ - "claude", "-p", full_prompt, - "--allowedTools", - "Read", "Glob", "Grep", - "mcp__slack-channel__reply", - "mcp__slack-channel__get_thread", - "mcp__slack-channel__read_history", - "mcp__slack-channel__add_reaction", - "mcp__claude_ai_Google_Calendar__*", - "mcp__claude_ai_Gmail__gmail_search_messages", - "mcp__claude_ai_Gmail__gmail_read_message", - "mcp__claude_ai_Gmail__gmail_read_thread", - "mcp__claude_ai_Gmail__gmail_get_profile", - "mcp__claude_ai_Gmail__gmail_list_labels", - "mcp__claude_ai_Notion__*", - ], - capture_output=True, text=True, timeout=300, - ) - output = result.stdout.strip() or "(no response)" - except sp.TimeoutExpired: - output = "(timed out)" - except Exception as e: - output = f"(error: {e})" - - # Remove eyes - try: - await _slack_client.reactions_remove(channel=channel, timestamp=msg_ts, name="eyes") - except Exception: - pass - - # Handle response modes - stripped = output.strip() - if stripped == NO_RESPONSE_SENTINEL: - logger.info("Cold reply: no response needed for %s", thread_ts) - return - - if stripped.startswith(REACT_PREFIX): - emoji = stripped[len(REACT_PREFIX):].strip() - logger.info("Cold reply: reacting with :%s: to %s", emoji, thread_ts) - try: - await _slack_client.reactions_add(channel=channel, timestamp=msg_ts, name=emoji) - except Exception: - pass - return - - # Send text reply - try: - if len(output) <= 4000: - await _slack_client.chat_postMessage( - channel=channel, text=output, thread_ts=thread_ts, - ) - else: - for i in range(0, len(output), 3900): - await _slack_client.chat_postMessage( - channel=channel, text=output[i:i + 3900], thread_ts=thread_ts, - ) - except Exception: - logger.error("Cold reply: failed to send response", exc_info=True) - - # --------------------------------------------------------------------------- # Event bus: leader writes, all instances read # --------------------------------------------------------------------------- -async def _watch_event_bus() -> None: - """Tail the shared event bus file and send notifications for owned threads. - - Polls every 0.5s for new lines. Each instance filters for its own threads. - """ - _STATE_DIR.mkdir(parents=True, exist_ok=True) - # Start from end of file (don't replay old events) - try: - pos = _EVENTS_FILE.stat().st_size - except OSError: - pos = 0 - - while True: - await anyio.sleep(0.5) - - try: - size = _EVENTS_FILE.stat().st_size - except OSError: - continue - - if size <= pos: - if size < pos: - pos = 0 # file was truncated - continue - - try: - with open(_EVENTS_FILE) as f: - f.seek(pos) - new_data = f.read() - pos = f.tell() - except OSError: - continue - - - # Reload threads from disk (another instance may have added threads via reply) - _owned_threads.update(_load_threads()) - - for line in new_data.strip().split("\n"): - if not line: - continue - try: - event = json.loads(line) - except json.JSONDecodeError: - continue - - thread_ts = event.get("thread_ts") - - if thread_ts and thread_ts in _owned_threads: - # Reply to an owned thread → notify this session - logger.info( - "Event bus → notify: thread=%s sender=%s", - thread_ts, event.get("_sender_name", "?"), - ) - msg_ts = event.get("ts", "") - msg_channel = event.get("channel", "") - if msg_ts and msg_channel: - try: - await _slack_client.reactions_add( - channel=msg_channel, timestamp=msg_ts, name="eyes", - ) - _pending_eyes[thread_ts] = msg_ts - except Exception: - pass - await _send_channel_notification(event) - elif _is_leader: - # Unowned message (including @mentions, new threads) → cold reply - msg_ts = event.get("ts", "") - msg_channel = event.get("channel", "") - if msg_channel: - import asyncio - asyncio.get_event_loop().create_task( - _cold_reply(event, msg_channel, thread_ts or msg_ts, msg_ts) - ) - - # --------------------------------------------------------------------------- # Slack Socket Mode listener (leader only) # --------------------------------------------------------------------------- -def _create_slack_app(bot_token: str) -> AsyncApp: - """Create the Slack Bolt async app with message handler.""" - app = AsyncApp(token=bot_token) - - @app.event("message") - async def on_message(event: dict, say: Any) -> None: - # Ignore subtypes (edits, deletes, joins, …) but allow file_share - subtype = event.get("subtype") - if subtype and subtype not in ("file_share",): - return - # Ignore bot messages - if event.get("bot_id"): - return - # Ignore our own messages - if event.get("user") == _bot_user_id: - return - # If a channel filter is set, only watch that channel - if CHANNEL_ID and event.get("channel") != CHANNEL_ID: - return - - # Resolve user name (cache it in the event for downstream) - sender = await _resolve_user(event.get("user", "")) - event["_sender_name"] = sender - - logger.info( - "Socket Mode event: user=%s channel=%s thread=%s", - sender, event.get("channel"), event.get("thread_ts"), - ) - - # Write to event bus — all instances will pick this up - _append_event(event) - - return app - - -async def _run_slack_listener(app: AsyncApp, app_token: str) -> None: - """Connect to Slack Socket Mode and listen forever (leader only).""" - global _bot_user_id - - try: - auth = await app.client.auth_test() - _bot_user_id = auth["user_id"] - logger.info("Leader: Slack connected — bot user %s (%s)", auth.get("user"), _bot_user_id) - except Exception: - logger.error("Slack auth_test failed", exc_info=True) - return - - handler = AsyncSocketModeHandler(app, app_token) - await handler.start_async() - - # --------------------------------------------------------------------------- # Entrypoint # --------------------------------------------------------------------------- async def _main() -> None: - global _session, _slack_client, _read_client, _is_leader, _bot_user_id + global _session, _slack_client, _read_client, _bot_user_id logging.basicConfig( level=logging.INFO, @@ -1142,19 +821,17 @@ async def _main() -> None: ) bot_token = os.environ.get("SLACK_BOT_TOKEN") - app_token = os.environ.get("SLACK_APP_TOKEN") if not bot_token: logger.error("SLACK_BOT_TOKEN is required") sys.exit(1) - # All instances need a Slack client for tools from slack_sdk.web.async_client import AsyncWebClient _slack_client = AsyncWebClient(token=bot_token) # Reads (history, channels, DMs) use the user token if provided so they see # everything Oded can see — including DMs the bot was never invited to. - # Writes/reactions/notifications stay on the bot token above (messages post as @Golem). + # Writes and reactions stay on the bot token above (messages post as @Golem). user_token = os.environ.get("SLACK_USER_TOKEN") _read_client = AsyncWebClient(token=user_token) if user_token else _slack_client if user_token: @@ -1162,42 +839,26 @@ async def _main() -> None: else: logger.info("SLACK_USER_TOKEN not set — reads fall back to bot token") - # Resolve bot user ID (needed for filtering even in non-leader instances) + # Resolve bot user ID (used to label our own messages) try: auth = await _slack_client.auth_test() _bot_user_id = auth["user_id"] except Exception: logger.warning("Could not resolve bot user ID", exc_info=True) - # Leader election: try to grab the Socket Mode lock - slack_app: AsyncApp | None = None - if app_token: - _is_leader = _try_become_leader() - if _is_leader: - logger.info("Elected as listener leader (pid=%d)", os.getpid()) - slack_app = _create_slack_app(bot_token) - _slack_client = slack_app.client # Use bolt's client (same token) - else: - logger.info("Another instance is the listener leader — watching event bus only") - else: - logger.warning("SLACK_APP_TOKEN not set — real-time notifications disabled") - async with stdio_server() as (read_stream, write_stream): init_options = InitializationOptions( server_name="slack-channel", server_version="0.1.0", capabilities=server.get_capabilities( notification_options=NotificationOptions(), - experimental_capabilities={"claude/channel": {}}, ), instructions=( - "Slack channel plugin. Listens for messages via Socket Mode and " - "emits notifications/claude/channel for replies to threads you " - "created and @mentions. Use the reply tool to send messages. " - "Threads are automatically tracked per conversation and survive " - "--resume. When you receive a channel notification, an eyes emoji " - "reaction is automatically added to the message and removed when " - "you reply in the thread." + "Slack channel plugin — outbound only. Use the reply tool to post to a " + "channel or thread, and read_history / get_thread to read. This server " + "does not listen for Slack messages and will never push anything to you; " + "inbound Slack is handled by QM. Threads you reply in are tracked per " + "conversation and survive --resume." ), ) @@ -1209,14 +870,6 @@ async def _main() -> None: _session = session async with anyio.create_task_group() as tg: - # Leader: run Socket Mode listener - if slack_app and app_token: - tg.start_soon(_run_slack_listener, slack_app, app_token) - - # ALL instances: watch event bus for notifications - if app_token: - tg.start_soon(_watch_event_bus) - # MCP message loop async for message in session.incoming_messages: tg.start_soon(