This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/nodes/dispatch_invoker.py

123 lines
4.6 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,
) -> 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.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():
_LOG.warning("dispatch_node: empty declared_scope from plan; parking")
return _parked
try:
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,
)
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
_LOG.info(
"dispatch_node: dispatched task %s to %s/%s (base=%s)",
thread_id,
owner,
repo,
base,
)
return {}
return dispatch_node