Files
cleveragents-core/tools/controller/worker/opencode_session.py
T
drew 0db0a15dad feat(controller): RUN_CI_LOCAL verdict source, ci-not-ready outcome, escalation hardening
Adds RUN_CI_LOCAL — an on-demand local-CI verdict source for when the
cluster's Forgejo CI is broken — plus robustness fixes, the telemetry
Live-tab rewrite, and PR-level cost attribution.

Controller:
- RUN_CI_LOCAL: the master swaps its Forgejo CI callbacks for local
  `forgejo-runner exec` runs (tools/run-ci-full-local.sh + local_ci.py).
  Async per-(owner,repo,SHA) on-disk job cache; preflights the
  forgejo-runner binary + Docker daemon at startup (fail loud, not a
  red verdict on every PR); GCs finished run dirs + per-run actcache.
- ci-not-ready implementer outcome + implementer_ci_not_ready event:
  an implementer that runs before the on-demand verdict exists parks
  in AWAITING_CI instead of dead-ending at STUCK; capped against
  ci_red ping-pong.
- ci_poll_exhaustion skips its sweep while a local CI run is in
  flight, so AWAITING_CI workflows queued behind on-demand CI are not
  STUCK'd by the remote-CI-sized timeout.
- Escalate the workflow after repeated worker-internal-error at a
  tier instead of retrying to pickup-exhaustion -> STUCK.
- forgejo_http: normalise Forgejo's per-gate `status` key to `state`
  so failing gates are actually counted (they previously all read as
  pending).
- Per-tier worker timeouts bumped +15 min; a timed-out attempt's
  dirty-worktree residue is preserved on auto-scratch/pr-<N> before
  the next attempt's reset.
- Implementer agents now verify only the CI-flagged gate(s) via a
  targeted re-run rather than the full local battery before claiming
  resolved/noop. Re-running the whole suite CI will run anyway was the
  #1 cause of implementer timeouts; CI remains the real gate and
  re-dispatches the implementer on red.

Telemetry:
- Live tab rebuilt on /api/live (controller DB run state + the live
  OpenCode session forest) after the live_log_writer sidecar was
  retired with the legacy dispatchers.
- Durable per-attempt input/output payloads surfaced in the Live
  drill-down, archived-session detail, and Workflows timeline.
- PR-level cost attribution: worker session tags carry -pr-<n>;
  backfill_llm_activity_pr.py repairs rows written before the fix.

Shared:
- tools/controller/session_tag.py — one canonical controller-tag
  parser shared by the telemetry server and the backfill.

Tests: new coverage for local_ci (state machine, log parsing,
_summarize_run, GC, in-flight probe, preflight), the ci-not-ready
path, the escalation/ci-not-ready SQL counters, the ci_poll
in-flight skip, and CI-status payload parsing across both sources.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-20 11:49:05 -04:00

266 lines
10 KiB
Python

"""OpenCode session adapter — wraps the existing
``_opencode_worker.run_session_blocking`` into the
``run_opencode_session`` protocol the controller's
``production_agent_runner`` expects.
Per plan v9: each attempt spawns a per-attempt MCP subprocess +
drives the LLM via OpenCode pointed at the matching agent. The
controller's responsibility is to map (role, tier) → agent name
and to build the prompt that tells the LLM "use the response-builder
MCP to construct your output."
Role-and-tier-to-agent mapping (matches existing pipeline names):
- implementer + tier 0 → task-implementor-tier-0
- implementer + tier 1 → task-implementor-tier-1
- implementer + tier 2 → task-implementor-tier-2
- reviewer → pr-review-worker
- estimator → estimator-implementation
- conflict_resolver → conflict-resolver-worker
- summarizer → controller-summarizer (new agent name; defined elsewhere)
This module is the wiring layer; the prompt-assembly logic itself
(turning input_payload + role into a worker-friendly text prompt
that instructs the LLM to use the MCP tools) lives in a follow-up
since it requires per-role prompt templates.
"""
from __future__ import annotations
import logging
import os
import sys
import time
from collections.abc import Callable
from pathlib import Path
from typing import Any
from .runner import WorkerError, WorkerLostLock
logger = logging.getLogger(__name__)
# Plan v9: role+tier → OpenCode agent name. Tier-aware for the
# implementer; flat for the rest.
#
# T5-4 (2026-05-19): reviewer routing branches on the input_payload's
# ``implementer_claim``. When the prior implementer attempt emitted
# ``outcome=dispute-reviewer``, the next reviewer attempt is a
# re-examination and must be at least as capable as the disputer
# (which is opus at tier-2). The dispute variant uses the
# ``pr-review-worker-dispute`` agent (also opus); the normal path
# uses ``pr-review-worker`` (sonnet baseline).
def agent_name_for(
role: str,
tier: int | None,
input_payload: dict[str, Any] | None = None,
) -> str:
if role == "implementer":
if tier is None:
raise ValueError("implementer role requires tier")
if tier not in {0, 1, 2}:
raise ValueError(f"invalid tier {tier!r}; must be 0/1/2")
return f"task-implementor-tier-{tier}"
if role == "reviewer":
if _is_dispute_reexamination(input_payload):
return "pr-review-worker-dispute"
return "pr-review-worker"
if role == "estimator":
return "estimator-implementation"
if role == "conflict_resolver":
return "conflict-resolver-worker"
if role == "summarizer":
return "controller-summarizer"
raise ValueError(f"unknown role {role!r}")
def _is_dispute_reexamination(input_payload: dict[str, Any] | None) -> bool:
"""Return True iff the reviewer's input_payload signals this attempt
is a re-examination of a prior implementer's dispute.
Cheap, defensive: ``input_payload`` may be None during test setup
or for early-init paths; either way, default to the non-dispute
branch.
"""
if not isinstance(input_payload, dict):
return False
claim = input_payload.get("implementer_claim")
if not isinstance(claim, dict):
return False
return claim.get("outcome") == "dispute-reviewer"
def _default_prompt_for(
role: str, tier: int | None, input_payload: dict[str, Any]
) -> str:
"""Per-role prompt builder; delegates to ``prompts.build_prompt``."""
from .prompts import build_prompt
return build_prompt(role, tier, input_payload)
# Default tier wallclock budgets (seconds). Per plan v9. Operator
# tunes via env at worker startup.
DEFAULT_TIER_TIMEOUT_S: dict[int | None, int] = {
0: int(os.environ.get("CONTROLLER_TIER_0_TIMEOUT_S", "600")),
1: int(os.environ.get("CONTROLLER_TIER_1_TIMEOUT_S", "1200")),
2: int(os.environ.get("CONTROLLER_TIER_2_TIMEOUT_S", "1800")),
None: int(os.environ.get("CONTROLLER_DEFAULT_AGENT_TIMEOUT_S", "600")),
}
def wire_opencode_session(
*,
opencode_server_url: str,
tag_prefix: str = "controller",
prompt_builder: Callable[[str, int | None, dict], str] | None = None,
run_session_blocking: Callable | None = None,
lost_lock_poll_interval_s: float = 5.0,
) -> Callable:
"""Returns a callable matching the production_agent_runner's
``run_opencode_session`` contract.
Args:
opencode_server_url: e.g. "http://localhost:4096".
tag_prefix: prefix for the OpenCode session title (operator-
visible in OpenCode's session list).
prompt_builder: function (role, tier, input_payload) → prompt
text. None uses the default ``_default_prompt_for``.
run_session_blocking: dependency-injection for the real
``_opencode_worker.run_session_blocking``. None uses the
real function. Tests inject a stub.
lost_lock_poll_interval_s: how often the OpenCode polling
callback checks lost_lock_check.
"""
if run_session_blocking is None:
# Lazy-import the real one. Done only when wire_opencode_session
# is called without an override (i.e., in production).
repo_root = Path(__file__).resolve().parents[3]
if str(repo_root) not in sys.path:
sys.path.insert(0, str(repo_root))
tools_dir = repo_root / "tools"
sys.path.insert(0, str(tools_dir))
from _opencode_worker import run_session_blocking as _real # type: ignore[import-not-found]
run_session_blocking = _real
builder = prompt_builder or _default_prompt_for
def run_opencode_session(
*,
role: str,
tier: int | None,
input_payload: dict[str, Any],
mcp_process,
attempt_id: int,
instance_id: str,
lost_lock_check: Callable[[], bool],
inline_output_callback: Callable[[dict[str, Any]], None] | None = None,
) -> None:
"""The injected callable. Drives the OpenCode session +
propagates lost_lock_check via the on_poll callback."""
agent = agent_name_for(role, tier, input_payload)
prompt = builder(role, tier, input_payload)
tag = f"{tag_prefix}-{role}-{attempt_id}"
# Append the PR component when this attempt is PR-bound. The
# session archive carries the tag verbatim, and the
# llm_activity scraper parses ``-pr-<n>`` out of it to attribute
# token cost to a PR in the telemetry cost dashboard. Without
# the suffix every controller turn lands as "unattributed".
# An estimator running on a pre-PR issue has no pr_number — the
# suffix is simply omitted and that row carries NULL, which is
# the correct "no PR to attribute to" outcome.
pr_number = (
input_payload.get("pr_number") if isinstance(input_payload, dict) else None
)
if isinstance(pr_number, int) and not isinstance(pr_number, bool):
tag = f"{tag}-pr-{pr_number}"
timeout = DEFAULT_TIER_TIMEOUT_S.get(tier, DEFAULT_TIER_TIMEOUT_S[None])
# The on_poll hook lets us bail out early if the heartbeat
# thread detects the lock has been reaped. Polling happens
# every poll_interval_seconds inside run_session_blocking.
def on_poll() -> None:
if lost_lock_check():
# Raise so run_session_blocking unwinds; the runner
# catches WorkerLostLock and aborts silently.
raise WorkerLostLock(
f"lost lock for attempt {attempt_id} during OpenCode poll"
)
try:
result = run_session_blocking(
server_url=opencode_server_url,
agent=agent,
tag=tag,
prompt=prompt,
timeout_seconds=timeout,
poll_interval_seconds=lost_lock_poll_interval_s,
on_poll=on_poll,
)
except WorkerLostLock:
raise
except Exception as exc:
raise WorkerError(
f"OpenCode run_session_blocking raised: {exc}",
outcome="worker-internal-error",
) from exc
# Phase 1k++++ trial-path: harvest the inline JSON the agent
# emitted as its final response. Legacy pipeline agents emit a
# single JSON object as their final message; OpenCode worker
# extracts it into ``SessionResult.parsed_json``. Since
# Phase 1m the response-builder MCPs ARE registered in
# opencode.json, so this inline channel is a fallback for
# agents that still emit chat-JSON. The agent_runner's
# callback skips the write if the MCP already wrote
# canonical V1 to the same path.
#
# PD16: re-raise WorkerLostLock from the callback so the
# runner's outer handler aborts properly. Other exceptions are
# converted to WorkerError so they classify as
# ``worker-internal-error`` instead of silently letting the
# canonical poller time out 30s later with no root cause.
parsed = getattr(result, "parsed_json", None)
if parsed is not None and inline_output_callback is not None:
try:
inline_output_callback(parsed)
except WorkerLostLock:
raise
except Exception as exc:
raise WorkerError(
f"inline_output_callback raised: {exc}",
outcome="worker-internal-error",
) from exc
# Inspect SessionResult.status.
status = getattr(result, "status", None)
if status == "completed":
# MCP's finalize emitted to the canonical-output file;
# production_agent_runner reads it after we return.
return
if status == "timeout":
raise WorkerError(
f"OpenCode session timed out after {timeout}s",
outcome="worker-internal-error",
)
if status == "transport-error":
error_kind = getattr(result, "error_kind", "transport-error")
raise WorkerError(
f"OpenCode transport error: {error_kind}",
outcome="worker-internal-error",
)
# Unknown status — be defensive.
raise WorkerError(
f"OpenCode session returned unexpected status: {status!r}",
outcome="worker-internal-error",
)
return run_opencode_session
__all__ = [
"DEFAULT_TIER_TIMEOUT_S",
"agent_name_for",
"wire_opencode_session",
]