febb352618
Round-3 deep pass identified four issues that would have prevented an
actual end-to-end pipeline run:
RB1 — DISCOVERED → ANALYZING never fired in production:
The state machine defines (DISCOVERED, discovery_picked_up) →
ANALYZING but NO production code fires the event. Workflows
created by discovery would sit in DISCOVERED forever.
Fix:
- New ``master/promote.py``: ``run_promote_discovered_tick`` scans
for DISCOVERED workflows + fires ``discovery_picked_up`` via
apply_event (state-machine invariants stay enforced) + emits a
``discovery-promoted`` controller_events row per transition.
- Composes with the master loop's other ticks; runs every iteration
(cheap — typically 0-1 row).
RB2 — scheduler.schedule_next_attempts never called from master loop:
The scheduler was exported by the master package but never invoked.
It creates the ``workflow_attempts`` rows that workers dequeue —
without it, workers would have nothing to pick up.
Fix:
- ``master/loop.py`` now accepts a ``prefetch: PrefetchCallback``
kwarg. When provided, the loop runs promote_discovered + scheduler
every iteration after tick/reaper/reconciliation.
- ``MasterTickReport`` gains ``promote_discovered`` and ``scheduler``
optional fields so on_iteration callbacks see both.
- ``master/__main__.py`` builds a ``PrefetchDataCallbacks`` from the
Forgejo callback bundle and constructs the production
``make_prefetch_callback(engine, callbacks)`` — wires through to
the loop's new prefetch kwarg.
RB3 — owner / repo missing from V1 input contracts:
The implementer / reviewer / estimator / conflict-resolver V1 inputs
had pr_number but not owner/repo. The OpenCode agent would have
had no way to know which Forgejo repo to clone — it would have had
to derive owner/repo from process env, coupling the worker to a
single repo.
Fix:
- ``contracts/v1.py``: added ``owner: str`` and ``repo: str``
(min_length=1) to ImplementerInputV1, ReviewerInputV1,
EstimatorInputV1, ConflictResolverInputV1.
- ``master/prefetch.py``: builders populate owner/repo from the
Workflow row (already known at prefetch time).
- Existing test fixtures in ``test_contracts_v1.py`` updated.
RB4 — input_payload.workspace_dir placeholder reached the agent:
Prefetch wrote ``workspace_dir = "<worker-injected>"`` as a
placeholder; the worker never patched it before invoking the
OpenCode session. The prompt builder rendered the literal
placeholder string into the agent's prompt — the agent had no idea
where to clone.
Fix:
- ``worker/agent_runner.py``: patches input_payload.workspace_dir
with the real path immediately before calling run_opencode_session.
Uses a shallow copy so the caller's dict isn't side-effected.
- ``worker/__main__.py``: workspace_dir naming convention is now
``pr-{owner}-{repo}-{pr_number}`` (matches workspace.py's
PerPRWorkspace convention) so the janitor's pr-* glob + the
agent's expected workspace location agree. Falls back to
``pr-attempt-{N}`` for legacy input_payloads missing owner/repo.
Tests:
- ``test_master_promote.py`` (NEW, +7 tests):
- empty DB no-op
- single workflow promoted
- multiple promoted in one tick
- only DISCOVERED targeted (non-DISCOVERED untouched)
- controller_events row emitted with correct shape
- idempotent after first promotion
- LoopIntegration end-to-end: DISCOVERED → ANALYZING → pending
estimator attempt visible in workflow_attempts (pins the entire
previously-broken pipeline from discovery to enqueue)
Total: 703 controller tests pass (+7 net), 0 regressions.
Without these four fixes, the pipeline would have looked alive in
unit tests but produced zero work in a real deployment.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
285 lines
11 KiB
Python
285 lines
11 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 .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")
|
|
)
|
|
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
|
|
|
|
|
|
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.
|
|
) -> 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
|
|
iteration = 0
|
|
|
|
logger.info(
|
|
"master loop starting: tick=%.1fs reaper=%.1fs reconcile=%.1fs "
|
|
"ci_poll_exh=%.1fs max_pickups=%d",
|
|
cfg.tick_interval_s, cfg.reaper_interval_s,
|
|
cfg.reconciliation_interval_s, cfg.ci_poll_exhaustion_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:
|
|
if len(reconciliation_args) == 5:
|
|
(owner, repo, get_pr_state, get_issue_state,
|
|
extra_kwargs) = reconciliation_args
|
|
else:
|
|
(owner, repo, get_pr_state,
|
|
get_issue_state) = reconciliation_args
|
|
extra_kwargs = {}
|
|
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"
|
|
)
|
|
|
|
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,
|
|
))
|
|
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",
|
|
]
|