Files
roboco/tests/unit/runtime/test_orchestrator_shutdown_drain.py
T
3ccc723cd4 v0.17.0 — Wave 3: sandbox DB, DB isolation, mobile UI, cloud auth, X account, roadmap engine (#303)
* feat(sandbox): throwaway per-agent Postgres/Redis sandbox containers

Orchestrator-provisioned sibling containers per agent spawn
(SandboxProvisioner, roboco/runtime/sandbox.py). Per-project opt-in via
projects.sandbox_services (migration 057); master switch
ROBOCO_SANDBOX_DB_ENABLED, default-off, armed in the NAS compose only.

When active, ROBOCO_TEST_DB_* / ROBOCO_TEST_REDIS_* point at the sandbox
and the prod-creds gate-env injection is suppressed (sandbox replaces,
never coexists). Sandbox lifetime tracks the agent container: teardown at
every removal path, orphan janitor at startup + each reaper tick with a
grace window for mid-flight spawns. The pre-spawn stale-clear spares the
just-provisioned sandbox; provision pre-clears stale same-named
containers from a crash-missed teardown.

Panel: per-project sandbox-service switches in the edit dialog + feature
flag card entry.

* docs: CLAUDE.md entry for the sandboxed dev DB/Redis subsystem

* feat(security): isolate prod Postgres/Redis from agent containers (roboco_data network)

Second user-defined bridge roboco_data carries postgres+redis only; the
orchestrator is multi-homed (default + data). Spawned agents and their
sandbox sidecars stay on roboco_default and can no longer resolve or
reach roboco-postgres:5432 / roboco-redis:6379 (redis has no auth —
membership is its only containment). Normal bridge, so host-published
ports (15432/16379) keep working. Applied to both build composes and
the registry compose; docker-compose.yml re-synced byte-identical with
docker-compose.yaml (it had drifted by the sandbox flag block).

ROBOCO_DB_NETWORK_ISOLATED (config default false, armed alongside the
topology) suppresses the legacy _append_gate_env prod-creds injection:
under isolation those creds dead-end, and unreachable creds are worse
than none. DB-needing projects opt into sandbox_services instead. The
flag is deliberately not a panel feature flag - it must travel with the
compose networks: stanzas.

Preserved by construction: agent<->agent A2A and orchestrator->agent SDK
polls on :9000, MCP->orchestrator on :8000, ollama reachability, docker
exec/inspect (daemon socket), host port publishing.

* feat(panel): full mobile responsiveness pass

Shared primitives: useIsMobile (useSyncExternalStore, hydration-safe,
memoized matchMedia subscribe), ResponsiveTable table->card switch below
md (single subtree mounted, no duplicated interactive rows), scrollable
snap TabsList in the base primitive (justify-center-safe so the first
tab stays reachable on overflow), persistent md:hidden bottom tab bar
(Overview/Tasks/Kanban/Chat, safe-area padded).

Applied: card lists for tasks/projects/products/work-sessions/sessions
+ the three raw metrics tables; CEO approval queue / release proposal /
playbook review action rows stack on narrow; command-center reorders
approvals above the fold on mobile; task-header metadata wraps;
Communications + A2A become URL-driven single-pane drill-downs below lg
(fixes the unconstrained-height ScrollArea bug) with dvh heights;
recharts label density/radius adapts via useIsMobile; git diff viewer
gets mobile font + wrap toggle; vh->dvh sweep; chat composers get
safe-area-inset padding; dashboard main p-4 md:p-6 + pb-20 for the bar.

Verified at 375px on the built app: bottom bar, drawer, approval-first
overview, swipeable kanban tab strip. All gates green (eslint, tsc,
vitest 249, next build 24/24 routes).

* feat(auth): cloud auth via FastAPI Users (default-off, single-user cookie session)

ROBOCO_CLOUD_AUTH_ENABLED (default off) lets the panel/API be exposed
beyond localhost without changing the CEO's local no-login flow while
off — get_agent_context and the WS gate are byte-for-byte unchanged in
off-mode. On: header-trust dies for humans — any agent-role claim (ceo
or a privileged PM/board role) with no valid HMAC token or session
cookie is 401, closing the header-spoof hole on the host-published
:8000 port for every role. The agent-fleet HMAC path and the system
self-PATCH keep working unmodified in both modes.

Single seeded CEO user (migration 058 users table, UserTable), no
registration router — idempotent env-driven upsert at startup by PK.
Cookie transport (httponly/secure/samesite=lax) + a JWTStrategy bound
to a fingerprint of the current password hash (rotating the password
invalidates every prior session). Sliding 30-day session: every
authenticated request re-mints the cookie, so an active session never
expires — no unexpected logouts.

Panel: (auth)/login page + proxy.ts (Next 16 rename of middleware; probes
/auth/status over the docker-internal URL, fails open to off) gate the
dashboard; client.ts gets withCredentials + 401->/login. nginx unchanged.

Review hardening: broadened the on-mode rejection from ceo-only to every
non-CEO role without a valid token (was only closed when
ROBOCO_AGENT_AUTH_REQUIRED was also armed); Next-16 proxy.ts rename to
clear the middleware deprecation warning.

* feat(x): RoboCo X account engine — HoM drafts, per-post CEO approval (default-off)

ROBOCO_X_ENGINE_ENABLED (default off, inert without creds). Mirrors the
ReleaseManagerEngine held-artifact shape: XEngine drafts a post when a
release publishes (via a draft_release_post seam on ReleaseProposalService
.approve) and drafts replies to meaningful mentions (dedicated poll loop,
x_seen_mentions dedup ledger, per-cycle/open caps). Drafting is
local-model-only, clamped to 280 chars. Nothing auto-posts — every tweet
is a held task (source x_post/x_reply, confirmed_by_human=False,
Secretary-owned, dispatcher-skipped) the CEO edits/approves/rejects in a
panel queue.

The four OAuth 1.0a secrets live Fernet-encrypted in a singleton
x_credentials row (migration 059, all-or-nothing, API returns only
has_credentials); decryption is server-side, agents never hold creds or
egress. Hand-rolled OAuth 1.0a HMAC-SHA1 signer, no new dependency;
NullXClient makes the unconfigured path a graceful no-op.

XPostService.approve (CEO-only) is the sole caller of post_tweet.

Review hardening: closed a double-post race — the approve path now
re-reads committed task state inside the Redis lock and commits COMPLETED
before releasing, so a concurrent approve that acquires the lock after the
winner released can't re-post (SET-NX is non-waiting, and the route-level
commit landed after the lock dropped). Added a regression test.

* feat(roadmap): board roadmap engine — PO proposes themed cycles, CEO approves per-item (default-off)

ROBOCO_ROADMAP_ENGINE_ENABLED (default off). Weekly, RoadmapEngine opens
ONE held exploration task (source=board_roadmap, confirmed_by_human=False,
Product-Owner-assigned), deduped to one open cycle. A dedicated one-shot
_dispatch_roadmap_exploration spawns the PO solo (not the two-reviewer
board path, which would also spawn HoM + fire Approve-&-Start). The PO
explores read-only (git/KB/metrics/releases/charter/web) and makes one
propose_roadmap call (PO-only content verb) authoring a themed cycle —
goal + 3-7 item drafts — persisted as a roadmap_cycle marker (no table,
no migration; head stays 059).

The CEO acts per-item in the panel roadmap queue: approve materializes a
BACKLOG task (source=roadmap, no assignee — never auto-starts), reject
records a reason; all-items-terminal completes the exploration task.
RoadmapService is idempotent per item. Dispatchers skip board_roadmap.

Includes a real SQLAlchemy dirty-check fix (deep-copy the JSON marker
before mutating, or the in-place edit + reassign compares equal to its
own baseline and the UPDATE is skipped).

Review hardening: create_task_from_draft now honors a draft-declared
source only from a {prompter, roadmap} whitelist — drafts are
LLM-authored, so an unbounded source could impersonate a privileged
origin (release_manager would even wedge that engine's dedup).

* chore(release): 0.17.0

Wave 3 — six default-off subsystems: sandboxed dev DB/Redis, prod
Postgres/Redis network isolation, full mobile UI pass, cloud auth
(FastAPI Users), the RoboCo X account engine, and the board roadmap
engine. Plus the waves 1+2 work already on master since 0.16.0.

Version bumped across the canonical set (config.py, __init__.py,
pyproject.toml, panel/package.json, uv.lock); CHANGELOG [Unreleased]
cut to [0.17.0]; docs/map delta added.

Compose: every optional feature armed :-true in the NAS composes, OFF
in the user-facing registry compose. Two opt-in exceptions default off
(CLOUD_AUTH — needs email/password/secret + TLS, would otherwise fail
startup; ROUTING_STRICT — fail-closed spawning). DB_NETWORK_ISOLATED
stays on in both (coupled to the roboco_data topology).

* chore(compose): arm cloud_auth + routing_strict ON in the NAS composes

Every feature defaults ON in the NAS composes per policy — these two
were wrongly left off. Both keep the ${VAR:-true} form so the operator
controls the real runtime via .env: cloud auth needs
ROBOCO_CLOUD_AUTH_EMAIL/_PASSWORD/_SECRET + TLS set there before a boot
(else startup fails loud), and routing_strict is fail-closed. Registry
compose keeps both off.

* fix(ci): reflow board.md prose (quality gate) + document v0.17.0 env creds

The roadmap section added hard-wrapped prose that failed the markdown
prose gate; reflowed (token-invariant). Also brought .env.example
current: cloud auth (now armed — needs SECRET or startup fails), routing
strict, the X engine (panel-entered OAuth), and web research.

* fix(ci): reduce cyclomatic complexity of five wave-3 blocks (xenon gate)

The wave-3 subagents introduced C-rank functions the CI xenon gate
rejects (my per-item reviews ran ruff/mypy/pytest but not xenon):
- sandbox.janitor_sweep -> extract _list_labeled_sandboxes /
  _list_live_agent_containers / _prune_grace
- x_client.fetch_mentions -> extract _parse_mention_items
- x_engine.run_cycle -> extract _process_mentions
- orchestrator._dispatch_pm_work -> extract the source-skip into a
  MODULE-level _is_held_ceo_source (module, not method, so the
  wholesale-mocked dispatcher unit tests exercise the real logic)
- auth/seed.ensure_seed_user -> extract _apply_seed_updates (module avg -> A)

Behavior-preserving; full suite green (11902), xenon clean.

* fix(ci): declare pyjwt + fastapi-users-db-sqlalchemy as direct deps (deptry)

The cloud-auth code imports jwt and fastapi_users_db_sqlalchemy directly
but they were only transitive deps (via fastapi-users), which deptry
(quality gate, DEP003) rejects. Declared explicitly; deptry roboco/ clean.
Missed originally because local make quality stopped at earlier gates
before reaching deptry.

* feat(x): gate mention replies behind ROBOCO_X_REPLIES_ENABLED (default off)

Per CEO decision: the X engine should only post about releases by
default. Reading mentions needs a paid X API tier, so the mention-reply
half is now a deliberate opt-in on top of release posting.

New default-off flag x_replies_enabled gates the mentions poll loop
(_x_mentions_poll_loop) and XEngine.run_cycle; release-post drafting
(the release-proposal approve hook) is unaffected and still runs when
x_engine_enabled + credentials are set. Added to FEATURE_FLAGS + the
panel card. Tests: release posting works with replies off; run_cycle +
the poll loop are no-ops with replies off.

* fix: 401 only redirects to /login when cloud auth is on; panel-token strips .env quotes

Two bugs that together dead-ended login in secure mode:
- client.ts redirected to /login on ANY 401, so a mismatched panel
  token (header-trust/secure mode, cloud auth off) bounced the user to a
  login page whose backend route isn't mounted -> 404. Now it probes
  /auth/status (bare fetch, no interceptor re-entry) and only redirects
  when cloud_auth_enabled.
- make panel-token read the .env secret with grep|cut without stripping
  surrounding quotes, so a quoted ROBOCO_AGENT_AUTH_SECRET produced a
  token signed with the quotes included — which never verifies against
  the orchestrator (docker-compose/pydantic unquote the secret). Now
  strips surrounding single/double quotes.

* fix: git-log 500 on '|' in commit message; X queue shows an empty state

- GET /api/git/log 500'd (ValueError: Invalid isoformat) when a commit
  SUBJECT contained a '|' (e.g. the 'curl|sh' lockdown commit): the
  fixed '|' field delimiter let the subject's pipe shift the split so
  author+date collapsed into one field. Switched to \x1f (Unit
  Separator), which can't appear in commit content. Regression test with
  a piped subject.
- The X Post Queue returned null when empty, so there was no visible
  place for the X drafts. It now renders a discoverable empty state
  pointing at Settings -> X credentials.

* docs: bring docs/rag + docs/map current for v0.17.0 (waves 1-3)

Agent-facing RAG corpus and codebase map updated for every feature in
the 0.17.0 span, code-verified:
- wave 3: sandbox DB, DB network isolation, cloud auth, X engine
  (+ x_replies_enabled sub-flag), board roadmap engine — new RAG
  architecture pages + role/tool/config-reference updates; new symbols,
  migrations 057-059, panel surfaces, and the get_agent_context
  dual-path across the map slices.
- waves 1-2: A2A live view + switchboard, prompter memory
  (search_past_tasks), Secretary edit access + PM-lighter scope, the
  PR-gate auto-submit turn cut (ROBOCO_PR_GATE_AUTO_SUBMIT_ENABLED).
- correctness fix: api-routes-schemas.md no longer claims the A2A admin
  routes are reachable by any authenticated agent — they carry a
  _require_ceo gate (wave 2c).

docs/internal, _front.md deltas, and the frozen _complete_map.md
snapshot untouched.

* fix(rag): atomic upsert for indexed-doc tracking (kills e2e segfault)

The indexed-document tracking write used check-then-insert in two paths
(IndexedDocumentRepository.upsert_batch and the file-source
_upsert_doc_record). Under concurrent indexing both callers saw no row
and both inserted, so the second violated uq_indexed_doc_source and
poisoned its transaction — surfacing in CI as the intermittent
_checkin_failed SIGSEGV on the failed connection's pool checkin.

Both paths now use INSERT ... ON CONFLICT DO UPDATE against the
constraint: coalesce keeps an existing title/preview when the new value
is empty (matching the old guards) and metadata is jsonb-merged. The
batch dedupes within itself first (ON CONFLICT can't touch a row twice
in one statement). expire_all after the Core upsert keeps same-session
ORM reads consistent with the merged DB row.

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2026-07-03 19:24:00 +02:00

250 lines
9.0 KiB
Python

"""Drain ``_bg_tasks`` on shutdown so fire-and-forget writes (respawn_tracker
upserts, audit-log writes, intake first-message delivery) are not abandoned.
Invariant: ``Orchestrator.stop()`` drains ``_bg_tasks`` with a bounded timeout —
short DB writes finish before the process exits (data preserved), a stuck task
is cancelled once the deadline passes (can't hang shutdown). The ``stop_agent``
loop is wrapped so one agent's stop error can't skip the drain.
"""
from __future__ import annotations
import asyncio
from typing import Any
from unittest.mock import AsyncMock, patch
import pytest
from roboco.models.runtime import AgentInstance
from roboco.runtime.orchestrator import (
_SHUTDOWN_DRAIN_TIMEOUT_SECONDS,
AgentOrchestrator,
)
# Floor encoding the logical-regression guard: a drain deadline below this
# would risk dropping a legitimate short DB write (an upsert that needs a
# second under load) before it commits — the exact data loss this fix targets.
# Named (not magic) for ruff PLR2004.
_MIN_DRAIN_TIMEOUT = 3.0
def _make_orchestrator() -> AgentOrchestrator:
"""AgentOrchestrator with constructor I/O skipped; stop() deps ready.
``stop()`` cancels the named loop tasks (all None here → no-op) and the
agents in ``_instances`` (empty here), then must drain ``_bg_tasks``.
"""
with patch.object(AgentOrchestrator, "__init__", return_value=None):
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch._instances = {}
orch._bg_tasks = set()
orch._pm_respawn_tracker = {} # #74: stop() flushes this after the drain
# Every named background loop ``stop()`` cancels — None makes each a no-op
# so the test exercises ONLY the _bg_tasks drain.
for attr in (
"_health_task",
"_dispatcher_task",
"_sweeper_task",
"_rate_limit_probe_task",
"_strategy_engine_task",
"_external_pr_poll_task",
"_self_heal_task",
"_ci_watch_task",
"_dep_update_task",
"_release_manager_task",
"_x_mentions_task",
"_roadmap_engine_task",
):
setattr(orch, attr, None)
return orch
def test_shutdown_drain_timeout_is_named_module_constant() -> None:
assert isinstance(_SHUTDOWN_DRAIN_TIMEOUT_SECONDS, int | float)
assert _SHUTDOWN_DRAIN_TIMEOUT_SECONDS > 0
def test_shutdown_drain_timeout_is_generous() -> None:
"""A short DB upsert under load can legitimately take a moment; the drain
deadline must not drop it. This guards the logical regression: a too-short
drain would silently lose the very writes it exists to preserve."""
assert _SHUTDOWN_DRAIN_TIMEOUT_SECONDS >= _MIN_DRAIN_TIMEOUT
@pytest.mark.asyncio
async def test_stop_drains_completing_bg_task_before_returning() -> None:
"""A bg task that finishes quickly MUST complete (its side effect observed)
before ``stop()`` returns. Without the drain, ``stop()`` returns immediately
and the task is abandoned mid-flight — the data-loss tail."""
orch = _make_orchestrator()
ran: list[bool] = []
async def _completes() -> None:
await asyncio.sleep(0.01)
ran.append(True)
orch._bg_tasks.add(asyncio.create_task(_completes()))
await asyncio.wait_for(orch.stop(), timeout=5.0)
assert ran == [True], "completing bg task was abandoned by stop()"
@pytest.mark.asyncio
async def test_stop_does_not_hang_on_stuck_bg_task(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""A bg task that never completes MUST NOT hang shutdown past the drain
deadline — it is cancelled once the drain times out. Without the drain,
a stuck bg task would let ``stop()`` (and thus the process) hang forever.
Deterministic: the drain deadline is patched tiny so a bounded fail-close is
asserted in well under a second, never relying on the real 5s default."""
monkeypatch.setattr(
"roboco.runtime.orchestrator._SHUTDOWN_DRAIN_TIMEOUT_SECONDS", 0.05
)
orch = _make_orchestrator()
async def _hangs() -> None:
await asyncio.Future() # never resolves
stuck = asyncio.create_task(_hangs())
orch._bg_tasks.add(stuck)
await asyncio.wait_for(orch.stop(), timeout=2.0)
assert stuck.cancelled(), "stuck bg task was not cancelled by the drain"
@pytest.mark.asyncio
async def test_stop_failing_agent_does_not_skip_drain() -> None:
"""If one agent's ``stop_agent`` raises, the drain must still run —
otherwise a single bad agent would re-introduce the data-loss tail for every
in-flight bg write. The completing bg task should still finish."""
orch = _make_orchestrator()
orch._instances["bad-agent"] = AgentInstance(agent_id="bad-agent")
async def _raises(_aid: str, **_kw: Any) -> None:
raise RuntimeError("boom")
patch.object(orch, "stop_agent", _raises).start()
ran: list[bool] = []
async def _completes() -> None:
await asyncio.sleep(0.01)
ran.append(True)
orch._bg_tasks.add(asyncio.create_task(_completes()))
await asyncio.wait_for(orch.stop(), timeout=5.0)
assert ran == [True], "failing stop_agent skipped the drain (data lost)"
@pytest.mark.asyncio
async def test_stop_is_idempotent_double_call_is_noop() -> None:
"""stop() is idempotent: the lifespan path and bootstrap's finally block both
call it, so the second call must be a clean no-op — not a re-drain or re-stop
of already-stopped agents — guarded by ``_stopped``."""
orch = _make_orchestrator()
real_drain = orch._drain_bg_tasks
drain_calls = 0
async def counting_drain() -> None:
nonlocal drain_calls
drain_calls += 1
await real_drain()
# Override via an Any-typed view so the assignment bypasses mypy's
# method-assign check while staying a plain attribute write (no setattr).
orch_any: Any = orch
orch_any._drain_bg_tasks = counting_drain
await orch.stop()
assert drain_calls == 1, "first stop() drained the bg tasks"
assert orch._stopped is True
await orch.stop() # safety-net double-call (lifespan already stopped it)
assert drain_calls == 1, "second stop() must not re-drain (idempotent no-op)"
assert orch._stopped is True
# ---------------------------------------------------------------------------
# #74: stop() flushes the authoritative in-memory respawn snapshot AFTER the
# bounded drain so a deadline-cancelled fire-and-forget persist can't leave the
# durable count lagging the in-memory counter (which would re-burn the strike
# threshold on the next restart).
# ---------------------------------------------------------------------------
_BE_PM_TID = "11111111-1111-1111-1111-111111111111"
_FE_PM_TID = "22222222-2222-2222-2222-222222222222"
_EXPECTED_FLUSHES = 2 # two seeded respawn rows → two shutdown persists
def _respawn_record(count: int) -> dict[str, Any]:
return {
"count": count,
"last_status": "blocked",
"last_check": None,
"tracing_resets": 0,
"notified": False,
}
@pytest.mark.asyncio
async def test_stop_flushes_respawn_tracker_after_drain() -> None:
"""#74: stop() writes every in-memory respawn row after the drain, carrying
the latest count so the durable counter matches memory on the next restart."""
orch = _make_orchestrator()
orch_any: Any = orch
orch_any._pm_respawn_tracker = {
("be-pm", _BE_PM_TID): _respawn_record(4),
("fe-pm", _FE_PM_TID): _respawn_record(2),
}
persist = AsyncMock()
orch_any._persist_respawn_record = persist
await asyncio.wait_for(orch.stop(), timeout=5.0)
assert persist.await_count == _EXPECTED_FLUSHES
keys = {(c.args[0], c.args[1]) for c in persist.await_args_list}
assert keys == {("be-pm", _BE_PM_TID), ("fe-pm", _FE_PM_TID)}
counts = {c.args[0]: c.args[2]["count"] for c in persist.await_args_list}
assert counts == {"be-pm": 4, "fe-pm": 2}
@pytest.mark.asyncio
async def test_flush_respawn_tracker_noop_when_empty() -> None:
"""An empty tracker flushes nothing — no DB churn on a clean shutdown."""
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch_any: Any = orch
orch_any._pm_respawn_tracker = {}
persist = AsyncMock()
orch_any._persist_respawn_record = persist
await orch._flush_respawn_tracker()
persist.assert_not_awaited()
@pytest.mark.asyncio
async def test_flush_respawn_tracker_swallows_row_errors() -> None:
"""#74: a row whose persist raises must not skip the remaining rows or crash
shutdown — the in-memory value is gone once the process exits anyway."""
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch_any: Any = orch
orch_any._pm_respawn_tracker = {
("be-pm", _BE_PM_TID): _respawn_record(4),
("fe-pm", _FE_PM_TID): _respawn_record(2),
}
async def _persist(slug: str, _tid: str, _record: dict[str, Any]) -> None:
if slug == "be-pm":
raise RuntimeError("db down")
# fe-pm succeeds
orch_any._persist_respawn_record = _persist
# Must not raise — the failing row is logged and the rest still flushed.
await orch._flush_respawn_tracker()