feat: added abort_task host tool for silent task abandonment
- Added an AbortController and wired it into task bindings so host tools can request an immediate worker teardown. - Implemented a new `abort_task` host tool that audits and logs an internal reason, marks the issue as abandoned, and triggers the abort controller. - Updated the worker RPC loop to treat abort-triggered process errors as a clean shutdown while still propagating other task failures.
This commit is contained in:
@@ -86,7 +86,7 @@ Lint + format: TypeScript via Biome (config in `biome.json`), Python via Ruff (c
|
||||
- `src/robomp/queue.py` — `WorkerPool` dispatcher and `_inflight` serialization.
|
||||
- `src/robomp/tasks.py` — the five task entry points the dispatcher calls.
|
||||
- `src/robomp/worker.py` — synchronous omp RPC driver, prompt assembly via `persona`.
|
||||
- `src/robomp/host_tools.py` — agent's GitHub surface; tool list: `classify_issue`, `set_issue_labels`, `gh_post_comment`, `repro_record`, `gh_push_branch`, `gh_open_pr`, `gh_request_review`, `mark_unable_to_reproduce`, `fetch_issue_thread`.
|
||||
- `src/robomp/host_tools.py` — agent's GitHub surface; tool list: `classify_issue`, `set_issue_labels`, `gh_post_comment`, `repro_record`, `gh_push_branch`, `gh_open_pr`, `gh_request_review`, `mark_unable_to_reproduce`, `abort_task`, `fetch_issue_thread`.
|
||||
- `src/robomp/sandbox.py` — clone pool + worktree lifecycle, `GitCommandError`, credential redaction.
|
||||
- `src/robomp/github_client.py` — typed httpx client; parses webhook payloads into `IssueInfo` / `CommentInfo` / `PullRequestInfo`.
|
||||
- `src/robomp/github_events.py` — routing and HMAC verification.
|
||||
|
||||
+11
-10
@@ -91,16 +91,17 @@ RUN pip install /tmp/wheels/omp_rpc-*.whl && rm -rf /tmp/wheels
|
||||
WORKDIR /app
|
||||
|
||||
# `omp` shim — calls into the mounted pi checkout via Bun.
|
||||
RUN cat > /usr/local/bin/omp <<'EOF' && chmod +x /usr/local/bin/omp
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
: "${PI_ROOT:=/work/pi}"
|
||||
if [ ! -d "$PI_ROOT/packages/coding-agent" ]; then
|
||||
echo "roboomp: PI_ROOT=$PI_ROOT does not look like a pi checkout" >&2
|
||||
exit 127
|
||||
fi
|
||||
exec bun "$PI_ROOT/packages/coding-agent/src/cli.ts" "$@"
|
||||
EOF
|
||||
RUN printf '%s\n' \
|
||||
'#!/usr/bin/env bash' \
|
||||
'set -euo pipefail' \
|
||||
': "${PI_ROOT:=/work/pi}"' \
|
||||
'if [ ! -d "$PI_ROOT/packages/coding-agent" ]; then' \
|
||||
' echo "roboomp: PI_ROOT=$PI_ROOT does not look like a pi checkout" >&2' \
|
||||
' exit 127' \
|
||||
'fi' \
|
||||
'exec bun "$PI_ROOT/packages/coding-agent/src/cli.ts" "$@"' \
|
||||
> /usr/local/bin/omp \
|
||||
&& chmod +x /usr/local/bin/omp
|
||||
|
||||
# roboomp itself. Drop the Vite-built dashboard into the package tree before
|
||||
# `pip install` so it lands in the installed wheel (`static/**/*` is declared
|
||||
|
||||
@@ -206,7 +206,7 @@ src/robomp/
|
||||
worker.py synchronous omp RPC driver, prompt assembly, env scrubbing
|
||||
host_tools.py classify_issue, set_issue_labels, gh_post_comment, repro_record,
|
||||
gh_push_branch, gh_open_pr, gh_request_review,
|
||||
mark_unable_to_reproduce, fetch_issue_thread
|
||||
mark_unable_to_reproduce, abort_task, fetch_issue_thread
|
||||
sandbox.py clone pool + worktree lifecycle
|
||||
github_client.py typed httpx client; webhook payload parsing
|
||||
proxy_client.py GitHubProxyClient + HMAC signer
|
||||
|
||||
@@ -11,7 +11,7 @@ import json
|
||||
import logging
|
||||
import subprocess
|
||||
import time
|
||||
from collections.abc import Mapping
|
||||
from collections.abc import Callable, Mapping
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
@@ -36,6 +36,34 @@ _PRE_PR_CHECK_MAX_OUTPUT = 12_000
|
||||
_PRE_PR_FIX_COMMIT_SUBJECT = "style: bun run fix"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class AbortController:
|
||||
"""Mutable handoff between the `abort_task` host tool and the worker.
|
||||
|
||||
`signal()` is called from the host-tool thread to request an irrecoverable
|
||||
teardown of the omp subprocess. The worker pre-populates `stop` with a
|
||||
thread-safe terminator (the same one used for queue cancellation and the
|
||||
hard-timeout watchdog), and inspects `triggered` after `prompt_and_wait`
|
||||
unblocks to decide whether the resulting `RpcError` is an intentional
|
||||
abort (swallow, mark event `done`) vs an actual failure (propagate).
|
||||
"""
|
||||
|
||||
triggered: bool = False
|
||||
reason: str = ""
|
||||
stop: Callable[[], None] | None = None
|
||||
|
||||
def signal(self, reason: str) -> None:
|
||||
# Idempotent. Only the first call records its reason; later calls are
|
||||
# silent no-ops so a retry inside the tool can't overwrite the
|
||||
# original diagnosis with a generic follow-up message.
|
||||
if self.triggered:
|
||||
return
|
||||
self.triggered = True
|
||||
self.reason = reason
|
||||
if self.stop is not None:
|
||||
self.stop()
|
||||
|
||||
|
||||
@dataclass(slots=True, frozen=True)
|
||||
class ToolBindings:
|
||||
"""Per-task closure that the host tools capture."""
|
||||
@@ -63,6 +91,10 @@ class ToolBindings:
|
||||
# itself does not carry triage labels.
|
||||
inbound_is_pr: bool = False
|
||||
slot_uid: int | None = None
|
||||
# Set by the worker before launching omp. Carries the abort-task signal
|
||||
# back out to the worker; `None` for unit tests that exercise tools
|
||||
# without a live RpcClient.
|
||||
abort: AbortController | None = None
|
||||
|
||||
@property
|
||||
def issue_key(self) -> str:
|
||||
@@ -837,6 +869,40 @@ def _build_mark_unable(bindings: ToolBindings) -> HostTool[Any, Any]:
|
||||
)
|
||||
|
||||
|
||||
# ---------- abort_task ----------
|
||||
def _build_abort_task(bindings: ToolBindings) -> HostTool[Any, Any]:
|
||||
def execute(args: dict[str, Any], _ctx: HostToolContext[Any]) -> str:
|
||||
reason = args.get("reason")
|
||||
if not isinstance(reason, str) or not reason.strip():
|
||||
_raise_command("abort_task requires a non-empty 'reason' string.")
|
||||
reason = reason.strip()
|
||||
# Audit FIRST so the diagnosis is durable even if anything below
|
||||
# races against the imminent omp teardown.
|
||||
_audit(bindings, "abort_task", args, result={"reason": reason})
|
||||
log.warning(
|
||||
"task_aborted",
|
||||
extra={"issue": bindings.issue_key, "reason": reason},
|
||||
)
|
||||
bindings.db.set_issue_state(bindings.issue_key, "abandoned")
|
||||
if bindings.abort is not None:
|
||||
bindings.abort.signal(reason)
|
||||
return "aborted"
|
||||
|
||||
return host_tool(
|
||||
name="abort_task",
|
||||
description=persona.host_tool_description("abort_task"),
|
||||
parameters={
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"reason": {"type": "string"},
|
||||
},
|
||||
"required": ["reason"],
|
||||
"additionalProperties": False,
|
||||
},
|
||||
execute=execute,
|
||||
)
|
||||
|
||||
|
||||
# ---------- fetch_issue_thread ----------
|
||||
def _build_fetch_thread(bindings: ToolBindings) -> HostTool[Any, Any]:
|
||||
def execute(args: dict[str, Any], _ctx: HostToolContext[Any]) -> str:
|
||||
@@ -1030,6 +1096,7 @@ def _build_classify_issue(bindings: ToolBindings) -> HostTool[Any, Any]:
|
||||
bindings.workspace,
|
||||
branch_slug,
|
||||
pr_number=existing.pr_number if existing is not None else None,
|
||||
slot_uid=bindings.slot_uid,
|
||||
)
|
||||
except ValueError as exc:
|
||||
_audit(bindings, "classify_issue", args, error=str(exc))
|
||||
@@ -1125,8 +1192,9 @@ def build(bindings: ToolBindings) -> tuple[HostTool[Any, Any], ...]:
|
||||
_build_request_review(bindings),
|
||||
_build_repro_record(bindings),
|
||||
_build_mark_unable(bindings),
|
||||
_build_abort_task(bindings),
|
||||
_build_fetch_thread(bindings),
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["ToolBindings", "build"]
|
||||
__all__ = ["AbortController", "ToolBindings", "build"]
|
||||
|
||||
@@ -32,6 +32,12 @@ reproduced = "True when the recorded run demonstrates the bug."
|
||||
[mark_unable_to_reproduce]
|
||||
description = "Close the loop without a PR: comment with diagnosis + info request, mark issue abandoned."
|
||||
|
||||
[abort_task]
|
||||
description = "Irrecoverably abandon this task WITHOUT posting any visible message. Use ONLY for orchestrator/environment defects you cannot work around (broken filesystem permissions, missing system tools, corrupted git metadata, harness bugs). NEVER for normal workflow problems — failed builds, missing repro info, unclear requests use `gh_post_comment` or `mark_unable_to_reproduce` instead. `reason` is audit-only and NEVER shown to the reporter."
|
||||
|
||||
[abort_task.parameters]
|
||||
reason = "Internal diagnosis for the operator. Concrete, specific, blameless. NEVER shown to the reporter."
|
||||
|
||||
[fetch_issue_thread]
|
||||
description = "Refetch the originating issue and its comments. Use sparingly."
|
||||
|
||||
|
||||
@@ -46,7 +46,7 @@ NEVER apply `provider` or `platform` speculatively. They REQUIRE explicit eviden
|
||||
9. **Publish.** Call `gh_push_branch`, then `gh_open_pr`. Both deterministically run `bun run fix` (auto-committing as `style: bun run fix`) then `bun check` before touching the remote. The same gate runs on every follow-up `gh_push_branch`. The tools also refuse dirty trees and commit-author mismatches.
|
||||
- `bun check` failed? Fix at the source, commit, call again.
|
||||
- **Escape hatch — `skip_checks=true`.** ONLY for breakage you have VERIFIED is pre-existing on the default branch. Verify by running the same command against the same paths on a clean checkout of the default branch and confirming the identical failure. NEVER use it to bypass a failure your diff introduced, and NEVER for transient or unclear failures. Document the bypass in the PR's `## Verification` section, one sentence: ``bun check` fails on `main` for unrelated reason X; skipped pre-publish gate.`
|
||||
- **NEVER tamper with git internals.** No editing `.git`/`gitdir:` pointers, no chown/chmod on worktree files, no `safe.directory` overrides, no pointing HEAD at a fabricated commit. Push refused for reasons you cannot resolve? Ask the maintainer via `gh_post_comment`, or use `mark_unable_to_reproduce`. NEVER improvise.
|
||||
- **NEVER tamper with git internals.** No editing `.git`/`gitdir:` pointers, no chown/chmod on worktree files, no `safe.directory` overrides, no pointing HEAD at a fabricated commit. Push refused for reasons you cannot resolve? Ask the maintainer via `gh_post_comment`, or use `mark_unable_to_reproduce`. Environmental/orchestrator defect that's not the reporter's problem (broken permissions, corrupted git metadata, missing tools)? Call `abort_task` with the diagnosis — silent abandonment, no comment leaked to the reporter. NEVER improvise.
|
||||
- **Two-strikes rule.** Two consecutive `gh_push_branch` rejections with the same error is a workflow bug. Fix the cause, use `skip_checks=true` with justification, or escalate via `gh_post_comment`. NEVER loop.
|
||||
10. **Link.** After the PR opens, one final `gh_post_comment` linking it.
|
||||
|
||||
|
||||
+18
-2
@@ -35,7 +35,7 @@ from robomp.config import Settings
|
||||
from robomp.db import Database, issue_key
|
||||
from robomp.github_backend import GitHubBackend
|
||||
from robomp.github_client import CommentInfo, IssueInfo, RepoInfo
|
||||
from robomp.host_tools import ToolBindings
|
||||
from robomp.host_tools import AbortController, ToolBindings
|
||||
from robomp.sandbox import GitTransport, Workspace, _prepare_slot_tmpdir
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
@@ -403,6 +403,8 @@ def _run_rpc_blocking(
|
||||
RpcProcessExitError("cancelled by operator")
|
||||
)
|
||||
|
||||
if bindings.abort is not None:
|
||||
bindings.abort.stop = _cancel_hook
|
||||
register_cancel_hook(_cancel_hook)
|
||||
try:
|
||||
client.install_headless_ui()
|
||||
@@ -473,7 +475,20 @@ def _run_rpc_blocking(
|
||||
hard_timer.daemon = True
|
||||
hard_timer.start()
|
||||
try:
|
||||
turn = client.prompt_and_wait(prompt, timeout=settings.task_timeout_seconds)
|
||||
try:
|
||||
turn = client.prompt_and_wait(prompt, timeout=settings.task_timeout_seconds)
|
||||
except (RpcError, RpcProcessExitError):
|
||||
# Did the agent intentionally pull the plug via `abort_task`?
|
||||
# If so, swallow — the abort path is a clean exit, not a
|
||||
# failure that should surface in the dashboard or trigger
|
||||
# a comment to the reporter. Anything else propagates.
|
||||
if bindings.abort is not None and bindings.abort.triggered:
|
||||
log.info(
|
||||
"rpc_aborted_by_tool",
|
||||
extra={"issue": bindings.issue_key, "task": task_kind, "reason": bindings.abort.reason},
|
||||
)
|
||||
return None
|
||||
raise
|
||||
finally:
|
||||
hard_timer.cancel()
|
||||
if hard_timeout_fired.is_set():
|
||||
@@ -517,6 +532,7 @@ async def run_task(
|
||||
inbound_thread_number=pr_number,
|
||||
inbound_is_pr=pr_number is not None,
|
||||
slot_uid=inputs.slot_uid,
|
||||
abort=AbortController(),
|
||||
)
|
||||
resuming = _has_prior_session(inputs.workspace.session_dir)
|
||||
prompt = _build_prompt(
|
||||
|
||||
@@ -14,7 +14,7 @@ from omp_rpc import HostToolContext, RpcCommandError
|
||||
|
||||
from robomp.db import Database
|
||||
from robomp.github_client import GitHubClient, IssueInfo, RepoInfo
|
||||
from robomp.host_tools import ToolBindings, build
|
||||
from robomp.host_tools import AbortController, ToolBindings, build
|
||||
from robomp.sandbox import LocalGitTransport, Workspace
|
||||
|
||||
|
||||
@@ -279,6 +279,78 @@ def test_mark_unable_posts_comment_and_abandons(db: Database, tmp_path: Path) ->
|
||||
assert issue and issue.state == "abandoned"
|
||||
|
||||
|
||||
def test_abort_task_signals_controller_and_abandons_without_comment(db: Database, tmp_path: Path) -> None:
|
||||
# Any HTTP call is a regression: abort_task MUST NOT touch GitHub.
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
raise AssertionError(f"abort_task issued an HTTP request to {request.url}")
|
||||
|
||||
bindings, loop, t = _bindings(db, tmp_path, httpx.MockTransport(handler))
|
||||
controller = AbortController()
|
||||
stops: list[None] = []
|
||||
controller.stop = lambda: stops.append(None)
|
||||
# Frozen dataclass — rebuild with the controller attached.
|
||||
bindings = ToolBindings(
|
||||
db=bindings.db,
|
||||
github=bindings.github,
|
||||
git_transport=bindings.git_transport,
|
||||
repo=bindings.repo,
|
||||
issue=bindings.issue,
|
||||
workspace=bindings.workspace,
|
||||
loop=bindings.loop,
|
||||
author_name=bindings.author_name,
|
||||
author_email=bindings.author_email,
|
||||
settings=bindings.settings,
|
||||
inbound_thread_number=bindings.inbound_thread_number,
|
||||
inbound_is_pr=bindings.inbound_is_pr,
|
||||
slot_uid=bindings.slot_uid,
|
||||
abort=controller,
|
||||
)
|
||||
try:
|
||||
tool = next(x for x in build(bindings) if x.name == "abort_task")
|
||||
result = tool.execute({"reason": "ref dir owned by foreign uid; git commit cannot lock HEAD"}, _ctx())
|
||||
finally:
|
||||
_stop_loop(loop, t)
|
||||
assert result == "aborted"
|
||||
assert controller.triggered
|
||||
assert "foreign uid" in controller.reason
|
||||
assert len(stops) == 1, "stop callback must fire exactly once"
|
||||
issue = db.get_issue(bindings.issue_key)
|
||||
assert issue and issue.state == "abandoned"
|
||||
# Audit row records the call. Use raw SQL because `Database` exposes a
|
||||
# writer but no reader for `tool_calls` — the dashboard reads via SQL too.
|
||||
with db._lock: # noqa: SLF001 - test-only inspection
|
||||
row = db._conn.execute( # noqa: SLF001
|
||||
"SELECT tool, args_json FROM tool_calls WHERE issue_key=? AND tool=?",
|
||||
(bindings.issue_key, "abort_task"),
|
||||
).fetchone()
|
||||
assert row is not None
|
||||
assert "foreign uid" in row["args_json"]
|
||||
|
||||
|
||||
def test_abort_task_rejects_empty_reason(db: Database, tmp_path: Path) -> None:
|
||||
bindings, loop, t = _bindings(db, tmp_path, httpx.MockTransport(lambda r: httpx.Response(500)))
|
||||
try:
|
||||
tool = next(x for x in build(bindings) if x.name == "abort_task")
|
||||
with pytest.raises(RpcCommandError):
|
||||
tool.execute({"reason": " "}, _ctx())
|
||||
finally:
|
||||
_stop_loop(loop, t)
|
||||
# No state change on rejected validation.
|
||||
issue = db.get_issue(bindings.issue_key)
|
||||
assert issue and issue.state == "reproducing"
|
||||
|
||||
|
||||
def test_abort_task_signal_is_idempotent(db: Database, tmp_path: Path) -> None:
|
||||
controller = AbortController()
|
||||
fires: list[None] = []
|
||||
controller.stop = lambda: fires.append(None)
|
||||
controller.signal("first")
|
||||
controller.signal("second")
|
||||
assert controller.triggered
|
||||
assert controller.reason == "first" # second call must not overwrite
|
||||
assert len(fires) == 1, "stop must not be called again after the first abort"
|
||||
|
||||
|
||||
def test_fetch_issue_thread_returns_markdown(db: Database, tmp_path: Path) -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
if request.url.path.endswith("/comments"):
|
||||
|
||||
@@ -124,6 +124,7 @@ def _make_inputs(
|
||||
repo=repo,
|
||||
issue=issue,
|
||||
issue_key=f"{repo.full_name}#{issue.number}",
|
||||
abort=None,
|
||||
)
|
||||
return inputs, bindings
|
||||
|
||||
|
||||
Reference in New Issue
Block a user