diff --git a/AGENTS.md b/AGENTS.md index 11f2f656..ffec329d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -84,7 +84,7 @@ There is intentionally no after-agent safety net that opens a PR for the agent. All tools live in `agent/tools/` and are flat-imported via `agent/tools/__init__.py`. The set is intentionally small and curated — see README "Tools — Curated, Not Accumulated". Wired into `get_agent`: -`http_request`, `fetch_url`, `web_search`, `linear_comment`, `linear_create_issue`, `linear_delete_issue`, `linear_get_issue`, `linear_get_issue_comments`, `linear_list_teams`, `linear_update_issue`, `request_pr_review`, `slack_read_thread_messages`, `slack_thread_reply`. +`http_request`, `fetch_url`, `web_search`, `linear_comment`, `linear_create_issue`, `linear_delete_issue`, `linear_get_issue`, `linear_get_issue_comments`, `linear_list_teams`, `linear_update_issue`, `request_pr_review`, `schedule_thread_wakeup`, `slack_read_thread_messages`, `slack_thread_reply`. Reviewer-only tools (in `agent/reviewer.py`): `add_finding`, `update_finding`, `list_findings`, `publish_review`. The review-style analyzer uses `save_review_style` (exported as `save_review_style_prompt`). diff --git a/CLAUDE.md b/CLAUDE.md index 7707bcf1..31cbc46c 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -80,7 +80,7 @@ There is intentionally no after-agent safety net that opens a PR for the agent. All tools live in `agent/tools/` and are flat-imported via `agent/tools/__init__.py`. The set is intentionally small and curated — see README "Tools — Curated, Not Accumulated". Wired into `get_agent`: -`http_request`, `fetch_url`, `web_search`, `linear_comment`, `linear_create_issue`, `linear_delete_issue`, `linear_get_issue`, `linear_get_issue_comments`, `linear_list_teams`, `linear_update_issue`, `request_pr_review`, `slack_read_thread_messages`, `slack_thread_reply`. +`http_request`, `fetch_url`, `web_search`, `linear_comment`, `linear_create_issue`, `linear_delete_issue`, `linear_get_issue`, `linear_get_issue_comments`, `linear_list_teams`, `linear_update_issue`, `request_pr_review`, `schedule_thread_wakeup`, `slack_read_thread_messages`, `slack_thread_reply`. Reviewer-only tools (in `agent/reviewer.py`): `add_finding`, `update_finding`, `list_findings`, `publish_review`. The review-style analyzer uses `save_review_style` (exported as `save_review_style_prompt`). diff --git a/agent/prompt.py b/agent/prompt.py index be90eabf..bee73037 100644 --- a/agent/prompt.py +++ b/agent/prompt.py @@ -176,6 +176,12 @@ Format messages using Slack's mrkdwn format, NOT standard Markdown. Do NOT use **bold**, [link](url), or other standard Markdown syntax. To mention/tag a user, use `<@USER_ID>` (e.g. `<@U06KD8BFY95>`). You can find user IDs in the conversation context next to display names (e.g. `@Name(U06KD8BFY95)`). +#### `request_pr_review` +Start the reviewer agent for a GitHub pull request URL. + +#### `schedule_thread_wakeup` +Schedule a one-shot re-trigger of the current thread after a delay. Pass `delay_minutes` (1–1440) and an optional `prompt` message. Use this to poll for updates — e.g. waiting for CI to finish, a deploy to complete, or an external process to settle. The thread will be re-invoked with the same run context (repo, source, Slack/Linear info) so you can continue where you left off. After the wakeup fires, the scheduled cron is automatically retired. + #### GitHub via `gh` Use `GH_TOKEN=dummy gh ` for GitHub operations: repository discovery, cloning, issues, pull requests, reviews, comments, labels, check status, and workflow operations. For local working-tree state, use `git` directly. Never pass a real GitHub token to `gh`.""" diff --git a/agent/server.py b/agent/server.py index f52deebe..c89b455a 100644 --- a/agent/server.py +++ b/agent/server.py @@ -73,6 +73,7 @@ from .tools import ( linear_update_issue, open_pull_request, request_pr_review, + schedule_thread_wakeup, slack_read_thread_messages, slack_thread_reply, web_search, @@ -732,6 +733,7 @@ async def get_agent(config: RunnableConfig) -> Pregel: linear_update_issue, open_pull_request, request_pr_review, + schedule_thread_wakeup, slack_read_thread_messages, slack_thread_reply, *corridor_tools, diff --git a/agent/tools/__init__.py b/agent/tools/__init__.py index 79f9e21a..041630a5 100644 --- a/agent/tools/__init__.py +++ b/agent/tools/__init__.py @@ -16,6 +16,7 @@ from .read_repo_file import read_repo_file from .reply_to_finding_thread import reply_to_finding_thread from .request_pr_review import request_pr_review from .resolve_finding_thread import resolve_finding_thread +from .schedule_thread_wakeup import schedule_thread_wakeup from .search_repo_code import search_repo_code from .slack_read_thread_messages import slack_read_thread_messages from .slack_thread_reply import slack_thread_reply @@ -41,6 +42,7 @@ __all__ = [ "request_pr_review", "reply_to_finding_thread", "resolve_finding_thread", + "schedule_thread_wakeup", "search_repo_code", "slack_read_thread_messages", "slack_thread_reply", diff --git a/agent/tools/schedule_thread_wakeup.py b/agent/tools/schedule_thread_wakeup.py new file mode 100644 index 00000000..85307a30 --- /dev/null +++ b/agent/tools/schedule_thread_wakeup.py @@ -0,0 +1,145 @@ +"""Tool that schedules a one-shot re-trigger of the current agent thread.""" + +from __future__ import annotations + +import asyncio +import logging +from datetime import UTC, datetime, timedelta +from typing import Any + +from langgraph.config import get_config +from langgraph_sdk import get_client + +from ..utils.thread_ops import langgraph_url + +logger = logging.getLogger(__name__) + +_AGENT_ASSISTANT_ID = "agent" +_MIN_DELAY_SECONDS = 60 +_MAX_DELAY_SECONDS = 86_400 +_END_TIME_PADDING_SECONDS = 90 + +_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 " + "waiting on and continue from there." +) + + +def _ceil_to_next_minute(value: datetime) -> datetime: + """Round a datetime up to the next whole minute.""" + rounded = value.replace(second=0, microsecond=0) + if rounded == value: + return rounded + return rounded + timedelta(minutes=1) + + +def _build_one_shot_cron(fire_time: datetime) -> str: + """Build a 5-field cron expression that fires at ``fire_time`` (UTC).""" + return " ".join( + [ + str(fire_time.minute), + str(fire_time.hour), + str(fire_time.day), + str(fire_time.month), + "*", + ] + ) + + +async def _create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], +) -> dict[str, Any]: + client = get_client(url=langgraph_url()) + schedule = _build_one_shot_cron(fire_time) + end_time = fire_time + timedelta(seconds=_END_TIME_PADDING_SECONDS) + run_config: dict[str, Any] = {"configurable": configurable} + cron = await client.crons.create_for_thread( + thread_id, + _AGENT_ASSISTANT_ID, + schedule=schedule, + input={"messages": [{"role": "user", "content": prompt}]}, + config=run_config, + end_time=end_time, + timezone="UTC", + metadata={ + "kind": "thread_wakeup", + "thread_id": thread_id, + }, + ) + cron_id = cron.get("cron_id") if isinstance(cron, dict) else getattr(cron, "cron_id", None) + return { + "success": True, + "cron_id": cron_id, + "scheduled_for": fire_time.isoformat(), + "thread_id": thread_id, + } + + +def schedule_thread_wakeup(delay_minutes: int, prompt: str | None = None) -> dict[str, Any]: + """Schedule a one-shot re-trigger of the current thread after a delay. + + Use this when you need to poll or check back on something later — e.g. + waiting for CI to finish, a deploy to complete, or an external process + to settle. The current thread will be re-invoked with the given prompt + (or a default wakeup message) after the specified delay. + + Args: + delay_minutes: How many minutes from now to wait before re-triggering. + Minimum 1 minute, maximum 1440 (24 hours). + prompt: Optional message to send to the thread when it wakes up. + If omitted, a default polling prompt is used. + + Returns a dict with ``success``, ``cron_id``, ``scheduled_for`` (ISO UTC), + and ``thread_id``. + """ + if not isinstance(delay_minutes, int) or delay_minutes < 1: + return {"success": False, "error": "delay_minutes must be a positive integer (>= 1)"} + delay_seconds = delay_minutes * 60 + if delay_seconds < _MIN_DELAY_SECONDS: + return {"success": False, "error": "delay must be at least 1 minute"} + if delay_seconds > _MAX_DELAY_SECONDS: + return {"success": False, "error": "delay must be at most 1440 minutes (24 hours)"} + + config = get_config() + configurable = config.get("configurable", {}) if isinstance(config, dict) else {} + thread_id = configurable.get("thread_id") + if not isinstance(thread_id, str) or not thread_id: + return {"success": False, "error": "No thread_id in current run config"} + + fire_time = _ceil_to_next_minute(datetime.now(UTC) + timedelta(seconds=delay_seconds)) + wakeup_prompt = ( + prompt.strip() if isinstance(prompt, str) and prompt.strip() else _DEFAULT_WAKEUP_PROMPT + ) + + passthrough_keys = ( + "repo", + "source", + "slack_thread", + "linear_issue", + "github_login", + "user_email", + "schedule_id", + ) + wakeup_configurable: dict[str, Any] = {"thread_id": thread_id} + for key in passthrough_keys: + value = configurable.get(key) + if value is not None: + wakeup_configurable[key] = value + + try: + return asyncio.run( + _create_wakeup_cron( + thread_id=thread_id, + fire_time=fire_time, + prompt=wakeup_prompt, + configurable=wakeup_configurable, + ) + ) + except Exception as exc: + logger.exception("Failed to schedule thread wakeup for %s", thread_id) + return {"success": False, "error": str(exc)} diff --git a/tests/test_schedule_thread_wakeup.py b/tests/test_schedule_thread_wakeup.py new file mode 100644 index 00000000..f06ec5d2 --- /dev/null +++ b/tests/test_schedule_thread_wakeup.py @@ -0,0 +1,237 @@ +from __future__ import annotations + +import importlib +from datetime import UTC, datetime +from typing import Any + +import pytest + +wakeup_tool = importlib.import_module("agent.tools.schedule_thread_wakeup") + + +def _config(**overrides: Any) -> dict[str, Any]: + base: dict[str, Any] = { + "configurable": { + "thread_id": "test-thread-123", + "source": "slack", + "repo": {"owner": "langchain-ai", "name": "open-swe"}, + "slack_thread": {"channel_id": "C1", "thread_ts": "1.0"}, + "github_login": "johannes117", + "user_email": "johannes@example.com", + } + } + base["configurable"].update(overrides) + return base + + +def test_schedule_thread_wakeup_rejects_zero_delay(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(wakeup_tool, "get_config", _config) + result = wakeup_tool.schedule_thread_wakeup(0) + assert result["success"] is False + assert "positive" in result["error"].lower() + + +def test_schedule_thread_wakeup_rejects_negative_delay(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(wakeup_tool, "get_config", _config) + result = wakeup_tool.schedule_thread_wakeup(-5) + assert result["success"] is False + + +def test_schedule_thread_wakeup_rejects_delay_over_24h(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(wakeup_tool, "get_config", _config) + result = wakeup_tool.schedule_thread_wakeup(1441) + assert result["success"] is False + assert "1440" in result["error"] + + +def test_schedule_thread_wakeup_rejects_missing_thread_id(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr( + wakeup_tool, + "get_config", + lambda: {"configurable": {"source": "slack"}}, + ) + result = wakeup_tool.schedule_thread_wakeup(5) + assert result["success"] is False + assert "thread_id" in result["error"].lower() + + +def test_schedule_thread_wakeup_creates_cron(monkeypatch: pytest.MonkeyPatch) -> None: + captured: dict[str, Any] = {} + + async def fake_create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], + ) -> dict[str, Any]: + captured.update( + { + "thread_id": thread_id, + "fire_time": fire_time, + "prompt": prompt, + "configurable": configurable, + } + ) + return { + "success": True, + "cron_id": "cron-abc", + "scheduled_for": fire_time.isoformat(), + "thread_id": thread_id, + } + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = wakeup_tool.schedule_thread_wakeup(10, prompt="Check CI status") + + assert result["success"] is True + assert result["cron_id"] == "cron-abc" + assert result["thread_id"] == "test-thread-123" + assert captured["thread_id"] == "test-thread-123" + assert captured["prompt"] == "Check CI status" + assert captured["configurable"]["thread_id"] == "test-thread-123" + assert captured["configurable"]["source"] == "slack" + assert captured["configurable"]["repo"] == {"owner": "langchain-ai", "name": "open-swe"} + assert captured["configurable"]["slack_thread"] == {"channel_id": "C1", "thread_ts": "1.0"} + assert captured["configurable"]["github_login"] == "johannes117" + + now = datetime.now(UTC) + delay = (captured["fire_time"] - now).total_seconds() + assert delay >= 600 + assert delay < 660 + assert captured["fire_time"].second == 0 + assert captured["fire_time"].microsecond == 0 + + +def test_schedule_thread_wakeup_uses_default_prompt_when_none( + monkeypatch: pytest.MonkeyPatch, +) -> None: + captured: dict[str, Any] = {} + + async def fake_create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], + ) -> dict[str, Any]: + captured["prompt"] = prompt + return {"success": True, "cron_id": "cron-1", "scheduled_for": "", "thread_id": thread_id} + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = wakeup_tool.schedule_thread_wakeup(5) + assert result["success"] is True + assert "automated re-trigger" in captured["prompt"].lower() + + +def test_schedule_thread_wakeup_uses_default_prompt_when_blank( + monkeypatch: pytest.MonkeyPatch, +) -> None: + captured: dict[str, Any] = {} + + async def fake_create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], + ) -> dict[str, Any]: + captured["prompt"] = prompt + return {"success": True, "cron_id": "cron-1", "scheduled_for": "", "thread_id": thread_id} + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = wakeup_tool.schedule_thread_wakeup(5, prompt=" ") + assert result["success"] is True + assert "automated re-trigger" in captured["prompt"].lower() + + +def test_schedule_thread_wakeup_returns_error_on_exception( + monkeypatch: pytest.MonkeyPatch, +) -> None: + async def fake_create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], + ) -> dict[str, Any]: + raise RuntimeError("connection refused") + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = wakeup_tool.schedule_thread_wakeup(5) + assert result["success"] is False + assert "connection refused" in result["error"] + + +def test_schedule_thread_wakeup_does_not_pass_none_configurable_keys( + monkeypatch: pytest.MonkeyPatch, +) -> None: + captured: dict[str, Any] = {} + + async def fake_create_wakeup_cron( + *, + thread_id: str, + fire_time: datetime, + prompt: str, + configurable: dict[str, Any], + ) -> dict[str, Any]: + captured["configurable"] = configurable + return {"success": True, "cron_id": "cron-1", "scheduled_for": "", "thread_id": thread_id} + + monkeypatch.setattr(wakeup_tool, "get_config", _config) + monkeypatch.setattr(wakeup_tool, "_create_wakeup_cron", fake_create_wakeup_cron) + + result = wakeup_tool.schedule_thread_wakeup(5) + assert result["success"] is True + cfg = captured["configurable"] + assert "linear_issue" not in cfg + assert "schedule_id" not in cfg + assert cfg["thread_id"] == "test-thread-123" + + +def test_ceil_to_next_minute_keeps_exact_minute() -> None: + value = datetime(2025, 1, 15, 14, 30, tzinfo=UTC) + assert wakeup_tool._ceil_to_next_minute(value) == value + + +def test_ceil_to_next_minute_rounds_up_with_seconds() -> None: + value = datetime(2025, 1, 15, 14, 30, 59, 123, tzinfo=UTC) + assert wakeup_tool._ceil_to_next_minute(value) == value.replace( + minute=31, second=0, microsecond=0 + ) + + +def test_ceil_to_next_minute_handles_day_boundary() -> None: + value = datetime(2025, 1, 31, 23, 59, 59, tzinfo=UTC) + assert wakeup_tool._ceil_to_next_minute(value) == value.replace( + year=2025, month=2, day=1, hour=0, minute=0, second=0 + ) + + +def test_build_one_shot_cron_format() -> None: + fire_time = datetime(2025, 1, 15, 14, 30, tzinfo=UTC) + cron = wakeup_tool._build_one_shot_cron(fire_time) + parts = cron.split(" ") + assert len(parts) == 5 + assert parts[0] == "30" + assert parts[1] == "14" + assert parts[2] == "15" + assert parts[3] == "1" + assert parts[4] == "*" + + +def test_build_one_shot_cron_handles_month_boundary() -> None: + fire_time = datetime(2025, 12, 31, 23, 59, tzinfo=UTC) + cron = wakeup_tool._build_one_shot_cron(fire_time) + parts = cron.split(" ") + assert parts[0] == "59" + assert parts[1] == "23" + assert parts[2] == "31" + assert parts[3] == "12"