diff --git a/agent/tools/schedule_thread_wakeup.py b/agent/tools/schedule_thread_wakeup.py index 88b558c5..78d42e06 100644 --- a/agent/tools/schedule_thread_wakeup.py +++ b/agent/tools/schedule_thread_wakeup.py @@ -18,6 +18,9 @@ _MIN_DELAY_SECONDS = 60 _MAX_DELAY_SECONDS = 86_400 _END_TIME_PADDING_SECONDS = 90 +_WAKEUP_KIND = "thread_wakeup" +_PURGE_PAGE_SIZE = 100 + _DEFAULT_WAKEUP_PROMPT = ( "This is an automated re-trigger of this thread. The agent scheduled this " "wakeup to poll for updates. Check the current state of whatever you were " @@ -46,6 +49,71 @@ def _build_one_shot_cron(fire_time: datetime) -> str: ) +def _parse_iso(value: Any) -> datetime | None: + if not isinstance(value, str) or not value: + return None + try: + return datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return None + + +async def find_expired_wakeup_cron_ids(client: Any, *, now: datetime) -> list[str]: + """Return the ids of ``thread_wakeup`` crons whose ``end_time`` has passed. + + Conservative: matches solely on ``metadata.kind == "thread_wakeup"`` AND a + past ``end_time``, so analyzer/dashboard crons are never selected. Paginates + fully before returning so the result is stable to delete afterwards. + """ + expired_ids: list[str] = [] + offset = 0 + while True: + page = await client.crons.search( + metadata={"kind": _WAKEUP_KIND}, + limit=_PURGE_PAGE_SIZE, + offset=offset, + ) + if not page: + break + for cron in page: + if not isinstance(cron, dict): + continue + end_time = _parse_iso(cron.get("end_time")) + cron_id = cron.get("cron_id") + if end_time is not None and end_time < now and isinstance(cron_id, str) and cron_id: + expired_ids.append(cron_id) + if len(page) < _PURGE_PAGE_SIZE: + break + offset += len(page) + return expired_ids + + +async def purge_expired_wakeup_crons(client: Any, *, now: datetime) -> int: + """Delete ``thread_wakeup`` crons whose ``end_time`` has already passed. + + Each wakeup is a thread-bound cron with an ``end_time`` (~90s past its fire) + that stops it re-firing, but the cron row itself is never removed, so dead + rows accumulate. This deletes only those dead rows. Returns the count deleted. + """ + expired_ids = await find_expired_wakeup_cron_ids(client, now=now) + deleted = 0 + for cron_id in expired_ids: + await client.crons.delete(cron_id) + deleted += 1 + return deleted + + +async def _purge_expired_wakeups_best_effort() -> None: + """Opportunistically purge expired wakeup crons; never raises.""" + try: + client = get_client(url=langgraph_url()) + deleted = await purge_expired_wakeup_crons(client, now=datetime.now(UTC)) + if deleted: + logger.info("Purged %d expired thread_wakeup cron(s)", deleted) + except Exception: + logger.warning("Failed to purge expired thread_wakeup crons", exc_info=True) + + async def _create_wakeup_cron( *, thread_id: str, @@ -130,6 +198,8 @@ async def schedule_thread_wakeup(delay_minutes: int, prompt: str | None = None) if value is not None: wakeup_configurable[key] = value + await _purge_expired_wakeups_best_effort() + try: return await _create_wakeup_cron( thread_id=thread_id, diff --git a/docs/upstream-sync/cherry-pick-runbook.md b/docs/upstream-sync/cherry-pick-runbook.md new file mode 100644 index 00000000..bd5312bc --- /dev/null +++ b/docs/upstream-sync/cherry-pick-runbook.md @@ -0,0 +1,95 @@ +# Cherry-picking upstream into the fork + +This repository is a long-lived fork of `langchain-ai/open-swe` (git remote `upstream`). +Upstream changes are brought in one commit at a time with `git cherry-pick`, and every +diverged commit is tracked in a triage ledger so a decision is made once and not revisited. +`dev` is the integration branch; the broader strategy lives in the fork-maintenance section +of `CLAUDE.md`. + +## Setup (once per clone) + + make install-hooks + +Installs the triage hooks and the `git cp` alias by pointing `core.hooksPath` at `.githooks/`. +Because that shadows the machine-global hook directory (`~/.config/git/hooks`, which holds the +mandatory security `pre-push`), `.githooks/pre-push` is a shim that re-execs the global hook, +and the installer verifies that delegation before it changes anything. `git` never auto-adopts +a repository's `core.hooksPath`, so this step cannot be skipped. + +Note: `core.hooksPath` applies repo-wide, but `.githooks/` is a tracked directory. The hooks +(and the security shim) only run on branches that actually contain `.githooks/`. Keep it present +on `dev` and `main` so no branch loses the security `pre-push`. + +## Finding what to pick + + make triage-sync # git fetch upstream, then append new dev..upstream/main commits + # to the ledger as `untriaged` (PR # + subject parsed from each) + +`triage-sync` is the discovery step: it records every diverged commit as `untriaged` and bumps +"Last synced" to the new `upstream/main` tip. Triage those rows (decide `deferred` / `wont-merge` +and which branch), then pick the ones you want. The underlying views if you prefer raw git: + + git fetch upstream + git log --oneline --no-merges dev..upstream/main # everything diverged + git show # inspect before deciding + +Cross-check candidates against the ledger first — most diverged commits already carry a +decision (already-in-dev, regression, deferred, or landed) and should not be re-examined. + +## Bringing in commits: `git cp` + + git cp -x # pre-check the ledger, cherry-pick -x, auto-reconcile + git cp -x ... # several, applied in the given order + git cp --continue # after resolving a conflict; also reconciles + git cp --force # override a SHA the ledger marks "Won't merge" + +`git cp` reads `docs/upstream-sync/triage.jsonl` before touching the tree and refuses a +known-reject SHA (override with `--force`). On success it runs `make triage-reconcile`, which +moves each applied SHA to Landed in the ledger and stages `triage.jsonl` + `triage.md` for you +to commit. + +Apply commits in upstream chronological order (oldest first), not the order you happen to list +them — a later commit often depends on an earlier one, and out-of-order picks conflict +needlessly: + + git log --reverse --topo-order --format=%h dev..upstream/main + +## The triage ledger + +`docs/upstream-sync/triage.jsonl` is the source of truth: one JSON row per upstream SHA, keyed +on the SHA (stable, unlike the local SHAs cherry-pick rewrites). `docs/upstream-sync/triage.md` +is generated from it and must not be hand-edited. Dispositions are `landed`, `wont-merge`, +`deferred`, `untriaged`. + + scripts/triage.py set --disposition deferred --branch slack-tooling --reason "..." + make triage-render # regenerate triage.md from the jsonl + make triage-check # CI gate: fail if triage.md is stale + +A SHA marked `wont-merge` is hard-blocked by both `git cp` and the `prepare-commit-msg` hook. +Override for a one-off re-evaluation with `git cp --force`, `SH_CHERRYPICK_ALLOW_REJECT=1`, or +`git config sh.cherrypick.blockRejects false`. + +## Branch layout + +Never cherry-pick onto `dev` directly. Work on a themed branch off `dev` and open a PR into +`dev`; the ledger's `branch` column records where each deferred commit is meant to land +(for example `slack-tooling`, `gateway-routing`, `plan-approval`, `durable-dispatch`). Keep +each PR to one theme so conflict resolution stays within one subsystem. + +## Raw `git cherry-pick` + +The hooks fire on a plain `git cherry-pick -x ` too: `post-commit` journals each applied +pick and `prepare-commit-msg` blocks known-rejects. Run `make triage-reconcile` once at the end +to land the picks in the ledger, then commit `triage.jsonl` + `triage.md`. + +## Conflicts + + # resolve the files, then: + git add + git cherry-pick --continue # or: git cp --continue + git cherry-pick --abort # bail out of the whole pick + git cherry-pick --skip # drop just this commit and continue the batch + +A commit that conflicts because `dev` already carries a newer version of the same code is a +regression, not a merge — skip it and record the decision as `wont-merge` in the ledger rather +than forcing it in. diff --git a/scripts/purge_wakeup_crons.py b/scripts/purge_wakeup_crons.py new file mode 100644 index 00000000..ad556593 --- /dev/null +++ b/scripts/purge_wakeup_crons.py @@ -0,0 +1,88 @@ +"""One-time backfill: delete expired ``thread_wakeup`` crons from a deployment. + +One-shot wakeup crons set an ``end_time`` that stops them re-firing, but the +cron row is never removed, so dead rows accumulate. The ``schedule_thread_wakeup`` +tool now purges these opportunistically; this script clears the backlog. + +Usage: + uv run python scripts/purge_wakeup_crons.py --dry-run + uv run python scripts/purge_wakeup_crons.py + +Resolves the deployment URL from ``--url`` or ``LANGGRAPH_URL`` / ``LANGGRAPH_URL_PROD``, +and the API key from ``LANGGRAPH_API_KEY`` / ``LANGSMITH_API_KEY`` / ``LANGSMITH_API_KEY_PROD``. +""" + +from __future__ import annotations + +import argparse +import asyncio +import logging +import os +from datetime import UTC, datetime + +from langgraph_sdk import get_client + +from agent.tools.schedule_thread_wakeup import ( + find_expired_wakeup_cron_ids, + purge_expired_wakeup_crons, +) + +logger = logging.getLogger(__name__) + + +def _load_dotenv_if_available() -> None: + try: + from dotenv import load_dotenv + except ImportError: + return + load_dotenv() + + +def _resolve_url(arg_url: str | None) -> str: + url = arg_url or os.environ.get("LANGGRAPH_URL") or os.environ.get("LANGGRAPH_URL_PROD") + if not url: + raise RuntimeError("Set --url or LANGGRAPH_URL / LANGGRAPH_URL_PROD") + return url + + +def _resolve_api_key() -> str | None: + return ( + os.environ.get("LANGGRAPH_API_KEY") + or os.environ.get("LANGSMITH_API_KEY") + or os.environ.get("LANGSMITH_API_KEY_PROD") + ) + + +async def _run(url: str, api_key: str | None, dry_run: bool) -> None: + client = get_client(url=url, api_key=api_key) + now = datetime.now(UTC) + if dry_run: + expired = await find_expired_wakeup_cron_ids(client, now=now) + logger.info("[dry-run] %d expired thread_wakeup cron(s) would be deleted", len(expired)) + for cron_id in expired: + logger.info(" %s", cron_id) + return + deleted = await purge_expired_wakeup_crons(client, now=now) + logger.info("Deleted %d expired thread_wakeup cron(s)", deleted) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Purge expired thread_wakeup crons.") + parser.add_argument("--url", default=None, help="Deployment URL (defaults to env).") + parser.add_argument( + "--dry-run", + action="store_true", + help="List the crons that would be deleted without deleting them.", + ) + return parser.parse_args() + + +def main() -> None: + _load_dotenv_if_available() + logging.basicConfig(level=logging.INFO, format="%(message)s") + args = parse_args() + asyncio.run(_run(_resolve_url(args.url), _resolve_api_key(), args.dry_run)) + + +if __name__ == "__main__": + main() diff --git a/tests/e2e/README.md b/tests/e2e/README.md index 990ef3cd..1469b129 100644 --- a/tests/e2e/README.md +++ b/tests/e2e/README.md @@ -12,17 +12,17 @@ This drives the **whole happy path** through two mock UIs: Only the **LLM** and the **external SaaS HTTP boundaries** are faked. All agent code runs for real. -| Piece | Real or fake | -|---|---| -| Slack webhook → `process_slack_mention` → run dispatch | **real** (`agent.webapp`) | -| `get_agent`, deepagents loop, tools, middleware, prompt | **real** | -| `open_pull_request`, `slack_thread_reply` tools | **real** | -| Sandbox | **real** `local` provider, rooted in a throwaway temp dir | -| Git remote ("GitHub") | **real git**, a local bare repo the agent clones/pushes | -| The LLM | **fake** — a scripted model (`fake_llm.py`) emitting a fixed tool sequence | -| `api.github.com` REST (PR create) | **fake** (`/fake-gh/...`), state rendered at `/mock/github` | -| `slack.com/api` (post message, etc.) | **fake** (`/fake-slack/...`), thread rendered at `/mock/slack` | -| GitHub App token mint, `api.github.com/user` identity | stubbed (offline) | +| Piece | Real or fake | +| ---------------------------------------------------------------- | -------------------------------------------------------------------------- | +| Slack webhook → `process_slack_mention` → run dispatch | **real** (`agent.webapp`) | +| `get_agent`, deepagents loop, tools, middleware, prompt | **real** | +| `open_pull_request`, `slack_thread_reply` tools | **real** | +| Sandbox | **real** `local` provider, rooted in a throwaway temp dir | +| Git remote ("GitHub") | **real git**, a local bare repo the agent clones/pushes | +| The LLM | **fake** — a scripted model (`fake_llm.py`) emitting a fixed tool sequence | +| `api.github.com` REST (PR create) + dashboard GitHub OAuth login | **fake** (`/fake-gh/...`), state rendered at `/mock/github` | +| `slack.com/api` (post message, etc.) | **fake** (`/fake-slack/...`), thread rendered at `/mock/slack` | +| GitHub App token mint, `api.github.com/user` identity | stubbed (offline) | The fake GitHub/Slack stores are the single source of truth the mock UIs render, so what Playwright asserts on is exactly what the real agent produced. diff --git a/tests/e2e/harness.py b/tests/e2e/harness.py index 1147300b..f774d546 100644 --- a/tests/e2e/harness.py +++ b/tests/e2e/harness.py @@ -18,8 +18,10 @@ import json import os import sys import time +from html import escape from pathlib import Path from typing import Any +from urllib.parse import quote sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) @@ -185,29 +187,43 @@ async def control_login_get(login: str = "", email: str = "", next_url: str = "" @app.get("/dashboard/api/auth/login") -async def mock_github_login(redirect_to: str = "", login: str = "") -> Response: - """Mock stand-in for GitHub OAuth: the dashboard's "Continue with GitHub" - button lands here. With no ``login``, render a picker of the fake GitHub - test users; once one is chosen, mint the real session cookie and redirect - back into the dashboard (``redirect_to``).""" +async def mock_github_login(redirect_to: str = "") -> Response: + """E2E stand-in for the dashboard OAuth start route. + + The real route would redirect to github.com. Keep the dashboard-facing URL + intact, then hand off to the fake GitHub simulator so Playwright exercises a + browser login flow instead of test code pre-minting a session cookie. + """ + ui = os.environ.get("DASHBOARD_BASE_URL", "").rstrip("/") + dest = redirect_to or (f"{ui}/agents" if ui else "/agents") + return RedirectResponse(f"/fake-gh/login/oauth/authorize?redirect_to={quote(dest)}", 302) + + +@app.get("/fake-gh/login/oauth/authorize") +async def fake_github_authorize(redirect_to: str = "", login: str = "") -> Response: + """Fake GitHub OAuth consent/login page for dashboard e2e tests.""" ui = os.environ.get("DASHBOARD_BASE_URL", "").rstrip("/") dest = redirect_to or (f"{ui}/agents" if ui else "/agents") if not login: options = "".join( - f'' for u in TEST_USERS + f'" + for u in TEST_USERS ) return HTMLResponse( - f"""Continue with GitHub (mock) + f"""GitHub · Authorize open-swe -

Continue with GitHub (mock)

-

Pick a fake GitHub account to sign in as.

-
- - - -
-

Tip: use a separate browser or profile per - user so their sessions don't overwrite each other.

+
+

Authorize open-swe

+

Pick a fake GitHub account to continue.

+
+ + + +
+
""" ) match = next((u for u in TEST_USERS if u["login"] == login), None) diff --git a/tests/e2e/tests/plan_review.spec.ts b/tests/e2e/tests/plan_review.spec.ts index e7637a7b..e7f92199 100644 --- a/tests/e2e/tests/plan_review.spec.ts +++ b/tests/e2e/tests/plan_review.spec.ts @@ -38,9 +38,14 @@ test.describe("Plan review (HTTP comments)", () => { // 1. A user asks the bot to PLAN something in Slack. await request.post("/control/reset"); const send = await request.post("/mock/slack/send", { - data: { text: "<@U0BOT> plan how to add a greet() helper", mention_bot: true }, + data: { + text: "<@U0BOT> plan how to add a greet() helper", + mention_bot: true, + }, }); - const { thread_id: threadId } = (await send.json()) as { thread_id: string }; + const { thread_id: threadId } = (await send.json()) as { + thread_id: string; + }; expect(threadId).toBeTruthy(); const planPath = `/agents/${threadId}/plan`; @@ -57,7 +62,9 @@ test.describe("Plan review (HTTP comments)", () => { }; return (state.values?.messages ?? []) .map((m) => - typeof m.content === "string" ? m.content : JSON.stringify(m.content), + typeof m.content === "string" + ? m.content + : JSON.stringify(m.content), ) .some((c) => c.includes("Plan mode is active")); }, @@ -67,13 +74,45 @@ test.describe("Plan review (HTTP comments)", () => { // 2. The agent shares the plan-review link, then announces the plan is ready. await expect - .poll(async () => (await botMessages(request)).join("\n"), { timeout: 60_000 }) + .poll(async () => (await botMessages(request)).join("\n"), { + timeout: 60_000, + }) .toMatch(/\/agents\/[^/]+\/plan\b/); await expect - .poll(async () => (await botMessages(request)).join("\n"), { timeout: 60_000 }) + .poll(async () => (await botMessages(request)).join("\n"), { + timeout: 60_000, + }) .toMatch(/ready for review/i); - // 3. The OWNER opens the conversation, follows the "Review plan" banner, and + // 3. A logged-out user follows the plan deep link, signs in through the fake + // GitHub OAuth simulator, and lands back on the same plan page. + const loggedOutCtx = await browser.newContext(); + const loggedOut = await loggedOutCtx.newPage(); + await loggedOut.goto(planPath); + await expect(loggedOut).toHaveURL( + new RegExp(`/login\\?redirect=.*${threadId}.*plan`), + ); + await expect(loggedOut.getByText("Sign in to open-swe")).toBeVisible({ + timeout: 30_000, + }); + await loggedOut.getByRole("link", { name: "Continue with GitHub" }).click(); + await expect(loggedOut).toHaveURL(/\/fake-gh\/login\/oauth\/authorize/); + await expect(loggedOut.getByTestId("fake-github-login")).toBeVisible(); + await loggedOut.getByLabel("GitHub user").selectOption(OWNER.login); + await loggedOut.getByRole("button", { name: "Authorize open-swe" }).click(); + await expect(loggedOut).toHaveURL(new RegExp(`/agents/${threadId}/plan$`)); + await expect(loggedOut.getByTestId("plan-review")).toBeVisible({ + timeout: 30_000, + }); + await expect(loggedOut.getByTestId("plan-document")).toContainText( + "greet", + { + timeout: 30_000, + }, + ); + await loggedOutCtx.close(); + + // 4. The OWNER opens the conversation, follows the "Review plan" banner, and // sees the rendered plan. const ownerCtx = await browser.newContext({ permissions: ["clipboard-read", "clipboard-write"], @@ -85,7 +124,9 @@ test.describe("Plan review (HTTP comments)", () => { await expect(reviewLink).toBeVisible({ timeout: 30_000 }); await reviewLink.click(); await expect(owner).toHaveURL(new RegExp(`/agents/${threadId}/plan$`)); - await expect(owner.getByTestId("plan-review")).toBeVisible({ timeout: 30_000 }); + await expect(owner.getByTestId("plan-review")).toBeVisible({ + timeout: 30_000, + }); await expect(owner.getByText("Back to conversation")).toBeVisible(); await expect(owner.getByTestId("plan-document")).toContainText("greet", { timeout: 30_000, @@ -98,7 +139,9 @@ test.describe("Plan review (HTTP comments)", () => { // Copy the whole plan as markdown. await owner.getByTestId("copy-plan").click(); await expect(owner.getByTestId("copy-plan")).toContainText("Copied!"); - const clipboard = await owner.evaluate(() => navigator.clipboard.readText()); + const clipboard = await owner.evaluate(() => + navigator.clipboard.readText(), + ); expect(clipboard).toContain("## Plan: Add greet() helper"); expect(clipboard).toContain("### Verification"); @@ -107,18 +150,24 @@ test.describe("Plan review (HTTP comments)", () => { await expect(owner.getByTestId("plan-comment")).toHaveCount(1); await expect(owner.getByTestId("reject-plan")).toBeEnabled(); - // 4. A COLLABORATOR opens the same plan: sees it AND the owner's comment + // 5. A COLLABORATOR opens the same plan: sees it AND the owner's comment // (fetched over HTTP), but has NO approve button. const collabCtx = await browser.newContext(); await collabCtx.request.post("/control/login", { data: COLLABORATOR }); const collab = await collabCtx.newPage(); await collab.goto(planPath); - await expect(collab.getByTestId("plan-review")).toBeVisible({ timeout: 30_000 }); + await expect(collab.getByTestId("plan-review")).toBeVisible({ + timeout: 30_000, + }); await expect(collab.getByTestId("plan-document")).toContainText("greet", { timeout: 30_000, }); - await expect(collab.getByTestId("plan-comment")).toHaveCount(1, { timeout: 30_000 }); - await expect(collab.getByTestId("plan-comment")).toContainText("looks solid"); + await expect(collab.getByTestId("plan-comment")).toHaveCount(1, { + timeout: 30_000, + }); + await expect(collab.getByTestId("plan-comment")).toContainText( + "looks solid", + ); await expect(collab.getByTestId("approve-plan")).toHaveCount(0); await expect(collab.getByTestId("reject-plan")).toBeVisible(); @@ -126,20 +175,27 @@ test.describe("Plan review (HTTP comments)", () => { await addComment(collab, "Reviewer: please also add a docstring."); await expect(collab.getByTestId("plan-comment")).toHaveCount(2); - // 5. The owner sees the collaborator's comment (polled), then approves. - await expect(owner.getByTestId("plan-comment")).toHaveCount(2, { timeout: 30_000 }); + // 6. The owner sees the collaborator's comment (polled), then approves and + // returns to the main conversation while implementation starts. + await expect(owner.getByTestId("plan-comment")).toHaveCount(2, { + timeout: 30_000, + }); await owner.getByTestId("approve-plan").click(); - await expect(owner.getByTestId("plan-decision")).toContainText(/implementing/i); + await expect(owner).toHaveURL(new RegExp(`/agents/${threadId}$`)); - // 6. The agent implements, opens a PR, and links it back in the Slack thread, + // 7. The agent implements, opens a PR, and links it back in the Slack thread, // echoing the reviewers' feedback — which proves the comments were stored // and harvested server-side on approve. await expect - .poll(async () => (await botMessages(request)).join("\n"), { timeout: 90_000 }) + .poll(async () => (await botMessages(request)).join("\n"), { + timeout: 90_000, + }) .toMatch(/\/pull\//); expect((await botMessages(request)).join("\n")).toMatch(/docstring/); - const prs = (await (await request.get("/mock/github/data")).json()) as Array; + const prs = (await ( + await request.get("/mock/github/data") + ).json()) as Array; expect(prs.length).toBeGreaterThan(0); await ownerCtx.close(); diff --git a/tests/test_dashboard_oauth_redirect.py b/tests/test_dashboard_oauth_redirect.py new file mode 100644 index 00000000..c19287ec --- /dev/null +++ b/tests/test_dashboard_oauth_redirect.py @@ -0,0 +1,30 @@ +from __future__ import annotations + +from agent.dashboard.oauth import sanitize_redirect_to + + +def test_sanitize_redirect_to_preserves_allowed_dashboard_target(monkeypatch) -> None: + monkeypatch.setenv("DASHBOARD_BASE_URL", "https://dashboard.example") + monkeypatch.setenv("DASHBOARD_ALLOWED_ORIGINS", "https://preview.example") + + target = "https://dashboard.example/agents/thread-1/plan?from=slack#review" + + assert sanitize_redirect_to(target) == target + + +def test_sanitize_redirect_to_preserves_allowed_preview_target(monkeypatch) -> None: + monkeypatch.setenv("DASHBOARD_BASE_URL", "https://dashboard.example") + monkeypatch.setenv("DASHBOARD_ALLOWED_ORIGINS", "https://preview.example") + + target = "https://preview.example/agents/thread-1/plan?from=slack#review" + + assert sanitize_redirect_to(target) == target + + +def test_sanitize_redirect_to_rejects_external_target(monkeypatch) -> None: + monkeypatch.setenv("DASHBOARD_BASE_URL", "https://dashboard.example") + monkeypatch.setenv("DASHBOARD_ALLOWED_ORIGINS", "https://preview.example") + + assert sanitize_redirect_to("https://evil.example/agents/thread-1/plan") == ( + "https://dashboard.example" + ) diff --git a/tests/test_schedule_thread_wakeup.py b/tests/test_schedule_thread_wakeup.py index 55b0c05e..50c4b1ad 100644 --- a/tests/test_schedule_thread_wakeup.py +++ b/tests/test_schedule_thread_wakeup.py @@ -1,13 +1,67 @@ from __future__ import annotations import importlib -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta from typing import Any import pytest wakeup_tool = importlib.import_module("agent.tools.schedule_thread_wakeup") +# Captured before the autouse stub replaces it, for the one test that needs the real wrapper. +_real_purge_best_effort = wakeup_tool._purge_expired_wakeups_best_effort + + +@pytest.fixture(autouse=True) +def _stub_purge(monkeypatch: pytest.MonkeyPatch) -> None: + """Keep the opportunistic purge from touching the network in every test.""" + + async def _noop() -> None: + return None + + monkeypatch.setattr(wakeup_tool, "_purge_expired_wakeups_best_effort", _noop) + + +class _FakeCrons: + def __init__(self, crons: list[dict[str, Any]]) -> None: + self._crons = list(crons) + self.deleted: list[str] = [] + self.search_calls: list[dict[str, Any]] = [] + + async def search( + self, + *, + metadata: dict[str, Any] | None = None, + limit: int = 10, + offset: int = 0, + **_: Any, + ) -> list[dict[str, Any]]: + self.search_calls.append({"metadata": metadata, "limit": limit, "offset": offset}) + items = [ + c + for c in self._crons + if not metadata + or all((c.get("metadata") or {}).get(k) == v for k, v in metadata.items()) + ] + return items[offset : offset + limit] + + async def delete(self, cron_id: str) -> None: + self.deleted.append(cron_id) + self._crons = [c for c in self._crons if c.get("cron_id") != cron_id] + + +class _FakeClient: + def __init__(self, crons: list[dict[str, Any]]) -> None: + self.crons = _FakeCrons(crons) + + +def _wakeup_cron(cron_id: str, end_time: datetime | None) -> dict[str, Any]: + return { + "cron_id": cron_id, + "end_time": end_time.isoformat() if end_time else None, + "metadata": {"kind": "thread_wakeup"}, + } + def _config(**overrides: Any) -> dict[str, Any]: base: dict[str, Any] = { @@ -241,3 +295,71 @@ def test_build_one_shot_cron_handles_month_boundary() -> None: assert parts[1] == "23" assert parts[2] == "31" assert parts[3] == "12" + + +async def test_purge_deletes_only_expired_wakeups() -> None: + now = datetime(2026, 6, 30, 22, 0, tzinfo=UTC) + client = _FakeClient( + [ + _wakeup_cron("expired-1", now - timedelta(hours=1)), + _wakeup_cron("expired-2", now - timedelta(days=1)), + _wakeup_cron("future-1", now + timedelta(hours=1)), + _wakeup_cron("no-end", None), + ] + ) + + deleted = await wakeup_tool.purge_expired_wakeup_crons(client, now=now) + + assert deleted == 2 + assert client.crons.deleted == ["expired-1", "expired-2"] + # Search is scoped to the thread_wakeup kind so other crons are never seen. + assert client.crons.search_calls[0]["metadata"] == {"kind": "thread_wakeup"} + + +async def test_purge_paginates(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(wakeup_tool, "_PURGE_PAGE_SIZE", 2) + now = datetime(2026, 6, 30, 22, 0, tzinfo=UTC) + client = _FakeClient([_wakeup_cron(f"expired-{i}", now - timedelta(hours=1)) for i in range(3)]) + + deleted = await wakeup_tool.purge_expired_wakeup_crons(client, now=now) + + assert deleted == 3 + assert sorted(client.crons.deleted) == ["expired-0", "expired-1", "expired-2"] + # Two pages fetched (offset 0 and 2), then a short final page ends the loop. + assert [c["offset"] for c in client.crons.search_calls] == [0, 2] + + +async def test_best_effort_purge_swallows_errors(monkeypatch: pytest.MonkeyPatch) -> None: + async def boom(*_: Any, **__: Any) -> int: + raise RuntimeError("search failed") + + monkeypatch.setattr(wakeup_tool, "purge_expired_wakeup_crons", boom) + monkeypatch.setattr(wakeup_tool, "get_client", lambda url: object()) + + # The real wrapper must never propagate — a purge failure can't block wakeups. + await _real_purge_best_effort() + + +async def test_schedule_purges_before_creating(monkeypatch: pytest.MonkeyPatch) -> None: + calls: list[str] = [] + + async def spy_purge() -> None: + calls.append("purge") + + async def fake_create_wakeup_cron(**kwargs: Any) -> dict[str, Any]: + calls.append("create") + return { + "success": True, + "cron_id": "cron-1", + "scheduled_for": "", + "thread_id": kwargs["thread_id"], + } + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_purge_expired_wakeups_best_effort", spy_purge) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = await wakeup_tool.schedule_thread_wakeup(5) + + assert result["success"] is True + assert calls == ["purge", "create"] diff --git a/ui/src/client.tsx b/ui/src/client.tsx index c6783110..5fe66e09 100644 --- a/ui/src/client.tsx +++ b/ui/src/client.tsx @@ -1,5 +1,5 @@ import { StartClient } from "@tanstack/react-start/client" -import { StrictMode, useEffect, useState } from "react" +import { useEffect, useState } from "react" import { hydrateRoot } from "react-dom/client" import { registerSW } from "virtual:pwa-register" @@ -35,8 +35,8 @@ function PwaUpdateProvider() { hydrateRoot( document, - + <> - + ) diff --git a/ui/src/components/agents/PlanReview.tsx b/ui/src/components/agents/PlanReview.tsx index 1220c4a3..c8ec2815 100644 --- a/ui/src/components/agents/PlanReview.tsx +++ b/ui/src/components/agents/PlanReview.tsx @@ -1,4 +1,5 @@ import { useCallback, useEffect, useState } from "react" +import { useNavigate } from "@tanstack/react-router" import type { PlanComment, PlanData } from "@/lib/plan" import { @@ -47,6 +48,7 @@ async function copyToClipboard(text: string): Promise { } export function PlanReview({ plan }: { plan: PlanData }) { + const navigate = useNavigate() const resolvedTheme = useResolvedTheme() const [comments, setComments] = useState>([]) const [draft, setDraft] = useState("") @@ -108,20 +110,23 @@ export function PlanReview({ plan }: { plan: PlanData }) { setBusy(kind) setError(null) try { - if (kind === "approve") await approvePlan(plan.threadId) - else await rejectPlan(plan.threadId) - setDecision( - kind === "approve" - ? "Plan approved — the agent is implementing it." - : "Changes requested — the agent is revising the plan." - ) + if (kind === "approve") { + await approvePlan(plan.threadId) + await navigate({ + to: "/agents/$threadId", + params: { threadId: plan.threadId }, + }) + return + } + await rejectPlan(plan.threadId) + setDecision("Changes requested — the agent is revising the plan.") } catch (e) { setError((e as Error).message) } finally { setBusy(null) } }, - [plan.threadId] + [navigate, plan.threadId] ) const copyPlan = useCallback(async () => { @@ -139,8 +144,8 @@ export function PlanReview({ plan }: { plan: PlanData }) { data-testid="plan-review" className="flex min-h-0 flex-1 flex-col bg-[var(--ui-bg)] text-[var(--ui-text)]" > -
-
+
+

Implementation plan

@@ -150,11 +155,11 @@ export function PlanReview({ plan }: { plan: PlanData }) { {plan.status}

-
+
{decision && ( {decision} @@ -194,9 +199,9 @@ export function PlanReview({ plan }: { plan: PlanData }) {
-
+
@@ -209,14 +214,14 @@ export function PlanReview({ plan }: { plan: PlanData }) { )}
-