Files
cleveragents-core/tools/controller/master/loop.py
T
drew f57d9f9478 fix(controller): batch L — round-4 trial-blockers (A1, P1, P3, P4, P5, T5)
Round-4 adversarial review found 5 trial-blockers + 1 silent-debt
item the post-round-3 deep pass missed. All fixed.

A1 — pre-clone the workspace so the agent has a worktree to operate on
``worker/__main__.py``: the agent_runner closure now constructs a
``PerPRWorkspace`` from input_payload.owner/repo/pr_number + the
FORGEJO_URL+FORGEJO_TOKEN env vars. Pre-flight:
- ``workspace.ensure_present()`` creates the dir skeleton.
- ``workspace.clone_if_absent()`` clones the repo into
  ``{workspace_dir}/worktree/`` if not already present (idempotent).
- ``workspace.fetch_and_validate(head_sha, head_ref)`` refreshes +
  verifies the workspace is at the expected head. ``StaleInputError``
  → ``WorkerError(outcome='stale-input')`` so the master re-prefetches
  without burning a pickup. ``RuntimeError`` → ``worker-internal-error``.
Previously the agent saw an empty workspace_dir + had no repo.

P1 — partial-write defense in the canonical-output poller
``worker/agent_runner.py:_wait_for_canonical_output`` now polls each
path with a two-pass quiescence check (size stable + content parses
as JSON) before returning. Partial writes (agent crashed mid-flush)
are skipped + the polling loop continues. The previous
``f.read().strip()`` returned partial JSON which then tripped
``ContractValidationError`` → ``worker-internal-error`` with no
record of WHICH path; now logs source path on every read.

P3 — TOCTOU defense in promote_discovered
``master/promote.py``: the UPDATE now filters
``current_state='DISCOVERED'``. If a concurrent reconciliation
moved the row off DISCOVERED between SELECT and UPDATE, rowcount=0
+ we skip the event-row write. No duplicate audit entry; no
overwriting a pause-by-label-removal.

P4 — explicit tuple-length validation in reconciliation_args + discovery_args
``master/loop.py``: previously a 6-tuple silently fell into the
``else`` 4-tuple unpack, raised ValueError("too many values"), got
swallowed by the per-iter ``except Exception``, and reconciliation
silently died forever. Now: ``elif n == 4`` + ``else: raise TypeError``.
The TypeError still hits the per-iter except (so the loop doesn't
crash) but ``logger.exception`` surfaces the actionable message in
journald. Operator sees "reconciliation_args must be a 4- or 5-tuple;
got length 6" instead of zero indication.

P5 — --tick-interval CLI flag preserves other config fields
``master/__main__.py``: replaced the manual ``MasterConfig(...)``
rebuild (which dropped reconciliation/ci_poll/discovery intervals)
with ``dataclasses.replace(cfg_loop, tick_interval_s=args.tick_interval)``.
Operators who pass --tick-interval no longer silently revert the
other intervals to defaults.

T5 — scheduler._commit_escalation uses safe_json_dumps
``master/scheduler.py``: the escalation event row's payload was the
only call site that bypassed safe_json_dumps. Now consistent — a
future contributor adding a datetime/Decimal field won't trip raw
json.dumps at runtime.

Tests (+4 net):
- ``test_worker_agent_runner.py::test_partial_write_not_read``: pins
  P1 (truncated fallback file + valid MCP output → MCP wins).
- ``test_master_promote.py::test_toctou_state_change_between_select_and_update``:
  pins P3 (steal state via monkey-patch → no double-promotion, no
  extra event row).
- ``test_master_loop.py::test_reconciliation_args_wrong_length_logs_not_silent``:
  pins P4 (6-tuple → logged error, not silent forever).
- ``test_entry_points.py::test_tick_interval_flag_preserves_other_cfg_fields``:
  pins P5 (env-set non-default intervals survive --tick-interval).

Total: 711 controller tests pass (+4 net), 0 regressions.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-18 17:23:35 -04:00

377 lines
16 KiB
Python

"""Master main loop — composes the per-tick work (tick + reaper +
pickup guard) and runs it on a configurable cadence until a stop
event fires.
Per plan v9: the master is a long-running singleton per (owner, repo).
This module is the orchestrator that ties together the deterministic
pieces shipped in Phase 1d-1/1d-2 + later additions.
Currently composes:
- ``tick.run_tick()`` — advance state for any completed attempts
- ``reaper.reap_stale_attempts()`` — reset stale-heartbeat in_progress rows
- ``pickup_guard.transition_exhausted_to_stuck()`` — STUCK workflows
whose attempts have been re-pended too many times
Deferred to Phase 1d-3+:
- Discovery (Forgejo poll for new PRs/issues + insert DISCOVERED rows)
- Per-workflow scheduling (after a state transition, enqueue the
next attempt's workflow_attempts row with input_payload prefetched)
- Forgejo writes (status comments, labels, merges)
- MERGING state's actual Forgejo merge call
- Periodic reconciliation tick
- Backfill at startup
- Operator CLI server (HTTP / unix socket for controller-cli)
"""
from __future__ import annotations
import logging
import os
import threading
from collections.abc import Callable
from dataclasses import dataclass
from sqlalchemy.engine import Engine
from ..pickup_guard import (
DEFAULT_MAX_PICKUPS,
PickupGuardReport,
transition_exhausted_to_stuck,
)
from ..reaper import ReaperReport, reap_stale_attempts
from .ci_poll import CIPollExhaustionReport, run_ci_poll_exhaustion_tick
from .discovery import DiscoveryReport, run_discovery
from .merging import MergeCallback, MergingHandlerReport, run_merging_tick
from .promote import PromoteDiscoveredReport, run_promote_discovered_tick
from .reconciliation import (
GetIssueStateCallback,
GetPRStateCallback,
ReconciliationReport,
run_reconciliation_tick,
)
from .scheduler import (
PrefetchCallback,
SchedulerReport,
schedule_next_attempts,
)
from .tick import TickReport, run_tick
logger = logging.getLogger(__name__)
@dataclass
class MasterConfig:
"""Per-master config; tunable via env vars."""
tick_interval_s: float = float(
os.environ.get("CONTROLLER_MASTER_TICK_INTERVAL_S", "30")
)
reaper_interval_s: float = float(
os.environ.get("CONTROLLER_REAPER_INTERVAL_S", "60")
)
reconciliation_interval_s: float = float(
os.environ.get("CONTROLLER_RECONCILIATION_INTERVAL_S", "300")
)
# CI poll-exhaustion runs on the same cadence as reconciliation
# by default — both are slow ticks that scan all non-terminal
# workflows.
ci_poll_exhaustion_interval_s: float = float(
os.environ.get("CONTROLLER_CI_POLL_EXHAUSTION_INTERVAL_S", "300")
)
# Periodic discovery — how often the master polls Forgejo for
# new PRs/issues. Without this, only PRs that existed at master
# startup get discovered (backfill is one-shot). Default 30s.
discovery_interval_s: float = float(
os.environ.get("CONTROLLER_DISCOVERY_INTERVAL_S", "30")
)
pickup_guard_max_pickups: int = DEFAULT_MAX_PICKUPS
@dataclass
class MasterTickReport:
"""Summary of one composite master tick."""
tick: TickReport
reaper: ReaperReport
pickup_guard: PickupGuardReport
reconciliation: ReconciliationReport | None = None
ci_poll_exhaustion: CIPollExhaustionReport | None = None
promote_discovered: PromoteDiscoveredReport | None = None
scheduler: SchedulerReport | None = None
merging: MergingHandlerReport | None = None
discovery: DiscoveryReport | None = None
def run_master_iteration(
engine: Engine, *, max_pickups: int = DEFAULT_MAX_PICKUPS,
) -> MasterTickReport:
"""Run one composite iteration: tick + reaper + pickup guard.
Order matters:
1. ``tick`` first: advance state machine for completed attempts;
may produce new transitions that the reaper / pickup guard
then notice.
2. ``reaper`` next: reset stale-heartbeat in_progress rows.
Post-reap, those attempts return to the pending pool +
pickup_count is preserved (the guard uses it).
3. ``pickup_guard`` last: STUCK any workflows whose pending
attempts have hit MAX_PICKUPS. Runs AFTER the reaper so a
just-reaped attempt's pickup_count is visible.
"""
return MasterTickReport(
tick=run_tick(engine),
reaper=reap_stale_attempts(engine),
pickup_guard=transition_exhausted_to_stuck(engine, max_pickups=max_pickups),
)
def master_main_loop(
engine: Engine,
*,
config: MasterConfig | None = None,
stop_event: threading.Event | None = None,
on_iteration: Callable[[MasterTickReport], None] | None = None,
reconciliation_args: tuple | None = None,
# If set: 4- OR 5-tuple — (owner, repo, get_pr_state, get_issue_state)
# OR (owner, repo, get_pr_state, get_issue_state, recon_kwargs_dict).
# The reconciliation tick fires every reconciliation_interval_s.
# None disables reconciliation (useful for tests that don't need it).
# recon_kwargs_dict (Phase 1k+) is passed through as keyword args
# to run_reconciliation_tick — used for opt_in_label /
# require_opt_in_label settings.
prefetch: PrefetchCallback | None = None,
# When set, the loop runs the DISCOVERED→ANALYZING promoter +
# the per-workflow scheduler every iteration. Without it, the
# scheduler is skipped — workflows would still transition between
# states but never get fresh worker attempts enqueued. Production
# __main__.py always passes a prefetch callback; tests may pass
# None to skip the scheduling layer.
merging_args: tuple | None = None,
# If set: (owner, repo, merge_callback). The MERGING handler
# fires every iteration (cheap if no workflows in MERGING).
# Without this, workflows that transition to MERGING never have
# the Forgejo merge call invoked — they sit in MERGING forever.
discovery_args: tuple | None = None,
# If set: 4- OR 5-tuple — (owner, repo, list_prs, list_issues)
# OR (owner, repo, list_prs, list_issues, discovery_kwargs).
# Periodic discovery fires every discovery_interval_s. Without
# this, only PRs that existed at master startup (via backfill)
# are ever managed — PRs created after master startup wait until
# the master restarts.
) -> None:
"""Run the master loop until ``stop_event`` is set.
Different ticks at different cadences:
- tick (state machine + transitions): every tick_interval_s (30s)
- reaper (stale-heartbeat reset): every reaper_interval_s (60s)
- reconciliation (Forgejo sync): every reconciliation_interval_s (300s)
- pickup guard: every iteration (cheap)
"""
cfg = config or MasterConfig()
stop = stop_event or threading.Event()
last_reap_at_iteration = 0
last_reconcile_at_iteration = 0
last_ci_poll_exhaustion_at_iteration = 0
last_discovery_at_iteration = 0
iteration = 0
logger.info(
"master loop starting: tick=%.1fs reaper=%.1fs reconcile=%.1fs "
"ci_poll_exh=%.1fs discovery=%.1fs max_pickups=%d",
cfg.tick_interval_s, cfg.reaper_interval_s,
cfg.reconciliation_interval_s, cfg.ci_poll_exhaustion_interval_s,
cfg.discovery_interval_s, cfg.pickup_guard_max_pickups,
)
try:
while not stop.is_set():
iteration += 1
# Always run tick + pickup guard. Run reaper + reconciliation
# less often per their own intervals.
tick_report = run_tick(engine)
should_reap = (
(iteration - last_reap_at_iteration) * cfg.tick_interval_s
>= cfg.reaper_interval_s
)
reaper_report = (
reap_stale_attempts(engine) if should_reap else ReaperReport()
)
if should_reap:
last_reap_at_iteration = iteration
should_reconcile = (
reconciliation_args is not None
and (iteration - last_reconcile_at_iteration) * cfg.tick_interval_s
>= cfg.reconciliation_interval_s
)
reconciliation_report: ReconciliationReport | None = None
if should_reconcile and reconciliation_args is not None:
try:
n = len(reconciliation_args)
if n == 5:
(owner, repo, get_pr_state, get_issue_state,
extra_kwargs) = reconciliation_args
elif n == 4:
(owner, repo, get_pr_state,
get_issue_state) = reconciliation_args
extra_kwargs = {}
else:
# R-round4 P4: explicit length check. A 6+ tuple
# used to silently fall to the `else` branch +
# raise ValueError("too many values to unpack")
# which the outer except Exception swallowed
# silently every iteration.
raise TypeError(
f"reconciliation_args must be a 4- or 5-tuple; "
f"got length {n}"
)
reconciliation_report = run_reconciliation_tick(
engine, owner=owner, repo=repo,
get_pr_state=get_pr_state,
get_issue_state=get_issue_state,
**(extra_kwargs or {}),
)
except Exception:
logger.exception("reconciliation tick raised; continuing")
last_reconcile_at_iteration = iteration
pickup_report = transition_exhausted_to_stuck(
engine, max_pickups=cfg.pickup_guard_max_pickups,
)
# CI poll-exhaustion: STUCK any AWAITING_CI workflow whose
# entered_state_at is older than the threshold. Runs on
# its own cadence (default = reconciliation cadence).
should_ci_poll_exh = (
(iteration - last_ci_poll_exhaustion_at_iteration)
* cfg.tick_interval_s
>= cfg.ci_poll_exhaustion_interval_s
)
ci_poll_report: CIPollExhaustionReport | None = None
if should_ci_poll_exh:
try:
ci_poll_report = run_ci_poll_exhaustion_tick(engine)
except Exception:
logger.exception(
"ci_poll_exhaustion tick raised; continuing"
)
last_ci_poll_exhaustion_at_iteration = iteration
# Phase 1k+++ (real-run): promote DISCOVERED → ANALYZING +
# schedule next attempts. Without these two ticks the
# master never enqueues anything for workers to pick up.
# Runs every iteration (both are cheap).
promote_report: PromoteDiscoveredReport | None = None
scheduler_report: SchedulerReport | None = None
if prefetch is not None:
try:
promote_report = run_promote_discovered_tick(engine)
except Exception:
logger.exception(
"promote_discovered tick raised; continuing"
)
try:
scheduler_report = schedule_next_attempts(
engine, prefetch=prefetch,
)
except Exception:
logger.exception(
"scheduler tick raised; continuing"
)
# Phase 1k+++ (real-run): MERGING handler. For workflows
# in MERGING state, call the Forgejo merge endpoint via
# the injected callback. Without this, workflows that
# transition to MERGING (via reviewer approval) never
# have the actual merge call invoked. Cheap when no
# workflows are in MERGING.
merging_report: MergingHandlerReport | None = None
if merging_args is not None:
try:
owner, repo, merge_cb = merging_args
merging_report = run_merging_tick(
engine, merge=merge_cb,
owner=owner, repo=repo,
)
except Exception:
logger.exception(
"merging tick raised; continuing"
)
# Phase 1k+++ (real-run): periodic discovery. Backfill is
# one-shot at startup; new PRs created later need this
# tick to be picked up. Runs every discovery_interval_s.
discovery_report: DiscoveryReport | None = None
should_discover = (
discovery_args is not None
and (iteration - last_discovery_at_iteration)
* cfg.tick_interval_s
>= cfg.discovery_interval_s
)
if should_discover and discovery_args is not None:
try:
n = len(discovery_args)
if n == 5:
(d_owner, d_repo, d_list_prs, d_list_issues,
d_kwargs) = discovery_args
elif n == 4:
(d_owner, d_repo, d_list_prs,
d_list_issues) = discovery_args
d_kwargs = {}
else:
raise TypeError(
f"discovery_args must be a 4- or 5-tuple; "
f"got length {n}"
)
discovery_report = run_discovery(
engine, owner=d_owner, repo=d_repo,
list_prs=d_list_prs,
list_issues=d_list_issues,
**(d_kwargs or {}),
)
except Exception:
logger.exception(
"discovery tick raised; continuing"
)
last_discovery_at_iteration = iteration
if on_iteration is not None:
try:
on_iteration(MasterTickReport(
tick=tick_report,
reaper=reaper_report,
pickup_guard=pickup_report,
reconciliation=reconciliation_report,
ci_poll_exhaustion=ci_poll_report,
promote_discovered=promote_report,
scheduler=scheduler_report,
merging=merging_report,
discovery=discovery_report,
))
except Exception:
logger.exception("on_iteration callback raised")
if (tick_report.transitions_applied
or reaper_report.rows_reaped
or pickup_report.workflows_stuck
or (reconciliation_report
and reconciliation_report.workflows_transitioned)):
logger.info(
"master iteration %d: transitions=%d reaped=%d "
"stuck=%d reconciled=%d",
iteration,
tick_report.transitions_applied,
reaper_report.rows_reaped,
pickup_report.workflows_stuck,
reconciliation_report.workflows_transitioned
if reconciliation_report else 0,
)
stop.wait(cfg.tick_interval_s)
finally:
logger.info("master loop stopped after %d iterations", iteration)
__all__ = [
"MasterConfig",
"MasterTickReport",
"master_main_loop",
"run_master_iteration",
]