open-swe/agent/dashboard/plan_api.py
Johannes du Plessis 209132d355
refactor: durable interrupt dispatch + completion webhook (#1621)
* wip(rebuild): core reliability spine

- remove PR-babysitting (ci_autofix + ci_monitor graph + webhook wiring)
- dispatch core: agent/dispatch.py with multitask_strategy=interrupt +
  durability=sync + completion webhook; reroute all webhook + plan triggers;
  drop the racy in-process lock + is_thread_active busy-check
- completion webhook: agent/completion.py + /webhooks/run-complete loopback
  route for failure/timeout replies (idempotent)

Co-authored-by: open-swe[bot]

* feat(rebuild): async tools, reconcile, shared http timeouts, assembly tuning

Parallel batch on top of the reliability spine:
- async-ify all 24 tools (drop asyncio.run; requests->httpx); re-implement the
  http_request/fetch_url SSRF + DNS-rebinding defense httpx-natively and harden
  the IP check to 'not is_global' (+ IPv4-mapped unwrap)
- reconcile.py: stale pending-run sweep (threads.search -> per-thread runs.list
  -> cancel_many), wired into the scheduler graph via task='reconcile'
- shared DEFAULT_HTTP_TIMEOUT (agent/utils/http.py) on every bare
  httpx.AsyncClient() across utils/dashboard/webapp/middleware
- run budget: MODEL_CALL_RECURSION_LIMIT 5000->250
- fix stale OpenAI->Anthropic fallback id (claude-opus-4-5 -> 4-8)
- drop redundant custom repair middleware (deepagents auto-adds PatchToolCalls)
- confirm tool-result eviction + summarization auto-wired via backend
- slim system prompt ~8% (full harness-profile rewrite deferred)

Co-authored-by: open-swe[bot]

* feat(rebuild): harness-profile prompt + split webhooks out of webapp

- prompt.py: own the system prompt via a registered harness profile
  (OPEN_SWE_SHARED_BASE, kept neutral so the read-only reviewer/analyzer that
  share it stay safe), registered across all 4 providers; per-thread values
  stay in construct_system_prompt. Assembled main-agent prompt ~6.8k -> ~3.1k
  tokens (~55% smaller); de-duped PR/commit/suite/force-push guidance; dropped
  ALL-CAPS markers.
- webapp.py 3325 -> 1890 LOC: moved 14 per-source handlers into
  agent/webhooks/{linear,slack,github}.py; webapp re-exports them for the
  routes + tests; moved handlers reach shared helpers via the webapp namespace
  to preserve the test suite's monkeypatch targets.

Full suite: 1168 passing, lint clean.

Co-authored-by: open-swe[bot]

* Restore MODEL_CALL_RECURSION_LIMIT to 5000 for long-running tasks

Reverts the 250 cap from the run-budget change — long-running tasks legitimately
need many model calls. The notify_step_limit_reached safety net still fires if a
run does hit the cap, so runs end with a signal either way.

Co-authored-by: open-swe[bot]

* fix: address PR review (auth, SSRF, interrupted status, redirect headers)

- completion.py: drop `interrupted` from failure statuses — with
  multitask_strategy=interrupt a follow-up ends the prior run as interrupted,
  which is healthy, not a failure to report. [open-swe]
- /webhooks/run-complete: shared-secret auth — dispatch appends ?token= when
  RUN_COMPLETE_WEBHOOK_SECRET is set; route verifies via hmac.compare_digest.
  [corridor-security]
- SSRF: extract the URL validator to agent/utils/url_safety.py and apply it
  before server-side image fetches in multimodal.fetch_image_block.
  [corridor-security]
- http_request: preserve caller headers/extensions across redirect hops instead
  of dropping them on the first hop. [open-swe]

Co-authored-by: open-swe[bot]

* chore: remove REBUILD_PLAN.md (planning doc, not needed in the repo)

Co-authored-by: open-swe[bot]

* fix: fail closed on run-complete webhook auth when secret unset

Corridor follow-up: verify_run_complete_token returns False (not True) when
RUN_COMPLETE_WEBHOOK_SECRET is unset, so the public route is never
unauthenticated. Logs a startup warning when the secret is absent, and dispatch
skips registering the webhook when there's no secret (no rejected callbacks).

Co-authored-by: open-swe[bot]

---------

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
2026-06-26 13:48:38 -07:00

263 lines
10 KiB
Python

"""REST API for the plan-review page: read the plan, comment, approve, or request
changes — all plain HTTP, no CRDT/WebSocket.
Reviewers leave whole-document comments via this API; they're stored server-side
and listed for everyone who can read the thread. On approve/reject the comments
are read back here, formatted, and handed to the agent as the instruction for the
follow-up run. The agent never sees comments during review — only this aggregated
feedback at the decision point.
Permissions: any authenticated org member can read a surfaced thread, comment, and
request changes (reject); only the thread owner can approve. A comment can be
deleted by its author or the thread owner.
"""
from __future__ import annotations
import logging
from typing import Any
from fastapi import APIRouter, Depends, HTTPException
from langgraph_sdk import get_client
from pydantic import BaseModel
from ..dispatch import dispatch_agent_run
from .oauth import require_same_origin_for_mutations, require_session
from .plan_store import (
PLAN_STATUS_APPROVED,
PLAN_STATUS_CANCELLED,
PLAN_STATUS_READY,
PLAN_STATUS_REVISING,
add_plan_comment,
delete_plan_comment,
get_plan_content,
list_plan_comments,
save_plan_content,
set_plan_status,
write_plan_to_sandbox,
)
from .thread_api import (
_repo_config_from_metadata,
_thread_is_readable,
_thread_source,
_user_owns_thread,
)
logger = logging.getLogger(__name__)
plan_router = APIRouter(
prefix="/dashboard/api/plan",
tags=["plan"],
dependencies=[Depends(require_same_origin_for_mutations)],
)
_SESSION_DEP = Depends(require_session)
class CommentBody(BaseModel):
body: str
class PlanUpdate(BaseModel):
markdown: str
async def _thread_metadata(thread_id: str) -> dict[str, Any]:
client = get_client()
try:
thread = await client.threads.get(thread_id)
except Exception as exc: # noqa: BLE001
raise HTTPException(404, "thread not found") from exc
metadata = (
thread.get("metadata") if isinstance(thread, dict) else getattr(thread, "metadata", None)
)
return metadata if isinstance(metadata, dict) else {}
@plan_router.get("/{thread_id}")
async def get_plan(thread_id: str, session: dict[str, Any] = _SESSION_DEP) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _thread_is_readable(metadata):
raise HTTPException(404, "thread not found")
login = session["sub"]
email = session.get("email")
content = await get_plan_content(thread_id) or {}
return {
"threadId": thread_id,
"status": content.get("status") or metadata.get("plan_status") or "planning",
"markdown": content.get("markdown", ""),
"isOwner": _user_owns_thread(metadata, login, email),
"user": {
"id": login,
"login": login,
"email": email,
"name": session.get("name") or login,
},
}
@plan_router.put("/{thread_id}")
async def update_plan(
thread_id: str, body: PlanUpdate, session: dict[str, Any] = _SESSION_DEP
) -> dict[str, Any]:
"""Owner-only manual edit of the plan markdown.
Re-publishes the edited plan as ``ready`` (and mirrors it into the sandbox
``plan.md``) while preserving reviewer comments, so the owner can refine the
plan before approving it."""
metadata = await _thread_metadata(thread_id)
if not _user_owns_thread(metadata, session["sub"], session.get("email")):
raise HTTPException(403, "only the plan owner can edit the plan")
markdown = body.markdown.strip()
if not markdown:
raise HTTPException(422, "plan markdown cannot be empty")
content = await get_plan_content(thread_id) or {}
status = content.get("status") or metadata.get("plan_status") or "planning"
if status in (PLAN_STATUS_APPROVED, PLAN_STATUS_CANCELLED):
raise HTTPException(409, f"cannot edit a {status} plan")
await save_plan_content(
thread_id, markdown=markdown, status=PLAN_STATUS_READY, clear_comments=False
)
await write_plan_to_sandbox(thread_id, markdown)
return {"status": PLAN_STATUS_READY, "markdown": markdown}
@plan_router.get("/{thread_id}/comments")
async def get_plan_comments(
thread_id: str, session: dict[str, Any] = _SESSION_DEP
) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _thread_is_readable(metadata):
raise HTTPException(404, "thread not found")
return {"comments": await list_plan_comments(thread_id)}
@plan_router.post("/{thread_id}/comments")
async def post_plan_comment(
thread_id: str, body: CommentBody, session: dict[str, Any] = _SESSION_DEP
) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _thread_is_readable(metadata):
raise HTTPException(404, "thread not found")
text = body.body.strip()
if not text:
raise HTTPException(422, "comment body cannot be empty")
login = session["sub"]
return await add_plan_comment(
thread_id, author=session.get("name") or login, author_login=login, body=text
)
@plan_router.delete("/{thread_id}/comments/{comment_id}")
async def remove_plan_comment(
thread_id: str, comment_id: str, session: dict[str, Any] = _SESSION_DEP
) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _thread_is_readable(metadata):
raise HTTPException(404, "thread not found")
comments = await list_plan_comments(thread_id)
target = next((c for c in comments if c.get("id") == comment_id), None)
if target is None:
raise HTTPException(404, "comment not found")
login = session["sub"]
is_owner = _user_owns_thread(metadata, login, session.get("email"))
if target.get("author_login") != login and not is_owner:
raise HTTPException(403, "only the author or the plan owner can delete a comment")
await delete_plan_comment(thread_id, comment_id)
return {"ok": True}
@plan_router.post("/{thread_id}/approve")
async def approve_plan(thread_id: str, session: dict[str, Any] = _SESSION_DEP) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _user_owns_thread(metadata, session["sub"], session.get("email")):
raise HTTPException(403, "only the plan owner can approve")
# Read the published plan + comments BEFORE mutating state: a store failure
# here aborts the decision (500) rather than dispatching without them. The
# published markdown may have been edited by the reviewer, so it is the
# source of truth handed to the agent (not its own stale history) — read it
# strictly so a transient failure can't silently drop the edit.
content = await get_plan_content(thread_id, raise_on_error=True) or {}
plan_markdown = str(content.get("markdown", "")).strip()
feedback = _format_comments(await list_plan_comments(thread_id, raise_on_error=True))
await set_plan_status(thread_id, PLAN_STATUS_APPROVED, plan_mode=False)
if plan_markdown:
text = (
"The plan has been approved. Implement it now exactly as written "
"below (it may have been edited by the reviewer, so treat this as "
f"the source of truth):\n\n{plan_markdown}"
)
else:
text = "The plan has been approved. Implement it now as described in the plan."
if feedback:
text += "\n\nAlso take this reviewer feedback into account:\n\n" + feedback
await _dispatch_followup(thread_id, metadata, text, plan_mode=False)
return {"status": PLAN_STATUS_APPROVED}
@plan_router.post("/{thread_id}/reject")
async def reject_plan(thread_id: str, session: dict[str, Any] = _SESSION_DEP) -> dict[str, Any]:
metadata = await _thread_metadata(thread_id)
if not _thread_is_readable(metadata):
raise HTTPException(404, "thread not found")
feedback = _format_comments(await list_plan_comments(thread_id, raise_on_error=True))
await set_plan_status(thread_id, PLAN_STATUS_REVISING, plan_mode=True)
text = (
"The plan needs changes before implementation. Address this reviewer "
"feedback and publish an updated plan with the save_plan tool:\n\n"
f"{feedback or '(no specific comments were left)'}"
)
await _dispatch_followup(thread_id, metadata, text, plan_mode=True)
return {"status": PLAN_STATUS_REVISING}
def _format_comments(comments: list[dict[str, Any]]) -> str:
lines: list[str] = []
index = 1
for comment in comments:
body = str(comment.get("body", "")).strip()
if not body:
continue
author = str(comment.get("author") or "reviewer").strip()
lines.append(f"{index}. {author}: {body}")
index += 1
return "\n".join(lines)
async def _dispatch_followup(
thread_id: str, metadata: dict[str, Any], text: str, *, plan_mode: bool
) -> None:
"""Continue the existing thread with a new instruction run.
Runs on the same LangGraph thread, so the agent resumes from the checkpoint
with the full planning history plus this instruction. The configurable is
rebuilt from the thread's stored owner/repo/Slack context so the agent can
push, open a PR, and reply in the original channel.
"""
configurable: dict[str, Any] = {
"thread_id": thread_id,
"source": _thread_source(metadata) or "slack",
}
email = metadata.get("triggering_user_email")
if isinstance(email, str) and email:
configurable["user_email"] = email
login = metadata.get("github_login")
if isinstance(login, str) and login:
configurable["github_login"] = login
repo = _repo_config_from_metadata(metadata)
if repo:
configurable["repo"] = repo
source_context = metadata.get("source_context")
if isinstance(source_context, dict):
slack_thread = source_context.get("slack_thread")
if isinstance(slack_thread, dict):
configurable["slack_thread"] = slack_thread
# Carry the decision to the follow-up run: approve continues out of plan
# mode (implement), reject stays in plan mode (revise the plan).
configurable["plan_mode"] = plan_mode
await dispatch_agent_run(
thread_id,
text,
configurable,
source=configurable["source"],
)