feat: add schedule_thread_wakeup tool for self-polling (#1592)

* feat: add schedule_thread_wakeup tool for self-polling

Add a new agent tool that schedules a one-shot re-trigger of the
current thread after a configurable delay (1–1440 minutes). Uses a
LangGraph cron with end_time to fire exactly once, then retire.

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>

* fix: prevent early thread wakeups

Round scheduled wakeup times up to the next whole minute so cron minute precision cannot fire before the requested delay.

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>

---------

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
This commit is contained in:
Johannes du Plessis 2026-06-23 11:06:01 -07:00 • committed by GitHub
parent 0bff510aae
commit 450a25d33f
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 394 additions and 2 deletions

View file

@ -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`).

View file

@ -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`).

View file

@ -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 <command>` 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`."""

View file

@ -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,

View file

@ -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",

View file

@ -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)}

View file

@ -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"