Files
roboco/tests/unit/runtime/test_intake_spawn.py
T
889f3689e7 MegaTask (#248)
* feat(batch): batch_id + collision descriptor columns

Sequenced batch intake ("Mega task") foundation: tasks.batch_id (indexed)
groups a batch of top-level tasks created together; intends_to_touch (text[]),
adds_migration and touches_shared (bool, NOT NULL default false) are the
per-task collision surface the SequencingService will read to wire dependency
waves. Mirrored on the Task model + TaskCreateRequest and wired through
TaskService.create. Migration 046 (real upgrade->downgrade->upgrade verified
vs a throwaway pgvector PG); a non-batch task declares no surface (defaults).

Task 1 of the 0.11.0 sequenced-batch-intake plan.

* feat(batch): flag + draft collision descriptors

Default-off ROBOCO_BATCH_INTAKE_ENABLED (config + FEATURE_FLAGS + panel card);
the propose_draft tool doc + the TS DraftProposal gain the per-task collision
surface intends_to_touch / adds_migration / touches_shared. The draft is a loose
dict so the descriptors ride it through the relay intact (test asserts the
forwarded payload); the analyzer (Task 3) reads them to wire dependency waves.

Task 2 of the 0.11.0 sequenced-batch-intake plan.

* feat(batch): deterministic collision-sequencing analyzer

SequencingService.analyze turns a batch's per-task collision surfaces into a
dependency DAG + execution waves — correctness in CODE, not agent judgment.
Rules in order: file overlap serializes (more-important first), migrations form
a serial chain (no concurrent Alembic heads), touches_shared runs last, cell
contention warns (never serializes); then dedupe, existence + cycle check, and
Kahn topological layering. Pure (no DB/services); SequencingError on a cycle or
out-of-range edge.

Golden test reproduces the CEO's hand-sequenced 4 waves of the 11-item
guard-core-app batch (the effort that deadlocked the Main PM): S6 alone last,
the R1/R3/R4 migration chain, R2/R3/S8 serialized on the shared threat service,
S1/S2/S7 in one parallel wave.

Task 3 of the 0.11.0 sequenced-batch-intake plan.

* chore(batch): brand the user-facing surfaces "MegaTask"

The user-facing name is MegaTask: the feature-flag label is "MegaTask intake",
the panel flag-card and the config description lead with MegaTask. Internal
names stay technical (batch_intake_enabled, batch_id, SequencingService).

* chore(batch): drop the feature flag — MegaTask is a core intake scope

MegaTask is additive and opt-in by its own nature (the Prompter proposes a
batch only when the CEO asks for several tasks; single-task intake is
unchanged), so there is no risk surface a flag protects — 'don't create a
MegaTask' is the off switch. Remove batch_intake_enabled from config, the
FEATURE_FLAGS registry, the panel flag card, and its tests. MegaTask will be
a third scope option in the Intake modal (single-cell / multi-project /
MegaTask), not a toggle.

* feat(batch): MegaTask identity predicate + orchestrator branchless recognition

The single source of truth for the umbrella's exemptions: pure
is_batch_umbrella / is_batch_root_subtask / is_branchless_coordination
(foundation/policy/batch.py) — an umbrella has a batch_id and is top-level; a
root-subtask shares the batch_id but is parented. The orchestrator's
_is_coordination_task now consults is_branchless_coordination, so a MegaTask
umbrella is recognized as doing no git of its own (git-exempt at spawn-readiness
/ stuck-detection) exactly like a product fan-out root. Non-batch behavior is
identical (the predicate reduces to the old no-project+product check; the
orchestrator coordination suite stays green), and the umbrella branch is inert
until the create path exists.

First slice of the MegaTask umbrella enforcement (branchless guard).

* feat(batch): branchless umbrella guard across the git-exemption sites

A MegaTask umbrella does no git of its own — every git-exemption site in
TaskService now consults the shared is_branchless_coordination predicate
instead of an inline product-only check, so the umbrella's exemptions
cannot drift between sites:

- the claimed->in_progress branch gate (GitContext.is_coordination) lets
  an unbranched umbrella reach in_progress and delegate;
- _ensure_branch_for_task short-circuits an umbrella to "" instead of the
  misconfigured raise (the claim path ignores the return, treating it as
  branchless);
- CEO-reject routing sends a rejected umbrella to the Main PM in PENDING
  (needs_revision is developer-claim-only and would deadlock it).

Covers both shapes via the predicate (product fan-out root OR umbrella);
a batch root-subtask keeps its own branch/PR. Adds orchestrator
recognition tests for the umbrella plus claim/branch/reject integration
tests.

* feat(batch): umbrella assembles no PR; completes branchless

submit_root now hard-rejects a MegaTask umbrella up front (a preflight
that also folds in the unknown-role refusal to stay within the
return-count budget): the umbrella spans many projects with no single
master, so each root-subtask opens and is reviewed on its own PR — the
umbrella never enters the in-path review gate. The Main PM completes it
directly once every root-subtask is terminal.

Umbrella completion needs no new code: it is branchless (no branch_name),
so _main_pm_complete_guard already accepts it from in_progress, checks
all_subtasks_terminal, and main_pm_complete walks it to awaiting_pm_review
and escalates to the CEO with no PR creation — exactly the product
fan-out root path. Adds the submit_root-reject and umbrella-completion
gateway tests; pins batch_id=None on the normal-root submit_root test
(a MagicMock auto-attr would otherwise read as an umbrella).

* feat(batch): MegaTask create path — umbrella + sequenced root-subtasks

PrompterService.confirm_live_batch turns N confirmed drafts into a real
MegaTask: it builds each draft's collision surface, runs the pure
SequencingService to get conflict-free waves, creates the branchless
umbrella (batch_id, no project/product), then one root-subtask per draft
(own project, parent=umbrella, sequence=wave index, descriptors), and
wires the analyzer's edges through add_dependency so the existing
dependency-gate runs the waves in order. The route picks the start path
like a single confirm: 'board' holds the root-subtasks in BACKLOG for the
batch review; 'main_pm' creates them PENDING so wave 0 dispatches at once.

create_task_from_draft gains a BatchPlacement (parent/batch/sequence/
team_override) and forwards the collision descriptors; the exactly-one-
target rule (here and the TaskService.create invariant) is relaxed for an
umbrella, which legitimately targets neither. New route
POST /live/{session}/confirm-batch + BatchConfirmRequest mirror the single
confirm. Adds the structural-invariant + board-hold + empty-batch tests.

* feat(batch): release MegaTask root-subtasks on CEO approval; board awareness

The board route holds a MegaTask's root-subtasks in BACKLOG so the work
waits for the batch review. approve_and_start (CEO gate #1, board->Main PM)
now releases them via _activate_batch_root_subtasks: each held child flips
BACKLOG -> PENDING + team=main_pm so the dependency-gate dispatches wave 0.
No-op for a non-umbrella; idempotent (children past BACKLOG untouched).

The Product Owner and Head of Marketing identity prompts gain a MegaTask
section so they review the whole batch + wave plan and adjust scope before
sign-off (they review drafts; the umbrella is their unit). Also extracts
the create() target invariant into _require_target_or_umbrella to keep the
method under the complexity gate after the umbrella exemption. Adds the
umbrella-approval activation test.

* feat(batch): multi-project intake scope for MegaTask

A MegaTask spans several possibly-unrelated repos, so the intake chat can
now be scoped to an explicit project list (not just one project or one
product). StartLiveRequest gains project_ids; /live/start threads it
through start/spawn_intake_session -> _spawn_intake_container ->
_clone_intake_scope. The multi-repo clone machinery already existed for
products; _intake_scope_slugs now also resolves an explicit project_ids
set (split into _slugs_for_project_ids / _slugs_for_product), cloning each
repo with the first as the primary cwd and the siblings readable. Scope
validation is now 'exactly one of project_slug / product_id / project_ids'
via the shared _require_one_intake_scope. Adds scope-resolution, spawn,
and route tests for the MegaTask path.

* feat(batch): propose_batch intake tool (MegaTask multi-draft hand-off)

The intake agent can now hand the panel a whole MegaTask in one tool call.
Both intake paths gain propose_batch alongside propose_draft:
- Claude (intake_driver): a propose_batch tool registered on the in-SDK
  MCP server + allowlisted; the driver intercepts the ToolUseBlock and
  emits ONE StreamChunk(kind="batch") carrying {drafts:[...], title}.
- grok (intake_server): a propose_batch tool that POSTs a "batch" relay
  event via the shared _post_event helper (post_draft/post_batch).

A batch carries N drafts, each the propose_draft shape PLUS its own
project_id (a MegaTask spans unrelated repos) and collision surface so the
analyzer sequences the waves. The prompter prompt documents the MegaTask
scope + when to call propose_batch. Adds Claude-normalize and grok-relay
tests for the batch path.

* feat(batch): MegaTask intake panel — third scope, batch review, waves

The panel now drives a MegaTask end to end. The intake modal gains a
third scope, 'MegaTask', beside Single cell and Board-led: a multi-project
checklist (a MegaTask spans several possibly-unrelated repos), validated
to at least two. start() sends project_ids; use-prompter accumulates the
agent's single propose_batch hand-off as a 'batch' SSE event into a
BatchProposal and lands in a new batch_preview state.

A new BatchReviewCard lists every proposed task with its target project +
collision-surface badges (migration / shared) and offers one start path
for the whole batch — Board review & Start or Approve & Start — wired to
confirmBatch → POST /confirm-batch. The success card shows the sequenced
result: N tasks in M waves (+ any advisory notes). prompter.ts gains the
DraftScale 'megatask' + the BatchConfirm payload/result types; the SSE
client allows the 'batch' kind. Panel typecheck + lint + 113 tests green.

* docs(batch): MegaTask across changelog, CLAUDE.md, site, and RAG

The four documentation obligations for the MegaTask feature:
- CHANGELOG: an Unreleased entry covering the umbrella model, sequencing,
  multi-project intake, propose_batch, and the create/approval path.
- CLAUDE.md: a MegaTask section (identity predicate, umbrella/root-subtask
  hierarchy, sequencing rules, intake + create path, board activation).
- Published site: a user-facing company/megatask.md (scopes, waves, the
  umbrella, the two start buttons) + nav entry; a pointer added to the
  intake chapter of the Tour.
- RAG corpus: workflows/megatask.md so the Main PM (and any agent) can
  retrieve the umbrella's branchless / no-PR / completion rules at runtime.

The runtime concurrent-migration guard is intentionally NOT added: the
analyzer already chains migration-adders into dependencies and the
dependency-gate serializes them, so a separate guard would be dead code.

* feat(batch): batch_id guardrail + wave preview + batch_id on TaskResponse

Guardrail (CEO): a batch_id is denied on any task that is not a well-formed
MegaTask member. is_valid_batch_shape permits batch_id only on an umbrella
(no parent → must target neither project nor product) or a root-subtask
(has a parent → exactly one target); TaskService.create enforces it AND
verifies a root-subtask's parent is the batch umbrella (same batch_id,
top-level). This closes a latent hole: is_batch_umbrella is true for a
batch_id + no-parent task even with a project, so a stray batch_id could
have spoofed the branchless branch-gate / no-PR exemption. (The public
task API never exposed batch_id for write; this guards the service layer.)

Wave preview: PrompterService.preview_batch + POST .../preview-batch
compute a MegaTask's waves from the proposed drafts WITHOUT creating
anything, so the panel can show the sequencing before confirm. Extracted
_sequence_drafts as the single source shared by preview and confirm, so
the previewed waves are exactly the ones wired.

TaskResponse now carries batch_id so the panel can badge the umbrella.

* feat(batch): MegaTask review — project editor, wave preview, persistence, badge

Closes the panel gaps in the MegaTask review experience:
- Per-task project editor: each proposed task gets an inline project
  Select (updateBatchDraftProject), so a task the agent put in the wrong
  or no repo can be fixed before launch — not only by re-chatting. Launch
  stays blocked until every task has a project.
- Wave preview: on a batch proposal the panel fetches POST .../preview-batch
  (no task created) and shows the conflict-free wave plan, so the human
  reviews the sequencing before confirming.
- Refresh durability: the MegaTask review (batch + waves + projectIds) is
  persisted, so a browser reload mid-review restores it like a single draft.
- MegaTask badge: TaskResponse exposes batch_id, the panel Task type
  carries it, and the task table badges the umbrella row 'MegaTask'.

Panel typecheck + lint + 113 tests green.

* test(batch): stub task carries batch_id for task_to_response

task_to_response now serializes batch_id (TaskResponse field), so the
_stub_task SimpleNamespace fixture must provide it — without it the reader
hit AttributeError, failing the 8 task-schema serialization/enrichment
tests. Test-only; the real TaskTable carries the column (migration 046).

* fix(batch): close MegaTask audit gaps — completion crash, analyzer cycle, guardrails

An adversarial multi-agent audit of the feature surfaced 20 verified gaps;
this closes the backend ones.

HIGH:
- Umbrella completion crashed. escalate_to_ceo hard-required a pr_number,
  which a branchless umbrella never has, so main_pm_complete dereferenced
  None. Both pr_number gates now waive a MegaTask umbrella (escalate_to_ceo
  + the awaiting_pm_review->awaiting_ceo_approval lifecycle gate via a new
  GitContext.is_umbrella), and main_pm_complete guards a None return. The
  completion test had mocked escalate_to_ceo, hiding it — now a real
  service test covers the waiver.
- The collision analyzer could fabricate a cycle (a touches_shared +
  adds_migration draft overlapping another migration draft) and raise
  SequencingError — a bare ValueError that escaped as an opaque 500. The
  migration chain is now shared-last-aware (never contradicts rule 3), and
  _sequence_drafts translates SequencingError to a clean 400.

MEDIUM:
- Collisions are now project-scoped: two repos can't collide on a
  coincidental path or serialize independent migrations (DraftSurface
  carries project_id; rules 1/2/3 respect it).
- The batch_id guardrail ran only at create. update() + the PATCH
  null-clear path now re-assert is_valid_batch_shape, so a mutation can't
  break a member's shape and spoof the branchless exemption.
- A draft missing title/acceptance_criteria now raises ValidationError
  (was a bare KeyError -> 500).
- confirm_live_batch re-asserts every draft targets a scoped project and
  the batch spans >=2 distinct projects (project_ids added to the request).
- Route-level tests for confirm-batch / preview-batch.

LOW: strict multi-repo clone (fail loud on any unresolvable project);
malformed/empty propose_batch surfaces an error chunk (Claude) / refuses
to POST (grok) instead of silently acking; dropped malformed drafts are
counted and surfaced; stale grok intake docstrings updated.

* fix(batch): MegaTask panel + doc audit gaps

Frontend half of the audit fixes:
- The confirm payload now carries project_ids (the schema requires it), and
  the panel re-checks every task targets one of the scoped repos before
  launching, naming the offending task.
- The Review-MegaTask project picker is filtered to the scoped repos and
  the per-task validity (border + launch gate) keys off scoped membership,
  so a task can only be (re)pointed at an in-scope project — also fixing the
  case where the agent emitted a non-UUID / unknown project.
- Dropped malformed drafts are surfaced as a chat error so the human knows
  the batch shrank instead of silently confirming fewer tasks.
- Doc wording: a wave releases on the previous wave's terminal state
  (normally a merge; a cancellation releases it too), not strictly 'merged'.

* test(batch): lock the CEO's EXACT 4-wave hand-sequencing as the golden bar

The golden test asserted the constraints (S6 last, the migration chain, the
shared-threats serialization, S1/S2/S7 parallel) but not the full wave
partition. The bar for MegaTask is 'reproduce my exact waves or it's not
done', so assert the exact 4-wave partition the analyzer produces for the
guard-core-app batch:
  wave 1: R1 R2 S1 S2 S3 S5 S7  ·  wave 2: R3  ·  wave 3: R4 S8  ·  wave 4: S6
Confirmed unchanged by the audit's analyzer fixes (no migration is shared;
single project).

* fix(batch): tolerate a stub task in assert_batch_shape_intact

The batch-shape re-validation read task.batch_id directly, but update()'s
partial-caller contract is exercised with a SimpleNamespace stub that has no
batch_id column → AttributeError. Use getattr(..., None) for batch_id and the
shape fields so the guard no-ops on any task lacking the column (a stub, or a
non-batch task) while still enforcing on a real batch member.

* fix(orchestrator): authenticate internal API self-calls with the system identity

The dispatcher httpx clients were built without an agent identity, so the
orchestrator's self-PATCHes to /api/tasks/{id} (auto-block, auto-resume,
auto-recover, SLA annotation) were rejected 401 "Missing X-Agent-ID" and
silently no-op'd. The auto-resume that lifts a PM's paused parent could never
write, so paused/blocked parents stayed wedged and stranded their dependents
(the fe-pm/be-pm respawn churn seen in prod).

Header propagation was inconsistent across the separate AsyncClient call-sites:
only the main dispatch client carried the system identity; the readiness and
sweep clients did not. Hoist the identity into a shared _SYSTEM_API_HEADERS
constant and apply it to every API-facing dispatcher client. The system role
holds TaskAction.ASSIGN, so it is authorized for the audited admin_set_status
path those write routes use. The external provider-recovery probe client is
intentionally left untouched.

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2026-06-24 01:15:57 +02:00

509 lines
19 KiB
Python

"""The persistent intake (prompter) live-session spawn/reap path.
The intake agent is not task-driven: ``spawn_intake_session`` launches a
long-lived Agent-SDK driver container (image ENTRYPOINT, NOT ``claude -p``),
clones the chat scope's repo(s), and registers the live relay session. These
tests cover the docker-command construction, scope resolution, and the
spawn/reap orchestration with docker + clone mocked (no daemon, no NAS).
"""
from __future__ import annotations
import asyncio
from pathlib import Path
from types import SimpleNamespace
from typing import Any
from unittest.mock import patch
from uuid import UUID
import pytest
from roboco.runtime.orchestrator import (
INTAKE_AGENT_ID,
AgentInstance,
AgentOrchestrator,
_IntakeRunSpec,
)
from roboco.services import prompter_live
def _make_minimal_orchestrator() -> AgentOrchestrator:
"""AgentOrchestrator with constructor I/O skipped; _instances ready."""
with patch.object(AgentOrchestrator, "__init__", return_value=None):
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch._instances = {}
orch._bg_tasks = set()
return orch
def _spec(**overrides: Any) -> _IntakeRunSpec:
base: dict[str, Any] = {
"container_name": "roboco-agent-intake-1",
"image": "roboco-agent-prompter",
"hosts": {
"claude": "/home/runner/.claude",
"prompt": "/data/prompts-generated/intake-1-prompt.md",
"workspaces": "/data/workspaces",
},
"session_id": "sess-abc",
"cwd": "/data/workspaces/roboco/board/intake-1",
"cli_model": "claude-opus-4-6",
"api_url": "http://roboco-orchestrator:8000",
"provider_base_url": None,
"provider_auth_token": None,
}
base.update(overrides)
return _IntakeRunSpec(**base)
@pytest.fixture(autouse=True)
def _fresh_registry() -> Any:
"""Isolate the process-wide live registry per test."""
prev = prompter_live._RegistryHolder.instance
prompter_live._RegistryHolder.instance = prompter_live.PrompterLiveRegistry()
yield
prompter_live._RegistryHolder.instance = prev
# ---------------------------------------------------------------------------
# _build_intake_run_cmd — the pure docker-argv builder.
# ---------------------------------------------------------------------------
class TestBuildIntakeRunCmd:
def test_image_is_last_and_no_claude_cli_args(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(_spec())
assert cmd[-1] == "roboco-agent-prompter"
# The image ENTRYPOINT is the driver — none of the claude CLI flags
# the task-driven path appends may appear here.
for flag in (
"-p",
"--model",
"--system-prompt-file",
"--mcp-config",
"--tools",
):
assert flag not in cmd, f"{flag} must not be in the intake run cmd"
def test_no_workdir_settings_or_manifest_mounts(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(_spec())
joined = " ".join(cmd)
assert "-w" not in cmd # driver sets cwd via ROBOCO_WORKSPACE/the SDK
assert "settings.json" not in joined # no hook mount (driver owns 9000)
assert "mcp-config.json" not in joined # MCP-free live agent
assert "tool-manifest.json" not in joined
def test_env_carries_session_workspace_and_api(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(_spec())
assert "ROBOCO_PROMPTER_SESSION_ID=sess-abc" in cmd
assert "ROBOCO_WORKSPACE=/data/workspaces/roboco/board/intake-1" in cmd
assert "ROBOCO_API_URL=http://roboco-orchestrator:8000" in cmd
assert "ROBOCO_AGENT_ID=intake-1" in cmd
assert "CLAUDE_CODE_SUBAGENT_MODEL=claude-opus-4-6" in cmd
def test_mounts_prompt_and_workspaces(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(_spec())
assert (
"/data/prompts-generated/intake-1-prompt.md:/app/system-prompt.md:ro" in cmd
)
assert "/data/workspaces:/data/workspaces" in cmd
def test_anthropic_default_omits_provider_env(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(_spec())
joined = " ".join(cmd)
assert "ANTHROPIC_BASE_URL" not in joined
assert "ANTHROPIC_AUTH_TOKEN" not in joined
def test_non_anthropic_injects_provider_env(self) -> None:
cmd = AgentOrchestrator._build_intake_run_cmd(
_spec(provider_base_url="http://ollama:11434/v1", provider_auth_token="tok")
)
assert "ANTHROPIC_BASE_URL=http://ollama:11434/v1" in cmd
assert "ANTHROPIC_AUTH_TOKEN=tok" in cmd
# ---------------------------------------------------------------------------
# _intake_scope_slugs — project XOR product resolution.
# ---------------------------------------------------------------------------
class TestIntakeScopeSlugs:
@pytest.mark.asyncio
async def test_project_scope_returns_single_slug(self) -> None:
slugs = await AgentOrchestrator._intake_scope_slugs(
db=object(), project_slug="roboco", product_id=None
)
assert slugs == ["roboco"]
@pytest.mark.asyncio
async def test_product_scope_resolves_distinct_projects_in_order(self) -> None:
# distinct_project_ids returns UUIDs in deterministic team order; the
# primary (cwd) is the first, so order must be preserved (not sorted).
ids = [
"11111111-1111-1111-1111-111111111111",
"22222222-2222-2222-2222-222222222222",
]
class _FakeProduct:
def __init__(self, _db: Any) -> None: ...
async def distinct_project_ids(self, _pid: Any) -> list[Any]:
return [UUID(i) for i in ids]
class _FakeProjectSvc:
async def get(self, pid: Any) -> Any:
return SimpleNamespace(slug=f"proj-{str(pid)[0]}")
with (
patch("roboco.services.product.ProductService", _FakeProduct),
patch(
"roboco.services.project.get_project_service",
lambda _db: _FakeProjectSvc(),
),
):
slugs = await AgentOrchestrator._intake_scope_slugs(
db=object(),
project_slug=None,
product_id="33333333-3333-3333-3333-333333333333",
)
assert slugs == ["proj-1", "proj-2"]
@pytest.mark.asyncio
async def test_product_with_no_projects_raises(self) -> None:
class _FakeProduct:
def __init__(self, _db: Any) -> None: ...
async def distinct_project_ids(self, _pid: Any) -> list[Any]:
return []
with (
patch("roboco.services.product.ProductService", _FakeProduct),
patch("roboco.services.project.get_project_service", lambda _db: object()),
pytest.raises(ValueError, match="no projects"),
):
await AgentOrchestrator._intake_scope_slugs(
db=object(),
project_slug=None,
product_id="33333333-3333-3333-3333-333333333333",
)
@pytest.mark.asyncio
async def test_megatask_scope_resolves_explicit_project_ids_in_order(self) -> None:
# A MegaTask spans an explicit set of (possibly unrelated) projects; the
# slugs are resolved in the given order (the first is the primary cwd).
ids = [
"11111111-1111-1111-1111-111111111111",
"22222222-2222-2222-2222-222222222222",
]
class _FakeProjectSvc:
async def get(self, pid: Any) -> Any:
return SimpleNamespace(slug=f"proj-{str(pid)[0]}")
with patch(
"roboco.services.project.get_project_service",
lambda _db: _FakeProjectSvc(),
):
slugs = await AgentOrchestrator._intake_scope_slugs(
db=object(),
project_slug=None,
product_id=None,
project_ids=ids,
)
assert slugs == ["proj-1", "proj-2"]
@pytest.mark.asyncio
async def test_megatask_scope_with_unresolvable_project_raises(self) -> None:
class _FakeProjectSvc:
async def get(self, _pid: Any) -> Any:
return None
with (
patch(
"roboco.services.project.get_project_service",
lambda _db: _FakeProjectSvc(),
),
pytest.raises(ValueError, match="not found"),
):
await AgentOrchestrator._intake_scope_slugs(
db=object(),
project_slug=None,
product_id=None,
project_ids=["11111111-1111-1111-1111-111111111111"],
)
@pytest.mark.asyncio
async def test_megatask_scope_with_one_unresolvable_id_raises(self) -> None:
# A PARTIAL failure (one of N ids invalid) must fail loud, not silently
# clone fewer repos than the agent was told it has.
good = "11111111-1111-1111-1111-111111111111"
bad = "22222222-2222-2222-2222-222222222222"
class _FakeProjectSvc:
async def get(self, pid: Any) -> Any:
return SimpleNamespace(slug="proj-a") if str(pid) == good else None
with (
patch(
"roboco.services.project.get_project_service",
lambda _db: _FakeProjectSvc(),
),
pytest.raises(ValueError, match="not found"),
):
await AgentOrchestrator._intake_scope_slugs(
db=object(),
project_slug=None,
product_id=None,
project_ids=[good, bad],
)
# ---------------------------------------------------------------------------
# spawn_intake_session / reap_intake_session — orchestration (docker mocked).
# ---------------------------------------------------------------------------
def _fake_route() -> SimpleNamespace:
return SimpleNamespace(
provider_type=SimpleNamespace(value="anthropic"),
model_name="opus",
base_url=None,
auth_token=None,
)
def _wire_spawn_mocks(
monkeypatch: pytest.MonkeyPatch,
orch: AgentOrchestrator,
run_calls: list[list[str]],
) -> None:
"""Patch every external boundary spawn_intake_session touches."""
async def _clone(_p: Any, _pr: Any, _pids: Any = None) -> tuple[str, list[str]]:
return "/data/workspaces/roboco/board/intake-1", [
"/data/workspaces/roboco/board/intake-1"
]
async def _route(_aid: str) -> Any:
return _fake_route()
async def _noop(*_a: Any, **_k: Any) -> None:
return None
async def _run(cmd: list[str]) -> str:
run_calls.append(cmd)
return "containerid0123456789"
monkeypatch.setattr(orch, "_clone_intake_scope", _clone)
monkeypatch.setattr(orch, "_resolve_agent_route", _route)
monkeypatch.setattr(orch, "_ensure_agent_image", _noop)
monkeypatch.setattr(orch, "_remove_container", _noop)
monkeypatch.setattr(orch, "_run_container_cmd", _run)
monkeypatch.setattr(orch, "_fire_audit", lambda **_k: None)
monkeypatch.setattr(
orch,
"_generate_composed_prompt",
lambda *_args, **_kwargs: Path("/tmp/intake-1-prompt.md"),
)
monkeypatch.setattr(
orch,
"_resolve_intake_host_paths",
lambda: {
"claude": "/home/runner/.claude",
"prompt": "/data/prompts-generated/intake-1-prompt.md",
"workspaces": "/data/workspaces",
},
)
class TestSpawnIntakeSession:
@pytest.mark.asyncio
async def test_spawn_registers_session_and_instance(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
run_calls: list[list[str]] = []
_wire_spawn_mocks(monkeypatch, orch, run_calls)
instance = await orch.spawn_intake_session("sess-1", project_slug="roboco")
# Live relay session opened for the container.
session = prompter_live.get_live_registry().get("sess-1")
assert session is not None
assert session.agent_id == INTAKE_AGENT_ID
# Orchestrator instance tracked and marked active.
assert orch._instances[INTAKE_AGENT_ID] is instance
assert instance.container_id == "containerid0123456789"
# The cloned cwd reached the docker cmd.
assert "ROBOCO_WORKSPACE=/data/workspaces/roboco/board/intake-1" in run_calls[0]
@pytest.mark.asyncio
async def test_scope_must_be_exactly_one(self) -> None:
orch = _make_minimal_orchestrator()
with pytest.raises(ValueError, match="exactly one"):
await orch.spawn_intake_session("s", project_slug="roboco", product_id="p")
with pytest.raises(ValueError, match="exactly one"):
await orch.spawn_intake_session("s")
# A MegaTask scope cannot combine with a single-project scope.
with pytest.raises(ValueError, match="exactly one"):
await orch.spawn_intake_session(
"s", project_slug="roboco", project_ids=["p1"]
)
@pytest.mark.asyncio
async def test_spawn_accepts_megatask_project_ids_scope(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
run_calls: list[list[str]] = []
_wire_spawn_mocks(monkeypatch, orch, run_calls)
instance = await orch.spawn_intake_session(
"sess-mega", project_ids=["11111111-1111-1111-1111-111111111111"]
)
assert orch._instances[INTAKE_AGENT_ID] is instance
assert run_calls # the container actually launched for the MegaTask scope
@pytest.mark.asyncio
async def test_spawn_reaps_prior_session_first(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
run_calls: list[list[str]] = []
_wire_spawn_mocks(monkeypatch, orch, run_calls)
stopped: list[str] = []
async def _stop(aid: str, **_kw: Any) -> None:
stopped.append(aid)
monkeypatch.setattr(orch, "stop_agent", _stop)
# A prior live container already registered for this agent.
orch._instances[INTAKE_AGENT_ID] = AgentInstance(agent_id=INTAKE_AGENT_ID)
await orch.spawn_intake_session("sess-2", project_slug="roboco")
assert stopped == [INTAKE_AGENT_ID] # the old one was reaped first
@pytest.mark.asyncio
async def test_initial_message_is_scheduled(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
run_calls: list[list[str]] = []
_wire_spawn_mocks(monkeypatch, orch, run_calls)
scheduled: list[tuple[str, str]] = []
monkeypatch.setattr(
orch,
"_schedule_intake_first_message",
lambda sid, text: scheduled.append((sid, text)),
)
await orch.spawn_intake_session(
"sess-3", project_slug="roboco", initial_message="build X"
)
assert scheduled == [("sess-3", "build X")]
class TestStartIntakeSession:
"""Non-blocking start: relay opens synchronously, spawn runs in the background."""
@pytest.mark.asyncio
async def test_opens_relay_now_and_schedules_spawn(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
spawned: list[str] = []
async def _spawn(session_id: str, **_kw: Any) -> Any:
spawned.append(session_id)
return AgentInstance(agent_id=INTAKE_AGENT_ID)
monkeypatch.setattr(orch, "_spawn_intake_container", _spawn)
await orch.start_intake_session("sess-A", project_slug="roboco")
# Relay is open the instant start returns — the SSE stream can connect
# before the (slow) container spawn finishes.
assert prompter_live.get_live_registry().get("sess-A") is not None
await asyncio.sleep(0) # let the scheduled bg spawn run
assert spawned == ["sess-A"]
@pytest.mark.asyncio
async def test_rejects_bad_scope(self) -> None:
orch = _make_minimal_orchestrator()
with pytest.raises(ValueError, match="exactly one"):
await orch.start_intake_session("s", project_slug="r", product_id="p")
class TestSpawnGuarded:
"""A background spawn failure surfaces on the relay instead of dying silently."""
@pytest.mark.asyncio
async def test_failure_pushes_error_and_closes(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
registry = prompter_live.get_live_registry()
registry.open("sess-B", INTAKE_AGENT_ID)
pushed: list[tuple[str, dict[str, Any]]] = []
closed: list[str] = []
async def _boom(_session_id: str, **_kw: Any) -> Any:
raise RuntimeError("clone exploded")
def _push(sid: str, ev: dict[str, Any]) -> bool:
pushed.append((sid, ev))
return True
monkeypatch.setattr(orch, "_spawn_intake_container", _boom)
monkeypatch.setattr(registry, "push", _push)
monkeypatch.setattr(registry, "close", closed.append)
await orch._spawn_intake_container_guarded(
"sess-B", project_slug="roboco", product_id=None, initial_message=None
)
assert len(pushed) == 1
assert pushed[0][1]["kind"] == "error"
assert "clone exploded" in pushed[0][1]["text"]
assert closed == ["sess-B"]
class TestReapIntakeSession:
@pytest.mark.asyncio
async def test_reap_closes_session_and_stops_container(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
stopped: list[str] = []
async def _stop(aid: str, **_kw: Any) -> None:
stopped.append(aid)
monkeypatch.setattr(orch, "stop_agent", _stop)
registry = prompter_live.get_live_registry()
registry.open("sess-x", INTAKE_AGENT_ID)
await orch.reap_intake_session("sess-x")
assert registry.get("sess-x") is None # relay session closed
assert stopped == [INTAKE_AGENT_ID]
class TestDeliverWhenReady:
@pytest.mark.asyncio
async def test_retries_until_receiver_is_up(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
orch = _make_minimal_orchestrator()
registry = prompter_live.get_live_registry()
succeed_on = 2 # fails once, then succeeds
attempts = {"n": 0}
async def _deliver(_sid: str, _text: str) -> bool:
attempts["n"] += 1
return attempts["n"] >= succeed_on
monkeypatch.setattr(registry, "deliver", _deliver)
await orch._deliver_when_ready("sess-y", "hi", attempts=5, delay=0)
assert attempts["n"] == succeed_on # stopped as soon as delivery succeeded