mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-10-02 12:03:15 +00:00
feat: port durable dispatch hardening and startup latency improvements (#160)
Some checks are pending
CI / Lint (push) Waiting to run
CI / Format check (push) Waiting to run
CI / Unit tests (push) Waiting to run
CI / Playwright E2E (push) Waiting to run
CI / Docker build smoke (push) Waiting to run
CI / Triage ledger up to date (push) Waiting to run
CI / ui bun.lock in sync (push) Waiting to run
Some checks are pending
CI / Lint (push) Waiting to run
CI / Format check (push) Waiting to run
CI / Unit tests (push) Waiting to run
CI / Playwright E2E (push) Waiting to run
CI / Docker build smoke (push) Waiting to run
CI / Triage ledger up to date (push) Waiting to run
CI / ui bun.lock in sync (push) Waiting to run
* 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 <adam@seahavenind.com>
This commit is contained in:
parent
f87847baa4
commit
0546085672
24 changed files with 998 additions and 181 deletions
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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}")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
82
agent/middleware/task_retry.py
Normal file
82
agent/middleware/task_retry.py
Normal file
|
|
@ -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,
|
||||
)
|
||||
64
agent/middleware/timeout_wrapup.py
Normal file
64
agent/middleware/timeout_wrapup.py
Normal file
|
|
@ -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 = """
|
||||
<time_limit_warning>
|
||||
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.
|
||||
</time_limit_warning>
|
||||
"""
|
||||
|
||||
|
||||
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))
|
||||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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]:
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
57
agent/utils/deferred_model.py
Normal file
57
agent/utils/deferred_model.py
Normal file
|
|
@ -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)
|
||||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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"}
|
||||
|
|
|
|||
|
|
@ -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 |
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
):
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
):
|
||||
|
|
|
|||
|
|
@ -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 "<https://ui/t1|Open SWE Web>" 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()
|
||||
|
|
|
|||
|
|
@ -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=[]
|
||||
),
|
||||
|
|
|
|||
124
tests/test_dispatch.py
Normal file
124
tests/test_dispatch.py
Normal file
|
|
@ -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"
|
||||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue