diff --git a/agent/analyzer.py b/agent/analyzer.py
index 13680de8..8f16934b 100644
--- a/agent/analyzer.py
+++ b/agent/analyzer.py
@@ -31,7 +31,11 @@ from langchain.agents.middleware import ModelCallLimitMiddleware
from .dashboard.team_settings import get_effective_gateway_enabled
from .integrations.langsmith import _configure_github_proxy
-from .middleware import SanitizeToolInputsMiddleware, ToolErrorMiddleware
+from .middleware import (
+ SanitizeToolInputsMiddleware,
+ TimeoutWrapupMiddleware,
+ ToolErrorMiddleware,
+)
from .review_style_guidance import REVIEWER_STYLE_THEMES
from .server import (
DEFAULT_LLM_MAX_TOKENS,
@@ -43,8 +47,9 @@ from .server import (
from .tools.read_finding_outcomes import read_finding_outcomes
from .tools.save_review_style import save_review_style_prompt
from .utils.analyzer_skills import SKILLS_ROUTE, skill_path_for_mode
+from .utils.deferred_model import make_model_or_defer
from .utils.github_app import get_github_app_installation_token
-from .utils.model import DEFAULT_LLM_REASONING, make_model, provider_model_kwargs
+from .utils.model import DEFAULT_LLM_REASONING, provider_model_kwargs
from .utils.sandbox_paths import aresolve_sandbox_work_dir
from .utils.sandbox_state import unwrap_sandbox_backend
from .utils.tracing import REVIEW_TRACING_PROJECT, traced_graph_factory
@@ -137,7 +142,7 @@ async def get_analyzer(config: RunnableConfig) -> Pregel:
system_prompt = f"{system_prompt}\n\n{user_context}"
return create_deep_agent(
- model=make_model(model_id, use_gateway=use_gateway, **model_kwargs),
+ model=make_model_or_defer(model_id, use_gateway=use_gateway, **model_kwargs),
system_prompt=system_prompt,
tools=[save_review_style_prompt, read_finding_outcomes],
backend=backend,
@@ -149,6 +154,7 @@ async def get_analyzer(config: RunnableConfig) -> Pregel:
exit_behavior="end",
),
ToolErrorMiddleware(),
+ TimeoutWrapupMiddleware(),
],
).with_config(config)
diff --git a/agent/completion.py b/agent/completion.py
index b711ce80..f2dfee9d 100644
--- a/agent/completion.py
+++ b/agent/completion.py
@@ -2,13 +2,17 @@
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
+ends in a failure state (``error`` / ``timeout``) 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.
+The reply is idempotent per run: the run id (or, for manual payloads without
+one, a run-distinguishing ``updated_at``/``created_at`` marker) is claimed in a
+bounded per-thread list *before* the post (claim-then-post), so a retried or
+concurrent completion webhook can't double-post. A payload that carries nothing
+to distinguish one run from another posts without deduping — a rare duplicate is
+preferred over silencing a later, genuinely different failed run.
"""
from __future__ import annotations
@@ -33,11 +37,13 @@ logger = logging.getLogger(__name__)
# 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"
+_FAILURE_REPLY_RUN_ID = "failure_reply_posted_run_id"
+_FAILURE_REPLY_RUN_IDS = "failure_reply_posted_run_ids"
+_MAX_FAILURE_REPLY_RUN_IDS = 20
class _ClaimFailed(Exception):
- """Raised when the dedup flag couldn't be claimed, so we skip the post."""
+ """Raised when the dedup key couldn't be claimed, so we skip the post."""
# Shared-secret bearer token proving a /webhooks/run-complete call came from our
@@ -65,9 +71,14 @@ def verify_run_complete_token(token: str | None) -> bool:
def _failure_text(status: str, dashboard_url: str | None = None) -> str:
- reason = "timed out" if status == "timeout" else "hit an unexpected error"
+ if status == "timeout":
+ reason = "timed out"
+ elif status == "interrupted":
+ reason = "was interrupted before it could finish"
+ else:
+ reason = "hit an unexpected error"
text = (
- f"⚠️ I wasn't able to finish that — the run {reason}. "
+ f"\u26a0\ufe0f I wasn't able to finish that \u2014 the run {reason}. "
"Send another message and I'll pick it back up."
)
if dashboard_url:
@@ -86,7 +97,7 @@ async def _post_failure_reply(
``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.
+ can't double-post and threads with no reply channel never burn the claim.
"""
source = metadata.get("source")
ctx = metadata.get("source_context")
@@ -131,6 +142,43 @@ async def _post_failure_reply(
return False
+def _dedup_key(payload: dict[str, Any]) -> str | None:
+ """A run-distinguishing key for dedupe, or None when nothing distinguishes runs.
+
+ Prefers the platform's ``run_id`` (always present on real completion
+ webhooks). Manual/legacy payloads without one fall back to a marker built
+ from ``updated_at``/``created_at`` so a retry of the *same* run still dedupes
+ while a *different* failed run on the same thread still gets its reply — the
+ old sticky per-thread boolean silenced every later run forever (SR160-03).
+ """
+ raw = payload.get("run_id")
+ if isinstance(raw, str) and raw:
+ return raw
+ for marker in ("updated_at", "created_at"):
+ value = payload.get(marker)
+ if isinstance(value, str) and value:
+ return f"{marker}:{value}"
+ return None
+
+
+def _posted_failure_run_ids(metadata: dict[str, Any]) -> list[str]:
+ raw = metadata.get(_FAILURE_REPLY_RUN_IDS)
+ ids = [item for item in raw if isinstance(item, str) and item] if isinstance(raw, list) else []
+ latest = metadata.get(_FAILURE_REPLY_RUN_ID)
+ if isinstance(latest, str) and latest and latest not in ids:
+ ids.append(latest)
+ return ids
+
+
+def _failure_reply_metadata(metadata: dict[str, Any], key: str) -> dict[str, Any]:
+ ids = [item for item in _posted_failure_run_ids(metadata) if item != key]
+ ids.append(key)
+ return {
+ _FAILURE_REPLY_RUN_ID: key,
+ _FAILURE_REPLY_RUN_IDS: ids[-_MAX_FAILURE_REPLY_RUN_IDS:],
+ }
+
+
async def handle_run_completion(payload: dict[str, Any]) -> dict[str, str]:
"""Handle a platform run-completion webhook POST.
@@ -153,17 +201,26 @@ async def handle_run_completion(payload: dict[str, Any]) -> dict[str, str]:
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"}
+ key = _dedup_key(payload)
+ if key is not None and key in _posted_failure_run_ids(metadata):
+ return {"status": "ignored", "reason": "failure reply already posted for run"}
- # 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.
+ # Claim-then-post: record the dedup key immediately before the actual post
+ # (via the claim callback) so a retried/concurrent completion webhook can't
+ # double-post (SR160-02). The key is only claimed on a branch that delivers,
+ # so a thread with no reply channel never burns it. If the claim itself fails
+ # we raise so we skip the post, leaving the key unclaimed for a later retry
+ # rather than reporting a clean success that invites a duplicate. A payload
+ # with no distinguishing key can't be recorded; it posts un-deduped rather
+ # than being permanently suppressed.
async def _claim() -> None:
+ if key is None:
+ return
try:
- await client.threads.update(thread_id=thread_id, metadata={_FAILURE_REPLY_FLAG: True})
+ await client.threads.update(
+ thread_id=thread_id,
+ metadata=_failure_reply_metadata(metadata, key),
+ )
except Exception as exc: # noqa: BLE001
logger.warning("run-complete: could not flag thread %s", thread_id, exc_info=True)
raise _ClaimFailed from exc
diff --git a/agent/dashboard/__init__.py b/agent/dashboard/__init__.py
index acf54f9e..e6a42bb4 100644
--- a/agent/dashboard/__init__.py
+++ b/agent/dashboard/__init__.py
@@ -1,5 +1,20 @@
-"""Dashboard backend: OAuth, profiles, and admin endpoints for the open-swe UI."""
+"""Dashboard backend: OAuth, profiles, and admin endpoints for the open-swe UI.
-from .routes import router
+``router`` is loaded lazily (PEP 562): importing any dashboard submodule
+(e.g. ``agent.dashboard.options`` from middleware) executes this __init__,
+and it must NOT drag in routes.py + FastAPI + every API/job module. Only the
+webapp, which actually mounts the router, pays that cost.
+"""
+
+from typing import Any
__all__ = ["router"]
+
+
+def __getattr__(name: str) -> Any:
+ if name == "router":
+ from .routes import router
+
+ globals()[name] = router
+ return router
+ raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
diff --git a/agent/dispatch.py b/agent/dispatch.py
index aa72225e..f818f8b8 100644
--- a/agent/dispatch.py
+++ b/agent/dispatch.py
@@ -17,7 +17,9 @@ from __future__ import annotations
import logging
import os
+import uuid
from typing import Any
+from urllib.parse import urlparse
from langgraph_sdk import get_client
from langgraph_sdk.client import LangGraphClient
@@ -25,22 +27,56 @@ from langgraph_sdk.client import LangGraphClient
logger = logging.getLogger(__name__)
ContentBlocks = str | list[dict[str, Any]]
+RunInput = dict[str, Any]
+RunConfig = 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.
+# FastAPI route the platform POSTs run completion/failure to. The platform
+# rejects loopback webhooks (relative URLs / localhost) — they bypass auth via
+# the in-process ASGI transport — so a loopback URL would 422 *every* run at
+# create time. COMPLETION_WEBHOOK_URL must therefore be the deployment's
+# absolute https URL (…/webhooks/run-complete). The route is fail-closed on
+# RUN_COMPLETE_WEBHOOK_SECRET, so we only attach 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). Secret unset, or URL relative/loopback
+# → no webhook attached (the completion reply is best-effort; it must never
+# break run creation).
_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 _is_loopback_webhook(url: str) -> bool:
+ """Whether a webhook URL is relative or points at localhost (platform-rejected)."""
+ parsed = urlparse(url)
+ if parsed.scheme not in ("http", "https") or not parsed.netloc:
+ return True # relative / schemeless
+ return (parsed.hostname or "").lower() in {"localhost", "127.0.0.1", "::1"}
+
+
+def _resolve_completion_webhook_url(base: str, secret: str | None) -> str | None:
+ """Resolve the completion webhook URL, or None to attach no webhook.
+
+ Degrades to None (with a warning) for a relative/loopback URL rather than
+ letting a rejected webhook poison every ``runs.create``.
+ """
+ if not secret:
+ return None
+ if _is_loopback_webhook(base):
+ logger.warning(
+ "RUN_COMPLETE_WEBHOOK_SECRET is set but COMPLETION_WEBHOOK_URL (%r) is relative "
+ "or loopback; the platform rejects such webhooks, so run-completion replies are "
+ "disabled. Set COMPLETION_WEBHOOK_URL to the deployment's absolute https URL "
+ "ending in /webhooks/run-complete to enable them.",
+ base,
+ )
+ return None
+ if "?" in base:
+ return base
+ return f"{base}?token={secret}"
+
+
+COMPLETION_WEBHOOK_URL: str | None = _resolve_completion_webhook_url(
+ _COMPLETION_WEBHOOK_BASE, _RUN_COMPLETE_SECRET
+)
def _langgraph_url() -> str:
@@ -53,6 +89,65 @@ def dispatch_client() -> LangGraphClient:
return get_client(url=_langgraph_url())
+def _config_with_prepare_run_id(
+ config: RunConfig | None,
+ metadata: dict[str, Any] | None,
+) -> RunConfig:
+ run_config = dict(config or {})
+ configurable = run_config.get("configurable")
+ configurable = dict(configurable) if isinstance(configurable, dict) else {}
+ configurable.setdefault("prepare_run_id", str(uuid.uuid4()))
+ run_config["configurable"] = configurable
+ if metadata is not None:
+ run_config["metadata"] = metadata
+ return run_config
+
+
+async def create_durable_run(
+ thread_id: str,
+ assistant_id: str,
+ *,
+ input: RunInput,
+ source: str,
+ config: RunConfig | None = None,
+ metadata: dict[str, Any] | None = None,
+ client: LangGraphClient | None = None,
+ multitask_strategy: str = "interrupt",
+ durability: str = "sync",
+ if_not_exists: str = "create",
+ stream_mode: Any | None = None,
+ stream_resumable: bool | None = None,
+ after_seconds: int | float | None = None,
+) -> dict[str, Any]:
+ """Create a run with Open SWE's durable LangGraph defaults."""
+ client = client or dispatch_client()
+ create_kwargs: dict[str, Any] = {
+ "input": input,
+ "config": _config_with_prepare_run_id(config, metadata),
+ "multitask_strategy": multitask_strategy,
+ "durability": durability,
+ "if_not_exists": if_not_exists,
+ }
+ if COMPLETION_WEBHOOK_URL:
+ create_kwargs["webhook"] = COMPLETION_WEBHOOK_URL
+ if stream_mode is not None:
+ create_kwargs["stream_mode"] = stream_mode
+ if stream_resumable is not None:
+ create_kwargs["stream_resumable"] = stream_resumable
+ if after_seconds is not None:
+ create_kwargs["after_seconds"] = after_seconds
+
+ run = await client.runs.create(thread_id, assistant_id, **create_kwargs)
+ 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
+
+
async def dispatch_agent_run(
thread_id: str,
content: ContentBlocks,
@@ -72,22 +167,13 @@ async def dispatch_agent_run(
``"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(
+ return await create_durable_run(
thread_id,
assistant_id,
input={"messages": [{"role": "user", "content": content}]},
- config={"configurable": configurable, "metadata": metadata or {}},
+ config={"configurable": configurable},
+ metadata=metadata or {},
+ source=source,
+ client=client or dispatch_client(),
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
diff --git a/agent/middleware/__init__.py b/agent/middleware/__init__.py
index d6d94403..19bc0466 100644
--- a/agent/middleware/__init__.py
+++ b/agent/middleware/__init__.py
@@ -1,22 +1,31 @@
-from .check_message_queue import check_message_queue_before_model
-from .ensure_no_empty_msg import ensure_no_empty_msg
-from .exclude_tools import ExcludeToolsMiddleware
-from .model_fallback import ModelFallbackMiddleware
-from .notify_step_limit import notify_step_limit_reached
-from .plan_mode import PlanModeMiddleware
-from .refresh_github_proxy import refresh_github_proxy_before_model
-from .refresh_slack_status import SlackAssistantStatusMiddleware
-from .repair_orphaned_tool_calls import RepairOrphanedToolCallsMiddleware
-from .sandbox_circuit_breaker import SandboxCircuitBreakerMiddleware
-from .sanitize_fireworks_messages import SanitizeFireworksMessagesMiddleware
-from .sanitize_openai_responses import SanitizeOpenAIResponsesMiddleware
-from .sanitize_thinking_blocks import SanitizeThinkingBlocksMiddleware
-from .sanitize_tool_inputs import SanitizeToolInputsMiddleware
-from .settle_review_check import settle_review_check_on_exit
-from .subdir_agents import SubdirAgentsReadMiddleware
-from .tool_artifact import ToolArtifactMiddleware
-from .tool_error_handler import ToolErrorMiddleware
-from .workflow_push_guard import WorkflowPushGuardMiddleware
+import sys
+from types import ModuleType
+from typing import TYPE_CHECKING, Any
+
+_MIDDLEWARE_MODULES = {
+ "check_message_queue_before_model": ".check_message_queue",
+ "ensure_no_empty_msg": ".ensure_no_empty_msg",
+ "ExcludeToolsMiddleware": ".exclude_tools",
+ "ModelFallbackMiddleware": ".model_fallback",
+ "notify_step_limit_reached": ".notify_step_limit",
+ "PlanModeMiddleware": ".plan_mode",
+ "refresh_github_proxy_before_model": ".refresh_github_proxy",
+ "RepairOrphanedToolCallsMiddleware": ".repair_orphaned_tool_calls",
+ "SlackAssistantStatusMiddleware": ".refresh_slack_status",
+ "SandboxCircuitBreakerMiddleware": ".sandbox_circuit_breaker",
+ "SanitizeFireworksMessagesMiddleware": ".sanitize_fireworks_messages",
+ "SanitizeOpenAIResponsesMiddleware": ".sanitize_openai_responses",
+ "SanitizeThinkingBlocksMiddleware": ".sanitize_thinking_blocks",
+ "SanitizeToolInputsMiddleware": ".sanitize_tool_inputs",
+ "settle_review_check_on_exit": ".settle_review_check",
+ "SubdirAgentsReadMiddleware": ".subdir_agents",
+ "task_on_failure": ".task_retry",
+ "task_retry_on": ".task_retry",
+ "TimeoutWrapupMiddleware": ".timeout_wrapup",
+ "ToolArtifactMiddleware": ".tool_artifact",
+ "ToolErrorMiddleware": ".tool_error_handler",
+ "WorkflowPushGuardMiddleware": ".workflow_push_guard",
+}
__all__ = [
"ExcludeToolsMiddleware",
@@ -30,6 +39,7 @@ __all__ = [
"SubdirAgentsReadMiddleware",
"ToolArtifactMiddleware",
"ToolErrorMiddleware",
+ "TimeoutWrapupMiddleware",
"WorkflowPushGuardMiddleware",
"SandboxCircuitBreakerMiddleware",
"SlackAssistantStatusMiddleware",
@@ -38,4 +48,51 @@ __all__ = [
"notify_step_limit_reached",
"refresh_github_proxy_before_model",
"settle_review_check_on_exit",
+ "task_on_failure",
+ "task_retry_on",
]
+
+if TYPE_CHECKING:
+ from .check_message_queue import check_message_queue_before_model
+ from .ensure_no_empty_msg import ensure_no_empty_msg
+ from .exclude_tools import ExcludeToolsMiddleware
+ from .model_fallback import ModelFallbackMiddleware
+ from .notify_step_limit import notify_step_limit_reached
+ from .plan_mode import PlanModeMiddleware
+ from .refresh_github_proxy import refresh_github_proxy_before_model
+ from .refresh_slack_status import SlackAssistantStatusMiddleware
+ from .repair_orphaned_tool_calls import RepairOrphanedToolCallsMiddleware
+ from .sandbox_circuit_breaker import SandboxCircuitBreakerMiddleware
+ from .sanitize_fireworks_messages import SanitizeFireworksMessagesMiddleware
+ from .sanitize_openai_responses import SanitizeOpenAIResponsesMiddleware
+ from .sanitize_thinking_blocks import SanitizeThinkingBlocksMiddleware
+ from .sanitize_tool_inputs import SanitizeToolInputsMiddleware
+ from .settle_review_check import settle_review_check_on_exit
+ from .subdir_agents import SubdirAgentsReadMiddleware
+ from .task_retry import task_on_failure, task_retry_on
+ from .timeout_wrapup import TimeoutWrapupMiddleware
+ from .tool_artifact import ToolArtifactMiddleware
+ from .tool_error_handler import ToolErrorMiddleware
+ from .workflow_push_guard import WorkflowPushGuardMiddleware
+
+
+def _load_export(name: str) -> Any:
+ module_name = _MIDDLEWARE_MODULES.get(name)
+ if module_name is None:
+ raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
+ from importlib import import_module
+
+ value = getattr(import_module(module_name, __name__), name)
+ globals()[name] = value
+ return value
+
+
+class _LazyMiddlewareModule(ModuleType):
+ def __getattribute__(self, name: str) -> Any:
+ module_map = ModuleType.__getattribute__(self, "__dict__").get("_MIDDLEWARE_MODULES", {})
+ if name in module_map:
+ return _load_export(name)
+ return ModuleType.__getattribute__(self, name)
+
+
+sys.modules[__name__].__class__ = _LazyMiddlewareModule
diff --git a/agent/middleware/task_retry.py b/agent/middleware/task_retry.py
new file mode 100644
index 00000000..d532bc2d
--- /dev/null
+++ b/agent/middleware/task_retry.py
@@ -0,0 +1,82 @@
+from __future__ import annotations
+
+import json
+
+_RETURN_TO_MODEL_CODES = frozenset({"invalid_prompt", "context_length_exceeded"})
+_RETURN_TO_MODEL_STATUS_CODES = frozenset({400, 422})
+_RETRY_HTTP_STATUS_CODES = frozenset({408, 409, 425, 429, 500, 502, 503, 504, 529})
+_TRANSIENT_ERROR_NAMES = frozenset(
+ {
+ "APIConnectionError",
+ "APITimeoutError",
+ "ConnectTimeout",
+ "ReadTimeout",
+ "TimeoutException",
+ "TransportError",
+ }
+)
+
+
+def _error_body(exc: Exception) -> dict[str, object]:
+ body = getattr(exc, "body", None)
+ if isinstance(body, dict):
+ nested = body.get("error")
+ return nested if isinstance(nested, dict) else body
+ return {}
+
+
+def _status_code(exc: Exception) -> int | None:
+ status = getattr(exc, "status_code", None)
+ if isinstance(status, int):
+ return status
+ response = getattr(exc, "response", None)
+ status = getattr(response, "status_code", None)
+ return status if isinstance(status, int) else None
+
+
+def _error_fields(exc: Exception) -> dict[str, object]:
+ body = _error_body(exc)
+ out: dict[str, object] = {}
+ status = _status_code(exc)
+ if status is not None:
+ out["status_code"] = status
+ for key in ("type", "code", "message"):
+ value = body.get(key)
+ if not isinstance(value, str) or not value:
+ value = getattr(exc, key, None)
+ if isinstance(value, str) and value:
+ out[key] = value
+ return out
+
+
+def _is_httpx_transport_error(exc: Exception) -> bool:
+ try:
+ import httpx
+ except ImportError: # pragma: no cover - dependency is declared in production
+ return False
+ return isinstance(exc, httpx.TransportError)
+
+
+def task_retry_on(exc: Exception) -> bool:
+ status = _status_code(exc)
+ if isinstance(status, int) and (status in _RETRY_HTTP_STATUS_CODES or status >= 500):
+ return True
+ return exc.__class__.__name__ in _TRANSIENT_ERROR_NAMES or _is_httpx_transport_error(exc)
+
+
+def task_on_failure(exc: Exception) -> str:
+ error = _error_fields(exc)
+ code = error.get("code")
+ status = error.get("status_code")
+ returnable = code in _RETURN_TO_MODEL_CODES or (
+ code is None
+ and error.get("type") == "invalid_request_error"
+ and isinstance(status, int)
+ and status in _RETURN_TO_MODEL_STATUS_CODES
+ )
+ if not returnable:
+ raise exc
+ return json.dumps(
+ {"status": "failed", "source": "subagent", "error": error},
+ sort_keys=True,
+ )
diff --git a/agent/middleware/timeout_wrapup.py b/agent/middleware/timeout_wrapup.py
new file mode 100644
index 00000000..78a0cd3e
--- /dev/null
+++ b/agent/middleware/timeout_wrapup.py
@@ -0,0 +1,64 @@
+from __future__ import annotations
+
+import os
+import time
+from collections.abc import Awaitable, Callable
+
+from langchain.agents.middleware.types import AgentMiddleware, ModelRequest, ModelResponse
+from langchain_core.messages import BaseMessage, SystemMessage
+
+_DEFAULT_TIMEOUT_SECONDS = 45 * 60
+_WRAPUP_INSTRUCTION = """
+
+You have been running for a long time. Wrap up immediately: finish the current
+step, save or report useful state, avoid starting new investigations, and end
+your turn with the best available result.
+
+"""
+
+
+def _configured_timeout_seconds() -> int:
+ raw = os.environ.get("OPEN_SWE_WRAPUP_TIMEOUT_SECONDS")
+ if not raw:
+ return _DEFAULT_TIMEOUT_SECONDS
+ try:
+ value = int(raw)
+ except ValueError:
+ return _DEFAULT_TIMEOUT_SECONDS
+ return value if value > 0 else _DEFAULT_TIMEOUT_SECONDS
+
+
+def _content_with_instruction(message: BaseMessage | None, instruction: str) -> str | list[object]:
+ if message is None:
+ return instruction
+ content = message.content
+ if isinstance(content, list):
+ return [*content, {"type": "text", "text": instruction}]
+ return f"{content}\n\n{instruction}" if content else instruction
+
+
+class TimeoutWrapupMiddleware(AgentMiddleware):
+ def __init__(self, timeout_seconds: int | None = None) -> None:
+ super().__init__()
+ self._timeout_seconds = timeout_seconds or _configured_timeout_seconds()
+ # Graph construction should create one middleware instance per run; start
+ # lazily so construction-time caching cannot age the run clock.
+ self._start: float | None = None
+
+ def _should_wrapup(self) -> bool:
+ if self._start is None:
+ self._start = time.monotonic()
+ return (time.monotonic() - self._start) >= self._timeout_seconds
+
+ def _apply(self, request: ModelRequest) -> ModelRequest:
+ if not self._should_wrapup():
+ return request
+ content = _content_with_instruction(request.system_message, _WRAPUP_INSTRUCTION)
+ return request.override(system_message=SystemMessage(content=content))
+
+ async def awrap_model_call(
+ self,
+ request: ModelRequest,
+ handler: Callable[[ModelRequest], Awaitable[ModelResponse]],
+ ) -> ModelResponse:
+ return await handler(self._apply(request))
diff --git a/agent/reviewer.py b/agent/reviewer.py
index 2bafe7ba..18151c0b 100644
--- a/agent/reviewer.py
+++ b/agent/reviewer.py
@@ -46,6 +46,7 @@ from .middleware import (
SanitizeThinkingBlocksMiddleware,
SanitizeToolInputsMiddleware,
SlackAssistantStatusMiddleware,
+ TimeoutWrapupMiddleware,
ToolErrorMiddleware,
check_message_queue_before_model,
refresh_github_proxy_before_model,
@@ -85,9 +86,10 @@ from .tools import (
)
from .utils.agents_md import fetch_agents_md
from .utils.api_standards_skill import fetch_api_standards_skill
+from .utils.deferred_model import make_model_or_defer
from .utils.github_app import get_github_app_installation_token_with_expiry
from .utils.github_token import cache_github_token_for_thread
-from .utils.model import DEFAULT_LLM_REASONING, make_model, provider_model_kwargs
+from .utils.model import DEFAULT_LLM_REASONING, provider_model_kwargs
from .utils.repo_prep import materialize_trusted_skills, prepare_review_repo
from .utils.sandbox_paths import aresolve_sandbox_work_dir
from .utils.tracing import REVIEW_TRACING_PROJECT, traced_graph_factory
@@ -826,7 +828,7 @@ async def _resolve_grouping_model(
max_tokens=DEFAULT_LLM_MAX_TOKENS,
openai_reasoning_default=DEFAULT_LLM_REASONING,
)
- return make_model(model_id, use_gateway=use_gateway, **model_kwargs)
+ return make_model_or_defer(model_id, use_gateway=use_gateway, **model_kwargs)
async def get_reviewer_agent(config: RunnableConfig) -> Pregel:
@@ -1162,8 +1164,8 @@ async def get_reviewer_agent(config: RunnableConfig) -> Pregel:
system_prompt = f"{system_prompt}\n\n{review_context}"
use_gateway = await get_effective_gateway_enabled()
- reviewer_model = make_model(model_id, use_gateway=use_gateway, **model_kwargs)
- reviewer_subagent_model = make_model(
+ reviewer_model = make_model_or_defer(model_id, use_gateway=use_gateway, **model_kwargs)
+ reviewer_subagent_model = make_model_or_defer(
subagent_model_id, use_gateway=use_gateway, **subagent_model_kwargs
)
@@ -1211,6 +1213,7 @@ async def get_reviewer_agent(config: RunnableConfig) -> Pregel:
refresh_github_proxy_before_model,
check_message_queue_before_model,
SlackAssistantStatusMiddleware(),
+ TimeoutWrapupMiddleware(),
SanitizeOpenAIResponsesMiddleware(),
SanitizeFireworksMessagesMiddleware(),
SanitizeThinkingBlocksMiddleware(),
diff --git a/agent/server.py b/agent/server.py
index 5e8af637..619d3cfa 100644
--- a/agent/server.py
+++ b/agent/server.py
@@ -27,7 +27,7 @@ from deepagents import create_deep_agent
from deepagents.backends import LangSmithSandbox
from deepagents.backends.protocol import SandboxBackendProtocol
from deepagents.middleware.subagents import GENERAL_PURPOSE_SUBAGENT, SubAgent
-from langchain.agents.middleware import ModelCallLimitMiddleware
+from langchain.agents.middleware import ModelCallLimitMiddleware, ToolRetryMiddleware
from langchain_core.language_models import BaseChatModel
from langsmith.sandbox import SandboxClientError
@@ -65,6 +65,7 @@ from .middleware import (
SanitizeToolInputsMiddleware,
SlackAssistantStatusMiddleware,
SubdirAgentsReadMiddleware,
+ TimeoutWrapupMiddleware,
ToolArtifactMiddleware,
ToolErrorMiddleware,
WorkflowPushGuardMiddleware,
@@ -72,6 +73,8 @@ from .middleware import (
ensure_no_empty_msg,
notify_step_limit_reached,
refresh_github_proxy_before_model,
+ task_on_failure,
+ task_retry_on,
)
from .prompt import construct_system_prompt
from .tools import (
@@ -103,6 +106,7 @@ from .utils.authorship import (
resolve_triggering_user_identity,
)
from .utils.dashboard_links import dashboard_plan_url, dashboard_thread_url
+from .utils.deferred_model import make_model_or_defer
from .utils.github_app import (
BASE_RUNTIME_PROXY_TOKEN_PERMISSIONS,
RUNTIME_PROXY_TOKEN_PERMISSIONS,
@@ -113,7 +117,6 @@ from .utils.github_proxy import record_proxy_token_expiry
from .utils.github_token import repo_cache_key
from .utils.model import (
fallback_model_id_for,
- make_model,
provider_model_kwargs,
)
from .utils.sandbox import create_sandbox
@@ -876,7 +879,7 @@ async def get_agent(config: RunnableConfig) -> Pregel:
)
fallback_middleware.append(
ModelFallbackMiddleware(
- make_model(fallback_model_id, use_gateway=use_gateway, **fallback_kwargs)
+ make_model_or_defer(fallback_model_id, use_gateway=use_gateway, **fallback_kwargs)
)
)
logger.info("Configured model fallback %s -> %s", model_id, fallback_model_id)
@@ -943,8 +946,10 @@ async def get_agent(config: RunnableConfig) -> Pregel:
notion_tools = []
logger.info("Returning agent with sandbox for thread %s", thread_id)
- main_model = make_model(model_id, use_gateway=use_gateway, **model_kwargs)
- subagent_model = make_model(subagent_model_id, use_gateway=use_gateway, **subagent_model_kwargs)
+ main_model = make_model_or_defer(model_id, use_gateway=use_gateway, **model_kwargs)
+ subagent_model = make_model_or_defer(
+ subagent_model_id, use_gateway=use_gateway, **subagent_model_kwargs
+ )
return create_deep_agent(
model=main_model,
system_prompt=construct_system_prompt(
@@ -994,11 +999,20 @@ async def get_agent(config: RunnableConfig) -> Pregel:
ModelCallLimitMiddleware(run_limit=MODEL_CALL_RECURSION_LIMIT, exit_behavior="end"),
ToolErrorMiddleware(),
SubdirAgentsReadMiddleware(),
+ ToolRetryMiddleware(
+ max_retries=2,
+ tools=["task"],
+ retry_on=task_retry_on,
+ on_failure=task_on_failure,
+ initial_delay=1.0,
+ max_delay=10.0,
+ ),
ToolArtifactMiddleware(),
WorkflowPushGuardMiddleware(),
refresh_github_proxy_before_model,
check_message_queue_before_model,
SlackAssistantStatusMiddleware(),
+ TimeoutWrapupMiddleware(),
ensure_no_empty_msg,
notify_step_limit_reached,
SandboxCircuitBreakerMiddleware(),
diff --git a/agent/tools/__init__.py b/agent/tools/__init__.py
index 8642474d..e48efdf7 100644
--- a/agent/tools/__init__.py
+++ b/agent/tools/__init__.py
@@ -1,32 +1,38 @@
-from .add_finding import add_finding
-from .enter_plan_mode import enter_plan_mode
-from .fetch_url import fetch_url
-from .http_request import http_request
-from .linear_comment import linear_comment
-from .linear_create_issue import linear_create_issue
-from .linear_delete_issue import linear_delete_issue
-from .linear_get_issue import linear_get_issue
-from .linear_get_issue_comments import linear_get_issue_comments
-from .linear_list_teams import linear_list_teams
-from .linear_update_issue import linear_update_issue
-from .list_findings import list_findings
-from .list_review_findings import list_review_findings
-from .open_pull_request import open_pull_request
-from .publish_review import publish_review
-from .read_repo_file import read_repo_file
-from .reply_to_finding_thread import reply_to_finding_thread
-from .report_platform_issue import report_platform_issue
-from .request_pr_review import request_pr_review
-from .resolve_finding_thread import resolve_finding_thread
-from .save_plan import save_plan
-from .schedule_thread_wakeup import schedule_thread_wakeup
-from .search_repo_code import search_repo_code
-from .slack_add_reaction import slack_add_reaction
-from .slack_read_thread_messages import slack_read_thread_messages
-from .slack_start_new_thread import slack_start_new_thread
-from .slack_thread_reply import slack_thread_reply
-from .update_finding import update_finding
-from .web_search import web_search
+import sys
+from types import ModuleType
+from typing import TYPE_CHECKING, Any
+
+_TOOL_MODULES = {
+ "add_finding": ".add_finding",
+ "enter_plan_mode": ".enter_plan_mode",
+ "fetch_url": ".fetch_url",
+ "http_request": ".http_request",
+ "linear_comment": ".linear_comment",
+ "linear_create_issue": ".linear_create_issue",
+ "linear_delete_issue": ".linear_delete_issue",
+ "linear_get_issue": ".linear_get_issue",
+ "linear_get_issue_comments": ".linear_get_issue_comments",
+ "linear_list_teams": ".linear_list_teams",
+ "linear_update_issue": ".linear_update_issue",
+ "list_findings": ".list_findings",
+ "list_review_findings": ".list_review_findings",
+ "open_pull_request": ".open_pull_request",
+ "publish_review": ".publish_review",
+ "read_repo_file": ".read_repo_file",
+ "report_platform_issue": ".report_platform_issue",
+ "request_pr_review": ".request_pr_review",
+ "reply_to_finding_thread": ".reply_to_finding_thread",
+ "resolve_finding_thread": ".resolve_finding_thread",
+ "save_plan": ".save_plan",
+ "schedule_thread_wakeup": ".schedule_thread_wakeup",
+ "search_repo_code": ".search_repo_code",
+ "slack_add_reaction": ".slack_add_reaction",
+ "slack_read_thread_messages": ".slack_read_thread_messages",
+ "slack_start_new_thread": ".slack_start_new_thread",
+ "slack_thread_reply": ".slack_thread_reply",
+ "update_finding": ".update_finding",
+ "web_search": ".web_search",
+}
__all__ = [
"add_finding",
@@ -59,3 +65,56 @@ __all__ = [
"update_finding",
"web_search",
]
+
+if TYPE_CHECKING:
+ from .add_finding import add_finding
+ from .enter_plan_mode import enter_plan_mode
+ from .fetch_url import fetch_url
+ from .http_request import http_request
+ from .linear_comment import linear_comment
+ from .linear_create_issue import linear_create_issue
+ from .linear_delete_issue import linear_delete_issue
+ from .linear_get_issue import linear_get_issue
+ from .linear_get_issue_comments import linear_get_issue_comments
+ from .linear_list_teams import linear_list_teams
+ from .linear_update_issue import linear_update_issue
+ from .list_findings import list_findings
+ from .list_review_findings import list_review_findings
+ from .open_pull_request import open_pull_request
+ from .publish_review import publish_review
+ from .read_repo_file import read_repo_file
+ from .reply_to_finding_thread import reply_to_finding_thread
+ from .report_platform_issue import report_platform_issue
+ from .request_pr_review import request_pr_review
+ from .resolve_finding_thread import resolve_finding_thread
+ from .save_plan import save_plan
+ from .schedule_thread_wakeup import schedule_thread_wakeup
+ from .search_repo_code import search_repo_code
+ from .slack_add_reaction import slack_add_reaction
+ from .slack_read_thread_messages import slack_read_thread_messages
+ from .slack_start_new_thread import slack_start_new_thread
+ from .slack_thread_reply import slack_thread_reply
+ from .update_finding import update_finding
+ from .web_search import web_search
+
+
+def _load_tool(name: str) -> Any:
+ module_name = _TOOL_MODULES.get(name)
+ if module_name is None:
+ raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
+ from importlib import import_module
+
+ value = getattr(import_module(module_name, __name__), name)
+ globals()[name] = value
+ return value
+
+
+class _LazyToolsModule(ModuleType):
+ def __getattribute__(self, name: str) -> Any:
+ tool_map = ModuleType.__getattribute__(self, "__dict__").get("_TOOL_MODULES", {})
+ if name in tool_map:
+ return _load_tool(name)
+ return ModuleType.__getattribute__(self, name)
+
+
+sys.modules[__name__].__class__ = _LazyToolsModule
diff --git a/agent/tools/request_pr_review.py b/agent/tools/request_pr_review.py
index d6f1bb54..e55a3506 100644
--- a/agent/tools/request_pr_review.py
+++ b/agent/tools/request_pr_review.py
@@ -2,8 +2,28 @@ from typing import Any
from langgraph.config import get_config
-from agent.utils.slack import parse_github_pr_url
-from agent.webapp import trigger_pr_review_from_ref
+from agent.utils.slack import GitHubPrRef, parse_github_pr_url
+
+
+async def trigger_pr_review_from_ref(
+ pr_ref: GitHubPrRef,
+ *,
+ source: str,
+ github_login: str = "",
+ github_user_id: int | None = None,
+ slack_channel_id: str = "",
+ slack_thread_ts: str = "",
+) -> dict[str, Any]:
+ from agent.webapp import trigger_pr_review_from_ref as _trigger_pr_review_from_ref
+
+ return await _trigger_pr_review_from_ref(
+ pr_ref,
+ source=source,
+ github_login=github_login,
+ github_user_id=github_user_id,
+ slack_channel_id=slack_channel_id,
+ slack_thread_ts=slack_thread_ts,
+ )
async def request_pr_review(pr_url: str) -> dict[str, Any]:
diff --git a/agent/tools/web_search.py b/agent/tools/web_search.py
index 18fd21e2..d6b99a05 100644
--- a/agent/tools/web_search.py
+++ b/agent/tools/web_search.py
@@ -3,8 +3,6 @@ import logging
import os
from typing import Any
-from exa_py import Exa
-
logger = logging.getLogger(__name__)
@@ -38,6 +36,8 @@ async def web_search(
}
async def _search() -> dict[str, Any]:
+ from exa_py import Exa # deferred: heavy import
+
client = Exa(api_key=api_key)
if include_contents:
result = await asyncio.to_thread(
diff --git a/agent/utils/deferred_model.py b/agent/utils/deferred_model.py
new file mode 100644
index 00000000..6511ff70
--- /dev/null
+++ b/agent/utils/deferred_model.py
@@ -0,0 +1,57 @@
+from __future__ import annotations
+
+import logging
+from typing import Any
+
+from langchain_core.language_models import BaseChatModel
+
+from .model import make_model
+
+logger = logging.getLogger(__name__)
+
+
+class DeferredErrorModel(BaseChatModel):
+ """Model placeholder that raises a stored setup error on first invocation."""
+
+ error_message: str
+ model_id: str | None = None
+
+ @property
+ def _llm_type(self) -> str:
+ return "deferred-error"
+
+ def _get_ls_params(self, stop: Any = None, **kwargs: Any) -> dict[str, Any]:
+ params = super()._get_ls_params(stop=stop, **kwargs)
+ if self.model_id:
+ params["ls_model_name"] = self.model_id
+ if ":" in self.model_id:
+ params["ls_provider"] = self.model_id.split(":", 1)[0]
+ return params
+
+ def bind_tools(self, tools: Any, **kwargs: Any) -> DeferredErrorModel:
+ return self
+
+ def _generate(self, messages: Any, stop: Any = None, run_manager: Any = None, **kwargs: Any):
+ raise ValueError(self.error_message)
+
+
+def make_deferred_error_model(
+ error: BaseException, *, model_id: str | None = None
+) -> BaseChatModel:
+ return DeferredErrorModel(error_message=f"{type(error).__name__}: {error}", model_id=model_id)
+
+
+def make_model_or_defer(model_id: str, **kwargs: Any) -> BaseChatModel:
+ """Call ``make_model`` and wrap any setup error in a ``DeferredErrorModel``.
+
+ This lets graph factories pass a model placeholder into the agent graph on
+ startup without crashing the process. The error is raised later inside the
+ agent run when the model is first invoked, so the caller can recover
+ gracefully (e.g. via fallback middleware or by surfacing the error to the
+ user).
+ """
+ try:
+ return make_model(model_id, **kwargs)
+ except Exception as e: # noqa: BLE001
+ logger.warning("Deferring model setup failure for %s", model_id, exc_info=True)
+ return make_deferred_error_model(e, model_id=model_id)
diff --git a/agent/utils/sandbox.py b/agent/utils/sandbox.py
index 0dc3e41d..9614420a 100644
--- a/agent/utils/sandbox.py
+++ b/agent/utils/sandbox.py
@@ -1,10 +1,9 @@
import os
from collections.abc import Callable
from importlib import import_module
+from typing import Any
-from deepagents.backends.protocol import SandboxBackendProtocol
-
-SandboxFactory = Callable[[str | None], SandboxBackendProtocol]
+SandboxFactory = Callable[..., Any]
SANDBOX_FACTORIES: dict[str, tuple[str, str]] = {
"langsmith": ("agent.integrations.langsmith", "create_langsmith_sandbox"),
@@ -31,7 +30,7 @@ def create_sandbox(
sandbox_id: str | None = None,
*,
snapshot_id: str | None = None,
-) -> SandboxBackendProtocol:
+) -> Any:
"""Create or reconnect to a sandbox using the configured provider.
The provider is selected via the SANDBOX_TYPE environment variable.
diff --git a/docs/upstream-sync/triage.jsonl b/docs/upstream-sync/triage.jsonl
index 4d340b8a..03d6f5e8 100644
--- a/docs/upstream-sync/triage.jsonl
+++ b/docs/upstream-sync/triage.jsonl
@@ -1,4 +1,4 @@
-{"_meta": {"last_synced": "fd2541ce", "last_synced_date": "2026-07-08"}}
+{"_meta": {"last_synced": "fd2541ce", "last_synced_date": "2026-07-09"}}
{"sha": "0b76afdc", "pr": 1653, "subject": "reviews block agenda, sticky headers, diff scroll", "disposition": "landed", "reason": "", "branch": "cherry-pick-upstream", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
{"sha": "7530653b", "pr": 1655, "subject": "ResizeObserver settle for review scroll-to", "disposition": "landed", "reason": "", "branch": "cherry-pick-upstream", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
{"sha": "23bd4a63", "pr": 1660, "subject": "top padding to sticky review block header", "disposition": "landed", "reason": "", "branch": "cherry-pick-upstream", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
@@ -18,9 +18,9 @@
{"sha": "e5a29eca", "pr": 1613, "subject": "plan links in PR descriptions", "disposition": "wont-merge", "reason": "regression — dev has async plan-ref", "branch": "", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
{"sha": "e1d85526", "pr": 1645, "subject": "switch ui to pnpm", "disposition": "wont-merge", "reason": "tooling — fork keeps bun", "branch": "", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
{"sha": "f5670f24", "pr": 1639, "subject": "require bun for ui agent work", "disposition": "landed", "reason": "cherry-picked (-x) in this PR (#132)", "branch": "PR1", "local_sha": null, "updated": "2026-07-08T23:17:34Z"}
-{"sha": "209132d3", "pr": 1621, "subject": "durable interrupt dispatch + completion webhook", "disposition": "deferred", "reason": "investigate first — may be applied", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
-{"sha": "02bb4dfd", "pr": 1658, "subject": "don't attach loopback run-complete webhooks", "disposition": "deferred", "reason": "", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
-{"sha": "29015fad", "pr": 1614, "subject": "gate workflow pushes with approval", "disposition": "deferred", "reason": "", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-02T00:00:00Z"}
+{"sha": "209132d3", "pr": 1621, "subject": "durable interrupt dispatch + completion webhook", "disposition": "landed", "reason": "investigate first — may be applied", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-09T17:22:52Z"}
+{"sha": "02bb4dfd", "pr": 1658, "subject": "don't attach loopback run-complete webhooks", "disposition": "landed", "reason": "", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-09T17:22:53Z"}
+{"sha": "29015fad", "pr": 1614, "subject": "gate workflow pushes with approval", "disposition": "landed", "reason": "", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-09T17:22:53Z"}
{"sha": "546042a4", "pr": 1652, "subject": "add workflow approval UI", "disposition": "landed", "reason": "", "branch": "plan-approval", "local_sha": null, "updated": "2026-07-09T17:10:22Z"}
{"sha": "ae04b72b", "pr": 1635, "subject": "publish plans from sandbox files", "disposition": "landed", "reason": "ported (adapted) in #128 (save_plan reads sandbox file)", "branch": "plan-approval", "local_sha": null, "updated": "2026-07-08T22:58:21Z"}
{"sha": "c03a6be7", "pr": 1634, "subject": "keep plan guidance high-level", "disposition": "landed", "reason": "", "branch": "plan-approval", "local_sha": null, "updated": "2026-07-09T17:10:22Z"}
@@ -64,9 +64,9 @@
{"sha": "290d0fee", "pr": 1669, "subject": "chore(deps): update langgraph-cli[inmem] requirement (#1669)", "disposition": "wont-merge", "reason": "already in dev — langgraph-cli[inmem] at 0.4.30", "branch": "", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
{"sha": "c9f6dd86", "pr": 1694, "subject": "fix: fall back on model stream transport errors (#1694)", "disposition": "deferred", "reason": "adds httpx.TransportError to transient set; small conflict w/ dev's diverged Bedrock fallback", "branch": "model-fallback", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
{"sha": "c9a9a7cd", "pr": 1695, "subject": "fix: retry model fallback exhaustion (#1695)", "disposition": "deferred", "reason": "alternating retry+backoff rewrite; reconcile by hand w/ dev's Bedrock + sync wrap path", "branch": "model-fallback", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
-{"sha": "e5dbc788", "pr": 1696, "subject": "fix: Harden durable agent runs (#1696)", "disposition": "deferred", "reason": "large durable-run hardening; 3 new modules dev lacks; rewrites fork dispatch/completion", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
+{"sha": "e5dbc788", "pr": 1696, "subject": "fix: Harden durable agent runs (#1696)", "disposition": "landed", "reason": "large durable-run hardening; 3 new modules dev lacks; rewrites fork dispatch/completion", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-09T17:22:53Z"}
{"sha": "52fe2916", "pr": 1698, "subject": "feat: add PR review link route (#1698)", "disposition": "landed", "reason": "cherry-picked (-x) in #127 (PR review link route)", "branch": "reviewer-misc", "local_sha": null, "updated": "2026-07-08T22:58:21Z"}
-{"sha": "5f7f5fbd", "pr": 1697, "subject": "fix: Reduce graph import and loader startup latency (#1697)", "disposition": "deferred", "reason": "import-hygiene refactor; cross-cutting, references many deferred upstream-only modules", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
+{"sha": "5f7f5fbd", "pr": 1697, "subject": "fix: Reduce graph import and loader startup latency (#1697)", "disposition": "landed", "reason": "import-hygiene refactor; cross-cutting, references many deferred upstream-only modules", "branch": "durable-dispatch", "local_sha": null, "updated": "2026-07-09T17:22:53Z"}
{"sha": "216cf181", "pr": 1699, "subject": "fix: keep workflow HITL without token downscoping (#1699)", "disposition": "landed", "reason": "DIVERGES-FROM-UPSTREAM: fork deliberately does NOT adopt #1699's standing-token workflows:write broadening. Security review (#159) BLOCKed it — the standing ALWAYS-ON proxy token carrying workflows:write turns the HITL guard's git-push-parser gaps (obfuscated-expansion push, `gh api` REST contents PUT, cross-branch refspecs) into live unapproved-workflow-push exploits. Fork keeps BASE without workflows:write and restores the transient per-approval elevation (_run_with_workflow_token mints WORKFLOW_RUNTIME_PROXY_TOKEN_PERMISSIONS around the approved, guard-normalized fixed_command, then downscopes to RUNTIME then BASE): the token scope is the backstop the parser relies on, so a bypass hits GitHub 403. HITL diff-preview/approval-URL/Slack-card additions from #159 retained; token-model divergence only.", "branch": "plan-approval", "local_sha": null, "updated": "2026-07-09T20:00:00Z"}
{"sha": "67abf5b0", "pr": 1659, "subject": "fix: surface attributed PR creation failures (#1659)", "disposition": "deferred", "reason": "PR-attribution-failure guard (new mw, safe imports); heavy conflict on diverged open_pull_request.py", "branch": "pr-attribution", "local_sha": null, "updated": "2026-07-08T20:14:42Z"}
{"sha": "3dbc0282", "pr": 1676, "subject": "fix: preserve plan redirects after login (#1676)", "disposition": "landed", "reason": "FLAG-HUMAN: follow-on to landed #1668 refining sanitize_redirect_to (open-redirect auth surface); not a dup", "branch": "plan-approval", "local_sha": null, "updated": "2026-07-09T17:10:22Z"}
diff --git a/docs/upstream-sync/triage.md b/docs/upstream-sync/triage.md
index 337184e6..078f4c5f 100644
--- a/docs/upstream-sync/triage.md
+++ b/docs/upstream-sync/triage.md
@@ -6,7 +6,7 @@ Commits on `upstream/main` (langchain-ai/open-swe) not yet in `dev`, and the dec
Rows key on the **upstream SHA** (stable across local cherry-picks). Deferred rows are provisional
— re-inspect before picking. See the fork-maintenance runbook in `CLAUDE.md`.
-**Last synced `upstream/main`:** `fd2541ce` (2026-07-08)
+**Last synced `upstream/main`:** `fd2541ce` (2026-07-09)
| sha | pr | subject | decision | why | branch |
|---|---|---|---|---|---|
@@ -21,6 +21,9 @@ Rows key on the **upstream SHA** (stable across local cherry-picks). Deferred ro
| `6575c327` | #1654 | disable React StrictMode | Landed | kept fork's `PwaUpdateProvider` | cherry-pick-upstream |
| `00906401` | #1610 | editable plan mode | Landed | re-implemented in fork via #130 (editable plan mode) | |
| `f5670f24` | #1639 | require bun for ui agent work | Landed | cherry-picked (-x) in this PR (#132) | PR1 |
+| `209132d3` | #1621 | durable interrupt dispatch + completion webhook | Landed | investigate first — may be applied | durable-dispatch |
+| `02bb4dfd` | #1658 | don't attach loopback run-complete webhooks | Landed | | durable-dispatch |
+| `29015fad` | #1614 | gate workflow pushes with approval | Landed | | durable-dispatch |
| `546042a4` | #1652 | add workflow approval UI | Landed | | plan-approval |
| `ae04b72b` | #1635 | publish plans from sandbox files | Landed | ported (adapted) in #128 (save_plan reads sandbox file) | plan-approval |
| `c03a6be7` | #1634 | keep plan guidance high-level | Landed | | plan-approval |
@@ -43,7 +46,9 @@ Rows key on the **upstream SHA** (stable across local cherry-picks). Deferred ro
| `4f913198` | #1647 | widen split review diffs | Landed | already present in dev; empty pick confirmed in #127 | reviewer-misc |
| `20f63e8c` | #1646 | install missing deps before verification | Landed | already in dev via #81 (upstream-sync); ledger was stale | prompt-tweaks |
| `2f237b53` | #1626 | fall back to vision model for image threads | Landed | ported (adapted to Bedrock/Fireworks vision) in #128 | gateway-routing |
+| `e5dbc788` | #1696 | fix: Harden durable agent runs (#1696) | Landed | large durable-run hardening; 3 new modules dev lacks; rewrites fork dispatch/completion | durable-dispatch |
| `52fe2916` | #1698 | feat: add PR review link route (#1698) | Landed | cherry-picked (-x) in #127 (PR review link route) | reviewer-misc |
+| `5f7f5fbd` | #1697 | fix: Reduce graph import and loader startup latency (#1697) | Landed | import-hygiene refactor; cross-cutting, references many deferred upstream-only modules | durable-dispatch |
| `216cf181` | #1699 | fix: keep workflow HITL without token downscoping (#1699) | Landed | DIVERGES-FROM-UPSTREAM: fork deliberately does NOT adopt #1699's standing-token workflows:write broadening. Security review (#159) BLOCKed it — the standing ALWAYS-ON proxy token carrying workflows:write turns the HITL guard's git-push-parser gaps (obfuscated-expansion push, `gh api` REST contents PUT, cross-branch refspecs) into live unapproved-workflow-push exploits. Fork keeps BASE without workflows:write and restores the transient per-approval elevation (_run_with_workflow_token mints WORKFLOW_RUNTIME_PROXY_TOKEN_PERMISSIONS around the approved, guard-normalized fixed_command, then downscopes to RUNTIME then BASE): the token scope is the backstop the parser relies on, so a bypass hits GitHub 403. HITL diff-preview/approval-URL/Slack-card additions from #159 retained; token-model divergence only. | plan-approval |
| `3dbc0282` | #1676 | fix: preserve plan redirects after login (#1676) | Landed | FLAG-HUMAN: follow-on to landed #1668 refining sanitize_redirect_to (open-redirect auth surface); not a dup | plan-approval |
| `bb104d93` | #1679 | fix: submit plan comments with cmd enter (#1679) | Landed | applies clean but edits fork-diverged PlanReview.tsx (#130); needs UI/e2e validation — separate PR | plan-approval |
@@ -67,9 +72,6 @@ Rows key on the **upstream SHA** (stable across local cherry-picks). Deferred ro
| `73a9e8b5` | #1693 | chore(deps): bump the minor-and-patch group across 1 directory with 19 updates (#1693) | Won't merge | dev at-or-ahead on 17/19; group fights dev's pinned langsmith==0.9.7 (#115) and carries an upstream plan-route test | |
| `9cd7e464` | #1700 | Fix workflow approval visibility (#1700) | Won't merge | superseded — dev's list_workflow_approvals_for_thread already enforces owner-only 403 | |
| `fd2541ce` | #1705 | fix: drop orphaned function_call items with stale OpenAI reasoning (#1705) | Won't merge | N/A — edits sanitize_openai_responses.py which dev deleted in the Bedrock/Fireworks migration (#62) | |
-| `209132d3` | #1621 | durable interrupt dispatch + completion webhook | Deferred | investigate first — may be applied | durable-dispatch |
-| `02bb4dfd` | #1658 | don't attach loopback run-complete webhooks | Deferred | | durable-dispatch |
-| `29015fad` | #1614 | gate workflow pushes with approval | Deferred | | durable-dispatch |
| `4cd5fa5c` | #1629 | avoid recapping Slack replies | Deferred | | slack-tooling |
| `baf0c248` | #1617 | filter & grouping menu in threads sidebar | Deferred | ~998 LOC | own branch |
| `f29868ff` | #1615 | recover thread work as patch | Deferred | ~495 LOC | own branch |
@@ -85,8 +87,6 @@ Rows key on the **upstream SHA** (stable across local cherry-picks). Deferred ro
| `fbc6de85` | #1667 | chore(deps): bump fireworks-ai from 1.2.0a75 to 1.2.0a86 (#1667) | Deferred | conflicts on pick — dev diverged to fireworks-ai a85; a86 still wanted, needs manual bump + uv lock | deps |
| `c9f6dd86` | #1694 | fix: fall back on model stream transport errors (#1694) | Deferred | adds httpx.TransportError to transient set; small conflict w/ dev's diverged Bedrock fallback | model-fallback |
| `c9a9a7cd` | #1695 | fix: retry model fallback exhaustion (#1695) | Deferred | alternating retry+backoff rewrite; reconcile by hand w/ dev's Bedrock + sync wrap path | model-fallback |
-| `e5dbc788` | #1696 | fix: Harden durable agent runs (#1696) | Deferred | large durable-run hardening; 3 new modules dev lacks; rewrites fork dispatch/completion | durable-dispatch |
-| `5f7f5fbd` | #1697 | fix: Reduce graph import and loader startup latency (#1697) | Deferred | import-hygiene refactor; cross-cutting, references many deferred upstream-only modules | durable-dispatch |
| `67abf5b0` | #1659 | fix: surface attributed PR creation failures (#1659) | Deferred | PR-attribution-failure guard (new mw, safe imports); heavy conflict on diverged open_pull_request.py | pr-attribution |
| `c75cbb1f` | #1677 | feat: re-add Fable 5 with an admin toggle to disable it (#1677) | Deferred | FLAG-HUMAN: re-adds Fable 5 via anthropic: — contradicts dev's deliberate hide (#1483) + Bedrock migration (#62); wont-merge candidate | fable-admin-toggle |
| `304032fa` | #1680 | chore: clarify question answering prompt (#1680) | Deferred | reword Slack info-only answer guidance; conflicts w/ fork's customized Slack prompt | prompt-tweaks |
diff --git a/tests/e2e/patches.py b/tests/e2e/patches.py
index a9fb0115..edfb16ee 100644
--- a/tests/e2e/patches.py
+++ b/tests/e2e/patches.py
@@ -29,7 +29,6 @@ def apply() -> None:
import importlib
- from agent import server
from agent.utils import auth, authorship
from agent.utils import slack as slack_utils
@@ -52,7 +51,11 @@ def apply() -> None:
def _fake_make_model(model_id: str, **kwargs: object): # noqa: ARG001
return FakeScriptedChatModel(script=build_script())
- server.make_model = _fake_make_model
+ # ``make_model_or_defer`` (used by server/reviewer/analyzer) resolves
+ # ``make_model`` via the module global at call time, so patch it there.
+ import agent.utils.deferred_model as deferred_model
+
+ deferred_model.make_model = _fake_make_model
async def _dummy_install_token_with_expiry() -> tuple[str, str | None]:
return "dummy-installation-token", None
diff --git a/tests/test_agent_assembly_context.py b/tests/test_agent_assembly_context.py
index 43668bf5..b6e0f36c 100644
--- a/tests/test_agent_assembly_context.py
+++ b/tests/test_agent_assembly_context.py
@@ -65,7 +65,7 @@ async def _capture_create_deep_agent_kwargs() -> dict[str, object]:
),
patch("agent.server.load_profile", new_callable=AsyncMock, return_value=None),
patch("agent.server.fallback_model_id_for", return_value=None),
- patch("agent.server.make_model", side_effect=[MagicMock(), MagicMock()]),
+ patch("agent.utils.deferred_model.make_model", side_effect=[MagicMock(), MagicMock()]),
patch("agent.server.construct_system_prompt", return_value="prompt"),
patch("agent.server.create_deep_agent", side_effect=fake_create_deep_agent),
):
diff --git a/tests/test_agent_subagent_models.py b/tests/test_agent_subagent_models.py
index 7d317f88..a3f5573c 100644
--- a/tests/test_agent_subagent_models.py
+++ b/tests/test_agent_subagent_models.py
@@ -66,7 +66,9 @@ async def test_agent_uses_profile_subagent_model_override() -> None:
},
),
patch("agent.server.fallback_model_id_for", return_value=None),
- patch("agent.server.make_model", side_effect=[main_model, subagent_model]) as make_model,
+ patch(
+ "agent.utils.deferred_model.make_model", side_effect=[main_model, subagent_model]
+ ) as make_model,
patch("agent.server.construct_system_prompt", return_value="prompt"),
patch("agent.server.create_deep_agent", side_effect=fake_create_deep_agent),
):
@@ -142,7 +144,9 @@ async def test_agent_subagent_inherits_profile_model_override_without_explicit_p
},
),
patch("agent.server.fallback_model_id_for", return_value=None),
- patch("agent.server.make_model", side_effect=[main_model, subagent_model]) as make_model,
+ patch(
+ "agent.utils.deferred_model.make_model", side_effect=[main_model, subagent_model]
+ ) as make_model,
patch("agent.server.construct_system_prompt", return_value="prompt"),
patch("agent.server.create_deep_agent", side_effect=fake_create_deep_agent),
):
diff --git a/tests/test_completion_webhook.py b/tests/test_completion_webhook.py
index 0966ed41..867ee84a 100644
--- a/tests/test_completion_webhook.py
+++ b/tests/test_completion_webhook.py
@@ -9,20 +9,27 @@ from agent import completion
class _FakeThreads:
- def __init__(self, metadata: dict[str, Any]) -> None:
+ def __init__(self, metadata: dict[str, Any], *, fail_updates: int = 0) -> None:
self._metadata = metadata
self.updates: list[dict[str, Any]] = []
+ self._fail_updates = fail_updates
async def get(self, thread_id: str) -> dict[str, Any]:
- return {"thread_id": thread_id, "metadata": self._metadata}
+ # Reflect prior claims so a duplicate/retried delivery sees them, exactly
+ # as the platform persists thread metadata across webhook deliveries.
+ return {"thread_id": thread_id, "metadata": dict(self._metadata)}
async def update(self, *, thread_id: str, metadata: dict[str, Any]) -> None:
+ if self._fail_updates > 0:
+ self._fail_updates -= 1
+ raise RuntimeError("simulated metadata write failure")
self.updates.append(metadata)
+ self._metadata.update(metadata)
class _FakeClient:
- def __init__(self, metadata: dict[str, Any]) -> None:
- self.threads = _FakeThreads(metadata)
+ def __init__(self, metadata: dict[str, Any], *, fail_updates: int = 0) -> None:
+ self.threads = _FakeThreads(metadata, fail_updates=fail_updates)
def _slack_metadata() -> dict[str, Any]:
@@ -42,7 +49,9 @@ async def test_error_status_posts_slack_failure_reply(monkeypatch: pytest.Monkey
completion, "dashboard_thread_url", lambda thread_id: f"https://ui/{thread_id}"
)
- result = await completion.handle_run_completion({"thread_id": "t1", "status": "error"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ )
assert result["status"] == "ok"
reply.assert_awaited_once()
@@ -50,7 +59,9 @@ async def test_error_status_posts_slack_failure_reply(monkeypatch: pytest.Monkey
assert args[0] == "C1"
assert args[1] == "123.45"
assert "" in args[2]
- assert client.threads.updates == [{"failure_reply_posted": True}]
+ assert client.threads.updates == [
+ {"failure_reply_posted_run_id": "run-1", "failure_reply_posted_run_ids": ["run-1"]}
+ ]
@pytest.mark.asyncio
@@ -60,7 +71,9 @@ async def test_success_status_is_ignored(monkeypatch: pytest.MonkeyPatch) -> Non
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
- result = await completion.handle_run_completion({"thread_id": "t1", "status": "success"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "success"}
+ )
assert result["status"] == "ignored"
reply.assert_not_called()
@@ -69,19 +82,46 @@ async def test_success_status_is_ignored(monkeypatch: pytest.MonkeyPatch) -> Non
@pytest.mark.asyncio
async def test_idempotent_when_already_replied(monkeypatch: pytest.MonkeyPatch) -> None:
metadata = _slack_metadata()
- metadata["failure_reply_posted"] = True
+ metadata["failure_reply_posted_run_ids"] = ["run-1"]
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"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "timeout"}
+ )
assert result["status"] == "ignored"
reply.assert_not_called()
assert client.threads.updates == []
+@pytest.mark.asyncio
+async def test_later_failed_run_posts_even_if_prior_run_replied(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ metadata = _slack_metadata()
+ metadata["failure_reply_posted_run_ids"] = ["run-1"]
+ 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", "run_id": "run-2", "status": "timeout"}
+ )
+
+ assert result["status"] == "ok"
+ reply.assert_awaited_once()
+ assert client.threads.updates == [
+ {
+ "failure_reply_posted_run_id": "run-2",
+ "failure_reply_posted_run_ids": ["run-1", "run-2"],
+ }
+ ]
+
+
@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"}}})
@@ -89,7 +129,9 @@ async def test_linear_source_comments_on_issue(monkeypatch: pytest.MonkeyPatch)
comment = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "comment_on_linear_issue", comment)
- result = await completion.handle_run_completion({"thread_id": "t1", "status": "timeout"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "timeout"}
+ )
assert result["status"] == "ok"
comment.assert_awaited_once()
@@ -98,46 +140,169 @@ async def test_linear_source_comments_on_issue(monkeypatch: pytest.MonkeyPatch)
@pytest.mark.asyncio
async def test_missing_thread_id_is_ignored() -> None:
- result = await completion.handle_run_completion({"status": "error"})
+ result = await completion.handle_run_completion({"run_id": "run-1", "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.
+async def test_no_run_id_and_no_marker_posts_without_dedupe(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ # A payload carrying nothing to distinguish one run from another posts
+ # un-deduped (no sticky flag written) rather than being suppressed. It
+ # prefers a rare duplicate over permanent silence (SR160-03).
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"
+
+ assert result["status"] == "ok"
+ reply.assert_awaited_once()
+ assert client.threads.updates == []
+
+
+@pytest.mark.asyncio
+async def test_legacy_sticky_flag_does_not_suppress_new_failure(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ # An old per-thread sticky boolean left by a prior code version must not
+ # permanently silence a later failed run (SR160-03).
+ 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", "run_id": "run-9", "status": "error"}
+ )
+
+ assert result["status"] == "ok"
+ reply.assert_awaited_once()
+
+
+@pytest.mark.asyncio
+async def test_no_run_id_dedupes_on_updated_at_marker(monkeypatch: pytest.MonkeyPatch) -> None:
+ # Without a run_id the same run's retry (identical updated_at) dedupes...
+ client = _FakeClient(_slack_metadata())
+ monkeypatch.setattr(completion, "langgraph_client", lambda: client)
+ reply = AsyncMock(return_value=True)
+ monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
+
+ first = await completion.handle_run_completion(
+ {"thread_id": "t1", "status": "error", "updated_at": "2026-07-09T00:00:00Z"}
+ )
+ retry = await completion.handle_run_completion(
+ {"thread_id": "t1", "status": "error", "updated_at": "2026-07-09T00:00:00Z"}
+ )
+
+ assert first["status"] == "ok"
+ assert retry["status"] == "ignored"
+ reply.assert_awaited_once()
+
+
+@pytest.mark.asyncio
+async def test_no_run_id_different_runs_each_reply(monkeypatch: pytest.MonkeyPatch) -> None:
+ # ...while two genuinely different failed runs (distinct updated_at) on one
+ # thread each get their reply — the crux of SR160-03.
+ client = _FakeClient(_slack_metadata())
+ monkeypatch.setattr(completion, "langgraph_client", lambda: client)
+ reply = AsyncMock(return_value=True)
+ monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
+
+ first = await completion.handle_run_completion(
+ {"thread_id": "t1", "status": "error", "updated_at": "2026-07-09T00:00:00Z"}
+ )
+ second = await completion.handle_run_completion(
+ {"thread_id": "t1", "status": "timeout", "updated_at": "2026-07-09T01:00:00Z"}
+ )
+
+ assert first["status"] == "ok"
+ assert second["status"] == "ok"
+ assert reply.await_count == 2
+
+
+@pytest.mark.asyncio
+async def test_claim_written_before_post(monkeypatch: pytest.MonkeyPatch) -> None:
+ # Claim-then-post: the dedup key is recorded before the network reply fires.
+ client = _FakeClient(_slack_metadata())
+ monkeypatch.setattr(completion, "langgraph_client", lambda: client)
+ order: list[str] = []
+
+ async def _update(*, thread_id: str, metadata: dict[str, Any]) -> None:
+ order.append("claim")
+ client.threads.updates.append(metadata)
+ client.threads._metadata.update(metadata)
+
+ async def _reply(*args: Any, **kwargs: Any) -> bool:
+ order.append("post")
+ return True
+
+ monkeypatch.setattr(client.threads, "update", _update)
+ monkeypatch.setattr(completion, "post_slack_thread_reply", _reply)
+
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ )
+
+ assert result["status"] == "ok"
+ assert order == ["claim", "post"]
+
+
+@pytest.mark.asyncio
+async def test_claim_write_failure_does_not_post_and_retry_posts_once(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ # SR160-02: if the claim write throws, we do NOT post and report an error so
+ # the retry re-attempts — the retry then claims and posts exactly once,
+ # rather than the old flow double-posting on retry.
+ client = _FakeClient(_slack_metadata(), fail_updates=1)
+ monkeypatch.setattr(completion, "langgraph_client", lambda: client)
+ reply = AsyncMock(return_value=True)
+ monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
+
+ first = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ )
+ assert first["status"] == "error"
reply.assert_not_called()
+ retry = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ )
+ assert retry["status"] == "ok"
+ reply.assert_awaited_once()
+
+
+@pytest.mark.asyncio
+async def test_duplicate_delivery_posts_exactly_once(monkeypatch: pytest.MonkeyPatch) -> None:
+ # A retried/duplicate delivery for one run_id posts exactly once because the
+ # first delivery's claim is persisted before the second reads metadata.
+ client = _FakeClient(_slack_metadata())
+ monkeypatch.setattr(completion, "langgraph_client", lambda: client)
+ reply = AsyncMock(return_value=True)
+ monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
+
+ payload = {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ first = await completion.handle_run_completion(dict(payload))
+ second = await completion.handle_run_completion(dict(payload))
+
+ assert first["status"] == "ok"
+ assert second["status"] == "ignored"
+ reply.assert_awaited_once()
+
@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"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "error"}
+ )
assert result["status"] == "ignored"
assert client.threads.updates == []
@@ -152,7 +317,9 @@ async def test_interrupted_status_is_ignored(monkeypatch: pytest.MonkeyPatch) ->
reply = AsyncMock(return_value=True)
monkeypatch.setattr(completion, "post_slack_thread_reply", reply)
- result = await completion.handle_run_completion({"thread_id": "t1", "status": "interrupted"})
+ result = await completion.handle_run_completion(
+ {"thread_id": "t1", "run_id": "run-1", "status": "interrupted"}
+ )
assert result["status"] == "ignored"
reply.assert_not_called()
diff --git a/tests/test_corridor_mcp.py b/tests/test_corridor_mcp.py
index bc248e9e..57cdfad3 100644
--- a/tests/test_corridor_mcp.py
+++ b/tests/test_corridor_mcp.py
@@ -154,7 +154,7 @@ async def test_get_agent_passes_corridor_prompt_state() -> None:
return_value=(("openai:gpt-5.5", "medium"), ("openai:gpt-5.5", "low")),
),
patch.object(server, "fallback_model_id_for", return_value=None),
- patch.object(server, "make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch.object(
server, "_load_observability_tools", new_callable=AsyncMock, return_value=[]
),
diff --git a/tests/test_dispatch.py b/tests/test_dispatch.py
new file mode 100644
index 00000000..cdf0e65f
--- /dev/null
+++ b/tests/test_dispatch.py
@@ -0,0 +1,124 @@
+from __future__ import annotations
+
+import importlib
+from typing import Any
+
+import pytest
+
+dispatch = importlib.import_module("agent.dispatch")
+
+_ABSOLUTE = "https://open-swe-v3-abc.us.langgraph.app/webhooks/run-complete"
+
+
+def test_is_loopback_webhook_relative() -> None:
+ assert dispatch._is_loopback_webhook("/webhooks/run-complete") is True
+
+
+def test_is_loopback_webhook_localhost() -> None:
+ assert dispatch._is_loopback_webhook("http://localhost:2024/webhooks/run-complete") is True
+ assert dispatch._is_loopback_webhook("http://127.0.0.1:8000/webhooks/run-complete") is True
+
+
+def test_is_loopback_webhook_absolute() -> None:
+ assert dispatch._is_loopback_webhook(_ABSOLUTE) is False
+
+
+def test_resolve_no_secret_attaches_nothing() -> None:
+ assert dispatch._resolve_completion_webhook_url(_ABSOLUTE, None) is None
+ assert dispatch._resolve_completion_webhook_url(_ABSOLUTE, "") is None
+
+
+def test_resolve_relative_url_degrades_to_none() -> None:
+ # Secret set but a loopback URL would 422 every run — attach nothing instead.
+ assert dispatch._resolve_completion_webhook_url("/webhooks/run-complete", "s3cret") is None
+
+
+def test_resolve_localhost_url_degrades_to_none() -> None:
+ assert dispatch._resolve_completion_webhook_url("http://localhost/x", "s3cret") is None
+
+
+def test_resolve_absolute_url_appends_token() -> None:
+ assert (
+ dispatch._resolve_completion_webhook_url(_ABSOLUTE, "s3cret") == f"{_ABSOLUTE}?token=s3cret"
+ )
+
+
+def test_resolve_absolute_url_with_existing_query_left_as_is() -> None:
+ url = f"{_ABSOLUTE}?token=preset"
+ assert dispatch._resolve_completion_webhook_url(url, "s3cret") == url
+
+
+class _FakeRuns:
+ def __init__(self) -> None:
+ self.created: list[dict[str, Any]] = []
+
+ async def create(self, thread_id: str, assistant_id: str, **kwargs: Any) -> dict[str, str]:
+ self.created.append({"thread_id": thread_id, "assistant_id": assistant_id, **kwargs})
+ return {"run_id": "run-1"}
+
+
+class _FakeClient:
+ def __init__(self) -> None:
+ self.runs = _FakeRuns()
+
+
+@pytest.mark.asyncio
+async def test_create_durable_run_applies_defaults(monkeypatch: pytest.MonkeyPatch) -> None:
+ client = _FakeClient()
+ monkeypatch.setattr(dispatch, "COMPLETION_WEBHOOK_URL", "https://app/webhooks/run-complete")
+
+ run = await dispatch.create_durable_run(
+ "thread-1",
+ "agent",
+ input={"messages": [{"role": "user", "content": "hi"}]},
+ source="test",
+ config={"configurable": {"thread_id": "thread-1"}, "metadata": {"kind": "test"}},
+ client=client,
+ )
+
+ assert run == {"run_id": "run-1"}
+ created = client.runs.created[0]
+ assert created["durability"] == "sync"
+ assert created["multitask_strategy"] == "interrupt"
+ assert created["if_not_exists"] == "create"
+ assert created["webhook"] == "https://app/webhooks/run-complete"
+ assert created["config"]["metadata"] == {"kind": "test"}
+ assert created["config"]["configurable"]["thread_id"] == "thread-1"
+ assert isinstance(created["config"]["configurable"]["prepare_run_id"], str)
+
+
+@pytest.mark.asyncio
+async def test_create_durable_run_no_webhook_when_none(monkeypatch: pytest.MonkeyPatch) -> None:
+ client = _FakeClient()
+ monkeypatch.setattr(dispatch, "COMPLETION_WEBHOOK_URL", None)
+
+ await dispatch.create_durable_run(
+ "thread-1",
+ "agent",
+ input={"messages": []},
+ source="test",
+ client=client,
+ )
+
+ created = client.runs.created[0]
+ assert "webhook" not in created
+
+
+@pytest.mark.asyncio
+async def test_dispatch_agent_run_preserves_multitask_strategy(
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ client = _FakeClient()
+ monkeypatch.setattr(dispatch, "COMPLETION_WEBHOOK_URL", None)
+
+ await dispatch.dispatch_agent_run(
+ "thread-1",
+ "hi",
+ {"key": "val"},
+ source="test",
+ client=client,
+ multitask_strategy="reject",
+ )
+
+ created = client.runs.created[0]
+ assert created["multitask_strategy"] == "reject"
diff --git a/tests/test_proxy_auth.py b/tests/test_proxy_auth.py
index 15ff5766..69ced5e6 100644
--- a/tests/test_proxy_auth.py
+++ b/tests/test_proxy_auth.py
@@ -333,7 +333,7 @@ class TestRefreshProxyOnSandboxReuse:
new_callable=AsyncMock,
return_value=mock_sandbox,
),
- patch("agent.server.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.server.construct_system_prompt", return_value="prompt"),
patch("agent.server.create_deep_agent", return_value=_DummyAgent()),
patch.dict(
@@ -386,7 +386,7 @@ class TestRefreshProxyOnSandboxReuse:
new_callable=AsyncMock,
return_value="/workspace",
),
- patch("agent.server.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.server.construct_system_prompt", return_value="prompt"),
patch("agent.server.create_deep_agent", return_value=_DummyAgent()),
patch.dict("agent.server.SANDBOX_BACKENDS", {}, clear=True),
diff --git a/tests/test_reviewer.py b/tests/test_reviewer.py
index b332502f..7f4c9264 100644
--- a/tests/test_reviewer.py
+++ b/tests/test_reviewer.py
@@ -204,7 +204,7 @@ async def test_reviewer_resolves_app_installation_token_at_run_start() -> None:
new_callable=AsyncMock,
return_value="/workspace",
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", return_value=dummy_agent) as create_agent,
):
await reviewer.get_reviewer_agent(config)
@@ -260,7 +260,7 @@ async def test_reviewer_reuses_app_token_for_sandbox_proxy() -> None:
new_callable=AsyncMock,
),
patch("agent.reviewer.fetch_agents_md", new_callable=AsyncMock, return_value=None),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", return_value=_DummyAgent()),
):
await reviewer.get_reviewer_agent(config)
@@ -335,7 +335,7 @@ async def test_reviewer_applies_eval_model_and_effort_overrides() -> None:
new_callable=AsyncMock,
return_value="/workspace",
),
- patch("agent.reviewer.make_model", return_value=MagicMock()) as make_model,
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()) as make_model,
patch("agent.reviewer.create_deep_agent", return_value=dummy_agent),
patch(
"agent.reviewer.fetch_agents_md",
@@ -383,7 +383,7 @@ async def test_reviewer_subagent_inherits_eval_model_without_explicit_override()
new_callable=AsyncMock,
return_value="/workspace",
),
- patch("agent.reviewer.make_model", return_value=MagicMock()) as make_model,
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()) as make_model,
patch("agent.reviewer.create_deep_agent", return_value=dummy_agent),
patch(
"agent.reviewer.fetch_agents_md",
@@ -441,7 +441,7 @@ async def test_reviewer_injects_repo_style_during_eval() -> None:
new_callable=AsyncMock,
return_value="Flag table rerender regressions.",
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
patch(
"agent.reviewer.fetch_agents_md",
@@ -507,7 +507,7 @@ async def test_reviewer_inlines_org_guidelines_into_system_prompt() -> None:
new_callable=AsyncMock,
return_value=None,
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -572,7 +572,7 @@ async def test_reviewer_inlines_agents_md_into_system_prompt() -> None:
new_callable=AsyncMock,
return_value="Always use the design system IconButton.",
) as mock_fetch_agents_md,
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -624,7 +624,7 @@ async def test_reviewer_inlines_claude_md_when_agents_md_absent() -> None:
new_callable=AsyncMock,
return_value="# CLAUDE.md\nUse semantic tokens only.",
) as mock_fetch_agents_md,
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -989,7 +989,7 @@ async def test_reviewer_injects_pr_review_threads_into_first_review_context() ->
new_callable=AsyncMock,
return_value=fake_threads,
) as mock_fetch_threads,
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1065,7 +1065,7 @@ async def test_reviewer_injects_pr_review_threads_into_re_review_context() -> No
new_callable=AsyncMock,
return_value=[],
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1123,7 +1123,7 @@ async def test_reviewer_omits_threads_block_when_fetch_returns_empty() -> None:
new_callable=AsyncMock,
return_value=[],
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1182,7 +1182,7 @@ async def test_reviewer_continues_when_thread_fetch_raises() -> None:
new_callable=AsyncMock,
side_effect=RuntimeError("network down"),
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1254,7 +1254,7 @@ async def test_reviewer_populates_diff_line_set_from_github_api() -> None:
new_callable=AsyncMock,
return_value=pr_diff,
) as mock_fetch_diff,
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1319,7 +1319,7 @@ async def test_reviewer_leaves_validation_disabled_when_diff_fetch_fails() -> No
new_callable=AsyncMock,
return_value=None,
),
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)
@@ -1385,7 +1385,7 @@ async def test_reviewer_injects_pr_title_and_body_into_context() -> None:
new_callable=AsyncMock,
return_value=("Add retry logic for uploads", "Retries flaky uploads up to 3 times."),
) as mock_fetch_metadata,
- patch("agent.reviewer.make_model", return_value=MagicMock()),
+ patch("agent.utils.deferred_model.make_model", return_value=MagicMock()),
patch("agent.reviewer.create_deep_agent", side_effect=fake_create_deep_agent),
):
await reviewer.get_reviewer_agent(config)