Files
cleveragents-core/tools/controller/master/prefetch.py
T
drew febb352618 fix(controller): pipeline run-blockers — promoter, scheduler, owner/repo, workspace_dir patch
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>
2026-05-18 16:43:36 -04:00

640 lines
23 KiB
Python

"""Prefetch callbacks — assemble V1 worker input from Forgejo + DB.
Plan v9 puts prefetch on the master side: at attempt-enqueue time the
scheduler calls prefetch (a pluggable callback) to build the worker's
input_payload from PR details + comments + reviews + prior attempts.
Workers then dequeue an already-assembled payload with NO extra
Forgejo I/O of their own.
This module ships the production prefetch implementation:
- ``PrefetchDataCallbacks`` — bundle of Forgejo fetchers (DI for tests).
- ``build_implementer_input`` / ``build_reviewer_input`` /
``build_estimator_input`` — per-role assembly into V1-shape dicts.
- ``make_prefetch_callback(engine, callbacks)`` — factory matching the
scheduler's ``PrefetchCallback`` protocol.
What the worker still patches in
--------------------------------
A few V1 input fields aren't known at prefetch time. The worker fills
these in after dequeue and before passing the payload to the agent:
- ``attempt_id`` — autoincrement PK, only known post-INSERT.
- ``attempt_number`` — known by the scheduler but not threaded through
the current ``PrefetchCallback`` signature; placeholder is 1.
- ``workspace_dir`` — per-worker filesystem path.
- ``wallclock_budget_s`` — worker-side config (default 600s).
For now these are placeholders the worker overwrites. The V1 contracts
admit the placeholders (``ge=0`` / ``ge=1`` / non-empty str) so the
schema parses; the worker's responsibility is to swap them for real
values before the agent ever sees the payload.
What Phase 1h DOES NOT yet produce
----------------------------------
- ``ci_summary`` / ``failing_gates`` — Phase 1j ships the deterministic
CI summarizer; until then these are None / empty.
- ``conflicted_files`` for the conflict_resolver role — needs a git
rebase + conflict-parse pass that lives on the worker side.
- ``implementer_claim`` for the reviewer role — derived from the most
recent implementer attempt's ``output_payload`` (we set it when the
payload validates as ImplementerOutputV1, else leave None).
"""
from __future__ import annotations
import json
import logging
from collections.abc import Callable
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Any
from sqlalchemy import text
from sqlalchemy.engine import Engine
from ..db.session import session_scope
logger = logging.getLogger(__name__)
# Default worker-side fields the prefetch fills with placeholders.
_PLACEHOLDER_ATTEMPT_ID = 0
_PLACEHOLDER_ATTEMPT_NUMBER = 1
_PLACEHOLDER_WORKSPACE_DIR = "<worker-injected>"
_DEFAULT_WALLCLOCK_BUDGET_S = 600
# Number of recent IMPLEMENTER outputs to include verbatim in
# PriorAttemptsBlock. Older attempts are summarized (Phase 2+).
_PRIOR_VERBATIM_LIMIT = 3
# ─── callback bundle ─────────────────────────────────────────────────
# Each callback returns Forgejo-shape data, or raises on transport
# failure. Tests inject synthetic callbacks; production wires them via
# ``forgejo_http.build_callbacks``.
#
# Signatures:
# get_pr_details(owner, repo, pr_number) -> dict | None
# Full PR object (head/base/title/body/etc.); None on 404.
# get_pr_diff(owner, repo, pr_number) -> str | None
# Raw unified diff text; None if unavailable.
# list_pr_reviews(owner, repo, pr_number) -> list[dict]
# Forgejo PR reviews (state/body/user/submitted_at/...).
# list_pr_comments(owner, repo, pr_number) -> list[dict]
# Issue-style PR comments (body/user/created_at/...).
GetPRDetailsCallback = Callable[[str, str, int], dict | None]
GetPRDiffCallback = Callable[[str, str, int], str | None]
ListPRReviewsCallback = Callable[[str, str, int], list[dict]]
ListPRCommentsCallback = Callable[[str, str, int], list[dict]]
@dataclass(frozen=True)
class PrefetchDataCallbacks:
"""Bundle of Forgejo fetchers the prefetch helpers depend on."""
get_pr_details: GetPRDetailsCallback
get_pr_diff: GetPRDiffCallback
list_pr_reviews: ListPRReviewsCallback
list_pr_comments: ListPRCommentsCallback
# ─── DB-side helpers ─────────────────────────────────────────────────
def _read_workflow(session, workflow_id: int) -> dict:
"""Read the workflow row as a plain dict. Raises if absent."""
row = session.execute(
text(
"SELECT workflow_id, kind, owner, repo, entity_number, "
" current_state, current_tier, tier_last_succeeded "
"FROM workflows WHERE workflow_id = :wf_id"
),
{"wf_id": workflow_id},
).first()
if row is None:
raise ValueError(f"workflow {workflow_id} not found")
return {
"workflow_id": row.workflow_id,
"kind": row.kind,
"owner": row.owner,
"repo": row.repo,
"entity_number": row.entity_number,
"current_state": row.current_state,
"current_tier": row.current_tier,
"tier_last_succeeded": row.tier_last_succeeded,
}
def _coerce_payload(raw: Any) -> dict | None:
"""Deserialize a JSON payload that may come back as text from raw
``text()`` SELECTs against the JSON column (the ORM's JSON type
processor only fires when the column type is known)."""
if isinstance(raw, dict):
return raw
if isinstance(raw, str):
try:
parsed = json.loads(raw)
except (TypeError, ValueError):
return None
return parsed if isinstance(parsed, dict) else None
return None
def _read_prior_implementer_attempts(
session, workflow_id: int,
) -> tuple[list[dict], int]:
"""Return (verbatim_outputs, total_attempts).
verbatim_outputs is the most-recent ``_PRIOR_VERBATIM_LIMIT``
implementer outputs whose payload parses as ImplementerOutputV1.
total_attempts is the count of all completed implementer attempts
(used by PriorAttemptsBlock for "older" bookkeeping).
"""
rows = session.execute(
text(
"SELECT attempt_number, output_payload, output_version "
"FROM workflow_attempts "
"WHERE workflow_id = :wf_id "
" AND role = 'implementer' "
" AND status = 'complete' "
" AND output_payload IS NOT NULL "
"ORDER BY attempt_number DESC"
),
{"wf_id": workflow_id},
).all()
total = len(rows)
verbatim: list[dict] = []
for r in rows[:_PRIOR_VERBATIM_LIMIT]:
payload = _coerce_payload(r.output_payload)
if payload is not None:
verbatim.append(payload)
# Caller expects oldest-first verbatim block:
verbatim.reverse()
return verbatim, total
def _read_most_recent_implementer_output(
session, workflow_id: int,
) -> dict | None:
"""Return the latest completed implementer output, or None."""
row = session.execute(
text(
"SELECT output_payload FROM workflow_attempts "
"WHERE workflow_id = :wf_id "
" AND role = 'implementer' "
" AND status = 'complete' "
" AND output_payload IS NOT NULL "
"ORDER BY attempt_number DESC LIMIT 1"
),
{"wf_id": workflow_id},
).first()
if row is None:
return None
return _coerce_payload(row.output_payload)
def _read_prior_reviewer_outputs(
session, workflow_id: int,
) -> list[dict]:
"""Return all completed reviewer outputs, oldest first."""
rows = session.execute(
text(
"SELECT output_payload FROM workflow_attempts "
"WHERE workflow_id = :wf_id "
" AND role = 'reviewer' "
" AND status = 'complete' "
" AND output_payload IS NOT NULL "
"ORDER BY attempt_number ASC"
),
{"wf_id": workflow_id},
).all()
out: list[dict] = []
for r in rows:
payload = _coerce_payload(r.output_payload)
if payload is not None:
out.append(payload)
return out
def _read_last_implementer_finished_at(
session, workflow_id: int,
) -> Any | None:
"""Timestamp of the most recently completed implementer attempt,
or None. Used to filter "comments since last attempt"."""
row = session.execute(
text(
"SELECT finished_at FROM workflow_attempts "
"WHERE workflow_id = :wf_id "
" AND role = 'implementer' "
" AND status = 'complete' "
"ORDER BY attempt_number DESC LIMIT 1"
),
{"wf_id": workflow_id},
).first()
return row.finished_at if row is not None else None
# ─── Forgejo-shape adapters ──────────────────────────────────────────
def _diff_summary_from_pr(pr: dict) -> str:
"""One-line diff summary derived from PR fields (additions /
deletions / changed_files). Falls back gracefully on missing keys."""
parts: list[str] = []
changed = pr.get("changed_files")
if isinstance(changed, int):
parts.append(f"{changed} files")
add = pr.get("additions")
if isinstance(add, int):
parts.append(f"+{add}")
delete = pr.get("deletions")
if isinstance(delete, int):
parts.append(f"-{delete}")
if not parts:
return "diff metadata unavailable"
return ", ".join(parts)
def _active_reviews_from_forgejo(reviews: list[dict]) -> list[dict]:
"""Project Forgejo reviews into the V1 ``Review`` shape.
Skips reviews that don't carry the required fields (review_id,
reviewer_login, state). Forgejo states map directly:
APPROVED / REQUEST_CHANGES / COMMENT / PENDING / DISMISSED.
"""
valid_states = {
"APPROVED", "REQUEST_CHANGES", "COMMENT", "PENDING", "DISMISSED",
}
out: list[dict] = []
for r in reviews:
if not isinstance(r, dict):
continue
review_id = r.get("id")
if not isinstance(review_id, int):
continue
user = r.get("user") if isinstance(r.get("user"), dict) else {}
login = user.get("login") if isinstance(user, dict) else None
if not isinstance(login, str) or not login:
continue
state = r.get("state")
if state not in valid_states:
continue
out.append({
"review_id": review_id,
"reviewer_login": login,
"state": state,
"body": r.get("body") or "",
"submitted_at": r.get("submitted_at"),
})
return out
def _comment_bodies_since(comments: list[dict], since: Any | None) -> list[str]:
"""Return comment bodies whose ``created_at`` is strictly after
``since``. If ``since`` is None, returns ALL comment bodies.
Compares as ``datetime`` (not string) so timezone-suffix variants
don't bifurcate the comparison: Forgejo emits
``2026-05-18T12:00:00Z``; Python ``datetime.isoformat()`` emits
``2026-05-18T12:00:00+00:00`` — lexicographic compare is wrong.
"""
since_dt = _to_aware_datetime(since)
out: list[str] = []
for c in comments:
if not isinstance(c, dict):
continue
body = c.get("body")
if not isinstance(body, str) or not body:
continue
if since_dt is None:
out.append(body)
continue
created_at_dt = _to_aware_datetime(c.get("created_at"))
if created_at_dt is None:
# Conservatively include — missing timestamp means we can't
# prove it's older than the cutoff. Better to show the
# worker a comment it might already have seen than to
# silently drop a relevant one.
out.append(body)
continue
if created_at_dt > since_dt:
out.append(body)
return out
def _to_aware_datetime(value: Any) -> datetime | None:
"""Coerce a value (str / datetime / None) to a timezone-aware
datetime. Returns None for None or unparseable values.
Accepts the trailing-Z variant Forgejo emits. Naive datetimes are
interpreted as UTC (the controller's invariant: all timestamps
are UTC at the boundary).
"""
if value is None:
return None
if isinstance(value, datetime):
return value if value.tzinfo is not None else value.replace(tzinfo=timezone.utc)
if isinstance(value, str):
if not value:
return None
candidate = value.replace("Z", "+00:00") if value.endswith("Z") else value
try:
parsed = datetime.fromisoformat(candidate)
except ValueError:
return None
return parsed if parsed.tzinfo is not None else parsed.replace(tzinfo=timezone.utc)
return None
# ─── per-role builders ───────────────────────────────────────────────
def build_implementer_input(
*, engine: Engine, workflow_id: int, tier: int,
callbacks: PrefetchDataCallbacks,
) -> dict:
"""Assemble an ImplementerInputV1-shape dict for a PR workflow.
Worker patches in attempt_id / attempt_number / workspace_dir
before validation.
"""
with session_scope(engine) as session:
wf = _read_workflow(session, workflow_id)
if wf["kind"] != "pr":
raise ValueError(
f"implementer prefetch needs kind='pr', got {wf['kind']!r}"
)
prior_verbatim, prior_total = _read_prior_implementer_attempts(
session, workflow_id,
)
since = _read_last_implementer_finished_at(session, workflow_id)
pr = callbacks.get_pr_details(wf["owner"], wf["repo"], wf["entity_number"])
if pr is None:
raise ValueError(
f"PR {wf['owner']}/{wf['repo']}#{wf['entity_number']} not found"
)
head_sha = pr.get("head", {}).get("sha") if isinstance(pr.get("head"), dict) else None
head_ref = pr.get("head", {}).get("ref") if isinstance(pr.get("head"), dict) else None
base_branch = pr.get("base", {}).get("ref") if isinstance(pr.get("base"), dict) else None
if not head_sha or not head_ref or not base_branch:
raise ValueError(
f"PR {wf['entity_number']} missing head/base metadata"
)
reviews = callbacks.list_pr_reviews(wf["owner"], wf["repo"], wf["entity_number"])
comments = callbacks.list_pr_comments(wf["owner"], wf["repo"], wf["entity_number"])
return {
"input_version": "V1",
"workflow_id": workflow_id,
"attempt_id": _PLACEHOLDER_ATTEMPT_ID,
"attempt_number": _PLACEHOLDER_ATTEMPT_NUMBER,
"owner": wf["owner"],
"repo": wf["repo"],
"pr_number": wf["entity_number"],
"tier": tier,
"head_sha": head_sha,
"head_ref": head_ref,
"base_branch": base_branch,
"ci_summary": None, # Phase 1j fills this.
"failing_gates": [], # Phase 1j fills this.
"active_reviews": _active_reviews_from_forgejo(reviews),
"pr_comments_since_last_attempt": _comment_bodies_since(comments, since),
"prior_attempts": {
"verbatim": prior_verbatim,
"older_summary": None,
"older_summary_covers_through_attempt": None,
"total_attempts": prior_total,
},
"allowed_files": None,
"diff_summary": _diff_summary_from_pr(pr),
"workspace_dir": _PLACEHOLDER_WORKSPACE_DIR,
"wallclock_budget_s": _DEFAULT_WALLCLOCK_BUDGET_S,
}
def build_reviewer_input(
*, engine: Engine, workflow_id: int,
callbacks: PrefetchDataCallbacks,
) -> dict:
"""Assemble a ReviewerInputV1-shape dict for a PR workflow."""
with session_scope(engine) as session:
wf = _read_workflow(session, workflow_id)
if wf["kind"] != "pr":
raise ValueError(
f"reviewer prefetch needs kind='pr', got {wf['kind']!r}"
)
prior_implementer_verbatim, prior_implementer_total = (
_read_prior_implementer_attempts(session, workflow_id)
)
prior_reviews = _read_prior_reviewer_outputs(session, workflow_id)
implementer_claim = _read_most_recent_implementer_output(
session, workflow_id,
)
pr = callbacks.get_pr_details(wf["owner"], wf["repo"], wf["entity_number"])
if pr is None:
raise ValueError(
f"PR {wf['owner']}/{wf['repo']}#{wf['entity_number']} not found"
)
head_sha = pr.get("head", {}).get("sha") if isinstance(pr.get("head"), dict) else None
if not head_sha:
raise ValueError(f"PR {wf['entity_number']} missing head.sha")
full_diff = callbacks.get_pr_diff(wf["owner"], wf["repo"], wf["entity_number"])
payload: dict[str, Any] = {
"input_version": "V1",
"workflow_id": workflow_id,
"attempt_id": _PLACEHOLDER_ATTEMPT_ID,
"attempt_number": _PLACEHOLDER_ATTEMPT_NUMBER,
"owner": wf["owner"],
"repo": wf["repo"],
"pr_number": wf["entity_number"],
"head_sha": head_sha,
"ci_summary": None, # Phase 1j fills this.
"diff_summary": _diff_summary_from_pr(pr),
"full_diff": full_diff,
"prior_reviews": prior_reviews,
"prior_implementer_attempts": {
"verbatim": prior_implementer_verbatim,
"older_summary": None,
"older_summary_covers_through_attempt": None,
"total_attempts": prior_implementer_total,
},
"implementer_claim": implementer_claim,
"workspace_dir": _PLACEHOLDER_WORKSPACE_DIR,
"wallclock_budget_s": _DEFAULT_WALLCLOCK_BUDGET_S,
}
return payload
def build_estimator_input(
*, engine: Engine, workflow_id: int,
callbacks: PrefetchDataCallbacks,
) -> dict:
"""Assemble an EstimatorInputV1-shape dict.
Works for both PR and issue workflows. For issues, head_sha /
diff_summary stay None (estimator runs pre-PR for issues).
"""
with session_scope(engine) as session:
wf = _read_workflow(session, workflow_id)
pr_number: int | None = None
head_sha: str | None = None
diff_summary: str | None = None
title = ""
body = ""
if wf["kind"] == "pr":
pr = callbacks.get_pr_details(wf["owner"], wf["repo"], wf["entity_number"])
if pr is None:
raise ValueError(
f"PR {wf['owner']}/{wf['repo']}#{wf['entity_number']} not found"
)
pr_number = wf["entity_number"]
head_sha = pr.get("head", {}).get("sha") if isinstance(pr.get("head"), dict) else None
diff_summary = _diff_summary_from_pr(pr)
title = pr.get("title") or ""
body = pr.get("body") or ""
else:
# Issue workflow: estimator runs pre-PR. We synthesize title /
# body from the issue payload via the same get_pr_details
# callback (Forgejo treats issues and PRs interchangeably at
# the /issues/{n} endpoint, but the controller binds prefetch
# callbacks to PRs; for issues we accept that title/body stay
# empty for now — Phase 1h+ adds a list_issue_details callback).
title = ""
body = ""
return {
"input_version": "V1",
"workflow_id": workflow_id,
"attempt_id": _PLACEHOLDER_ATTEMPT_ID,
"owner": wf["owner"],
"repo": wf["repo"],
"pr_number": pr_number,
"head_sha": head_sha,
"ci_summary": None,
"diff_summary": diff_summary,
"pr_title": title,
"pr_body": body,
"workspace_dir": _PLACEHOLDER_WORKSPACE_DIR,
"wallclock_budget_s": _DEFAULT_WALLCLOCK_BUDGET_S,
}
def build_conflict_resolver_input(
*, engine: Engine, workflow_id: int, tier: int,
callbacks: PrefetchDataCallbacks,
) -> dict:
"""Stub conflict-resolver input.
Conflicted-file extraction needs a git rebase pass that lives on
the worker side. Phase 1h returns an empty conflicted_files list +
placeholder base_sha; the worker fills both in.
"""
with session_scope(engine) as session:
wf = _read_workflow(session, workflow_id)
if wf["kind"] != "pr":
raise ValueError(
f"conflict_resolver prefetch needs kind='pr', got {wf['kind']!r}"
)
prior_verbatim, prior_total = _read_prior_implementer_attempts(
session, workflow_id,
)
pr = callbacks.get_pr_details(wf["owner"], wf["repo"], wf["entity_number"])
if pr is None:
raise ValueError(
f"PR {wf['owner']}/{wf['repo']}#{wf['entity_number']} not found"
)
head_sha = pr.get("head", {}).get("sha") if isinstance(pr.get("head"), dict) else None
base_branch = pr.get("base", {}).get("ref") if isinstance(pr.get("base"), dict) else None
base_sha = pr.get("base", {}).get("sha") if isinstance(pr.get("base"), dict) else None
if not head_sha or not base_branch or not base_sha:
raise ValueError(
f"PR {wf['entity_number']} missing head/base metadata for conflict resolver"
)
return {
"input_version": "V1",
"workflow_id": workflow_id,
"attempt_id": _PLACEHOLDER_ATTEMPT_ID,
"owner": wf["owner"],
"repo": wf["repo"],
"pr_number": wf["entity_number"],
"head_sha": head_sha,
"base_branch": base_branch,
"base_sha": base_sha,
"conflicted_files": [], # Worker fills via git rebase.
"prior_implementer_outputs": {
"verbatim": prior_verbatim,
"older_summary": None,
"older_summary_covers_through_attempt": None,
"total_attempts": prior_total,
},
"workspace_dir": _PLACEHOLDER_WORKSPACE_DIR,
"wallclock_budget_s": _DEFAULT_WALLCLOCK_BUDGET_S,
}
# ─── PrefetchCallback factory ────────────────────────────────────────
def make_prefetch_callback(
engine: Engine, callbacks: PrefetchDataCallbacks,
) -> Callable[[int, str, int | None], tuple[dict, str]]:
"""Build a ``PrefetchCallback`` matching the scheduler protocol.
Routes by role to the per-role builder. Returns ``(payload, "V1")``.
"""
def prefetch(
workflow_id: int, role: str, tier: int | None,
) -> tuple[dict, str]:
if role == "implementer":
payload = build_implementer_input(
engine=engine, workflow_id=workflow_id,
tier=int(tier) if tier is not None else 0,
callbacks=callbacks,
)
elif role == "reviewer":
payload = build_reviewer_input(
engine=engine, workflow_id=workflow_id,
callbacks=callbacks,
)
elif role == "estimator":
payload = build_estimator_input(
engine=engine, workflow_id=workflow_id,
callbacks=callbacks,
)
elif role == "conflict_resolver":
payload = build_conflict_resolver_input(
engine=engine, workflow_id=workflow_id,
tier=int(tier) if tier is not None else 0,
callbacks=callbacks,
)
else:
raise ValueError(f"unknown role for prefetch: {role!r}")
return payload, "V1"
return prefetch
__all__ = [
"GetPRDetailsCallback",
"GetPRDiffCallback",
"ListPRReviewsCallback",
"ListPRCommentsCallback",
"PrefetchDataCallbacks",
"build_conflict_resolver_input",
"build_estimator_input",
"build_implementer_input",
"build_reviewer_input",
"make_prefetch_callback",
]