Compare commits

..

1 commit

Author SHA1 Message Date
dependabot[bot]
3fb2bf28d2
chore(deps): bump the minor-and-patch group across 1 directory with 18 updates
Bumps the minor-and-patch group with 14 updates in the / directory:

| Package | From | To |
| --- | --- | --- |
| [deepagents](https://github.com/langchain-ai/deepagents) | `0.6.8` | `0.6.12` |
| [fastapi](https://github.com/fastapi/fastapi) | `0.136.3` | `0.138.2` |
| [uvicorn](https://github.com/Kludex/uvicorn) | `0.48.0` | `0.49.0` |
| [markdownify](https://github.com/matthewwithanm/python-markdownify) | `1.2.2` | `1.2.3` |
| [langsmith](https://github.com/langchain-ai/langsmith-sdk) | `0.9.3` | `0.9.4` |
| [langchain-openai](https://github.com/langchain-ai/langchain) | `1.2.2` | `1.3.3` |
| [langchain-fireworks](https://github.com/langchain-ai/langchain) | `1.4.2` | `1.4.3` |
| [langchain-daytona](https://github.com/langchain-ai/deepagents) | `0.0.6` | `0.0.7` |
| [langchain-modal](https://github.com/langchain-ai/deepagents) | `0.0.4` | `0.0.5` |
| [langchain-runloop](https://github.com/langchain-ai/deepagents) | `0.0.5` | `0.0.6` |
| [exa-py](https://github.com/exa-labs/exa-py) | `2.13.0` | `2.15.0` |
| [langchain-mcp-adapters](https://github.com/langchain-ai/langchain-mcp-adapters) | `0.2.2` | `0.3.0` |
| [pytest](https://github.com/pytest-dev/pytest) | `9.0.3` | `9.1.1` |
| [ruff](https://github.com/astral-sh/ruff) | `0.15.15` | `0.15.20` |



Updates `deepagents` from 0.6.8 to 0.6.12
- [Release notes](https://github.com/langchain-ai/deepagents/releases)
- [Commits](https://github.com/langchain-ai/deepagents/compare/deepagents==0.6.8...deepagents==0.6.12)

Updates `fastapi` from 0.136.3 to 0.138.2
- [Release notes](https://github.com/fastapi/fastapi/releases)
- [Commits](https://github.com/fastapi/fastapi/compare/0.136.3...0.138.2)

Updates `uvicorn` from 0.48.0 to 0.49.0
- [Release notes](https://github.com/Kludex/uvicorn/releases)
- [Changelog](https://github.com/Kludex/uvicorn/blob/main/docs/release-notes.md)
- [Commits](https://github.com/Kludex/uvicorn/compare/0.48.0...0.49.0)

Updates `langchain` from 1.3.9 to 1.3.11
- [Release notes](https://github.com/langchain-ai/langchain/releases)
- [Commits](https://github.com/langchain-ai/langchain/compare/langchain==1.3.9...langchain==1.3.11)

Updates `langgraph` from 1.2.4 to 1.2.7
- [Release notes](https://github.com/langchain-ai/langgraph/releases)
- [Commits](https://github.com/langchain-ai/langgraph/compare/1.2.4...1.2.7)

Updates `markdownify` from 1.2.2 to 1.2.3
- [Release notes](https://github.com/matthewwithanm/python-markdownify/releases)
- [Commits](https://github.com/matthewwithanm/python-markdownify/compare/1.2.2...1.2.3)

Updates `langchain-anthropic` from 1.4.6 to 1.4.8
- [Release notes](https://github.com/langchain-ai/langchain/releases)
- [Commits](https://github.com/langchain-ai/langchain/compare/langchain-anthropic==1.4.6...langchain-anthropic==1.4.8)

Updates `langsmith` from 0.9.3 to 0.9.4
- [Release notes](https://github.com/langchain-ai/langsmith-sdk/releases)
- [Commits](https://github.com/langchain-ai/langsmith-sdk/compare/v0.9.3...v0.9.4)

Updates `langchain-openai` from 1.2.2 to 1.3.3
- [Release notes](https://github.com/langchain-ai/langchain/releases)
- [Commits](https://github.com/langchain-ai/langchain/compare/langchain-openai==1.2.2...langchain-openai==1.3.3)

Updates `langchain-fireworks` from 1.4.2 to 1.4.3
- [Release notes](https://github.com/langchain-ai/langchain/releases)
- [Commits](https://github.com/langchain-ai/langchain/compare/langchain-fireworks==1.4.2...langchain-fireworks==1.4.3)

Updates `langchain-daytona` from 0.0.6 to 0.0.7
- [Release notes](https://github.com/langchain-ai/deepagents/releases)
- [Commits](https://github.com/langchain-ai/deepagents/compare/langchain-daytona==0.0.6...langchain-daytona==0.0.7)

Updates `langchain-modal` from 0.0.4 to 0.0.5
- [Release notes](https://github.com/langchain-ai/deepagents/releases)
- [Commits](https://github.com/langchain-ai/deepagents/compare/langchain-modal==0.0.4...langchain-modal==0.0.5)

Updates `langchain-runloop` from 0.0.5 to 0.0.6
- [Release notes](https://github.com/langchain-ai/deepagents/releases)
- [Commits](https://github.com/langchain-ai/deepagents/compare/langchain-runloop==0.0.5...langchain-runloop==0.0.6)

Updates `exa-py` from 2.13.0 to 2.15.0
- [Release notes](https://github.com/exa-labs/exa-py/releases)
- [Changelog](https://github.com/exa-labs/exa-py/blob/master/CHANGELOG.md)
- [Commits](https://github.com/exa-labs/exa-py/commits)

Updates `langchain-google-genai` from 4.2.4 to 4.2.6
- [Release notes](https://github.com/langchain-ai/langchain-google/releases)
- [Commits](https://github.com/langchain-ai/langchain-google/compare/libs/genai/v4.2.4...libs/genai/v4.2.6)

Updates `langchain-mcp-adapters` from 0.2.2 to 0.3.0
- [Release notes](https://github.com/langchain-ai/langchain-mcp-adapters/releases)
- [Commits](https://github.com/langchain-ai/langchain-mcp-adapters/compare/langchain-mcp-adapters==0.2.2...langchain-mcp-adapters==0.3.0)

Updates `pytest` from 9.0.3 to 9.1.1
- [Release notes](https://github.com/pytest-dev/pytest/releases)
- [Changelog](https://github.com/pytest-dev/pytest/blob/main/CHANGELOG.rst)
- [Commits](https://github.com/pytest-dev/pytest/compare/9.0.3...9.1.1)

Updates `ruff` from 0.15.15 to 0.15.20
- [Release notes](https://github.com/astral-sh/ruff/releases)
- [Changelog](https://github.com/astral-sh/ruff/blob/main/CHANGELOG.md)
- [Commits](https://github.com/astral-sh/ruff/compare/0.15.15...0.15.20)

---
updated-dependencies:
- dependency-name: deepagents
  dependency-version: 0.6.12
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: fastapi
  dependency-version: 0.138.2
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: uvicorn
  dependency-version: 0.49.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: langchain
  dependency-version: 1.3.11
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langgraph
  dependency-version: 1.2.7
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: markdownify
  dependency-version: 1.2.3
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-anthropic
  dependency-version: 1.4.8
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langsmith
  dependency-version: 0.9.4
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-openai
  dependency-version: 1.3.3
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: langchain-fireworks
  dependency-version: 1.4.3
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-daytona
  dependency-version: 0.0.7
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-modal
  dependency-version: 0.0.5
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-runloop
  dependency-version: 0.0.6
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: exa-py
  dependency-version: 2.15.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: langchain-google-genai
  dependency-version: 4.2.6
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
- dependency-name: langchain-mcp-adapters
  dependency-version: 0.3.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: pytest
  dependency-version: 9.1.1
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: minor-and-patch
- dependency-name: ruff
  dependency-version: 0.15.20
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: minor-and-patch
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-06-30 20:53:13 +00:00
35 changed files with 2137 additions and 3102 deletions

View file

@ -40,7 +40,7 @@ The FastAPI app is `agent.webapp:app`.
- **`agent/server.py` → `get_agent(config)`** — main graph factory. Called per-thread. Resolves the GitHub token, gets-or-creates the sandbox for the thread, resolves the team/profile/per-thread model + effort, then constructs a fresh `create_deep_agent(...)` with the curated tool list and middleware stack. The agent itself is stateless — all per-thread state lives in the sandbox + thread metadata.
- **`agent/reviewer.py` → `get_reviewer_agent(config)`** — reviewer graph factory. Shares `ensure_sandbox_for_thread` with the main agent but wires a reviewer-only toolset (`add_finding`, `update_finding`, `list_findings`, `publish_review`, `web_search`, `fetch_url`, `http_request`) and a different system prompt that pins the single-evolving-findings model and the diff-anchored bar for filing a finding. Read-only: no commit/push/PR-opening tools.
- **`agent/analyzer.py` → `get_analyzer(config)`** — small graph that emits a per-repo style prompt via the `save_review_style_prompt` tool, consumed by the reviewer as a "repository-specific review style" appendix. It runs in one of two modes (`analyzer_mode` in `configurable`): **bootstrap** (cold-start: crawl historical PR reviews) and **continual** (nightly: refine using this reviewer's own finding outcomes via `read_finding_outcomes`). Each mode's procedure lives in a deepagents **skill** (`agent/skills/bootstrap-repo-analysis/`, `agent/skills/continual-learning/`) served as virtual files via a `CompositeBackend` `/skills/` route + `StateBackend` (seeded into the run's `files` channel by the launcher — never written to the sandbox). Launchers and the per-repo nightly cron live in `agent/dashboard/review_style_jobs.py` and `agent/dashboard/analyzer_cron.py`; the cron is registered when bootstrap completes.
- **`agent/webapp.py`** — thin FastAPI routing layer mounted alongside the LangGraph server. Defines the webhook routes (GitHub, Linear, Slack) plus `/webhooks/run-complete`, and keeps the shared helpers/constants; the per-source handlers live in **`agent/webhooks/{github,slack,linear}.py`** (re-exported from `webapp` so existing call sites and tests keep working). Each webhook resolves a deterministic `thread_id` (so follow-up messages route to the same agent run) and triggers a run through the single durable dispatch contract in **`agent/dispatch.py`** (`dispatch_agent_run`: `multitask_strategy="interrupt"` + `durability="sync"` + completion webhook); `agent/completion.py` posts a failure reply if a run dies, and `agent/reconcile.py` (a `scheduler`-graph sweep) catches stragglers. The GitHub handler also auto-reviews PRs on `opened` / `ready_for_review` and drives the CI auto-fix flow (`agent/ci_autofix.py`).
- **`agent/webapp.py`** — custom FastAPI routes mounted alongside the LangGraph server. Webhooks land here (GitHub, Linear, Slack). Each webhook resolves a deterministic `thread_id` (so follow-up messages route to the same agent run) and triggers/streams a run via the `langgraph_sdk` client. Also auto-reviews PRs on `opened` / `ready_for_review` events when the repo+author opt in.
- **`agent/dashboard/`** — `router` mounted under the FastAPI app at startup (`app.include_router(dashboard_router)`). Owns GitHub OAuth, per-user profiles, admin endpoints, team defaults, enabled-repo lists, review-style management, and the Agents chat thread API used by the UI in `ui/`.
### Sandbox lifecycle (the tricky part)

View file

@ -25,7 +25,6 @@ from langgraph_sdk import get_client
from .dashboard.agent_overrides import load_profile, resolve_login_from_email_async
from .dashboard.autofix_state import is_pr_autofix_disabled
from .dashboard.enabled_repos import is_review_repo_enabled
from .dispatch import dispatch_agent_run
from .reviewer_findings import REVIEWER_THREAD_KIND
from .utils.dashboard_links import dashboard_thread_url
from .utils.github_app import get_github_app_installation_token
@ -41,7 +40,7 @@ from .utils.github_ci import (
)
from .utils.github_org_membership import INTERNAL_BOT_LOGINS
from .utils.thread_ops import (
get_thread_active_status,
is_thread_active,
langgraph_client,
)
@ -272,25 +271,17 @@ async def _mark_pending_autofix_event(thread_id: str, reason: str, detail: str =
async def _dispatch_or_batch(
thread_id: str, prompt: str, *, configurable: dict[str, Any], reason: str, detail: str = ""
) -> str:
# Deliberate skip-rule: batch auto-fix events while the agent thread is
# actively running so we don't interrupt an in-progress fix. ``interrupt``
# is fine for human follow-ups but undesirable for autofix, so we keep the
# busy-check here even though the webhook hot-path no longer needs one.
if await get_thread_active_status(thread_id) is True:
if await is_thread_active(thread_id):
logger.info("Agent thread %s busy; batching auto-fix event %s", thread_id, reason)
await _mark_pending_autofix_event(thread_id, reason, detail)
return "batched"
# The busy-check above has a TOCTOU window (the dedupe SHA is only recorded
# after dispatch), so a burst of near-simultaneous CI events for one head SHA
# can all pass the gate. Dispatch with ``reject`` — matching ``dev``'s prior
# platform default — so the platform drops the duplicate concurrent creates
# instead of letting them interrupt each other.
await dispatch_agent_run(
client = langgraph_client()
await client.runs.create(
thread_id,
prompt,
configurable,
source=str(configurable.get("source") or "github_autofix"),
multitask_strategy="reject",
"agent",
input={"messages": [{"role": "user", "content": prompt}]},
config={"configurable": configurable},
if_not_exists="create",
)
logger.info(
"Created auto-fix run for thread %s (source=%s)", thread_id, configurable.get("source")

View file

@ -1,173 +0,0 @@
"""Run-completion webhook handler — guarantees every run ends with a signal.
The platform POSTs a run-completion payload to ``/webhooks/run-complete`` (wired
as the ``webhook`` on every dispatched run, see ``agent.dispatch``). When a run
ends in a failure state (``error`` / ``timeout`` / ``interrupted``) we post a
short failure reply to the originating channel, so a run that died on a server
recycle or hit a limit never leaves the user in silence.
This decouples "the user gets an answer" from "the agent remembered to reply."
The reply is idempotent: a per-thread metadata flag prevents double-posting when
the platform retries the webhook or a checkpoint replays.
"""
from __future__ import annotations
import hmac
import logging
import os
from collections.abc import Awaitable, Callable
from typing import Any
from .utils.github_app import get_github_app_installation_token
from .utils.github_comments import post_github_comment
from .utils.linear import comment_on_linear_issue
from .utils.slack import post_slack_thread_reply
from .utils.thread_ops import langgraph_client
logger = logging.getLogger(__name__)
# Run statuses that mean the user will otherwise get nothing back. "interrupted"
# is intentionally excluded: with multitask_strategy="interrupt", a normal
# follow-up halts the prior run (status "interrupted") while its replacement
# carries on — that's healthy, not a failure worth a "couldn't finish" reply.
_TERMINAL_FAILURE_STATUSES = frozenset({"error", "timeout"})
_FAILURE_REPLY_FLAG = "failure_reply_posted"
class _ClaimFailed(Exception):
"""Raised when the dedup flag couldn't be claimed, so we skip the post."""
# Shared-secret bearer token proving a /webhooks/run-complete call came from our
# own dispatch (which appends ?token= when this is set) rather than from an
# attacker hitting the public route. Fail closed when unset: the route rejects
# every call, so completion replies stay off until the secret is configured.
RUN_COMPLETE_WEBHOOK_SECRET = os.environ.get("RUN_COMPLETE_WEBHOOK_SECRET")
if not RUN_COMPLETE_WEBHOOK_SECRET:
logger.warning(
"RUN_COMPLETE_WEBHOOK_SECRET is not set; /webhooks/run-complete is fail-closed "
"(all calls rejected) and run-failure replies are disabled. Set it to enable them."
)
def verify_run_complete_token(token: str | None) -> bool:
"""Return whether a run-completion webhook token is acceptable.
Fail closed: with no secret configured, reject every call rather than accept
unauthenticated requests on a publicly reachable route.
"""
secret = RUN_COMPLETE_WEBHOOK_SECRET
if not secret:
return False
return token is not None and hmac.compare_digest(token, secret)
def _failure_text(status: str) -> str:
reason = "timed out" if status == "timeout" else "hit an unexpected error"
return (
f"⚠️ I wasn't able to finish that — the run {reason}. "
"Send another message and I'll pick it back up."
)
async def _post_failure_reply(
thread_id: str,
metadata: dict[str, Any],
status: str,
*,
claim: Callable[[], Awaitable[None]],
) -> bool:
"""Post a failure reply to the run's originating channel. Best-effort.
``claim`` is awaited immediately before the network post (claim-then-post),
only on a branch that actually delivers, so a retried/concurrent webhook
can't double-post and threads with no channel never burn the flag.
"""
source = metadata.get("source")
ctx = metadata.get("source_context")
ctx = ctx if isinstance(ctx, dict) else {}
text = _failure_text(status)
if source == "slack":
slack_thread = ctx.get("slack_thread")
if isinstance(slack_thread, dict):
channel_id = slack_thread.get("channel_id")
thread_ts = slack_thread.get("thread_ts")
if channel_id and thread_ts:
await claim()
return await post_slack_thread_reply(channel_id, thread_ts, text)
return False
if source == "linear":
linear_issue = ctx.get("linear_issue")
if isinstance(linear_issue, dict):
issue_id = linear_issue.get("id")
if issue_id:
await claim()
return await comment_on_linear_issue(issue_id, text)
return False
if source in ("github", "github_issue"):
repo_config = metadata.get("repo")
number = ctx.get("pr_number")
if number is None:
github_issue = ctx.get("github_issue")
if isinstance(github_issue, dict):
number = github_issue.get("number")
if isinstance(repo_config, dict) and isinstance(number, int):
token = await get_github_app_installation_token()
if token:
await claim()
return await post_github_comment(repo_config, number, text, token=token)
return False
logger.info("No failure-reply channel for thread %s (source=%s)", thread_id, source)
return False
async def handle_run_completion(payload: dict[str, Any]) -> dict[str, str]:
"""Handle a platform run-completion webhook POST.
Posts a failure reply only when the run ended in a failure state and we
haven't already replied for this thread.
"""
status = payload.get("status")
thread_id = payload.get("thread_id")
if not isinstance(thread_id, str) or not thread_id:
return {"status": "ignored", "reason": "missing thread_id"}
if status not in _TERMINAL_FAILURE_STATUSES:
return {"status": "ignored", "reason": f"non-failure status: {status}"}
client = langgraph_client()
try:
thread = await client.threads.get(thread_id)
except Exception: # noqa: BLE001
logger.warning("run-complete: could not load thread %s", thread_id, exc_info=True)
return {"status": "error", "reason": "thread fetch failed"}
metadata = thread.get("metadata") if isinstance(thread, dict) else None
metadata = metadata if isinstance(metadata, dict) else {}
if metadata.get(_FAILURE_REPLY_FLAG):
return {"status": "ignored", "reason": "failure reply already posted"}
# Claim-then-post: set the dedup flag immediately before the actual post (via
# the claim callback) so a retried/concurrent completion webhook can't
# double-post. The flag is only claimed on a branch that delivers, so a
# thread with no reply channel never burns it. If the claim itself fails we
# skip the post, leaving the flag unset so a later retry can try again.
async def _claim() -> None:
try:
await client.threads.update(thread_id=thread_id, metadata={_FAILURE_REPLY_FLAG: True})
except Exception as exc: # noqa: BLE001
logger.warning("run-complete: could not flag thread %s", thread_id, exc_info=True)
raise _ClaimFailed from exc
try:
posted = await _post_failure_reply(thread_id, metadata, status, claim=_claim)
except _ClaimFailed:
return {"status": "error", "reason": "could not claim failure reply"}
if not posted:
return {"status": "ignored", "reason": "no reply posted"}
logger.info("Posted failure reply for thread %s (status=%s)", thread_id, status)
return {"status": "ok", "reason": "failure reply posted"}

View file

@ -21,7 +21,6 @@ from fastapi import APIRouter, Depends, HTTPException
from langgraph_sdk import get_client
from pydantic import BaseModel
from ..dispatch import dispatch_agent_run
from .oauth import require_same_origin_for_mutations, require_session
from .plan_store import (
PLAN_STATUS_APPROVED,
@ -213,9 +212,11 @@ async def _dispatch_followup(
# mode (implement), reject stays in plan mode (revise the plan).
configurable["plan_mode"] = plan_mode
await dispatch_agent_run(
client = get_client()
await client.runs.create(
thread_id,
text,
configurable,
source=configurable["source"],
"agent",
input={"messages": [{"role": "user", "content": text}]},
config={"configurable": configurable},
if_not_exists="create",
)

View file

@ -3,7 +3,6 @@
from __future__ import annotations
import logging
import re
import uuid
from datetime import UTC, datetime
from typing import Any
@ -23,7 +22,6 @@ SCHEDULES_NAMESPACE: list[str] = ["agent_schedules"]
_AGENT_ASSISTANT_ID = "agent"
_SCHEDULER_ASSISTANT_ID = "scheduler"
_CRON_FIELD_RANGES = ((0, 59), (0, 23), (1, 31), (1, 12), (0, 7))
_SLACK_CHANNEL_ID_RE = re.compile(r"^[A-Z][A-Z0-9]{5,}$")
class ScheduleCreateBody(BaseModel):
@ -33,18 +31,12 @@ class ScheduleCreateBody(BaseModel):
repo: str | None = None
model_id: str | None = None
effort: str | None = None
slack_report_channel: str | None = Field(default=None, max_length=120)
@field_validator("schedule")
@classmethod
def _valid_schedule(cls, value: str) -> str:
return normalize_cron_schedule(value)
@field_validator("slack_report_channel")
@classmethod
def _valid_slack_report_channel(cls, value: str | None) -> str | None:
return normalize_slack_channel_id(value)
class ScheduleUpdateBody(BaseModel):
prompt: str | None = Field(default=None, min_length=1, max_length=20_000)
@ -54,18 +46,12 @@ class ScheduleUpdateBody(BaseModel):
model_id: str | None = None
effort: str | None = None
enabled: bool | None = None
slack_report_channel: str | None = Field(default=None, max_length=120)
@field_validator("schedule")
@classmethod
def _valid_schedule(cls, value: str | None) -> str | None:
return normalize_cron_schedule(value) if value is not None else None
@field_validator("slack_report_channel")
@classmethod
def _valid_slack_report_channel(cls, value: str | None) -> str | None:
return normalize_slack_channel_id(value)
def _client():
return langgraph_client()
@ -123,18 +109,6 @@ def normalize_cron_schedule(raw: str) -> str:
return value
def normalize_slack_channel_id(raw: str | None) -> str | None:
"""Normalize a Slack channel ID; blank becomes None (no report channel)."""
if raw is None:
return None
value = raw.strip().lstrip("#")
if not value:
return None
if not _SLACK_CHANNEL_ID_RE.match(value):
raise ValueError("slack_report_channel must be a Slack channel ID (e.g. C0123ABCD)")
return value
def _derive_name(prompt: str) -> str:
return prompt.strip().splitlines()[0][:80] or "Scheduled agent"
@ -157,7 +131,6 @@ def _schedule_summary(record: dict[str, Any]) -> dict[str, Any]:
"repo": _repo_full_name(repo),
"model": record.get("model"),
"effort": record.get("effort"),
"slackReportChannel": record.get("slack_report_channel"),
"enabled": bool(record.get("enabled")),
"cronId": record.get("cron_id"),
"lastThreadId": record.get("last_thread_id"),
@ -301,7 +274,6 @@ async def create_agent_schedule(
"repo": repo,
"model": chosen_model or profile.get("default_model") or "Default",
"effort": chosen_effort or profile.get("reasoning_effort"),
"slack_report_channel": body.slack_report_channel,
"base_branch": profile.get("base_branch") or "main",
"branch_prefix": profile.get("branch_prefix"),
"enabled": True,
@ -350,8 +322,6 @@ async def update_agent_schedule(
patch["effort"] = effort
if body.enabled is not None:
patch["enabled"] = body.enabled
if "slack_report_channel" in body.model_fields_set:
patch["slack_report_channel"] = body.slack_report_channel
updated = {**existing, **patch}
schedule_changed = updated.get("schedule") != existing.get("schedule")
@ -420,9 +390,6 @@ def _agent_run_config(record: dict[str, Any], thread_id: str) -> dict[str, Any]:
if model and effort:
configurable["agent_model_id"] = model
configurable["agent_effort"] = effort
report_channel = record.get("slack_report_channel")
if isinstance(report_channel, str) and report_channel.strip():
configurable["slack_thread"] = {"channel_id": report_channel.strip()}
return {"configurable": configurable, "metadata": _agent_version_metadata()}

View file

@ -1,93 +0,0 @@
"""Single durable dispatch contract behind every agent/reviewer run trigger.
Replaces the per-site ``runs.create`` calls (plus the ``is_thread_active``
busy-check and the custom store-queue) with one function that always uses:
- ``multitask_strategy="interrupt"`` — a follow-up halts the active run
(progress preserved by the sync checkpoint) and resumes the agent with full
history + the new message; on an idle thread it just starts. This is the
platform-native, cross-process replacement for the racy busy-check + queue.
- ``durability="sync"`` — checkpoint before each step so a crash/recycle
resumes from the last checkpoint instead of losing all work.
- ``webhook=COMPLETION_WEBHOOK_URL`` — the platform calls us on completion or
failure so every run ends with a signal even if the agent died.
"""
from __future__ import annotations
import logging
import os
from typing import Any
from langgraph_sdk import get_client
from langgraph_sdk.client import LangGraphClient
logger = logging.getLogger(__name__)
ContentBlocks = str | list[dict[str, Any]]
# Same-server FastAPI route the platform POSTs run completion/failure to. A
# relative URL loopback-posts into this app (no SSRF/loopback config needed);
# override with an absolute URL via env for split deployments. The route is
# fail-closed on RUN_COMPLETE_WEBHOOK_SECRET, so only register the webhook when
# the secret is set, appending it as ?token= so the route can verify the call
# came from us (completion.verify_run_complete_token). Unset → no webhook.
_COMPLETION_WEBHOOK_BASE = os.environ.get("COMPLETION_WEBHOOK_URL") or "/webhooks/run-complete"
_RUN_COMPLETE_SECRET = os.environ.get("RUN_COMPLETE_WEBHOOK_SECRET")
COMPLETION_WEBHOOK_URL: str | None
if not _RUN_COMPLETE_SECRET:
COMPLETION_WEBHOOK_URL = None
elif "?" in _COMPLETION_WEBHOOK_BASE:
COMPLETION_WEBHOOK_URL = _COMPLETION_WEBHOOK_BASE
else:
COMPLETION_WEBHOOK_URL = f"{_COMPLETION_WEBHOOK_BASE}?token={_RUN_COMPLETE_SECRET}"
def _langgraph_url() -> str:
return os.environ.get("LANGGRAPH_URL") or os.environ.get(
"LANGGRAPH_URL_PROD", "http://localhost:2024"
)
def dispatch_client() -> LangGraphClient:
return get_client(url=_langgraph_url())
async def dispatch_agent_run(
thread_id: str,
content: ContentBlocks,
configurable: dict[str, Any],
*,
source: str,
assistant_id: str = "agent",
metadata: dict[str, Any] | None = None,
client: LangGraphClient | None = None,
multitask_strategy: str = "interrupt",
) -> dict[str, Any]:
"""Create (or interrupt-and-resume) a run for ``thread_id``.
Routes every Slack / Linear / GitHub / dashboard trigger through one
contract. ``source`` is for logging/metadata only; ``assistant_id`` selects
the graph (``"agent"`` or ``"reviewer"``). ``multitask_strategy`` defaults to
``"interrupt"`` (human follow-ups halt + resume); autofix passes ``"reject"``
so a burst of concurrent CI events for one head SHA can't interrupt each other.
"""
client = client or dispatch_client()
run = await client.runs.create(
thread_id,
assistant_id,
input={"messages": [{"role": "user", "content": content}]},
config={"configurable": configurable, "metadata": metadata or {}},
multitask_strategy=multitask_strategy,
durability="sync",
webhook=COMPLETION_WEBHOOK_URL,
if_not_exists="create",
)
logger.info(
"Dispatched %s run on thread %s (source=%s, run=%s)",
assistant_id,
thread_id,
source,
run.get("run_id") if isinstance(run, dict) else None,
)
return run

View file

@ -368,23 +368,6 @@ ALWAYS_CREATE_PR_SECTION = """---
The user's dashboard setting **Always Create PRs** is enabled. For code-change tasks, always open or update a draft pull request after committing and pushing the branch. This does not apply to questions, explanations, status checks, or other information-only requests where no files are changed."""
def _render_scheduled_report_section(channel_id: str | None) -> str:
if not channel_id or not channel_id.strip():
return ""
return (
"---\n\n"
"### Scheduled Run Report\n\n"
"This is a scheduled (automated) run with a configured Slack report channel. "
"When you finish, post your final summary to that channel by calling "
"`slack_thread_reply` with your report — it posts a top-level message to the "
f"configured channel (`{channel_id.strip()}`) as the bot. Post exactly one "
"final report. If the run produced a pull request, include its link. If "
"`slack_thread_reply` reports a failure (for example `not_in_channel`, meaning "
"the bot is not a member of the channel), do not retry repeatedly — surface the "
"error in your final output instead."
)
def _render_repo_instructions_section(instructions: str | None) -> str:
if not instructions or not instructions.strip():
return ""
@ -415,7 +398,6 @@ SYSTEM_PROMPT_TEMPLATE = (
+ EXTERNAL_UNTRUSTED_COMMENTS_SECTION
+ COMMIT_PR_SECTION
+ "{pr_policy_override_section}"
+ "{scheduled_report_section}"
+ "{collaboration_section}"
+ "{repo_instructions_section}"
)
@ -433,7 +415,6 @@ def construct_system_prompt(
repo_custom_instructions: str | None = None,
thread_url: str | None = None,
corridor_enabled: bool = False,
slack_report_channel: str | None = None,
) -> str:
default_prompt_section = _load_default_prompt()
if default_repo and default_repo.get("owner") and default_repo.get("name"):
@ -463,7 +444,6 @@ def construct_system_prompt(
default_prompt_section=default_prompt_section,
corridor_prompt_section=CORRIDOR_PROMPT if corridor_enabled else "",
pr_policy_override_section=ALWAYS_CREATE_PR_SECTION if create_prs else "",
scheduled_report_section=_render_scheduled_report_section(slack_report_channel),
collaboration_section=_render_collaboration_section(triggering_user_identity, thread_url),
repo_instructions_section=_render_repo_instructions_section(repo_custom_instructions),
commit_identity_name=commit_identity_name,

View file

@ -1,121 +0,0 @@
"""Reconciliation sweep: cancel runs stuck in ``pending`` past their deadline.
The durable-dispatch contract relies on the platform's completion webhook to
end every run. When that webhook never fires (crash, lost delivery), a run can
sit in ``pending`` forever and hold its thread ``busy``. This sweep is the
safety net: find busy threads, look for stale ``pending`` runs on them, and
cancel the ones older than ``max_age_seconds`` so the thread frees up.
"""
from __future__ import annotations
import logging
from datetime import UTC, datetime
from typing import Any
from .utils.thread_ops import langgraph_client
logger = logging.getLogger(__name__)
_SEARCH_PAGE_SIZE = 100
def _parse_created_at(value: Any) -> datetime | None:
"""Parse a run's ``created_at`` into an aware UTC datetime, or None."""
if isinstance(value, datetime):
return value if value.tzinfo else value.replace(tzinfo=UTC)
if not isinstance(value, str) or not value:
return None
text = value.strip()
if text.endswith("Z"):
text = f"{text[:-1]}+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
return parsed if parsed.tzinfo else parsed.replace(tzinfo=UTC)
async def reconcile_stale_runs(*, max_age_seconds: int = 1800) -> dict[str, int]:
"""Cancel ``pending`` runs older than ``max_age_seconds`` on busy threads.
Walks every ``busy`` thread (paginated), lists its ``pending`` runs, and
cancels those whose ``created_at`` is older than the cutoff. Per-thread work
is wrapped in try/except so one bad thread never aborts the sweep.
Returns counts: ``{"threads_checked", "stale_runs", "cancelled"}``.
"""
client = langgraph_client()
now = datetime.now(UTC)
threads_checked = 0
stale_runs = 0
cancelled = 0
offset = 0
while True:
try:
threads = await client.threads.search(
metadata=None,
status="busy",
limit=_SEARCH_PAGE_SIZE,
offset=offset,
)
except Exception:
logger.exception("Reconcile sweep: thread search failed at offset %d", offset)
break
if not threads:
break
for thread in threads:
thread_id = thread.get("thread_id") if isinstance(thread, dict) else None
if not thread_id:
continue
threads_checked += 1
try:
runs = await client.runs.list(thread_id, status="pending")
stale_run_ids: list[str] = []
for run in runs:
created = _parse_created_at(run.get("created_at"))
if created is None:
logger.warning(
"Reconcile sweep: unparseable created_at on run %s (thread %s)",
run.get("run_id"),
thread_id,
)
continue
if (now - created).total_seconds() <= max_age_seconds:
continue
run_id = run.get("run_id")
if run_id:
stale_run_ids.append(run_id)
if not stale_run_ids:
continue
stale_runs += len(stale_run_ids)
await client.runs.cancel_many(
thread_id=thread_id,
run_ids=stale_run_ids,
action="interrupt",
)
cancelled += len(stale_run_ids)
logger.info(
"Reconcile sweep: cancelled %d stale pending run(s) on thread %s",
len(stale_run_ids),
thread_id,
)
except Exception:
logger.exception("Reconcile sweep: failed to reconcile thread %s", thread_id)
continue
if len(threads) < _SEARCH_PAGE_SIZE:
break
offset += _SEARCH_PAGE_SIZE
counts = {
"threads_checked": threads_checked,
"stale_runs": stale_runs,
"cancelled": cancelled,
}
logger.info("Reconcile sweep complete: %s", counts)
return counts

View file

@ -9,22 +9,17 @@ from langgraph.graph import END, START, StateGraph
from langgraph.graph.state import RunnableConfig
from .dashboard.schedules import launch_scheduled_agent_run
from .reconcile import reconcile_stale_runs
logger = logging.getLogger(__name__)
class SchedulerState(TypedDict, total=False):
schedule_id: str
task: str
result: dict[str, Any]
async def _launch(state: SchedulerState, config: RunnableConfig) -> dict[str, Any]:
configurable = config.get("configurable") or {}
task = state.get("task") or configurable.get("task")
if task == "reconcile":
return {"result": await reconcile_stale_runs()}
schedule_id = state.get("schedule_id") or configurable.get("schedule_id")
if not isinstance(schedule_id, str) or not schedule_id:
logger.warning("Scheduled agent tick missing schedule_id")

View file

@ -647,17 +647,6 @@ def _get_cached_sandbox_backend(thread_id: str) -> SandboxBackendProtocol:
return sandbox_backend
def _scheduled_report_channel(configurable: dict[str, Any]) -> str | None:
"""Slack channel a scheduled run should post its final report to, if configured."""
if (configurable or {}).get("source") != "schedule":
return None
slack_thread = (configurable or {}).get("slack_thread") or {}
channel_id = slack_thread.get("channel_id")
if isinstance(channel_id, str) and channel_id.strip():
return channel_id.strip()
return None
async def _observability_authorized(config: RunnableConfig, profile_login: str | None) -> bool:
"""Whether the triggering user may use the team observability tools.
@ -939,7 +928,6 @@ async def get_agent(config: RunnableConfig) -> Pregel:
repo_custom_instructions=repo_custom_instructions,
thread_url=dashboard_thread_url(thread_id),
corridor_enabled=bool(corridor_tools),
slack_report_channel=_scheduled_report_channel(configurable),
),
tools=[
http_request,

View file

@ -2,6 +2,7 @@
from __future__ import annotations
import asyncio
import logging
from typing import Annotated
@ -22,7 +23,7 @@ _ENTERED_MESSAGE = (
)
async def enter_plan_mode(tool_call_id: Annotated[str, InjectedToolCallId]) -> Command:
def enter_plan_mode(tool_call_id: Annotated[str, InjectedToolCallId]) -> Command:
"""Activate plan mode mid-run.
Call this when you believe the task would benefit from a structured
@ -40,7 +41,7 @@ async def enter_plan_mode(tool_call_id: Annotated[str, InjectedToolCallId]) -> C
thread_id = _thread_id_from_config()
if thread_id:
try:
await set_plan_status(thread_id, PLAN_STATUS_PLANNING, plan_mode=True)
asyncio.run(set_plan_status(thread_id, PLAN_STATUS_PLANNING, plan_mode=True))
except Exception:
logger.warning("Failed to persist plan-mode entry for %s", thread_id, exc_info=True)
return Command(

View file

@ -1,3 +1,4 @@
import asyncio
from typing import Any
from langgraph.config import get_config
@ -6,7 +7,7 @@ from agent.utils.slack import parse_github_pr_url
from agent.webapp import trigger_pr_review_from_ref
async def request_pr_review(pr_url: str) -> dict[str, Any]:
def request_pr_review(pr_url: str) -> dict[str, Any]:
"""Start the reviewer agent for a GitHub pull request URL."""
pr_ref = parse_github_pr_url(pr_url)
if not pr_ref:
@ -18,11 +19,13 @@ async def request_pr_review(pr_url: str) -> dict[str, Any]:
configurable = get_config().get("configurable", {})
source = configurable.get("source") or "agent"
slack_thread = configurable.get("slack_thread") or {}
return await trigger_pr_review_from_ref(
pr_ref,
source=source,
github_login=configurable.get("github_login", ""),
github_user_id=configurable.get("github_user_id"),
slack_channel_id=slack_thread.get("channel_id", ""),
slack_thread_ts=slack_thread.get("thread_ts", ""),
return asyncio.run(
trigger_pr_review_from_ref(
pr_ref,
source=source,
github_login=configurable.get("github_login", ""),
github_user_id=configurable.get("github_user_id"),
slack_channel_id=slack_thread.get("channel_id", ""),
slack_thread_ts=slack_thread.get("thread_ts", ""),
)
)

View file

@ -8,6 +8,7 @@ changes. Available in plan mode (it does not modify the repository under review)
from __future__ import annotations
import asyncio
import logging
from typing import Any
@ -20,7 +21,7 @@ logger = logging.getLogger(__name__)
PLAN_FILE_PATH = "plan.md"
async def save_plan(plan_markdown: str) -> dict[str, Any]:
def save_plan(plan_markdown: str) -> dict[str, Any]:
"""Write your implementation plan as a markdown file and publish it for review.
Use this in plan mode once your plan is ready. The plan is saved as
@ -53,7 +54,7 @@ async def save_plan(plan_markdown: str) -> dict[str, Any]:
return {"success": False, "error": "no thread_id in run config"}
try:
path = await _save(str(thread_id), content)
path = asyncio.run(_save(str(thread_id), content))
except Exception as exc: # noqa: BLE001
logger.exception("save_plan failed for thread %s", thread_id)
return {"success": False, "error": f"failed to save plan: {exc}"}

View file

@ -1,6 +1,5 @@
import json
import os
from collections import OrderedDict
from typing import Any
from langgraph.config import get_config
@ -9,7 +8,6 @@ from langgraph_sdk import get_client
from ..utils.slack import (
convert_mentions_to_slack_format,
post_slack_thread_reply_with_ts,
post_slack_top_level_message_with_ts,
store_slack_message_run_mapping,
)
@ -17,12 +15,6 @@ LANGGRAPH_URL = os.environ.get("LANGGRAPH_URL") or os.environ.get(
"LANGGRAPH_URL_PROD", "http://localhost:2024"
)
# Runs that have already posted their single top-level (channel) message. Scheduled
# runs seed `slack_thread` with a channel but no `thread_ts`, so every reply would
# otherwise spray a new top-level message into the report channel; cap it at one.
_MAX_TRACKED_RUNS = 2048
_top_level_posts: "OrderedDict[str, None]" = OrderedDict()
async def slack_thread_reply(
message: str,
@ -58,33 +50,17 @@ async def slack_thread_reply(
channel_id = slack_thread.get("channel_id")
thread_ts = slack_thread.get("thread_ts")
if not channel_id:
if not channel_id or not thread_ts:
return {
"success": False,
"error": "Missing slack_thread.channel_id in config",
"error": "Missing slack_thread.channel_id or slack_thread.thread_ts in config",
}
if not message.strip():
return {"success": False, "error": "Message cannot be empty"}
top_level = not thread_ts
run_key = _run_key(config) if top_level else None
if top_level and run_key is not None and run_key in _top_level_posts:
return {
"success": False,
"error": "A message was already posted to this channel for this run",
"hint": (
"Only one top-level message per run is allowed for the configured "
"report channel; post a single final report and do not call this again."
),
}
message = convert_mentions_to_slack_format(message)
if top_level:
# Interactive blocks (options / plan_approval) render dead buttons in a
# report channel where no run is driving the approval/option flow.
slack_blocks = blocks
elif plan_approval:
if plan_approval:
slack_blocks = _build_plan_approval_blocks(message)
else:
slack_blocks = blocks or _build_option_blocks(message, options)
@ -99,25 +75,9 @@ async def slack_thread_reply(
"message_chars": len(message),
"hint": _slack_reply_failure_hint(slack_error),
}
if top_level and run_key is not None:
_top_level_posts[run_key] = None
if len(_top_level_posts) > _MAX_TRACKED_RUNS:
_top_level_posts.popitem(last=False)
return {"success": True}
def _run_key(config: dict[str, Any]) -> str | None:
candidates = [config.get("run_id")]
configurable = config.get("configurable")
if isinstance(configurable, dict):
candidates.append(configurable.get("run_id"))
candidates.append(configurable.get("thread_id"))
for candidate in candidates:
if isinstance(candidate, str) and candidate:
return candidate
return None
def _build_option_blocks(message: str, options: list[str] | None) -> list[dict[str, Any]] | None:
if not options:
return None
@ -228,15 +188,11 @@ def _slack_reply_failure_hint(slack_error: str | None) -> str:
async def _post_and_store_mapping(
channel_id: str,
thread_ts: str | None,
thread_ts: str,
message: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
if not thread_ts:
# Top-level report posts are fire-and-forget: a scheduled run is one-shot, so
# there is no live run to route channel replies back to (no mapping stored).
return await post_slack_top_level_message_with_ts(channel_id, message, blocks=blocks)
message_ts, slack_error = await post_slack_thread_reply_with_ts(
channel_id, thread_ts, message, blocks=blocks
)

View file

@ -1,16 +1,12 @@
"""Shared LangGraph thread helpers for the dashboard.
The webhook triggers (Slack / Linear / GitHub) dispatch through
``agent.dispatch.dispatch_agent_run`` with ``multitask_strategy="interrupt"``,
so they no longer need a busy-check or an in-process lock. The store-queue
below is retained for the dashboard's deliberate "inject a follow-up into a
run that's already in flight" path (``thread_api.send_dashboard_message``).
"""
"""Shared LangGraph thread helpers for webhooks and the dashboard."""
from __future__ import annotations
import asyncio
import logging
import os
from collections.abc import AsyncIterator
from contextlib import asynccontextmanager
from typing import Any
from langgraph_sdk import get_client
@ -19,6 +15,25 @@ logger = logging.getLogger(__name__)
MAX_QUEUED_MESSAGES = 100
_THREAD_RUN_LOCKS: dict[str, asyncio.Lock] = {}
def get_thread_run_lock(thread_id: str) -> asyncio.Lock:
"""Return a per-thread-id asyncio.Lock, creating one lazily if needed."""
lock = _THREAD_RUN_LOCKS.get(thread_id)
if lock is None:
lock = asyncio.Lock()
_THREAD_RUN_LOCKS[thread_id] = lock
return lock
@asynccontextmanager
async def thread_run_lock(thread_id: str) -> AsyncIterator[None]:
"""Serialize run dispatch for a thread."""
lock = get_thread_run_lock(thread_id)
async with lock:
yield
def langgraph_url() -> str:
return os.environ.get("LANGGRAPH_URL") or os.environ.get(
@ -42,14 +57,15 @@ async def get_thread_active_status(thread_id: str) -> bool | None:
return None
async def is_thread_active(thread_id: str) -> bool:
"""Return whether the thread currently has a running run."""
return await get_thread_active_status(thread_id) is True
async def queue_message_for_thread(
thread_id: str, message_content: str | list[dict[str, Any]] | dict[str, Any]
) -> bool:
"""Queue a follow-up message for a busy thread (FIFO store namespace).
Used by the dashboard to inject a follow-up into a run that's already in
flight; webhook triggers use ``multitask_strategy="interrupt"`` instead.
"""
"""Queue a follow-up message for a busy thread (FIFO store namespace)."""
client = langgraph_client()
try:
namespace = ("queue", thread_id)

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -1,235 +0,0 @@
"""Linear webhook handler — moved out of webapp.py (behavior-identical).
Helpers and constants stay in webapp.py; they are accessed through the module
object (``webapp.X``) so tests that monkeypatch them keep working.
"""
from typing import Any
import httpx
from langchain_core.messages.content import create_text_block
from agent import webapp
async def process_linear_issue( # noqa: PLR0912, PLR0915
issue_data: dict[str, Any], repo_config: dict[str, str]
) -> None:
"""Process a Linear issue by creating a new LangGraph thread and run.
Args:
issue_data: The Linear issue data from webhook (basic info only).
repo_config: The repo configuration with owner and name.
"""
issue_id = issue_data.get("id", "")
webapp.logger.info(
"Processing Linear issue %s for repo %s/%s",
issue_id,
repo_config.get("owner"),
repo_config.get("name"),
)
triggering_comment_id = issue_data.get("triggering_comment_id", "")
if triggering_comment_id:
await webapp.react_to_linear_comment(triggering_comment_id, "👀")
thread_id = webapp.generate_thread_id_from_issue(issue_id)
full_issue = await webapp.fetch_linear_issue_details(issue_id)
if not full_issue:
full_issue = issue_data
user_email = None
user_name = None
comment_author = issue_data.get("comment_author", {})
if comment_author:
user_email = comment_author.get("email")
user_name = comment_author.get("name")
if not user_email:
creator = full_issue.get("creator", {})
if creator:
user_email = creator.get("email")
user_name = user_name or creator.get("name")
if not user_email:
assignee = full_issue.get("assignee", {})
if assignee:
user_email = assignee.get("email")
user_name = user_name or assignee.get("name")
webapp.logger.info("User email for issue %s: %s", issue_id, user_email)
title = full_issue.get("title", "No title")
description = full_issue.get("description") or "No description"
image_urls: list[str] = []
description_image_urls = webapp.extract_image_urls(description)
if description_image_urls:
image_urls.extend(description_image_urls)
webapp.logger.debug(
"Found %d image URL(s) in issue description",
len(description_image_urls),
)
comments = full_issue.get("comments", {}).get("nodes", [])
comments_text = ""
triggering_comment = issue_data.get("triggering_comment", "")
triggering_comment_id = issue_data.get("triggering_comment_id", "")
bot_message_prefixes = (
"🔐 **GitHub Authentication Required**",
"✅ **Pull Request Created**",
"✅ **Pull Request Updated**",
"**Pull Request Created**",
"**Pull Request Updated**",
"🤖 **Agent Response**",
"❌ **Agent Error**",
)
comment_ids: set[str] = set()
comment_id_to_index: dict[str, int] = {}
if comments:
for i, comment in enumerate(comments):
comment_id = comment.get("id", "")
if comment_id:
comment_ids.add(comment_id)
comment_id_to_index[comment_id] = i
relevant_comments = []
trigger_index = None
if triggering_comment_id:
trigger_index = comment_id_to_index.get(triggering_comment_id)
if trigger_index is not None:
relevant_comments = comments[trigger_index:]
webapp.logger.debug(
"Using triggering comment index %d to build relevant comments",
trigger_index,
)
else:
relevant_comments = webapp.get_recent_comments(comments, bot_message_prefixes)
if relevant_comments:
comments_text = "\n\n## Comments:\n"
for comment in relevant_comments:
user = comment.get("user") or {}
author = user.get("name", "User")
body = comment.get("body", "")
body_image_urls = webapp.extract_image_urls(body)
if body_image_urls:
image_urls.extend(body_image_urls)
webapp.logger.debug(
"Found %d image URL(s) in comment by %s",
len(body_image_urls),
author,
)
if any(body.startswith(prefix) for prefix in bot_message_prefixes):
continue
comments_text += f"\n**{author}:** {body}\n"
if triggering_comment and triggering_comment_id not in comment_ids:
if not comments_text:
comments_text = "\n\n## Comments:\n"
trigger_author = comment_author.get("name", "Unknown")
trigger_body = triggering_comment
trigger_image_urls = webapp.extract_image_urls(trigger_body)
if trigger_image_urls:
image_urls.extend(trigger_image_urls)
webapp.logger.debug(
"Found %d image URL(s) in triggering comment by %s",
len(trigger_image_urls),
trigger_author,
)
comments_text += f"\n**{trigger_author}:** {trigger_body}\n"
webapp.logger.debug(
"Appended triggering comment %s not present in issue comments list",
triggering_comment_id or "<missing-id>",
)
identifier = full_issue.get("identifier", "") or issue_data.get("identifier", "")
triggered_by_line = f"## Triggered by: {user_name}\n\n" if user_name else ""
tag_instruction = (
f"When calling linear_comment, tag @{user_name} if you are asking them a question, need their input, or are notifying them of something important (e.g. a completed PR). For simple answers, tagging is not required."
if user_name
else ""
)
prompt = (
f"Please work on the following issue:\n\n"
f"## Repository: {repo_config.get('owner')}/{repo_config.get('name')}\n\n"
f"## Title: {title}\n\n"
f"{triggered_by_line}"
f"## Linear Ticket: {identifier} - Ticket ID: {issue_id}\n\n"
f"## Description:\n{description}\n"
f"{comments_text}\n\n"
f"Please analyze this issue and implement the necessary changes. "
f"When you're done, commit and push your changes. {tag_instruction}"
)
content_blocks: list[dict[str, Any]] = [create_text_block(prompt)]
if image_urls:
image_urls = webapp.dedupe_urls(image_urls)
linear_login = (
await webapp.resolve_login_from_email_async(user_email) if user_email else None
)
resolved_model_id = await webapp.resolve_agent_model_id(linear_login)
if webapp.model_supports_images(resolved_model_id):
webapp.logger.info("Preparing %d image(s) for multimodal content", len(image_urls))
webapp.logger.debug("Image URLs: %s", image_urls)
async with httpx.AsyncClient(timeout=webapp.DEFAULT_HTTP_TIMEOUT) as client:
for image_url in image_urls:
image_block = await webapp.fetch_image_block(image_url, client)
if image_block:
content_blocks.append(image_block)
webapp.logger.info("Built %d content block(s) for prompt", len(content_blocks))
else:
webapp.logger.warning(
"Skipping %d image(s) for Linear issue: model %s does not support images",
len(image_urls),
resolved_model_id,
)
prompt += webapp.vision_not_supported_warning(resolved_model_id, len(image_urls))
content_blocks[0] = create_text_block(prompt)
image_urls = []
linear_project_id = ""
linear_issue_number = ""
if identifier and "-" in identifier:
parts = identifier.split("-", 1)
linear_project_id = parts[0]
linear_issue_number = parts[1]
configurable: dict[str, Any] = {
"repo": repo_config,
"linear_issue": {
"id": issue_id,
"title": title,
"url": full_issue.get("url", "") or issue_data.get("url", ""),
"identifier": identifier,
"linear_project_id": linear_project_id,
"linear_issue_number": linear_issue_number,
"triggering_user_name": user_name or "",
},
"user_email": user_email,
"source": "linear",
}
await webapp.upsert_agent_thread_owner_metadata(
thread_id,
source="linear",
repo_config=repo_config,
user_email=user_email or "",
title=title or identifier or "Linear issue",
source_context={"linear_issue": configurable["linear_issue"]},
)
run = await webapp.dispatch_agent_run(
thread_id,
content_blocks,
configurable,
source="linear",
metadata=webapp._AGENT_VERSION_METADATA,
)
webapp.logger.info(
"LangGraph run dispatched for thread %s (run=%s)",
thread_id,
run.get("run_id") if isinstance(run, dict) else None,
)
await webapp.post_linear_trace_comment(issue_id, thread_id, triggering_comment_id)

View file

@ -1,269 +0,0 @@
"""Slack webhook handler — moved out of webapp.py (behavior-identical).
Helpers and constants stay in webapp.py; they are accessed through the module
object (``webapp.X``) so tests that monkeypatch them keep working.
"""
from typing import Any
import httpx
from langchain_core.messages.content import create_text_block
from agent import webapp
async def process_slack_mention(event_data: dict[str, Any], repo_config: dict[str, str]) -> None:
"""Process a Slack app mention by creating a run or queuing a mid-run message."""
channel_id = event_data.get("channel_id", "")
thread_ts = event_data.get("thread_ts", "")
event_ts = event_data.get("event_ts", "")
user_id = event_data.get("user_id", "")
text = event_data.get("text", "")
bot_user_id = event_data.get("bot_user_id", "")
if not channel_id or not thread_ts or not event_ts:
webapp.logger.warning(
"Missing Slack event fields (channel_id=%s, thread_ts=%s, event_ts=%s)",
channel_id,
thread_ts,
event_ts,
)
return
await webapp.set_slack_assistant_status(channel_id, thread_ts)
thread_id = webapp.generate_thread_id_from_slack_thread(channel_id, thread_ts)
# Prime the user-mapping cache so login/email/slack-id lookups below are warm.
try:
await webapp.refresh_user_mapping_cache()
except Exception: # noqa: BLE001
webapp.logger.debug("Could not refresh user mapping cache for Slack mention", exc_info=True)
user_email = None
user_name = ""
if user_id:
slack_user = await webapp.get_slack_user_info(user_id)
if slack_user:
profile = slack_user.get("profile", {})
if isinstance(profile, dict):
user_email = profile.get("email")
user_name = (
profile.get("display_name")
or profile.get("real_name")
or slack_user.get("real_name")
or slack_user.get("name")
or ""
)
thread_messages = await webapp.fetch_slack_thread_messages(channel_id, thread_ts)
if not any(str(message.get("ts")) == str(event_ts) for message in thread_messages):
thread_messages.append({"ts": event_ts, "text": text, "user": user_id})
context_messages, context_mode = webapp.select_slack_context_messages(
thread_messages, event_ts, bot_user_id, webapp.SLACK_BOT_USERNAME
)
context_user_ids = [
value
for value in (message.get("user") for message in context_messages)
if isinstance(value, str) and value
]
user_names_by_id = await webapp.get_slack_user_names(context_user_ids)
if user_id and user_name and user_id not in user_names_by_id:
user_names_by_id[user_id] = user_name
context_text = webapp.format_slack_messages_for_prompt(
context_messages,
user_names_by_id,
bot_user_id=bot_user_id,
bot_username=webapp.SLACK_BOT_USERNAME,
)
context_source = (
"the previous message where I was tagged"
if context_mode == "last_mention"
else "the beginning of the thread"
)
clean_text = (
webapp.strip_bot_mention(text, bot_user_id, bot_username=webapp.SLACK_BOT_USERNAME)
or "(no text in mention)"
)
trigger_user = user_name or (f"<@{user_id}>" if user_id else "Unknown user")
# Auto-resolve cross-posted Slack message links in context
resolved_links_section, image_urls_from_links = await webapp.resolve_slack_links_in_context(
context_messages, user_names_by_id
)
prompt = (
"You were mentioned in Slack.\n\n"
"## Default Repository Hint\n"
f"{repo_config.get('owner')}/{repo_config.get('name')}\n"
"Use this only if the Slack conversation does not identify a different repository.\n\n"
f"## Triggered by\n{trigger_user}\n\n"
f"## Slack Thread\n- Channel: {channel_id}\n- Thread TS: {thread_ts}\n"
f"- Context starts at: {context_source}\n\n"
f"## Conversation Context\n{context_text}\n\n"
f"## Latest Mention Request\n{clean_text}\n\n"
+ (f"{resolved_links_section}\n\n" if resolved_links_section else "")
+ "Use `slack_thread_reply` to communicate in this Slack thread for clarifications, "
"status updates, and final summaries. Use `slack_read_thread_messages` to read any "
"Slack messages by providing channel_id and message_ts."
)
content_blocks: list[dict[str, Any]] = [create_text_block(prompt)]
image_urls = webapp.dedupe_urls(
[url for msg in context_messages for url in webapp.extract_image_urls(msg.get("text", ""))]
+ [
f["url_private"]
for msg in context_messages
for f in msg.get("files", [])
if isinstance(f, dict)
and f.get("mimetype", "").startswith("image/")
and f.get("url_private")
]
+ image_urls_from_links
)
mapped_login = await webapp.login_for_slack_id(user_id)
if not mapped_login and user_email:
mapped_login = await webapp.login_for_email(user_email)
if image_urls:
resolved_model_id = await webapp.resolve_agent_model_id(mapped_login)
if webapp.model_supports_images(resolved_model_id):
webapp.logger.info("Preparing %d image(s) for Slack mention", len(image_urls))
async with httpx.AsyncClient(timeout=webapp.DEFAULT_HTTP_TIMEOUT) as http_client:
for image_url in image_urls:
image_block = await webapp.fetch_image_block(image_url, http_client)
if image_block:
content_blocks.append(image_block)
else:
webapp.logger.warning(
"Skipping %d image(s) for Slack mention: model %s does not support images",
len(image_urls),
resolved_model_id,
)
prompt += webapp.vision_not_supported_warning(resolved_model_id, len(image_urls))
content_blocks[0] = create_text_block(prompt)
image_urls = []
# Open SWE opens PRs as the triggering user, so a run only proceeds when we
# have a valid user GitHub token. Users who have never signed in with
# GitHub, and users whose stored authorization is no longer usable, are
# blocked and prompted to set up via the dashboard. Bot-token-only
# deployments are exempt — they run on the installation token.
user_token: str | None = None
if mapped_login:
try:
user_token = await webapp.get_valid_access_token(mapped_login)
except Exception: # noqa: BLE001
webapp.logger.debug(
"Failed to resolve GitHub token for %s; treating as unauthenticated",
mapped_login,
exc_info=True,
)
user_token = None
has_valid_user_token = bool(user_token)
if not has_valid_user_token and not webapp.is_bot_token_only_mode():
# A stored-but-unusable token means "sign in again"; no record at all
# means the user has never connected GitHub + Slack via the dashboard.
# Guard the store read like token resolution above so a transient
# failure still yields an actionable prompt and clears the status.
has_token_record = False
if mapped_login:
try:
has_token_record = await webapp.has_access_token_record(mapped_login)
except Exception: # noqa: BLE001
webapp.logger.debug(
"Failed to check GitHub token record for %s; prompting sign-in",
mapped_login,
exc_info=True,
)
reason = "revoked" if has_token_record else "unlinked"
webapp.logger.info(
"Blocking Slack run for thread %s: no valid user GitHub token (%s)",
thread_id,
reason,
)
if user_id:
await webapp._post_account_link_prompt(
channel_id, thread_ts, user_id, user_email, reason=reason
)
await webapp.set_slack_assistant_status(channel_id, thread_ts, status="")
return
configurable: dict[str, Any] = {
"repo": repo_config,
"slack_thread": {
"channel_id": channel_id,
"thread_ts": thread_ts,
"triggering_user_id": user_id,
"triggering_user_name": user_name,
"triggering_user_email": user_email,
"triggering_event_ts": event_ts,
},
"user_email": user_email,
"source": "slack",
}
if mapped_login:
configurable["github_login"] = mapped_login
thread_plan_mode = await webapp._get_thread_plan_mode(thread_id)
if thread_plan_mode is not None:
configurable["plan_mode"] = thread_plan_mode
langgraph_client = webapp.get_client(url=webapp.LANGGRAPH_URL)
is_first_mention = not await webapp._thread_exists(thread_id)
await webapp._upsert_slack_thread_repo_metadata(thread_id, repo_config, langgraph_client)
# Pass the login resolved above (from the stable Slack user id) so the thread is
# always tagged with github_login — the key the dashboard searches by. Without
# it, upsert re-resolves from the Slack profile email, which can miss.
await webapp.upsert_agent_thread_owner_metadata(
thread_id,
source="slack",
repo_config=repo_config,
github_login=mapped_login or "",
user_email=user_email or "",
title=clean_text if is_first_mention else "",
source_context={"slack_thread": configurable["slack_thread"]},
)
run = await webapp.dispatch_agent_run(
thread_id,
content_blocks,
configurable,
source="slack",
metadata=webapp._AGENT_VERSION_METADATA,
client=langgraph_client,
)
webapp.logger.info(
"Slack LangGraph run %s dispatched for thread %s",
webapp._run_id_for_logging(run),
thread_id,
)
run_id = run.get("run_id")
if is_first_mention:
trace_message_ts = await webapp.post_slack_trace_reply(channel_id, thread_ts, thread_id)
await webapp.set_slack_assistant_status(channel_id, thread_ts)
if isinstance(run_id, str) and run_id:
await webapp.store_slack_run_mapping(
langgraph_client,
channel_id,
thread_ts,
run_id,
message_ts=trace_message_ts,
triggering_user_id=user_id,
)
else:
webapp.logger.info(
"Skipping Slack trace reply for thread %s — agent will reply when run completes",
thread_id,
)
if isinstance(run_id, str) and run_id:
await webapp.store_slack_run_mapping(
langgraph_client,
channel_id,
thread_ts,
run_id,
triggering_user_id=user_id,
)

View file

@ -303,20 +303,20 @@ Mappings can't be fully created from the dashboard → must write the Store dire
- **Fix (Phase B):** add the `work_email` field to the Admin mappings form so mappings are fully creatable from the UI.
### Fix #5 — GitHub webhook path doesn't refresh the user-mapping cache (multi-replica break)
**Files:** `agent/webhooks/github.py` (GitHub handlers: `process_github_pr_comment`, `process_github_issue`) vs `agent/webhooks/slack.py` (`process_slack_mention`). (Pre-modular-refactor these all lived in `agent/webapp.py`.)
**Files:** `agent/webapp.py:3052` (GitHub path) vs `agent/webapp.py:1091` (Slack path).
The Slack path refreshes before lookup:
```python
# agent/webhooks/slack.py (process_slack_mention)
await webapp.refresh_user_mapping_cache()
# agent/webapp.py:1089-1093 (Slack)
await refresh_user_mapping_cache()
...
```
The GitHub path historically did **not** — it called `email = await email_for_login(github_login)` cold. The cache (`user_mappings.py` `_ensure_cache_loaded`) is **one-shot per process** (`_cache_loaded` flag). On self-host single-process this was fine; on managed's **multi-replica autoscaling**, a freshly-added mapping isn't seen by a replica whose cache loaded earlier — until restart.
- **Fix (applied):** the GitHub issue and PR-comment handlers now call `webapp.refresh_user_mapping_cache()` before email resolution, mirroring the Slack path.
- **Generalize (Phase B5):** audit ALL in-process caches for the single-process → multi-replica assumption — `SANDBOX_BACKENDS` dict (`agent/utils/sandbox_state.py`), `_by_login`/`_by_email`/`_by_slack_id` (`user_mappings.py`). Sandbox affinity is already thread-keyed + persisted in thread metadata (`sandbox_id`), so it's the cache state that needs the multi-replica review. (The legacy in-process thread lock has been removed: webhook triggers now serialize through `dispatch_agent_run`'s `multitask_strategy="interrupt"` instead.)
The GitHub path does **not** — it calls `email = await email_for_login(github_login)` (`webapp.py:3052`, again at `:3331`) cold. The cache (`user_mappings.py` `_ensure_cache_loaded`, line 197) is **one-shot per process** (`_cache_loaded` flag). On self-host single-process this was fine; on managed's **multi-replica autoscaling**, a freshly-added mapping isn't seen by a replica whose cache loaded earlier — until restart.
- **Fix:** refresh-before-lookup on the GitHub path (mirror the Slack path), or add a TTL / cross-replica invalidation to the cache.
- **Generalize (Phase B5):** audit ALL in-process caches for the single-process → multi-replica assumption — `SANDBOX_BACKENDS` dict (`agent/utils/sandbox_state.py`), `_THREAD_RUN_LOCKS` (`thread_ops.py:18`), `_by_login`/`_by_email`/`_by_slack_id` (`user_mappings.py:67-69`). Sandbox affinity is already thread-keyed + persisted in thread metadata (`sandbox_id`), so it's the cache/lock state that needs the multi-replica review.
### Fix #6 — Slow custom-app import (~8s startup)
**Symptom:** "exceeded expected startup time" → risks the deployment being marked unhealthy / slow to scale out.
- **Fix (Phase B6):** lazy imports / reduce import-time work in `agent/webapp.py` (+ `agent/webhooks/*.py`) and the graph factories. Profile with `FF_PROFILE_IMPORTS` (the import-profiling flag) to find the heavy modules.
- **Fix (Phase B6):** lazy imports / reduce import-time work in `agent/webapp.py` and the graph factories. Profile with `FF_PROFILE_IMPORTS` (the import-profiling flag) to find the heavy modules.
---

View file

@ -7,8 +7,7 @@
"reviewer": "agent.reviewer:traced_reviewer_agent",
"analyzer": "agent.analyzer:traced_analyzer",
"chat": "agent.chat:traced_chat_agent",
"scheduler": "agent.scheduler:get_scheduler",
"ci_monitor": "agent.ci_monitor:get_ci_monitor"
"scheduler": "agent.scheduler:get_scheduler"
},
"dependencies": [
"."

View file

@ -136,96 +136,6 @@ def test_cron_validation_accepts_steps_ranges_and_lists() -> None:
assert body.schedule == "*/15 9-17 * * 1,3,5"
def test_slack_report_channel_normalizes_and_validates() -> None:
body = ScheduleCreateBody(
prompt="hello", schedule="0 9 * * 1", slack_report_channel=" #C0123ABCD "
)
assert body.slack_report_channel == "C0123ABCD"
blank = ScheduleCreateBody(prompt="hello", schedule="0 9 * * 1", slack_report_channel=" ")
assert blank.slack_report_channel is None
with pytest.raises(ValidationError):
ScheduleCreateBody(prompt="hello", schedule="0 9 * * 1", slack_report_channel="not a chan")
with pytest.raises(ValidationError):
ScheduleCreateBody(prompt="hello", schedule="0 9 * * 1", slack_report_channel="123456")
async def test_create_agent_schedule_persists_slack_report_channel(fake_client, auth) -> None: # noqa: ANN001, ARG001
body = ScheduleCreateBody(
name="Daily report",
prompt="Summarize merged PRs",
schedule="0 9 * * 1-5",
slack_report_channel="C0123ABCD",
)
result = await schedules.create_agent_schedule("alice", body, email="alice@example.com")
assert result["slackReportChannel"] == "C0123ABCD"
stored = fake_client.store.items[(tuple(schedules.SCHEDULES_NAMESPACE), result["id"])]
assert stored["slack_report_channel"] == "C0123ABCD"
async def test_update_agent_schedule_clears_slack_report_channel(fake_client) -> None: # noqa: ANN001
record = {
"id": "sched_1",
"name": "Daily",
"prompt": "Run daily",
"schedule": "0 9 * * *",
"repo": None,
"model": "Default",
"effort": None,
"slack_report_channel": "C0123ABCD",
"enabled": True,
"cron_id": "cron_old",
"created_by": "alice",
"user_email": "alice@example.com",
"created_at": "2026-01-01T00:00:00+00:00",
"updated_at": "2026-01-01T00:00:00+00:00",
}
await fake_client.store.put_item(schedules.SCHEDULES_NAMESPACE, "sched_1", record)
result = await schedules.update_agent_schedule(
"sched_1",
"alice",
ScheduleUpdateBody(slack_report_channel=""),
email="alice@example.com",
)
assert result["slackReportChannel"] is None
def test_agent_run_config_seeds_slack_thread_channel() -> None:
record = {
"id": "sched_1",
"model": "Default",
"effort": None,
"created_by": "alice",
"user_email": "alice@example.com",
"slack_report_channel": "C0123ABCD",
}
config = schedules._agent_run_config(record, "thread_1")
assert config["configurable"]["slack_thread"] == {"channel_id": "C0123ABCD"}
def test_agent_run_config_omits_slack_thread_without_channel() -> None:
record = {
"id": "sched_1",
"model": "Default",
"effort": None,
"created_by": "alice",
"user_email": "alice@example.com",
"slack_report_channel": None,
}
config = schedules._agent_run_config(record, "thread_1")
assert "slack_thread" not in config["configurable"]
async def test_create_agent_schedule_registers_scheduler_cron(fake_client, auth) -> None: # noqa: ANN001, ARG001
body = ScheduleCreateBody(
name="Daily report",

View file

@ -7,7 +7,6 @@ from unittest.mock import AsyncMock, patch
import pytest
from agent import webapp
from agent.webhooks import github as webhooks_github
def test_parse_autofix_command() -> None:
@ -120,7 +119,7 @@ async def test_process_github_ci_event_dispatches() -> None:
},
}
handle = AsyncMock(return_value="dispatched")
with patch.object(webhooks_github, "handle_ci_failure", handle):
with patch.object(webapp, "handle_ci_failure", handle):
await webapp.process_github_ci_event(payload, "check_run")
handle.assert_awaited_once()
kwargs = handle.await_args.kwargs
@ -136,7 +135,7 @@ async def test_process_github_ci_event_ignores_success() -> None:
"check_run": {"status": "completed", "conclusion": "success", "head_sha": "s"},
}
handle = AsyncMock()
with patch.object(webhooks_github, "handle_ci_failure", handle):
with patch.object(webapp, "handle_ci_failure", handle):
await webapp.process_github_ci_event(payload, "check_run")
handle.assert_not_called()
@ -150,7 +149,7 @@ async def test_process_autofix_command_sets_flag() -> None:
}
setter = AsyncMock()
with (
patch.object(webhooks_github, "set_pr_autofix_disabled", setter),
patch.object(webapp, "set_pr_autofix_disabled", setter),
patch.object(webapp, "get_github_app_installation_token", AsyncMock(return_value="")),
):
await webapp.process_github_autofix_command(payload, "issue_comment", disabled=True)
@ -165,7 +164,7 @@ async def test_autofix_review_dispatches_for_writer() -> None:
"review": {"body": "rename to userId", "user": {"login": "alice"}},
}
handle = AsyncMock(return_value="dispatched")
with patch.object(webhooks_github, "handle_review_feedback", handle):
with patch.object(webapp, "handle_review_feedback", handle):
await webapp.process_github_autofix_review(payload, "pull_request_review")
handle.assert_awaited_once()
@ -178,7 +177,7 @@ async def test_autofix_review_delegates_permission_check_to_core() -> None:
"review": {"body": "inject code", "user": {"login": "attacker"}},
}
handle = AsyncMock(return_value="reviewer_no_write_permission")
with patch.object(webhooks_github, "handle_review_feedback", handle):
with patch.object(webapp, "handle_review_feedback", handle):
await webapp.process_github_autofix_review(payload, "pull_request_review")
handle.assert_awaited_once()

View file

@ -29,12 +29,9 @@ def happy(monkeypatch: pytest.MonkeyPatch) -> dict[str, Any]:
threads_update = AsyncMock()
store_client = MagicMock()
store_client.threads.update = threads_update
# Auto-fix runs now dispatch through the durable dispatch_agent_run contract
# rather than a raw runs.create; assert against that.
dispatch_run = AsyncMock(return_value={"run_id": "r1"})
mocks: dict[str, Any] = {
"runs_create": dispatch_run,
"runs_create": runs_create,
"threads_update": threads_update,
"status_check": AsyncMock(return_value=True),
"store_put": store_put,
@ -61,10 +58,9 @@ def happy(monkeypatch: pytest.MonkeyPatch) -> dict[str, Any]:
monkeypatch.setattr(
ci_autofix, "head_commit_author_login", AsyncMock(return_value="open-swe[bot]")
)
monkeypatch.setattr(ci_autofix, "get_thread_active_status", AsyncMock(return_value=False))
monkeypatch.setattr(ci_autofix, "is_thread_active", AsyncMock(return_value=False))
monkeypatch.setattr(ci_autofix, "post_autofix_status_check", mocks["status_check"])
monkeypatch.setattr(ci_autofix, "langgraph_client", lambda: lg_client)
monkeypatch.setattr(ci_autofix, "dispatch_agent_run", mocks["runs_create"])
monkeypatch.setattr(ci_autofix, "get_client", lambda: store_client)
return mocks
@ -89,18 +85,9 @@ async def test_dispatch_happy_path(happy: dict[str, Any]) -> None:
happy["status_check"].assert_awaited()
@pytest.mark.asyncio
async def test_autofix_dispatch_uses_reject_strategy(happy: dict[str, Any]) -> None:
# A burst of concurrent CI events for one head SHA can slip past the busy-check
# before the dedupe SHA is recorded; dispatching with "reject" lets the platform
# drop the duplicate concurrent creates instead of interrupting each other.
await _run()
assert happy["runs_create"].await_args.kwargs["multitask_strategy"] == "reject"
@pytest.mark.asyncio
async def test_batches_when_thread_busy(happy: dict[str, Any], monkeypatch) -> None:
monkeypatch.setattr(ci_autofix, "get_thread_active_status", AsyncMock(return_value=True))
monkeypatch.setattr(ci_autofix, "is_thread_active", AsyncMock(return_value=True))
result = await _run()
assert result == "batched"
happy["store_put"].assert_awaited()
@ -207,7 +194,7 @@ async def test_review_feedback_skips_user_disabled(happy: dict[str, Any], monkey
@pytest.mark.asyncio
async def test_review_feedback_batches_when_thread_busy(happy: dict[str, Any], monkeypatch) -> None:
monkeypatch.setattr(ci_autofix, "has_repo_write_permission", AsyncMock(return_value=True))
monkeypatch.setattr(ci_autofix, "get_thread_active_status", AsyncMock(return_value=True))
monkeypatch.setattr(ci_autofix, "is_thread_active", AsyncMock(return_value=True))
result = await ci_autofix.handle_review_feedback(
repo_config={"owner": "o", "name": "r"},
pr_number=5,

View file

@ -1,168 +0,0 @@
from __future__ import annotations
from typing import Any
from unittest.mock import AsyncMock
import pytest
from agent import completion
class _FakeThreads:
def __init__(self, metadata: dict[str, Any]) -> None:
self._metadata = metadata
self.updates: list[dict[str, Any]] = []
async def get(self, thread_id: str) -> dict[str, Any]:
return {"thread_id": thread_id, "metadata": self._metadata}
async def update(self, *, thread_id: str, metadata: dict[str, Any]) -> None:
self.updates.append(metadata)
class _FakeClient:
def __init__(self, metadata: dict[str, Any]) -> None:
self.threads = _FakeThreads(metadata)
def _slack_metadata() -> dict[str, Any]:
return {
"source": "slack",
"source_context": {"slack_thread": {"channel_id": "C1", "thread_ts": "123.45"}},
}
@pytest.mark.asyncio
async def test_error_status_posts_slack_failure_reply(monkeypatch: pytest.MonkeyPatch) -> None:
client = _FakeClient(_slack_metadata())
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "error"})
assert result["status"] == "ok"
reply.assert_awaited_once()
args = reply.await_args.args
assert args[0] == "C1"
assert args[1] == "123.45"
assert client.threads.updates == [{"failure_reply_posted": True}]
@pytest.mark.asyncio
async def test_success_status_is_ignored(monkeypatch: pytest.MonkeyPatch) -> None:
client = _FakeClient(_slack_metadata())
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "success"})
assert result["status"] == "ignored"
reply.assert_not_called()
@pytest.mark.asyncio
async def test_idempotent_when_already_replied(monkeypatch: pytest.MonkeyPatch) -> None:
metadata = _slack_metadata()
metadata["failure_reply_posted"] = True
client = _FakeClient(metadata)
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "timeout"})
assert result["status"] == "ignored"
reply.assert_not_called()
assert client.threads.updates == []
@pytest.mark.asyncio
async def test_linear_source_comments_on_issue(monkeypatch: pytest.MonkeyPatch) -> None:
client = _FakeClient({"source": "linear", "source_context": {"linear_issue": {"id": "iss_1"}}})
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
comment = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "comment_on_linear_issue", comment)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "timeout"})
assert result["status"] == "ok"
comment.assert_awaited_once()
assert comment.await_args.args[0] == "iss_1"
@pytest.mark.asyncio
async def test_missing_thread_id_is_ignored() -> None:
result = await completion.handle_run_completion({"status": "error"})
assert result["status"] == "ignored"
@pytest.mark.asyncio
async def test_claims_flag_before_posting(monkeypatch: pytest.MonkeyPatch) -> None:
# Claim-then-post: the dedup flag must be set before the reply is posted so a
# retried/concurrent webhook can't double-post the canned failure message.
client = _FakeClient(_slack_metadata())
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
async def _reply(*_args: Any, **_kwargs: Any) -> bool:
assert client.threads.updates == [{"failure_reply_posted": True}]
return True
monkeypatch.setattr(completion, "post_slack_thread_reply", AsyncMock(side_effect=_reply))
result = await completion.handle_run_completion({"thread_id": "t1", "status": "error"})
assert result["status"] == "ok"
@pytest.mark.asyncio
async def test_does_not_post_when_claim_fails(monkeypatch: pytest.MonkeyPatch) -> None:
client = _FakeClient(_slack_metadata())
client.threads.update = AsyncMock(side_effect=RuntimeError("boom"))
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "error"})
assert result["status"] == "error"
reply.assert_not_called()
@pytest.mark.asyncio
async def test_no_reply_channel_does_not_flag(monkeypatch: pytest.MonkeyPatch) -> None:
client = _FakeClient({"source": "schedule"})
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "error"})
assert result["status"] == "ignored"
assert client.threads.updates == []
@pytest.mark.asyncio
async def test_interrupted_status_is_ignored(monkeypatch: pytest.MonkeyPatch) -> None:
# Follow-ups use multitask_strategy="interrupt", so an interrupted run is a
# healthy hand-off, not a failure to report.
client = _FakeClient(_slack_metadata())
monkeypatch.setattr(completion, "langgraph_client", lambda: client)
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
result = await completion.handle_run_completion({"thread_id": "t1", "status": "interrupted"})
assert result["status"] == "ignored"
reply.assert_not_called()
assert client.threads.updates == []
def test_verify_run_complete_token(monkeypatch: pytest.MonkeyPatch) -> None:
# No secret configured: fail closed (reject everything).
monkeypatch.setattr(completion, "RUN_COMPLETE_WEBHOOK_SECRET", None)
assert completion.verify_run_complete_token(None) is False
assert completion.verify_run_complete_token("whatever") is False
# Secret configured: require an exact match.
monkeypatch.setattr(completion, "RUN_COMPLETE_WEBHOOK_SECRET", "s3cret")
assert completion.verify_run_complete_token("s3cret") is True
assert completion.verify_run_complete_token("wrong") is False
assert completion.verify_run_complete_token(None) is False

View file

@ -323,6 +323,9 @@ def test_process_github_review_finding_reply_uses_rereview_config(monkeypatch) -
captured["interaction"] = (finding_id, interaction)
return {}
async def fake_is_thread_active(_thread_id: str) -> bool:
return False
async def fake_store_current_run_id(_thread_id: str, _run: object) -> None:
return None
@ -345,6 +348,7 @@ def test_process_github_review_finding_reply_uses_rereview_config(monkeypatch) -
monkeypatch.setattr(webapp, "reconcile_findings_with_review_threads", fake_reconcile)
monkeypatch.setattr(webapp, "list_reviewer_findings", fake_list_findings)
monkeypatch.setattr(webapp, "append_finding_interaction", fake_append_interaction)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "_store_current_reviewer_run_id", fake_store_current_run_id)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
@ -377,7 +381,7 @@ def test_process_github_review_finding_reply_uses_rereview_config(monkeypatch) -
assert config["finding_reply_id"] == "f_1"
def test_process_github_review_finding_reply_dispatches_sanitized_reply_body(monkeypatch) -> None:
def test_process_github_review_finding_reply_queues_reply_body_when_active(monkeypatch) -> None:
captured: dict[str, object] = {}
async def fake_get_thread_metadata_safe(_thread_id: str) -> dict[str, object]:
@ -403,16 +407,15 @@ def test_process_github_review_finding_reply_dispatches_sanitized_reply_body(mon
) -> dict[str, object]:
return {}
async def fake_store_current_run_id(_thread_id: str, _run: object) -> None:
return None
async def fake_is_thread_active(_thread_id: str) -> bool:
return True
class _FakeRunsClient:
async def create(self, thread_id: str, graph: str, **kwargs) -> dict[str, str]:
captured["kwargs"] = kwargs
return {"run_id": "run-1"}
async def fake_queue_message_for_thread(thread_id: str, message_content: object) -> bool:
captured["queued"] = {"thread_id": thread_id, "message_content": message_content}
return True
class _FakeLangGraphClient:
runs = _FakeRunsClient()
def fail_get_client(*_args: object, **_kwargs: object) -> None:
raise AssertionError("active reviewer thread should not create a new run")
monkeypatch.setattr(webapp, "_get_thread_metadata_safe", fake_get_thread_metadata_safe)
monkeypatch.setattr(
@ -423,8 +426,9 @@ def test_process_github_review_finding_reply_dispatches_sanitized_reply_body(mon
monkeypatch.setattr(webapp, "reconcile_findings_with_review_threads", fake_reconcile)
monkeypatch.setattr(webapp, "list_reviewer_findings", fake_list_findings)
monkeypatch.setattr(webapp, "append_finding_interaction", fake_append_interaction)
monkeypatch.setattr(webapp, "_store_current_reviewer_run_id", fake_store_current_run_id)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "queue_message_for_thread", fake_queue_message_for_thread)
monkeypatch.setattr(webapp, "get_client", fail_get_client)
asyncio.run(
webapp.process_github_review_finding_reply(
@ -447,9 +451,9 @@ def test_process_github_review_finding_reply_dispatches_sanitized_reply_body(mon
)
)
kwargs = captured["kwargs"]
assert isinstance(kwargs, dict)
message_content = kwargs["input"]["messages"][0]["content"]
queued = captured["queued"]
assert isinstance(queued, dict)
message_content = queued["message_content"]
assert isinstance(message_content, str)
assert "Open SWE finding f_1" in message_content
assert "untrusted data from GitHub" in message_content
@ -811,6 +815,10 @@ def test_process_github_pr_ready_creates_reviewer_run(monkeypatch) -> None:
captured["cache_token"] = token
captured["cache_expires_at"] = expires_at
async def fake_is_thread_active(thread_id: str) -> bool:
captured["active_thread_id"] = thread_id
return False
class _FakeRunsClient:
async def create(self, thread_id: str, graph: str, **kwargs) -> None:
captured["thread_id"] = thread_id
@ -840,6 +848,7 @@ def test_process_github_pr_ready_creates_reviewer_run(monkeypatch) -> None:
return 1
monkeypatch.setattr(webapp, "cache_github_token_for_thread", fake_cache_github_token)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "set_reviewer_thread_metadata", fake_set_reviewer_thread_metadata)
monkeypatch.setattr(webapp, "post_review_started_comment", fake_post_review_started_comment)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
@ -904,6 +913,10 @@ def test_trigger_pr_review_from_ref_creates_reviewer_run(monkeypatch) -> None:
captured["cache_token"] = token
captured["cache_expires_at"] = expires_at
async def fake_is_thread_active(thread_id: str) -> bool:
captured["active_thread_id"] = thread_id
return False
class _FakeRunsClient:
async def create(self, thread_id: str, graph: str, **kwargs) -> None:
captured["thread_id"] = thread_id
@ -937,6 +950,7 @@ def test_trigger_pr_review_from_ref_creates_reviewer_run(monkeypatch) -> None:
monkeypatch.setattr(webapp, "fetch_github_pr_metadata", fake_fetch_github_pr_metadata)
monkeypatch.setattr(webapp, "cache_github_token_for_thread", fake_cache_github_token)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "set_reviewer_thread_metadata", fake_set_reviewer_thread_metadata)
monkeypatch.setattr(webapp, "post_review_started_comment", fake_post_review_started_comment)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
@ -1015,7 +1029,7 @@ def test_trigger_pr_review_from_ref_respects_dashboard_opt_in(monkeypatch) -> No
assert called is False
async def test_request_pr_review_tool_uses_shared_trigger(monkeypatch) -> None:
def test_request_pr_review_tool_uses_shared_trigger(monkeypatch) -> None:
captured: dict[str, object] = {}
async def fake_trigger_pr_review_from_ref(
@ -1051,7 +1065,7 @@ async def test_request_pr_review_tool_uses_shared_trigger(monkeypatch) -> None:
},
)
result = await request_pr_review_tool("https://github.com/langchain-ai/open-swe/pull/1244")
result = request_pr_review_tool("https://github.com/langchain-ai/open-swe/pull/1244")
pr_ref = captured["pr_ref"]
assert isinstance(pr_ref, GitHubPrRef)
@ -1142,6 +1156,9 @@ def test_process_github_issue_uses_resolved_user_token_for_reaction(monkeypatch)
captured["fetch_token"] = token
return []
async def fake_is_thread_active(thread_id: str) -> bool:
return False
class _FakeRunsClient:
async def create(self, *args, **kwargs) -> None:
captured["run_created"] = True
@ -1158,6 +1175,7 @@ def test_process_github_issue_uses_resolved_user_token_for_reaction(monkeypatch)
monkeypatch.setattr(webapp, "_thread_exists", lambda thread_id: asyncio.sleep(0, result=False))
monkeypatch.setattr(webapp, "react_to_github_comment", fake_react_to_github_comment)
monkeypatch.setattr(webapp, "fetch_issue_comments", fake_fetch_issue_comments)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
monkeypatch.setattr(
webapp,
@ -1221,6 +1239,9 @@ def test_process_github_issue_existing_thread_uses_followup_prompt(monkeypatch)
async def fake_thread_exists(thread_id: str) -> bool:
return True
async def fake_is_thread_active(thread_id: str) -> bool:
return False
class _FakeRunsClient:
async def create(self, *args, **kwargs) -> None:
captured["prompt"] = kwargs["input"]["messages"][0]["content"]
@ -1237,6 +1258,7 @@ def test_process_github_issue_existing_thread_uses_followup_prompt(monkeypatch)
monkeypatch.setattr(webapp, "_thread_exists", fake_thread_exists)
monkeypatch.setattr(webapp, "react_to_github_comment", fake_react_to_github_comment)
monkeypatch.setattr(webapp, "fetch_issue_comments", fake_fetch_issue_comments)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClient())
monkeypatch.setattr(
webapp,

View file

@ -169,7 +169,7 @@ def test_plan_mode_guidance_section_present_when_enabled() -> None:
assert "Plan Mode (ACTIVE)" in prompt
async def test_enter_plan_mode_tool_returns_command() -> None:
def test_enter_plan_mode_tool_returns_command() -> None:
from langchain_core.messages import ToolMessage
from langchain_core.tools import tool as as_tool
from langgraph.types import Command
@ -178,7 +178,7 @@ async def test_enter_plan_mode_tool_returns_command() -> None:
# Wrap as the agent does so the InjectedToolCallId is supplied from the call.
wrapped = as_tool(enter_plan_mode)
result = await wrapped.ainvoke(
result = wrapped.invoke(
{"name": "enter_plan_mode", "args": {}, "id": "call-1", "type": "tool_call"}
)
assert isinstance(result, Command)

View file

@ -95,19 +95,19 @@ async def test_clear_plan_comments_deletes_each(monkeypatch: pytest.MonkeyPatch)
assert deleted == ["a", "b"]
async def test_save_plan_requires_run_context() -> None:
def test_save_plan_requires_run_context() -> None:
from agent.tools.save_plan import save_plan
# No LangGraph run context → no thread_id → graceful error, not a crash.
result = await save_plan("## Plan")
result = save_plan("## Plan")
assert result["success"] is False
assert "thread_id" in result["error"]
async def test_save_plan_rejects_empty_markdown() -> None:
def test_save_plan_rejects_empty_markdown() -> None:
from agent.tools.save_plan import save_plan
result = await save_plan(" ")
result = save_plan(" ")
assert result["success"] is False
assert "empty" in result["error"]

View file

@ -1,158 +0,0 @@
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from typing import Any
from unittest.mock import AsyncMock
import pytest
from agent import reconcile
def _run(run_id: str, thread_id: str, age_seconds: float) -> dict[str, Any]:
created = datetime.now(UTC) - timedelta(seconds=age_seconds)
return {
"run_id": run_id,
"thread_id": thread_id,
"status": "pending",
"created_at": created.isoformat(),
}
class _FakeThreads:
def __init__(self, pages: list[list[dict[str, Any]]]) -> None:
self._pages = pages
self.search_calls: list[dict[str, Any]] = []
async def search(self, **kwargs: Any) -> list[dict[str, Any]]:
self.search_calls.append(kwargs)
offset = kwargs.get("offset", 0)
limit = kwargs.get("limit", 100)
index = offset // limit if limit else 0
if index < len(self._pages):
return self._pages[index]
return []
class _FakeRuns:
def __init__(self, runs_by_thread: dict[str, Any]) -> None:
self._runs_by_thread = runs_by_thread
self.cancel_many = AsyncMock(return_value=None)
self.list_calls: list[tuple[str, dict[str, Any]]] = []
async def list(self, thread_id: str, **kwargs: Any) -> list[dict[str, Any]]:
self.list_calls.append((thread_id, kwargs))
value = self._runs_by_thread.get(thread_id, [])
if isinstance(value, Exception):
raise value
return value
class _FakeClient:
def __init__(self, threads: _FakeThreads, runs: _FakeRuns) -> None:
self.threads = threads
self.runs = runs
def _patch(monkeypatch: pytest.MonkeyPatch, client: _FakeClient) -> None:
monkeypatch.setattr(reconcile, "langgraph_client", lambda: client)
@pytest.mark.asyncio
async def test_cancels_only_stale_pending_runs(monkeypatch: pytest.MonkeyPatch) -> None:
threads = _FakeThreads([[{"thread_id": "t1"}]])
runs = _FakeRuns(
{
"t1": [
_run("old1", "t1", age_seconds=4000),
_run("fresh1", "t1", age_seconds=60),
_run("old2", "t1", age_seconds=10000),
]
}
)
_patch(monkeypatch, _FakeClient(threads, runs))
counts = await reconcile.reconcile_stale_runs(max_age_seconds=1800)
assert counts == {"threads_checked": 1, "stale_runs": 2, "cancelled": 2}
runs.cancel_many.assert_awaited_once()
kwargs = runs.cancel_many.await_args.kwargs
assert kwargs["thread_id"] == "t1"
assert sorted(kwargs["run_ids"]) == ["old1", "old2"]
@pytest.mark.asyncio
async def test_no_stale_runs_means_no_cancel(monkeypatch: pytest.MonkeyPatch) -> None:
threads = _FakeThreads([[{"thread_id": "t1"}]])
runs = _FakeRuns({"t1": [_run("fresh1", "t1", age_seconds=30)]})
_patch(monkeypatch, _FakeClient(threads, runs))
counts = await reconcile.reconcile_stale_runs(max_age_seconds=1800)
assert counts == {"threads_checked": 1, "stale_runs": 0, "cancelled": 0}
runs.cancel_many.assert_not_awaited()
@pytest.mark.asyncio
async def test_bad_thread_does_not_abort_sweep(monkeypatch: pytest.MonkeyPatch) -> None:
threads = _FakeThreads([[{"thread_id": "bad"}, {"thread_id": "good"}]])
runs = _FakeRuns(
{
"bad": RuntimeError("runs.list exploded"),
"good": [_run("old1", "good", age_seconds=5000)],
}
)
_patch(monkeypatch, _FakeClient(threads, runs))
counts = await reconcile.reconcile_stale_runs(max_age_seconds=1800)
# Both threads counted; the good thread is still reconciled despite the bad one.
assert counts == {"threads_checked": 2, "stale_runs": 1, "cancelled": 1}
runs.cancel_many.assert_awaited_once()
assert runs.cancel_many.await_args.kwargs["thread_id"] == "good"
assert runs.cancel_many.await_args.kwargs["run_ids"] == ["old1"]
@pytest.mark.asyncio
async def test_paginates_busy_threads(monkeypatch: pytest.MonkeyPatch) -> None:
full_page = [{"thread_id": f"t{i}"} for i in range(reconcile._SEARCH_PAGE_SIZE)]
second_page = [{"thread_id": "tail"}]
threads = _FakeThreads([full_page, second_page])
runs_by_thread: dict[str, Any] = {t["thread_id"]: [] for t in full_page}
runs_by_thread["tail"] = [_run("old", "tail", age_seconds=9000)]
runs = _FakeRuns(runs_by_thread)
_patch(monkeypatch, _FakeClient(threads, runs))
counts = await reconcile.reconcile_stale_runs(max_age_seconds=1800)
assert counts["threads_checked"] == reconcile._SEARCH_PAGE_SIZE + 1
assert counts["cancelled"] == 1
# Two search calls: first full page triggers a second page fetch.
assert len(threads.search_calls) == 2
assert threads.search_calls[0]["offset"] == 0
assert threads.search_calls[1]["offset"] == reconcile._SEARCH_PAGE_SIZE
assert threads.search_calls[0]["status"] == "busy"
@pytest.mark.asyncio
async def test_unparseable_created_at_is_skipped(monkeypatch: pytest.MonkeyPatch) -> None:
threads = _FakeThreads([[{"thread_id": "t1"}]])
runs = _FakeRuns(
{
"t1": [
{
"run_id": "bad",
"thread_id": "t1",
"status": "pending",
"created_at": "not-a-date",
},
_run("old", "t1", age_seconds=5000),
]
}
)
_patch(monkeypatch, _FakeClient(threads, runs))
counts = await reconcile.reconcile_stale_runs(max_age_seconds=1800)
assert counts == {"threads_checked": 1, "stale_runs": 1, "cancelled": 1}
assert runs.cancel_many.await_args.kwargs["run_ids"] == ["old"]

View file

@ -454,6 +454,10 @@ def _setup_slack_mention_fakes(
captured["user_names_by_id"] = user_names_by_id
return "", []
async def fake_is_thread_active(thread_id: str) -> bool:
captured["active_thread_id"] = thread_id
return False
async def fake_post_slack_trace_reply(channel_id: str, thread_ts: str, thread_id: str) -> None:
captured["trace_reply"] = {
"channel_id": channel_id,
@ -501,6 +505,7 @@ def _setup_slack_mention_fakes(
async def fake_post_prompt(*args, **kwargs) -> None:
captured["prompt"] = {"args": args, "kwargs": kwargs}
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "post_slack_trace_reply", fake_post_slack_trace_reply)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClientForProcess())
monkeypatch.setattr(webapp, "login_for_slack_id", fake_login_for_slack_id)
@ -542,6 +547,7 @@ def test_process_slack_mention_creates_thread_first_run_with_trace_reply(
assert captured["thread_exists_check"] == expected_thread_id
assert captured["fetch_thread"] == {"channel_id": "C123", "thread_ts": thread_ts}
assert captured["active_thread_id"] == expected_thread_id
assert captured["metadata_update"] == {
"thread_id": expected_thread_id,
"metadata": {"repo": {"owner": "langchain-ai", "name": "open-swe"}},
@ -558,8 +564,7 @@ def test_process_slack_mention_creates_thread_first_run_with_trace_reply(
assert run_create["graph"] == "agent"
kwargs = run_create["kwargs"]
assert kwargs["if_not_exists"] == "create"
assert kwargs["multitask_strategy"] == "interrupt"
assert kwargs["durability"] == "sync"
assert "multitask_strategy" not in kwargs
assert kwargs["config"]["configurable"]["slack_thread"]["thread_ts"] == thread_ts
prompt_block = kwargs["input"]["messages"][0]["content"][0]
assert "## Default Repository Hint\nlangchain-ai/open-swe" in prompt_block["text"]
@ -610,6 +615,217 @@ def test_process_slack_mention_skips_trace_reply_on_followup_mention(
assert run_create["thread_id"] == expected_thread_id
def test_process_slack_mention_queues_active_thread_message(
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, object] = {}
async def fake_get_slack_user_info(user_id: str) -> dict:
return {
"profile": {
"email": "mason@example.com",
"display_name": "Mason",
}
}
async def fake_fetch_slack_thread_messages(channel_id: str, thread_ts: str) -> list[dict]:
return [
{"ts": "1700000000.000100", "text": "<@UBOT> first request", "user": "U123"},
{
"ts": "1700000000.000200",
"text": "<@UBOT> include this screenshot https://example.com/image.png",
"user": "U123",
},
]
async def fake_get_slack_user_names(user_ids: list[str]) -> dict[str, str]:
captured["user_ids"] = user_ids
return {"U123": "Mason"}
async def fake_resolve_slack_links_in_context(
context_messages: list[dict], user_names_by_id: dict[str, str]
) -> tuple[str, list[str]]:
captured["context_messages"] = context_messages
return "", []
async def fake_fetch_image_block(image_url: str, http_client: object) -> None:
captured["image_url"] = image_url
return None
async def fake_is_thread_active(thread_id: str) -> bool:
captured["active_thread_id"] = thread_id
return True
async def fake_queue_message_for_thread(thread_id: str, message_content: object) -> bool:
captured["queued"] = {"thread_id": thread_id, "message_content": message_content}
return True
async def fake_post_slack_trace_reply(*args, **kwargs) -> None:
raise AssertionError("trace reply should not be posted for queued mid-run Slack messages")
async def fake_thread_exists(thread_id: str) -> bool:
return True
class _FakeRunsClient:
async def create(self, *args, **kwargs) -> None:
raise AssertionError("run should not be created for active Slack threads")
class _FakeThreadsClientForProcess:
async def update(self, *, thread_id: str, metadata: dict) -> None:
captured["metadata_update"] = {"thread_id": thread_id, "metadata": metadata}
class _FakeLangGraphClientForProcess:
runs = _FakeRunsClient()
threads = _FakeThreadsClientForProcess()
monkeypatch.setattr(webapp, "SLACK_BOT_USERNAME", "open-swe")
monkeypatch.setattr(webapp, "get_slack_user_info", fake_get_slack_user_info)
monkeypatch.setattr(webapp, "fetch_slack_thread_messages", fake_fetch_slack_thread_messages)
monkeypatch.setattr(webapp, "get_slack_user_names", fake_get_slack_user_names)
monkeypatch.setattr(
webapp, "resolve_slack_links_in_context", fake_resolve_slack_links_in_context
)
monkeypatch.setattr(webapp, "fetch_image_block", fake_fetch_image_block)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "queue_message_for_thread", fake_queue_message_for_thread)
monkeypatch.setattr(webapp, "post_slack_trace_reply", fake_post_slack_trace_reply)
monkeypatch.setattr(webapp, "_thread_exists", fake_thread_exists)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClientForProcess())
async def fake_login_for_slack_id(slack_user_id):
return "mason-gh"
async def fake_login_for_email(email):
return None
async def fake_refresh_cache() -> list:
return []
async def fake_get_valid_access_token(login):
return "user-token"
monkeypatch.setattr(webapp, "login_for_slack_id", fake_login_for_slack_id)
monkeypatch.setattr(webapp, "login_for_email", fake_login_for_email)
monkeypatch.setattr(webapp, "refresh_user_mapping_cache", fake_refresh_cache)
monkeypatch.setattr(webapp, "get_valid_access_token", fake_get_valid_access_token)
async def fake_resolve_agent_model_id(github_login, per_thread_model_id=None):
return "bedrock_converse:us.anthropic.claude-opus-4-8"
monkeypatch.setattr(webapp, "resolve_agent_model_id", fake_resolve_agent_model_id)
thread_ts = "1700000000.000100"
event_ts = "1700000000.000200"
expected_thread_id = generate_thread_id_from_slack_thread("C123", thread_ts)
asyncio.run(
webapp.process_slack_mention(
{
"channel_id": "C123",
"thread_ts": thread_ts,
"event_ts": event_ts,
"user_id": "U123",
"text": "<@UBOT> include this screenshot https://example.com/image.png",
"bot_user_id": "UBOT",
},
{"owner": "langchain-ai", "name": "open-swe"},
)
)
assert captured["active_thread_id"] == expected_thread_id
assert captured["queued"]["thread_id"] == expected_thread_id
queued_payload = captured["queued"]["message_content"]
assert queued_payload["image_urls"] == ["https://example.com/image.png"]
assert "## Latest Mention Request\ninclude this screenshot" in queued_payload["text"]
def test_process_slack_mention_serializes_concurrent_run_dispatch(
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, object] = {}
_setup_slack_mention_fakes(monkeypatch, captured)
thread_ts = "1700000001.000100"
expected_thread_id = generate_thread_id_from_slack_thread("C123", thread_ts)
first_active_started = asyncio.Event()
finish_first_active = asyncio.Event()
active_calls: list[str] = []
run_creates: list[dict[str, object]] = []
queued_messages: list[dict[str, object]] = []
async def fake_thread_exists(thread_id: str) -> bool:
return False
async def fake_is_thread_active(thread_id: str) -> bool:
active_calls.append(thread_id)
if len(active_calls) == 1:
first_active_started.set()
await finish_first_active.wait()
return bool(run_creates)
async def fake_queue_message_for_thread(thread_id: str, message_content: object) -> bool:
queued_messages.append({"thread_id": thread_id, "message_content": message_content})
return True
class _FakeRunsClient:
async def create(self, thread_id: str, graph: str, **kwargs) -> dict[str, str]:
run_creates.append({"thread_id": thread_id, "graph": graph, "kwargs": kwargs})
return {"run_id": f"run-{len(run_creates)}"}
class _FakeThreadsClientForProcess:
async def update(self, *, thread_id: str, metadata: dict) -> None:
captured["metadata_update"] = {"thread_id": thread_id, "metadata": metadata}
class _FakeLangGraphClientForProcess:
runs = _FakeRunsClient()
threads = _FakeThreadsClientForProcess()
monkeypatch.setattr(webapp, "_thread_exists", fake_thread_exists)
monkeypatch.setattr(webapp, "is_thread_active", fake_is_thread_active)
monkeypatch.setattr(webapp, "queue_message_for_thread", fake_queue_message_for_thread)
monkeypatch.setattr(webapp, "get_client", lambda url: _FakeLangGraphClientForProcess())
async def run_concurrent_mentions() -> None:
first = asyncio.create_task(
webapp.process_slack_mention(
{
"channel_id": "C123",
"thread_ts": thread_ts,
"event_ts": "1700000000.000200",
"user_id": "U123",
"text": "<@UBOT> first request",
"bot_user_id": "UBOT",
},
{"owner": "langchain-ai", "name": "open-swe"},
)
)
await first_active_started.wait()
second = asyncio.create_task(
webapp.process_slack_mention(
{
"channel_id": "C123",
"thread_ts": thread_ts,
"event_ts": "1700000000.000300",
"user_id": "U123",
"text": "<@UBOT> second request",
"bot_user_id": "UBOT",
},
{"owner": "langchain-ai", "name": "open-swe"},
)
)
await asyncio.sleep(0.05)
assert active_calls == [expected_thread_id]
finish_first_active.set()
await asyncio.gather(first, second)
asyncio.run(run_concurrent_mentions())
assert active_calls == [expected_thread_id, expected_thread_id]
assert len(run_creates) == 1
assert run_creates[0]["thread_id"] == expected_thread_id
assert queued_messages[0]["thread_id"] == expected_thread_id
def test_process_slack_mention_unmapped_user_blocked_and_prompted(
monkeypatch: pytest.MonkeyPatch,
) -> None:

View file

@ -169,157 +169,3 @@ async def test_slack_thread_reply_builds_option_blocks(monkeypatch: pytest.Monke
assert actions["type"] == "actions"
assert [button["text"]["text"] for button in actions["elements"]] == ["A", "B"]
assert actions["elements"][0]["action_id"] == "open_swe_option_select"
def _channel_only_config() -> dict[str, Any]:
return {"configurable": {"slack_thread": {"channel_id": "C9"}}}
def _channel_run_config() -> dict[str, Any]:
return {
"configurable": {
"slack_thread": {"channel_id": "C9"},
"thread_id": "run-1",
}
}
@pytest.fixture(autouse=True)
def _reset_top_level_posts() -> Any:
slack_reply_tool._top_level_posts.clear()
yield
slack_reply_tool._top_level_posts.clear()
async def test_slack_thread_reply_requires_channel_id(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setattr(slack_reply_tool, "get_config", lambda: {"configurable": {}})
result = await slack_reply_tool.slack_thread_reply("hello")
assert result["success"] is False
assert result["error"] == "Missing slack_thread.channel_id in config"
async def test_slack_thread_reply_posts_top_level_when_no_thread_ts(
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, Any] = {}
async def fake_top_level(
channel_id: str,
text: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
captured.update({"channel_id": channel_id, "text": text, "blocks": blocks})
return "3.0", None
async def fail_thread_reply(*args: Any, **kwargs: Any) -> tuple[str | None, str | None]:
raise AssertionError("should not post a thread reply without thread_ts")
monkeypatch.setattr(slack_reply_tool, "get_config", _channel_only_config)
monkeypatch.setattr(slack_reply_tool, "post_slack_top_level_message_with_ts", fake_top_level)
monkeypatch.setattr(slack_reply_tool, "post_slack_thread_reply_with_ts", fail_thread_reply)
result = await slack_reply_tool.slack_thread_reply("Scheduled report")
assert result == {"success": True}
assert captured["channel_id"] == "C9"
assert captured["text"] == "Scheduled report"
async def test_slack_thread_reply_top_level_surfaces_not_in_channel(
monkeypatch: pytest.MonkeyPatch,
) -> None:
async def fake_top_level(
channel_id: str,
text: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
return None, "not_in_channel"
monkeypatch.setattr(slack_reply_tool, "get_config", _channel_only_config)
monkeypatch.setattr(slack_reply_tool, "post_slack_top_level_message_with_ts", fake_top_level)
result = await slack_reply_tool.slack_thread_reply("Scheduled report")
assert result["success"] is False
assert result["slack_error"] == "not_in_channel"
assert "do not retry" in result["hint"]
async def test_slack_thread_reply_top_level_drops_interactive_blocks(
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, Any] = {}
async def fake_top_level(
channel_id: str,
text: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
captured["blocks"] = blocks
return "3.0", None
monkeypatch.setattr(slack_reply_tool, "get_config", _channel_only_config)
monkeypatch.setattr(slack_reply_tool, "post_slack_top_level_message_with_ts", fake_top_level)
options_result = await slack_reply_tool.slack_thread_reply("Pick", options=["A", "B"])
assert options_result == {"success": True}
assert captured["blocks"] is None
approval_result = await slack_reply_tool.slack_thread_reply("Plan?", plan_approval=True)
assert approval_result == {"success": True}
assert captured["blocks"] is None
async def test_slack_thread_reply_allows_only_one_top_level_post_per_run(
monkeypatch: pytest.MonkeyPatch,
) -> None:
calls = 0
async def fake_top_level(
channel_id: str,
text: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
nonlocal calls
calls += 1
return "3.0", None
monkeypatch.setattr(slack_reply_tool, "get_config", _channel_run_config)
monkeypatch.setattr(slack_reply_tool, "post_slack_top_level_message_with_ts", fake_top_level)
first = await slack_reply_tool.slack_thread_reply("First")
second = await slack_reply_tool.slack_thread_reply("Second")
assert first == {"success": True}
assert second["success"] is False
assert "Only one top-level message per run" in second["hint"]
assert calls == 1
async def test_slack_thread_reply_failed_top_level_post_does_not_consume_slot(
monkeypatch: pytest.MonkeyPatch,
) -> None:
results = iter([(None, "rate_limited"), ("3.0", None)])
async def fake_top_level(
channel_id: str,
text: str,
*,
blocks: list[dict[str, Any]] | None = None,
) -> tuple[str | None, str | None]:
return next(results)
monkeypatch.setattr(slack_reply_tool, "get_config", _channel_run_config)
monkeypatch.setattr(slack_reply_tool, "post_slack_top_level_message_with_ts", fake_top_level)
first = await slack_reply_tool.slack_thread_reply("First")
second = await slack_reply_tool.slack_thread_reply("Retry")
assert first["success"] is False
assert second == {"success": True}

View file

@ -44,7 +44,9 @@ function scheduleToSelection(
(model) =>
model.id === schedule.model && model.efforts.includes(schedule.effort!)
)
return supported ? { modelId: schedule.model, effort: schedule.effort } : null
return supported
? { modelId: schedule.model, effort: schedule.effort }
: null
}
export function AutomationEditor({
@ -71,9 +73,6 @@ export function AutomationEditor({
)
const [repo, setRepo] = useState<string | null>(schedule?.repo ?? null)
const [enabled, setEnabled] = useState(schedule?.enabled ?? true)
const [slackReportChannel, setSlackReportChannel] = useState(
schedule?.slackReportChannel ?? ""
)
// undefined = untouched (derive from the schedule / default as models load).
const [selectionOverride, setSelectionOverride] = useState<
ModelSelection | null | undefined
@ -116,7 +115,6 @@ export function AutomationEditor({
repo,
model_id: modelId,
effort,
slack_report_channel: slackReportChannel.trim() || null,
},
{ onSuccess: () => navigate({ to: "/agents/automations" }) }
)
@ -134,7 +132,6 @@ export function AutomationEditor({
model_id: modelId,
effort,
enabled,
slack_report_channel: slackReportChannel.trim() || null,
},
},
{ onSuccess: () => navigate({ to: "/agents/automations" }) }
@ -267,21 +264,6 @@ export function AutomationEditor({
</div>
</div>
<SectionLabel>Report to Slack channel</SectionLabel>
<div className="rounded-xl border border-[var(--ui-border)] bg-[var(--ui-surface)] p-3">
<input
value={slackReportChannel}
onChange={(e) => setSlackReportChannel(e.target.value)}
placeholder="Channel ID (e.g. C0123ABCD)"
className="w-full bg-transparent font-mono text-sm text-[var(--ui-text)] outline-none placeholder:text-[var(--ui-text-dim)]"
/>
<p className="mt-2 text-xs text-[var(--ui-text-dim)]">
When set, each run posts its final summary to this channel as the
Open SWE bot. The bot must be a member of the channel. Leave blank
for no Slack report.
</p>
</div>
{errorMessage && (
<p className="mt-4 text-xs text-[var(--ui-danger)]">{errorMessage}</p>
)}

View file

@ -27,7 +27,6 @@ export interface ScheduleCreateRequest {
repo?: string | null
model_id?: string | null
effort?: string | null
slack_report_channel?: string | null
}
export interface ScheduleUpdateRequest {
@ -38,7 +37,6 @@ export interface ScheduleUpdateRequest {
model_id?: string | null
effort?: string | null
enabled?: boolean | null
slack_report_channel?: string | null
}
export interface ThreadPrDiffFile {

View file

@ -164,7 +164,6 @@ export interface AgentSchedule {
repo: string | null
model: string
effort?: string | null
slackReportChannel?: string | null
enabled: boolean
cronId?: string | null
lastThreadId?: string | null