Files
cleveragents-core/tools/controller/worker/opencode_session.py
T
drew bb8d872ba7 feat(controller): Phase 1d-4 — OpenCode session adapter
The production wrapper that wires the existing
_opencode_worker.run_session_blocking into the controller's
run_opencode_session protocol that production_agent_runner expects.

tools/controller/worker/opencode_session.py:

- agent_name_for(role, tier): role+tier → OpenCode agent name.
  - implementer + tier {0,1,2} → task-implementor-tier-{0,1,2}
  - reviewer → pr-review-worker
  - estimator → estimator-implementation
  - conflict_resolver → conflict-resolver-worker
  - summarizer → controller-summarizer (new agent name)

- DEFAULT_TIER_TIMEOUT_S: per-tier wallclock budgets via env vars.
  Defaults match plan v9 (tier 0: 600s, tier 1: 1200s, tier 2: 1800s).
  Reviewer/estimator/conflict use the default-agent timeout (600s).

- wire_opencode_session(opencode_server_url, ...) → callable matching
  the production_agent_runner's run_opencode_session contract.
  - Resolves role+tier to agent name + per-tier timeout
  - Builds the prompt (default stub or custom builder; per-role prompt
    templates are a follow-up)
  - Calls run_session_blocking with an on_poll callback that raises
    WorkerLostLock if the controller's lost_lock_check returns True
    mid-session (heartbeat thread reaper detected stolen lock)
  - Routes SessionResult.status:
    - 'completed' → return None (MCP's canonical output is in the
      tempfile; production_agent_runner reads it after we return)
    - 'timeout' → raise WorkerError(worker-internal-error)
    - 'transport-error' → raise WorkerError(worker-internal-error,
      including error_kind for forensics)
    - unknown → raise WorkerError defensively
  - Any unexpected exception from run_session_blocking wrapped as
    WorkerError(worker-internal-error).

- Production wiring: dependency injection. wire_opencode_session()
  lazy-imports the real run_session_blocking when not overridden.
  Tests inject a stub.

22 new tests in test_worker_opencode_session.py:
- agent_name_for: parametrized per role (4 implementer-tier rows +
  4 flat-role rows) + invalid tier + unknown role
- DEFAULT_TIER_TIMEOUT_S: each tier has a timeout + monotonic
- wire_opencode_session: completed/timeout/transport-error/unknown-
  status routing; run_session_blocking raises wrapped as WorkerError;
  on_poll propagates lost_lock; on_poll no-op when lock held;
  custom prompt_builder used; default prompt includes role + tier
  + finalize() instruction + rendered input_payload; tag prefix
  customizable.

Total: 398 controller tests; full auto_agents suite 2760 pass.
2026-05-18 14:09:34 -04:00

202 lines
7.7 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.
def agent_name_for(role: str, tier: int | 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":
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 _default_prompt_for(role: str, tier: int | None, input_payload: dict[str, Any]) -> str:
"""Minimal v1 prompt: tells the LLM to use the builder MCP.
Real per-role prompt assembly (with the worker contract fields
rendered for the LLM's consumption) is a follow-up. This stub
keeps the adapter complete — the agent's own system prompt
(in `.opencode/agents/{name}.md`) carries the instructions for
HOW to use the MCP.
"""
import json
head = f"You are running as the {role} worker"
if tier is not None:
head += f" at tier {tier}"
head += ".\n\n"
head += (
f"Use the response-builder MCP tools to construct your "
f"{role}OutputV1 output. Call {role}_finalize() when done.\n\n"
)
head += "Input payload (V1 contract):\n"
head += "```json\n"
head += json.dumps(input_payload, indent=2, sort_keys=True, default=str)[:8000]
head += "\n```"
return head
# 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],
) -> None:
"""The injected callable. Drives the OpenCode session +
propagates lost_lock_check via the on_poll callback."""
agent = agent_name_for(role, tier)
prompt = builder(role, tier, input_payload)
tag = f"{tag_prefix}-{role}-{attempt_id}"
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
# 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",
]