End-to-end validation surfaced that auto-dispatch parked every task at "empty declared_scope": dispatch_node required plan["scope"], but the planner emits only summary+phases (never a scope) and config.allowed_scope defaults to None, so no node ever populated it. (The earlier dispatch test was operator-initiated with an explicit scope; the auto planner->build->dispatch path was never exercised until box-side App dispatch went live.) Fix: when no planner-/operator-declared scope is present, dispatch_node derives declared_scope from the candidate diff's own touched paths (ci_gate.diff_touched_paths). This supplies the missing scope without relaxing any CI trust control — the apply/verify workflow still INDEPENDENTLY re-checks the materialized diff against the denylist + '..'-escape + this scope + the diff-hash binding, and the agent-apply environment's required reviewer remains the human gate. A non-empty diff that parses to zero touched paths still parks (fail closed). Tests: the three old park-on-missing-scope cases now assert scope-from-diff dispatch; added a park case for a diff with no parseable paths. 1527 pass.
155 lines
6.3 KiB
Python
155 lines
6.3 KiB
Python
"""Dispatch-node factory: carry an approved diff into org CI (§3.3, §7.1 P3).
|
|
|
|
:mod:`agent_team.nodes.builders` proposes the diff;
|
|
:mod:`agent_team.nodes.verifier` clears the pure-code gate; this node is the
|
|
final in-graph step that calls
|
|
:func:`agent_team.dispatcher.dispatch_apply_verify` to push the head branch and
|
|
fire ``workflow_dispatch``.
|
|
|
|
Model/role: this node executes only when the verifier emits the
|
|
``APPROVED_ROUTE`` signal. It reads the already-validated candidate diff and
|
|
task context from graph state, assembles the dispatch inputs (``owner``/``repo``
|
|
injected at coordinator startup), and calls the dispatcher's seams. The actual
|
|
diff was hashed and scope-checked by the verifier; this node transports — it
|
|
makes no new trust decision.
|
|
|
|
INERT unless configured: the node is built only when the coordinator threads a
|
|
``dispatch_node_wiring`` factory through
|
|
:func:`~agent_team.coordinator.Coordinator`. With no factory,
|
|
:func:`agent_team.graph.build_graph` leaves ``APPROVED_ROUTE → END`` unchanged
|
|
and no dispatch ever fires. This keeps the default path inert and the P3
|
|
subgraph opt-in, exactly as the verifier node.
|
|
|
|
Fail-safe: any dispatch error (invalid state, empty diff, owner/repo
|
|
misconfigured, network/subprocess failure) parks the task rather than crashing
|
|
the graph. The verifier already validated the diff hash; a dispatch error is an
|
|
infrastructure problem, not a security bypass.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from collections.abc import Callable
|
|
from typing import Any
|
|
|
|
__all__ = ["DispatchNodeFactory", "make_dispatch_node"]
|
|
|
|
_LOG = logging.getLogger("agent_team.nodes.dispatch_invoker")
|
|
|
|
# A dispatch-node factory type: takes no args, returns the LangGraph node
|
|
# callable. Mirrors the other node-factory types in coordinator.py.
|
|
DispatchNodeFactory = Callable[[], "Callable[[Any], Any]"]
|
|
|
|
|
|
def make_dispatch_node(
|
|
*,
|
|
owner: str,
|
|
repo: str,
|
|
base: str = "main",
|
|
pusher: Any = None,
|
|
dispatcher: Any = None,
|
|
locator: Any = None,
|
|
) -> Callable[[Any], Any]:
|
|
"""Build a LangGraph dispatch node for ``owner``/``repo``.
|
|
|
|
Returns a single-arg ``(state) -> dict`` node. At runtime it reads
|
|
``thread_id``, ``candidate_diff``, and ``plan.scope`` from ``state``, then
|
|
calls :func:`agent_team.dispatcher.dispatch_apply_verify` with the injected
|
|
``pusher``/``dispatcher`` seams (default: real git/gh subprocess paths).
|
|
|
|
Fail-safe: any :class:`~agent_team.dispatcher.DispatcherError` or unexpected
|
|
exception parks the task (returns ``status=PARKED``); the caller retains the
|
|
full graph state, so the coordinator can ALARM and a human can inspect.
|
|
"""
|
|
# Deferred import: no orchestrator / subprocess module at module load.
|
|
from agent_team.ci_gate import diff_touched_paths
|
|
from agent_team.dispatcher import DispatcherError, dispatch_apply_verify
|
|
from agent_team.task_model import Phase, TaskStatus
|
|
|
|
def dispatch_node(state: Any) -> Any:
|
|
thread_id: str = state.get("thread_id") or ""
|
|
diff_text: str = state.get("candidate_diff") or ""
|
|
plan: Any = state.get("plan") or {}
|
|
scope_list: list[Any] = (
|
|
plan.get("scope") or [] if isinstance(plan, dict) else []
|
|
)
|
|
declared_scope: str = "\n".join(str(s) for s in scope_list if s)
|
|
|
|
_parked: dict[str, Any] = {
|
|
"status": TaskStatus.PARKED.value,
|
|
"current_phase": Phase.PARKED.value,
|
|
}
|
|
|
|
if not thread_id or not diff_text.strip():
|
|
_LOG.warning("dispatch_node: missing thread_id or candidate_diff; parking")
|
|
return _parked
|
|
if not declared_scope.strip():
|
|
# No planner-/operator-declared scope (the planner emits only
|
|
# summary+phases, never a scope) — derive an HONEST, non-empty
|
|
# declared_scope from the candidate diff's own touched paths. The
|
|
# workflow's guard still INDEPENDENTLY re-checks the materialized diff
|
|
# against the denylist + '..'-escape + this scope + the diff-hash
|
|
# binding, and the `agent-apply` environment's required reviewer remains
|
|
# the human gate — so this only supplies the scope that was missing, it
|
|
# does not relax any CI trust control.
|
|
declared_scope = "\n".join(diff_touched_paths(diff_text))
|
|
if not declared_scope.strip():
|
|
_LOG.warning(
|
|
"dispatch_node: no declared scope and diff touches no paths; parking"
|
|
)
|
|
return _parked
|
|
|
|
try:
|
|
result = dispatch_apply_verify(
|
|
owner=owner,
|
|
repo=repo,
|
|
task_id=thread_id,
|
|
diff_text=diff_text,
|
|
declared_scope=declared_scope,
|
|
base=base,
|
|
pusher=pusher,
|
|
dispatcher=dispatcher,
|
|
locator=locator,
|
|
)
|
|
except DispatcherError as exc:
|
|
_LOG.error(
|
|
"dispatch_node: DispatcherError for task %s: %s; parking",
|
|
thread_id,
|
|
exc,
|
|
)
|
|
return _parked
|
|
except Exception as exc: # noqa: BLE001
|
|
_LOG.error(
|
|
"dispatch_node: unexpected error for task %s (%s); parking",
|
|
thread_id,
|
|
type(exc).__name__,
|
|
)
|
|
return _parked
|
|
|
|
if not result.run_id:
|
|
# Fired, but the run could not be correlated. Persist the watermark
|
|
# anyway and let the verifier gate fail closed (no authenticated run
|
|
# to read -> BLOCK/park) rather than fabricating progress.
|
|
_LOG.warning(
|
|
"dispatch_node: task %s dispatched but run_id unresolved; "
|
|
"downstream verify will fail closed",
|
|
thread_id,
|
|
)
|
|
|
|
_LOG.info(
|
|
"dispatch_node: dispatched task %s to %s/%s (base=%s, run_id=%s)",
|
|
thread_id,
|
|
owner,
|
|
repo,
|
|
base,
|
|
result.run_id,
|
|
)
|
|
# Persist the located run identity so the verifier's read-only fetcher
|
|
# polls THIS task's run and the pure-code gate binds its verdict to it.
|
|
return {
|
|
"run_id": result.run_id,
|
|
"dispatched_at": result.dispatched_at,
|
|
"ci_correlation_tag": result.correlation_tag,
|
|
}
|
|
|
|
return dispatch_node
|