feat: add Slack breakout thread tool (#1638)

* feat: add Slack breakout thread tool

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>

* chore: make fake LLM scripts declarative

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>

* fix: exclude slack_start_new_thread from plan mode

The breakout tool can dispatch a fresh agent run that starts outside the
current plan-mode state, bypassing the approval flow. Add it to
PLAN_MODE_EXCLUDED_TOOLS so it's hidden alongside the other mutating
tools while planning.

---------

Co-authored-by: Ramon Nogueira <270434257+ramon-langchain@users.noreply.github.com>
Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
(cherry picked from commit 747ce4bbe5)
This commit is contained in:
Ramon Nogueira 2026-06-29 16:16:54 -04:00 • committed by Adam Moussa
parent 5dca57146b
commit b2b49ed944
No known key found for this signature in database
9 changed files with 569 additions and 7 deletions

View file

@ -83,6 +83,7 @@ OPEN_SWE_SHARED_BASE = """You are **Open SWE**, an open-source agent built on La
### Communication
- Focus on the substance and keep summaries brief. Use light markdown (`###`/`####` headings, bold, code) — avoid `#`/`##` titles.
- In Slack, when a user asks to “break out,” “split out,” or “start a separate thread” for part of the work, summarize the requested aspect and relevant context into self-contained instructions, then call `slack_start_new_thread` instead of only replying in the current thread.
- When you post to Slack with `slack_thread_reply`, do not repeat that text in a later assistant message; the user can already see the Slack message.
- When delegated work to a subagent: the calling agent only sees your final message, so make it the complete answer.

View file

@ -86,6 +86,7 @@ from .tools import (
save_plan,
schedule_thread_wakeup,
slack_read_thread_messages,
slack_start_new_thread,
slack_thread_reply,
web_search,
)
@ -624,6 +625,7 @@ PLAN_MODE_EXCLUDED_TOOLS: frozenset[str] = frozenset(
"http_request",
"open_pull_request",
"request_pr_review",
"slack_start_new_thread",
"linear_create_issue",
"linear_update_issue",
"linear_delete_issue",
@ -958,6 +960,7 @@ async def get_agent(config: RunnableConfig) -> Pregel:
request_pr_review,
schedule_thread_wakeup,
slack_read_thread_messages,
slack_start_new_thread,
slack_thread_reply,
*corridor_tools,
*observability_tools,

View file

@ -21,6 +21,7 @@ from .save_plan import save_plan
from .schedule_thread_wakeup import schedule_thread_wakeup
from .search_repo_code import search_repo_code
from .slack_read_thread_messages import slack_read_thread_messages
from .slack_start_new_thread import slack_start_new_thread
from .slack_thread_reply import slack_thread_reply
from .update_finding import update_finding
from .web_search import web_search
@ -49,6 +50,7 @@ __all__ = [
"schedule_thread_wakeup",
"search_repo_code",
"slack_read_thread_messages",
"slack_start_new_thread",
"slack_thread_reply",
"update_finding",
"web_search",

View file

@ -0,0 +1,260 @@
import os
import re
from typing import Any
from langgraph.config import get_config
from langgraph_sdk import get_client
from ..dispatch import dispatch_agent_run
from ..utils.dashboard_links import dashboard_thread_url
from ..utils.slack import (
post_slack_top_level_message_with_ts,
post_slack_trace_reply,
store_slack_run_mapping,
)
from ..utils.thread_ids import generate_thread_id_from_slack_thread
LANGGRAPH_URL = os.environ.get("LANGGRAPH_URL") or os.environ.get(
"LANGGRAPH_URL_PROD", "http://localhost:2024"
)
_TITLE_MAX_CHARS = 160
_INSTRUCTIONS_MAX_CHARS = 12000
_VISIBLE_INSTRUCTIONS_MAX_CHARS = 2800
_REPO_RE = re.compile(r"^[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+$")
def _failure_hint(slack_error: str | None) -> str:
if slack_error == "msg_too_long":
return "Slack rejected the message as too long; retry with shorter title or instructions."
if slack_error in {"channel_not_found", "not_in_channel"}:
return "Slack rejected the channel; do not retry with another channel."
if slack_error and slack_error.startswith("rate_limited"):
retry_after = slack_error.partition(":")[2].strip()
if retry_after:
return f"Slack rate limited the request; wait at least {retry_after}s before retrying."
return "Slack rate limited the request; wait before retrying."
if slack_error == "missing_slack_bot_token":
return "Slack bot token is missing; do not retry."
if slack_error and slack_error.startswith("http_error:"):
return "Slack posting hit an HTTP error; retry once."
return "Slack post failed; retry once with concise instructions."
def _validate_text(value: str, *, field: str, max_chars: int) -> str | dict[str, Any]:
text = value.strip() if isinstance(value, str) else ""
if not text:
return {"success": False, "error": f"{field} is required"}
if len(text) > max_chars:
return {
"success": False,
"error": f"{field} is too long",
"max_chars": max_chars,
"actual_chars": len(text),
}
return text
def _resolve_repo(configurable: dict[str, Any], default_repo: str | None) -> dict[str, str] | None:
if default_repo and default_repo.strip():
candidate = default_repo.strip()
if not _REPO_RE.fullmatch(candidate):
return None
owner, name = candidate.split("/", 1)
return {"owner": owner, "name": name}
repo = configurable.get("repo")
if isinstance(repo, dict):
owner = repo.get("owner")
name = repo.get("name")
if isinstance(owner, str) and owner.strip() and isinstance(name, str) and name.strip():
return {"owner": owner.strip(), "name": name.strip()}
return None
def _truncate_for_slack(text: str) -> str:
if len(text) <= _VISIBLE_INSTRUCTIONS_MAX_CHARS:
return text
omitted = len(text) - _VISIBLE_INSTRUCTIONS_MAX_CHARS
return f"{text[:_VISIBLE_INSTRUCTIONS_MAX_CHARS].rstrip()}\n\n…truncated {omitted} chars; the new Open SWE thread received the full instructions."
def _visible_message(title: str, instructions: str, repo: dict[str, str] | None) -> str:
repo_line = f"\n*Repository:* `{repo['owner']}/{repo['name']}`" if repo else ""
return (
f"*Open SWE breakout thread:* {title}{repo_line}\n\n"
f"*Instructions for the new thread:*\n{_truncate_for_slack(instructions)}"
)
def _run_prompt(
title: str,
instructions: str,
repo: dict[str, str] | None,
original_slack_thread: dict[str, Any],
) -> str:
repo_text = f"{repo['owner']}/{repo['name']}" if repo else "(no repository specified)"
channel_id = original_slack_thread.get("channel_id", "")
thread_ts = original_slack_thread.get("thread_ts", "")
return (
"You were started from another Open SWE Slack thread as a breakout task.\n\n"
f"## Breakout Title\n{title}\n\n"
f"## Default Repository Hint\n{repo_text}\n"
"Use this repository unless the instructions below clearly identify a different repository.\n\n"
"## Source Slack Thread\n"
f"- Channel: {channel_id}\n"
f"- Thread TS: {thread_ts}\n\n"
"## Breakout Instructions\n"
f"{instructions}\n\n"
"Use `slack_thread_reply` to communicate in this new Slack thread for clarifications, "
"status updates, and final summaries."
)
def _new_slack_thread_context(
original: dict[str, Any],
*,
channel_id: str,
thread_ts: str,
) -> dict[str, Any]:
return {
"channel_id": channel_id,
"thread_ts": thread_ts,
"triggering_user_id": original.get("triggering_user_id", ""),
"triggering_user_name": original.get("triggering_user_name", ""),
"triggering_user_email": original.get("triggering_user_email", ""),
"triggering_event_ts": thread_ts,
}
async def slack_start_new_thread(
title: str,
instructions: str,
default_repo: str | None = None,
) -> dict[str, Any]:
"""Start a new Open SWE thread in a top-level Slack message in the current channel."""
config = get_config()
configurable = config.get("configurable", {})
current_slack_thread = configurable.get("slack_thread")
if not isinstance(current_slack_thread, dict):
return {"success": False, "error": "Missing slack_thread config"}
channel_id = current_slack_thread.get("channel_id")
current_thread_ts = current_slack_thread.get("thread_ts")
if not isinstance(channel_id, str) or not channel_id.strip():
return {"success": False, "error": "Missing slack_thread.channel_id in config"}
clean_title = _validate_text(title, field="title", max_chars=_TITLE_MAX_CHARS)
if isinstance(clean_title, dict):
return clean_title
clean_instructions = _validate_text(
instructions, field="instructions", max_chars=_INSTRUCTIONS_MAX_CHARS
)
if isinstance(clean_instructions, dict):
return clean_instructions
repo = _resolve_repo(configurable, default_repo)
if default_repo and default_repo.strip() and repo is None:
return {
"success": False,
"error": "default_repo must be a simple owner/name repository string",
}
message_ts, slack_error = await post_slack_top_level_message_with_ts(
channel_id.strip(),
_visible_message(clean_title, clean_instructions, repo),
unfurl_links=False,
unfurl_media=False,
)
if message_ts is None:
return {
"success": False,
"error": slack_error or "post failed",
"slack_error": slack_error,
"hint": _failure_hint(slack_error),
}
thread_id = generate_thread_id_from_slack_thread(channel_id.strip(), message_ts)
new_slack_thread = _new_slack_thread_context(
current_slack_thread,
channel_id=channel_id.strip(),
thread_ts=message_ts,
)
breakout_from = {
"channel_id": channel_id.strip(),
"thread_ts": current_thread_ts or "",
"message_ts": current_slack_thread.get("triggering_event_ts", ""),
}
metadata: dict[str, Any] = {
"source": "slack",
"title": clean_title[:80],
"source_context": {
"slack_thread": new_slack_thread,
"breakout_from": breakout_from,
},
}
if repo:
metadata.update(
{
"repo": repo,
"repo_owner": repo["owner"],
"repo_name": repo["name"],
}
)
github_login = configurable.get("github_login")
if isinstance(github_login, str) and github_login:
metadata["github_login"] = github_login
user_email = configurable.get("user_email")
if isinstance(user_email, str) and user_email:
metadata["triggering_user_email"] = user_email.strip().lower()
new_configurable: dict[str, Any] = {
"slack_thread": new_slack_thread,
"source": "slack",
}
if repo:
new_configurable["repo"] = repo
for key in ("user_email", "github_login", "agent_model_id", "agent_effort"):
value = configurable.get(key)
if value:
new_configurable[key] = value
client = get_client(url=LANGGRAPH_URL)
await client.threads.create(thread_id=thread_id, if_exists="do_nothing", metadata=metadata)
await client.threads.update(thread_id=thread_id, metadata=metadata)
run = await dispatch_agent_run(
thread_id,
_run_prompt(clean_title, clean_instructions, repo, current_slack_thread),
new_configurable,
source="slack",
client=client,
)
run_id = run.get("run_id") if isinstance(run, dict) else None
trace_message_ts = await post_slack_trace_reply(channel_id.strip(), message_ts, thread_id)
if isinstance(run_id, str) and run_id:
await store_slack_run_mapping(
client,
channel_id.strip(),
message_ts,
run_id,
message_ts=message_ts,
triggering_user_id=new_slack_thread.get("triggering_user_id") or None,
)
if trace_message_ts:
await store_slack_run_mapping(
client,
channel_id.strip(),
message_ts,
run_id,
message_ts=trace_message_ts,
triggering_user_id=new_slack_thread.get("triggering_user_id") or None,
)
return {
"success": True,
"thread_id": thread_id,
"thread_ts": message_ts,
"dashboard_url": dashboard_thread_url(thread_id),
}

View file

@ -125,6 +125,7 @@ from .utils.slack_feedback import (
process_slack_reaction_added,
process_slack_reaction_removed,
)
from .utils.thread_ids import generate_thread_id_from_slack_thread
logger = logging.getLogger(__name__)
@ -378,13 +379,6 @@ def generate_thread_id_from_github_issue(issue_id: str) -> str:
)
def generate_thread_id_from_slack_thread(channel_id: str, thread_id: str) -> str:
"""Generate a deterministic thread ID from a Slack thread identifier."""
composite = f"{channel_id}:{thread_id}"
md5_hex = hashlib.md5(composite.encode("utf-8")).hexdigest()
return str(uuid.UUID(hex=md5_hex))
def generate_reviewer_thread_id(owner: str, repo: str, pr_number: int) -> str:
stable_key = f"{owner}/{repo}/pr/{pr_number}/reviewer"
return str(uuid.uuid5(uuid.NAMESPACE_URL, stable_key))

View file

@ -308,6 +308,23 @@ SCRIPT_LIBRARY: dict[str, tuple[StepSpec, ...]] = {
),
_dynamic_step(_reply_step),
),
"breakout": (
_tool_step(
"Starting a separate Slack thread for the breakout task.",
"slack_start_new_thread",
{
"title": "Add greet() helper",
"instructions": "Please add a greet() helper and open a draft PR in the default repository. Use the current Slack request as context, and report progress in this new thread.",
},
"call-breakout",
),
_tool_step(
"Confirming the breakout thread was started.",
"slack_thread_reply",
{"message": "I started a separate Open SWE thread for that aspect."},
"call-breakout-reply",
),
),
"plan": (
_tool_step(
"This is worth planning first — entering plan mode.",
@ -330,6 +347,11 @@ def _is_plan_request(text: str) -> bool:
return "plan" in text.lower()
def _is_breakout_request(text: str) -> bool:
t = text.lower()
return "break out" in t or "separate thread" in t or "split out" in t
def _is_approval(text: str) -> bool:
t = text.lower()
return "approved" in t and "implement" in t
@ -344,6 +366,9 @@ SCRIPT_RULES: tuple[ScriptRule, ...] = (
ScriptRule("implement", lambda ctx: _is_approval(ctx.last_text)),
ScriptRule("plan", lambda ctx: _is_revision(ctx.last_text)),
ScriptRule("plan", lambda ctx: ctx.human_count <= 1 and _is_plan_request(ctx.first_text)),
ScriptRule(
"breakout", lambda ctx: ctx.human_count <= 1 and _is_breakout_request(ctx.first_text)
),
ScriptRule("implement", lambda ctx: ctx.human_count <= 1),
ScriptRule("followup", lambda _ctx: True),
)

View file

@ -37,6 +37,26 @@ test.describe("Open SWE full flow", () => {
await expect(page.locator('.pr[data-pr="1"]')).toContainText("greet.py");
});
test("Slack breakout request starts a new top-level Open SWE thread", async ({ page }) => {
await page.locator("#text").fill("<@U0BOT> please break out adding a greet() helper into a separate thread");
await page.locator("#send").click();
const breakout = page
.locator(".msg.bot")
.filter({ hasText: /Open SWE breakout thread:\* Add greet\(\) helper/ });
await expect(breakout).toBeVisible({ timeout: 60_000 });
const breakoutThreadTs = await breakout.getAttribute("data-thread-ts");
expect(breakoutThreadTs).toBeTruthy();
const breakoutThreadMessages = page.locator(`.msg.bot[data-thread-ts="${breakoutThreadTs}"]`);
await expect(breakoutThreadMessages.locator('a[href*="/agents/"]')).toBeVisible({
timeout: 60_000,
});
await expect(
page.locator(".msg.bot").filter({ hasText: "I started a separate Open SWE thread" }),
).toBeVisible({ timeout: 60_000 });
});
test("a message that does not mention the bot produces no run and no PR", async ({ page }) => {
await page.locator("#mention").uncheck();
await page.locator("#text").fill("just chatting with the team, nothing for the bot");

View file

@ -29,6 +29,7 @@ def test_plan_mode_excluded_tools_cover_mutating_tools() -> None:
"task",
"open_pull_request",
"request_pr_review",
"slack_start_new_thread",
"linear_create_issue",
"linear_update_issue",
"linear_delete_issue",

View file

@ -0,0 +1,256 @@
from __future__ import annotations
import importlib
from typing import Any
import pytest
from agent.utils.thread_ids import generate_thread_id_from_slack_thread
slack_breakout_tool = importlib.import_module("agent.tools.slack_start_new_thread")
def _config() -> dict[str, Any]:
return {
"configurable": {
"repo": {"owner": "langchain-ai", "name": "open-swe"},
"github_login": "alice",
"user_email": "alice@example.com",
"agent_model_id": "anthropic:claude-sonnet-4-5",
"agent_effort": "high",
"slack_thread": {
"channel_id": "C1",
"thread_ts": "1700000000.000001",
"triggering_user_id": "U1",
"triggering_user_name": "Alice",
"triggering_user_email": "alice@example.com",
"triggering_event_ts": "1700000000.000002",
},
}
}
class _FakeThreadsClient:
def __init__(self, captured: dict[str, Any]) -> None:
self.captured = captured
async def create(self, *, thread_id: str, if_exists: str, metadata: dict[str, Any]) -> None:
self.captured["thread_create"] = {
"thread_id": thread_id,
"if_exists": if_exists,
"metadata": metadata,
}
async def update(self, *, thread_id: str, metadata: dict[str, Any]) -> None:
self.captured["thread_update"] = {"thread_id": thread_id, "metadata": metadata}
class _FakeClient:
def __init__(self, captured: dict[str, Any]) -> None:
self.threads = _FakeThreadsClient(captured)
async def test_slack_start_new_thread_success(monkeypatch: pytest.MonkeyPatch) -> None:
captured: dict[str, Any] = {"stored_mappings": []}
new_ts = "1700000000.111111"
trace_ts = "1700000000.222222"
async def fake_post_top_level(
channel_id: str,
text: str,
*,
unfurl_links: bool = True,
unfurl_media: bool = True,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
captured["top_level_post"] = {
"channel_id": channel_id,
"text": text,
"unfurl_links": unfurl_links,
"unfurl_media": unfurl_media,
"blocks": blocks,
}
return new_ts, None
async def fake_dispatch_agent_run(
thread_id: str,
content: str,
configurable: dict[str, Any],
*,
source: str,
client: Any,
**kwargs: Any,
) -> dict[str, str]:
captured["dispatch"] = {
"thread_id": thread_id,
"content": content,
"configurable": configurable,
"source": source,
"client": client,
"kwargs": kwargs,
}
return {"run_id": "run-123"}
async def fake_post_trace(channel_id: str, thread_ts: str, thread_id: str) -> str:
captured["trace"] = {
"channel_id": channel_id,
"thread_ts": thread_ts,
"thread_id": thread_id,
}
return trace_ts
async def fake_store_mapping(
client: Any,
channel_id: str,
thread_ts: str,
run_id: str,
*,
message_ts: str | None = None,
triggering_user_id: str | None = None,
) -> None:
captured["stored_mappings"].append(
{
"client": client,
"channel_id": channel_id,
"thread_ts": thread_ts,
"run_id": run_id,
"message_ts": message_ts,
"triggering_user_id": triggering_user_id,
}
)
fake_client = _FakeClient(captured)
monkeypatch.setattr(slack_breakout_tool, "get_config", _config)
monkeypatch.setattr(slack_breakout_tool, "get_client", lambda url: fake_client)
monkeypatch.setattr(
slack_breakout_tool, "post_slack_top_level_message_with_ts", fake_post_top_level
)
monkeypatch.setattr(slack_breakout_tool, "dispatch_agent_run", fake_dispatch_agent_run)
monkeypatch.setattr(slack_breakout_tool, "post_slack_trace_reply", fake_post_trace)
monkeypatch.setattr(slack_breakout_tool, "store_slack_run_mapping", fake_store_mapping)
monkeypatch.setattr(
slack_breakout_tool,
"dashboard_thread_url",
lambda thread_id: f"https://dashboard.example/agents/{thread_id}",
)
result = await slack_breakout_tool.slack_start_new_thread(
"Investigate follow-up",
"Use the same repo and investigate the follow-up aspect in detail.",
)
expected_thread_id = generate_thread_id_from_slack_thread("C1", new_ts)
assert result == {
"success": True,
"thread_id": expected_thread_id,
"thread_ts": new_ts,
"dashboard_url": f"https://dashboard.example/agents/{expected_thread_id}",
}
assert captured["top_level_post"]["channel_id"] == "C1"
assert "Investigate follow-up" in captured["top_level_post"]["text"]
assert "langchain-ai/open-swe" in captured["top_level_post"]["text"]
assert captured["top_level_post"]["unfurl_links"] is False
assert captured["thread_create"]["if_exists"] == "do_nothing"
assert captured["thread_create"]["thread_id"] == expected_thread_id
metadata = captured["thread_update"]["metadata"]
assert metadata["source"] == "slack"
assert metadata["repo"] == {"owner": "langchain-ai", "name": "open-swe"}
assert metadata["github_login"] == "alice"
assert metadata["triggering_user_email"] == "alice@example.com"
assert metadata["source_context"]["slack_thread"]["thread_ts"] == new_ts
assert metadata["source_context"]["slack_thread"]["triggering_user_id"] == "U1"
assert metadata["source_context"]["breakout_from"] == {
"channel_id": "C1",
"thread_ts": "1700000000.000001",
"message_ts": "1700000000.000002",
}
dispatch = captured["dispatch"]
assert dispatch["thread_id"] == expected_thread_id
assert dispatch["source"] == "slack"
assert dispatch["configurable"]["slack_thread"]["thread_ts"] == new_ts
assert dispatch["configurable"]["repo"] == {"owner": "langchain-ai", "name": "open-swe"}
assert dispatch["configurable"]["github_login"] == "alice"
assert dispatch["configurable"]["agent_model_id"] == "anthropic:claude-sonnet-4-5"
assert "Breakout Instructions" in dispatch["content"]
assert captured["trace"] == {
"channel_id": "C1",
"thread_ts": new_ts,
"thread_id": expected_thread_id,
}
assert [item["message_ts"] for item in captured["stored_mappings"]] == [new_ts, trace_ts]
assert all(item["triggering_user_id"] == "U1" for item in captured["stored_mappings"])
async def test_slack_start_new_thread_requires_slack_config(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(slack_breakout_tool, "get_config", lambda: {"configurable": {}})
result = await slack_breakout_tool.slack_start_new_thread("Title", "Instructions")
assert result == {"success": False, "error": "Missing slack_thread config"}
@pytest.mark.parametrize(
("title", "instructions", "error"),
[
("", "Instructions", "title is required"),
("Title", "", "instructions is required"),
("x" * 161, "Instructions", "title is too long"),
("Title", "x" * 12001, "instructions is too long"),
],
)
async def test_slack_start_new_thread_validates_text(
title: str,
instructions: str,
error: str,
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(slack_breakout_tool, "get_config", _config)
result = await slack_breakout_tool.slack_start_new_thread(title, instructions)
assert result["success"] is False
assert result["error"] == error
async def test_slack_start_new_thread_rejects_invalid_repo_override(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(slack_breakout_tool, "get_config", _config)
result = await slack_breakout_tool.slack_start_new_thread(
"Title", "Instructions", default_repo="https://github.com/langchain-ai/open-swe"
)
assert result == {
"success": False,
"error": "default_repo must be a simple owner/name repository string",
}
async def test_slack_start_new_thread_returns_slack_failure_without_dispatch(
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, bool] = {"dispatched": False}
async def fake_post_top_level(*args: Any, **kwargs: Any) -> tuple[str | None, str | None]:
return None, "msg_too_long"
async def fake_dispatch_agent_run(*args: Any, **kwargs: Any) -> dict[str, str]:
captured["dispatched"] = True
return {"run_id": "run-123"}
monkeypatch.setattr(slack_breakout_tool, "get_config", _config)
monkeypatch.setattr(
slack_breakout_tool, "post_slack_top_level_message_with_ts", fake_post_top_level
)
monkeypatch.setattr(slack_breakout_tool, "dispatch_agent_run", fake_dispatch_agent_run)
result = await slack_breakout_tool.slack_start_new_thread("Title", "Instructions")
assert result["success"] is False
assert result["error"] == "msg_too_long"
assert result["slack_error"] == "msg_too_long"
assert "shorter" in result["hint"]
assert captured["dispatched"] is False