WS Slack-UX Feature 1. A /new-task task now maps to ONE Slack thread instead of
several top-level messages.
- /new-task posts an immediate root "📥 Task received: …" ack and captures its
ts (root_ts); this is the instant acknowledgement.
- root_ts is plumbed into start: new PipelineState/TaskRecord channel
slack_thread_ts, seeded by graph.start_task and threaded through
Coordinator.start_task. The NewTaskCallback is now (task_text, via, root_ts).
- All clarifier questions for the task post as THREADED REPLIES under root_ts
(chat.postMessage thread_ts=root_ts), and each question's ledger channel_ref
is set to root_ts (NOT the reply's own ts). Because answer-mapping resolves a
reply via find_open_question_by_channel_ref(thread_ts), a reply in the root
thread (thread_ts==root_ts) maps to the task's currently-open question with NO
change to the mapping logic or the first-answer-wins CAS. The open-only
partial-unique index still holds (one open question per task at a time).
- Lifecycle milestones (parked / plan-ready / needs-input) and follow-up
questions thread under root_ts too; the notify sink gained an optional
thread_ts kwarg (degrades to top-level on a sink that doesn't accept it).
notify failures still never break tick.
- SlackTransport.post_question + the live poster accept/forward thread_ts.
- No root_ts (non-/new-task origin) ⇒ top-level posts exactly as before.
AUTHZ-01 (owner-allowlist-first, fail-closed) and the atomic open→answered
compare-and-set are unchanged.
Adds plumbing for the inbound-ack reactor seam used by Feature 2 (dormant until
a reactor is injected). Tests cover thread_ts forwarding, channel_ref=root_ts,
graph seeding, and coordinator threading.
229 lines
8.4 KiB
Python
229 lines
8.4 KiB
Python
"""Activation-wiring tests (integration branch): prove the WS seams that the
|
|
``serve`` path flips ON are actually wired, without a live Slack socket.
|
|
|
|
Covers:
|
|
* run-team ``_build_context_provider`` returns the handbook loader (WS5 / D10).
|
|
* run-team ``_build_coordinator`` threads that context_provider into the planner
|
|
node factory.
|
|
* run-team ``_build_coordinator`` sets a ``/new-task`` callback that starts a
|
|
task on the SAME coordinator with transport_name="slack" (WS2).
|
|
* ``default_slack_listener_factory`` forwards ``new_task_callback`` to the
|
|
SlackListener (WS2).
|
|
* ``Coordinator`` stores ``new_task_callback`` / ``set_new_task_callback`` and
|
|
forwards it when it builds the default listener.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import importlib.util
|
|
import sys
|
|
from pathlib import Path
|
|
from types import SimpleNamespace
|
|
from typing import Any
|
|
from unittest.mock import MagicMock
|
|
|
|
|
|
_AGENT_TEAM_DIR = Path(__file__).resolve().parents[1]
|
|
if str(_AGENT_TEAM_DIR) not in sys.path:
|
|
sys.path.insert(0, str(_AGENT_TEAM_DIR))
|
|
|
|
|
|
def _load_run_team():
|
|
"""Import run-team.py (hyphenated, so loaded by path) as a module."""
|
|
cli_path = _AGENT_TEAM_DIR / "run-team.py"
|
|
spec = importlib.util.spec_from_file_location("run_team_cli", cli_path)
|
|
assert spec is not None and spec.loader is not None
|
|
module = importlib.util.module_from_spec(spec)
|
|
spec.loader.exec_module(module)
|
|
return module
|
|
|
|
|
|
def _dry_args(tmp_path: Path) -> SimpleNamespace:
|
|
return SimpleNamespace(
|
|
db=str(tmp_path / "agent_team.sqlite"),
|
|
transport="slack",
|
|
dry_run=True,
|
|
)
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# WS5: context_provider = handbook loader, threaded into the plan node
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_build_context_provider_returns_handbook_loader() -> None:
|
|
cli = _load_run_team()
|
|
provider = cli._build_context_provider()
|
|
assert callable(provider)
|
|
# Zero-arg and returns a string (fail-safe: "" when no handbook dir).
|
|
result = provider()
|
|
assert isinstance(result, str)
|
|
|
|
|
|
def test_build_coordinator_threads_context_provider_into_plan_node(
|
|
tmp_path: Path, monkeypatch: Any
|
|
) -> None:
|
|
cli = _load_run_team()
|
|
from agent_team import coordinator as coord_mod
|
|
|
|
captured: dict[str, Any] = {}
|
|
|
|
def _spy_plan_factory(context_provider: Any = None):
|
|
captured["context_provider"] = context_provider
|
|
return lambda state: state
|
|
|
|
monkeypatch.setattr(coord_mod, "default_plan_node_factory", _spy_plan_factory)
|
|
|
|
coordinator = cli._build_coordinator(_dry_args(tmp_path))
|
|
# The plan node factory the coordinator holds is the run-team lambda; calling
|
|
# it must invoke default_plan_node_factory WITH a non-None context_provider.
|
|
coordinator._build_plan_node()
|
|
assert "context_provider" in captured
|
|
assert callable(captured["context_provider"])
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# WS2: /new-task callback wired to this coordinator's start_task
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_build_coordinator_wires_new_task_callback_to_start_task(
|
|
tmp_path: Path, monkeypatch: Any
|
|
) -> None:
|
|
cli = _load_run_team()
|
|
coordinator = cli._build_coordinator(_dry_args(tmp_path))
|
|
|
|
# Replace start_task so we can observe the callback routing without running
|
|
# the real graph.
|
|
calls: dict[str, Any] = {}
|
|
|
|
def _fake_start_task(
|
|
*, task_text: str, transport_name: str, slack_thread_ts: str = ""
|
|
) -> str:
|
|
calls["task_text"] = task_text
|
|
calls["transport_name"] = transport_name
|
|
calls["slack_thread_ts"] = slack_thread_ts
|
|
return "thread-xyz"
|
|
|
|
monkeypatch.setattr(coordinator, "start_task", _fake_start_task)
|
|
|
|
cb = coordinator._new_task_callback
|
|
assert cb is not None
|
|
# The 3-arg callback (one-thread-per-task): (task_text, via, root_ts). The
|
|
# root_ts is forwarded into start_task so the task threads under the root.
|
|
thread_id = cb("fix the flaky test", "slack", "1700000000.000100")
|
|
assert thread_id == "thread-xyz"
|
|
assert calls == {
|
|
"task_text": "fix the flaky test",
|
|
"transport_name": "slack",
|
|
"slack_thread_ts": "1700000000.000100",
|
|
}
|
|
|
|
|
|
def test_set_new_task_callback_overrides() -> None:
|
|
from agent_team.coordinator import Coordinator
|
|
from agent_team.transport.base import Transport
|
|
|
|
class _T(Transport):
|
|
def post_question(self, **kw: Any) -> str: # type: ignore[override]
|
|
return "q"
|
|
|
|
def parse_answer(self, raw: Any): # type: ignore[override]
|
|
raise NotImplementedError
|
|
|
|
coord = Coordinator(db_path=":memory:", transport=_T())
|
|
assert coord._new_task_callback is None
|
|
sentinel = lambda t, s, r: "tid" # noqa: E731
|
|
coord.set_new_task_callback(sentinel)
|
|
assert coord._new_task_callback is sentinel
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# WS2: factory + coordinator forward new_task_callback to the SlackListener
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_slack_listener_factory_forwards_new_task_callback(
|
|
tmp_path: Path, monkeypatch: Any
|
|
) -> None:
|
|
import agent_team.coordinator as coord_mod
|
|
|
|
captured: dict[str, Any] = {}
|
|
|
|
class _FakeListener:
|
|
def __init__(self, *args: Any, **kwargs: Any) -> None:
|
|
captured["new_task_callback"] = kwargs.get("new_task_callback")
|
|
|
|
# Patch the lazily-imported SlackListener symbol.
|
|
import agent_team.transport.slack_listener as sl_mod
|
|
|
|
monkeypatch.setattr(sl_mod, "SlackListener", _FakeListener)
|
|
|
|
sentinel = lambda t, s, r: "tid" # noqa: E731
|
|
coord_mod.default_slack_listener_factory(
|
|
transport=MagicMock(),
|
|
db_path=tmp_path / "x.sqlite",
|
|
enqueue_resume=lambda _x: None,
|
|
new_task_callback=sentinel,
|
|
)
|
|
assert captured["new_task_callback"] is sentinel
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# WS5 (G4): context_provider is threaded into the CLARIFIER too, not just plan
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_build_coordinator_threads_context_provider_into_clarify_node(
|
|
tmp_path: Path, monkeypatch: Any
|
|
) -> None:
|
|
cli = _load_run_team()
|
|
from agent_team import coordinator as coord_mod
|
|
|
|
captured: dict[str, Any] = {}
|
|
|
|
def _spy_clarify_factory(context_provider: Any = None):
|
|
captured["context_provider"] = context_provider
|
|
return lambda state: state
|
|
|
|
monkeypatch.setattr(coord_mod, "default_clarify_node_factory", _spy_clarify_factory)
|
|
|
|
coordinator = cli._build_coordinator(_dry_args(tmp_path))
|
|
coordinator._build_clarify_node()
|
|
assert "context_provider" in captured
|
|
assert callable(captured["context_provider"])
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# WS1 (G1): the in-process cross-reviewer bootstraps the orchestrator root onto
|
|
# sys.path before importing `models` (else the daemon silently REQUEST_CHANGES).
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_cross_reviewer_invoker_bootstraps_orchestrator_path(monkeypatch: Any) -> None:
|
|
import types
|
|
|
|
from agent_team.invoker_multi import _ensure_orchestrator_on_path # noqa: F401
|
|
from agent_team.nodes.review_loop_llm import make_cross_reviewer_invoker
|
|
|
|
# The orchestrator root is parents[2] of invoker_multi.py.
|
|
import agent_team.invoker_multi as im
|
|
|
|
root = str(Path(im.__file__).resolve().parents[2])
|
|
|
|
# Simulate the daemon: root NOT on sys.path. Inject a fake `models` so the
|
|
# deferred import resolves without real provider keys — the point is to
|
|
# prove the bootstrap runs (root re-added) BEFORE the import.
|
|
monkeypatch.setattr(sys, "path", [p for p in sys.path if p != root])
|
|
fake_models = types.ModuleType("models")
|
|
fake_reviewer = MagicMock()
|
|
fake_reviewer.invoke.return_value = MagicMock(content="APPROVE")
|
|
fake_models.get_cross_reviewer = lambda: fake_reviewer # type: ignore[attr-defined]
|
|
monkeypatch.setitem(sys.modules, "models", fake_models)
|
|
|
|
invoker = make_cross_reviewer_invoker()
|
|
out = invoker("review this plan")
|
|
assert out == "APPROVE"
|
|
assert root in sys.path, (
|
|
"invoker must bootstrap the orchestrator root onto sys.path"
|
|
)
|