From 0546085672c26877747825b5072f47f759e5e999 Mon Sep 17 00:00:00 2001 From: "seahaven-openswe[bot]" <296972425+seahaven-openswe[bot]@users.noreply.github.com> Date: Thu, 9 Jul 2026 17:11:25 -0400 Subject: [PATCH] feat: port durable dispatch hardening and startup latency improvements (#160) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: port plan-review & workflow-approval UX (#135) Port six upstream commits onto dev: - c03a6be7 (already ported): keep plan guidance high-level - 546042a4: add workflow approval UI with diff preview, approval URLs, web review links, and polling for approval status during active runs - 216cf181: remove workflow token elevation; approved pushes pass through directly without proxy token rewriting - 3dbc0282: preserve plan redirects after login by accepting relative same-origin redirect_to values and rejecting blocked paths - bb104d93: submit plan comments with cmd+enter - 90cb6caa: terse Slack replies, shared content via save_plan outside plan mode (PLAN_STATUS_SHARED), reject shared-content mutations Refs: #135 * feat: port durable dispatch hardening and startup latency improvements Port five upstream PRs onto dev: - #1621 / #1658: durable dispatch with loopback webhook defense, create_durable_run helper, _config_with_prepare_run_id, degradation to None for relative/loopback completion webhook URLs - #1696: run-level completion webhook deduplication (replace claim-then-post with post-then-flag per run_id), DeferredErrorModel for graph-factory resilience, ToolRetryMiddleware for task subagents, TimeoutWrapupMiddleware for all three graphs - #1697: lazy-load __init__.py for agent.middleware, agent.tools, agent.dashboard (PEP 562); defer heavy imports (exa_py in web_search, agent.webapp in request_pr_review, deepagents in sandbox.py); add ttl_cache.py with stale-while-revalidate for tool loaders Refs: #137 * fix: restore login page render and clear CI lint/format The plan-review port removed the authRedirectUrl import from login.tsx but left its call site, crashing the login page at runtime (blank page, no 'Sign in to open-swe'). Pass the relative path straight to loginUrl, matching the plan route and the backend relative-redirect handling. Also drop an unused os import in the guard test and reformat workflow_push_guard.py to satisfy ruff. * fix: restore RepairOrphaned middleware export and repoint model fake to deferred_model boundary * fix: restore RepairOrphanedToolCallsMiddleware, fix E2E model-fake patch, drop dead ttl_cache - Re-add RepairOrphanedToolCallsMiddleware to the lazy middleware __init__ (_MIDDLEWARE_MODULES, __all__, TYPE_CHECKING) so agent.reviewer can import it. - Reroute E2E model patching to deferred_model.make_model so make_model_or_defer (used by all three graph factories) returns the scripted fake instead of building a real model with fake credentials. - Drop unused agent/utils/ttl_cache.py — no agent module imports it. - Fix import ordering in agent/reviewer.py and agent/analyzer.py (ruff I001). - Format tests/test_dispatch.py. * fix: claim-then-post run-level failure dedup; stop permanent suppression --------- Co-authored-by: amoussa1229 <166072409+amoussa1229@users.noreply.github.com> Co-authored-by: Adam Moussa --- agent/analyzer.py | 12 +- agent/completion.py | 89 ++++++++-- agent/dashboard/__init__.py | 19 ++- agent/dispatch.py | 140 ++++++++++++---- agent/middleware/__init__.py | 95 ++++++++--- agent/middleware/task_retry.py | 82 ++++++++++ agent/middleware/timeout_wrapup.py | 64 ++++++++ agent/reviewer.py | 11 +- agent/server.py | 24 ++- agent/tools/__init__.py | 117 ++++++++++---- agent/tools/request_pr_review.py | 24 ++- agent/tools/web_search.py | 4 +- agent/utils/deferred_model.py | 57 +++++++ agent/utils/sandbox.py | 7 +- docs/upstream-sync/triage.jsonl | 12 +- docs/upstream-sync/triage.md | 12 +- tests/e2e/patches.py | 7 +- tests/test_agent_assembly_context.py | 2 +- tests/test_agent_subagent_models.py | 8 +- tests/test_completion_webhook.py | 233 +++++++++++++++++++++++---- tests/test_corridor_mcp.py | 2 +- tests/test_dispatch.py | 124 ++++++++++++++ tests/test_proxy_auth.py | 4 +- tests/test_reviewer.py | 30 ++-- 24 files changed, 998 insertions(+), 181 deletions(-) create mode 100644 agent/middleware/task_retry.py create mode 100644 agent/middleware/timeout_wrapup.py create mode 100644 agent/utils/deferred_model.py create mode 100644 tests/test_dispatch.py 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)