Files
roboco/tests/unit/services/test_heartbeat_mutex.py
e9d0e0bd48 feat(video): 0.19.0 video engine (Remotion) + preview auth + render persistence (#307)
* feat(video): Phase A — VideoEngine origination spine + held-source gates

New default-off engine skeleton: opens a UX/UI authoring task (source=video, assigned to a ux-dev, LOW complexity to clear the dev-needs-subtasks guard) and materializes a held CEO-approval draft (source=video_post). Excludes video_post from all three held-source skip sites; adds the video_draft marker, six config flags, and the feature-flag entries. Origination + gate behavior unit-tested.

* refactor(orchestrator): fold _dispatch_dev_work skip chain into a helper

The per-source if/continue chain grew past xenon's --max-absolute B when the video_post held source joined it. Extract _is_non_dev_dispatch_source (every held-CEO source plus the two Board exploration sources) so the dev loop's skip is one flat call. Behavior-identical.

* feat(video): Phase B — propose_video do-tool (metadata-only, team-gated)

UX/UI dev records a video's composition ref + per-platform captions onto the authoring task's video_draft marker. Team-gated at runtime via _caller_team (Role.DEVELOPER can't tell a ux-dev from a be-dev). Resolves the caller's ACTIVE task via get_active_task_for_agent, not an oldest-first scan that would clobber a second open video task. Metadata only, no render. Wired through do_server + route + schema; added to _DEV_DO.

* feat(video): Phase D — render loop + RemotionRenderer client

Orchestrator-async _video_render_loop renders a completed authoring task's merged composition to MP4 (vertical + square) via the remotion-renderer sidecar and materializes the held video_post draft. RemotionRenderer tars the read-clone's motion/ source, POSTs it, and saves the returned MP4 bytes to a TASK-scoped local path (no shared volume; a composition is reused across videos so a composition-scoped path would clobber an earlier draft). Render failures bounded-retry (read-clone catch-up window, transient sidecar) up to a cap, then terminal-fail. Client tested vs a mock transport; loop vs a mock renderer + real DB.

* feat(video): Phase C — release / spotlight / on-demand video triggers

Three entry points open a UX/UI video-authoring task via VideoEngine.open_video_task: (1) a published release drafts a companion video — best-effort in ReleaseProposalService.approve, never fails the publish; script from the CHANGELOG via the local model with a template fallback. (2) propose_feature_spotlight gains optional wants_video/video_script — best-effort, gated on video_on_spotlight, default-off leaves the spotlight flow byte-for-byte unchanged. (3) POST /video/request (CEO-only) for an on-demand brief, with clean disabled/not_opened responses. All gated on video_engine_enabled.

* fix(video): savepoint-isolate video-task inserts (F042 poisoned session)

The best-effort try/except around open_video_task (release-publish + spotlight hooks) swallowed the Python exception, but a DBAPI error at the insert flush left the shared session must-rollback — so the caller's next commit (release finalize / request boundary) threw PendingRollbackError: the release stuck 'pending' after actually publishing, or the spotlight draft + HTTP response were lost. Wrap both inserts (open_video_task, _originate_video_post) in a begin_nested savepoint (the repo's established F042 pattern) so a DB error rolls back only the insert. open_video_task returns None (every caller already handles it); _originate_video_post propagates to the render loop's handler. Regression test: an insert FK error returns None with the session left usable. Dormant while the flags were off; armed on the NAS.

* feat(video): Phase G — motion/ package + remotion-renderer sidecar + compose

In-repo Remotion v4 motion/ package (ReleaseAnnouncement composition; calculateMetadata returns 1080x1920 vertical / 1080x1080 square from inputProps.orientation) + a credential-free remotion-renderer sidecar: untar the POSTed motion/ source, bundle (LRU-cached per source sha), selectComposition + renderMedia h264, stream the MP4 bytes back — matching the RemotionRenderer client contract. docker/remotion.Dockerfile on Debian (Chrome apt deps, build-time Chrome pre-warm, ffmpeg bundled in @remotion/renderer). Wired into both compose files (roboco_default only, shm_size 1gb, /health check) + the release publish matrix. Verified via a real local render of both cuts; the Debian docker build is the CEO's to run.

* chore(video): D-hardening — video_post source_task_id + render-loop docstring

Add a source_task_id back-reference to the video_post held-draft marker (traceability from a draft to its authoring task; also makes the render loop's two-key idempotency check wireable later). Fix the render-loop test's stale docstring ('never retried' -> bounded-retry). Both from the Phase D critic's non-blocking follow-ups.

* feat(video): Phase E1 — VideoPostService + heartbeat mutex (approve->post)

CEO-approve->post service: heartbeat-renewed Redis mutex (fail-closed, grace=ttl-2*heartbeat), re-read-in-lock double-post guard, per-platform durable commits (asyncio.shield-ed, settle-before-rollback on lock-loss), all writes inside the lock (captions validated pre-lock, applied in-lock — no stale whole-column clobber), idempotent, per-platform retry-skip. Poster interfaces (X/TikTok, mocked here). Reject + list-held-drafts. Survived 3 adversarial rounds; residual = a crash in the poster->commit window (CEO-gated low-freq, documented).

* fix(video): G-hardening — renderer leaks + Share Tech Mono brand font

Sidecar: give bundle() an explicit outDir tracked + deleted on LRU eviction (was leaking ~19MB remotion-webpack-bundle-* per source); res.on('close') cleanup so an aborted/retried download no longer leaks its remotion-out-* MP4 dir. Fonts: vendor Share Tech Mono (roboco-website brand font) as the display face (self-hosted woff2, 400-weight, headline fontWeight 700->400 to avoid faux-bold) + self-hosted Inter body — no gstatic fetch at render time (lsof-verified). Extras: composition_id whitelist (400) + Multer error middleware (400/413).

* feat(video): Phase E2 — X v2 + TikTok posters, tiktok_credentials, routes

LiveXVideoPoster (X v2 chunked media upload: init/append/finalize/STATUS-poll -> tweet w/ media_ids, OAuth1 signer reused). LiveTikTokPoster (OAuth2 inbox: init -> chunked PUT with asymmetric final chunk -> status-fetch; 401 -> refresh_token grant, rotated token persisted). tiktok_credentials Fernet singleton + migration 062 (single head). Routes: CEO approve/reject + list held drafts + write-only tiktok creds, wiring real posters into VideoPostService. Residual: a lock-loss right after a token-refresh flush can discard the rotated token (same rare CEO-gated class as the documented post->commit window).

* feat(video): Phase F — panel video-post queue + TikTok creds card + flags

video-post-queue.tsx: <video> MP4 preview with 9:16/1:1 cut switch, per-platform editable captions (280/2200 counters, over-limit disables approve), approve/reject, Request-a-video dialog. tiktok-credentials-card.tsx (4 write-only OAuth2 fields). feature-flags-card inlines TikTokCredentialsForm under video_engine_enabled. Mounted in command-center. tsc/eslint clean, 273 panel tests green. NOTE: needs the GET /video/posts/{id}/media route + mp4_paths on VideoPostResponse (folded into H) for the preview source.

* feat(video): Phase H — media route + e2e smoke + NAS arming + docs

GET /video/posts/{id}/media?cut= (CEO-gated FileResponse of the rendered MP4; closes the panel preview gap) + mp4_paths on VideoPostResponse. e2e smoke tests/e2e_smoke/test_video_pipeline.py (full flow, sidecar+X/TikTok mocked; asserts dispatcher skips, render-loop materialize, propose_video team-gate, approve idempotency). NAS arming: docker-compose.yml/.yaml ROBOCO_VIDEO_ENGINE_ENABLED/ON_RELEASE/ON_SPOTLIGHT default-on (.yaml resynced to .yml); registry stays off. CLAUDE.md video-engine section + CHANGELOG. Fixed 2 pre-existing route-test pollution leaks. Full suite 11763 passed.

* fix(video): auth-carrying preview, media route confinement, VideoPost type drift

Three fixes along the video preview path:

1. panel video preview auth: the <video> element was pointed straight at
   GET /video/posts/{id}/media, but a native <video src> GET carries none
   of axios's X-Agent-ID/X-Agent-Role headers — so in the default
   header-trust deployment the request 401s. Fetch the cut via
   videoApi.getMediaBlob (axios, responseType: blob) and drive <video>
   off a URL.createObjectURL result instead. The object URL is revoked
   on cut-change (the previous cut's URL) and on unmount, so neither
   cut switches nor row teardown leak blob URLs.

2. backend media route confinement: GET /video/posts/{id}/media now
   resolves mp4_path and refuses it with 404 when it falls outside
   settings.video_output_dir. Defense-in-depth against any future
   writer of mp4_paths serving files from arbitrary disk locations.

3. panel VideoPost type/comment drift: added mp4_paths to the
   VideoPost interface (the committed VideoPostResponse already
   carries it), and corrected the stale comment on videoMediaUrl
   that claimed no route served the rendered bytes — the route has
   existed since the media endpoint landed; the comment now describes
   why getMediaBlob exists instead of a direct <video src>.

* Persist rendered videos to data in physical storage.

* ++

* docs(video): 0.18.0 CHANGELOG entry + RAG + map reference for video engine

- Move the video engine bullet from [Unreleased] into [0.18.0] and note
  the ROBOCO_VIDEO_OUTPUT_DIR bind-mount persistence.
- Add docs/rag/architecture/video-engine.md (mirrors x-engine.md shape:
  enable/disable, three triggers, render loop + sidecar, CEO gate, media
  route confinement, credentials).
- Reference the video render loop in docs/map/orchestrator.md's engine list.

* chore(video): re-bump to 0.19.0 + sync registry compose defaults

Version was wrongly bumped to 0.18.0; 0.18.0 is an already-released
section. Restore its 2026-07-04 date and move the video-engine CHANGELOG
bullet into a new [0.19.0] - 2026-07-05 section above it. Bump
pyproject.toml, roboco/__init__.py, roboco/config.py (app_version),
panel/package.json, and the motion/README inputProps example to 0.19.0.

docker-compose.registry.yml: add ROBOCO_VIDEO_ENGINE_ENABLED /
_VIDEO_ON_RELEASE / _VIDEO_ON_SPOTLIGHT defaulted false (NAS arms them
true), and comment out the video-renders bind mount with a short note
so the public registry image ships video off by default. Structural
sync with docker-compose.yml maintained.

* fix(video): rate-limit /render + reflow motion/README

CodeQL flagged js/missing-rate-limiting on the renderer /render route.
The sidecar is container-network-only with one trusted caller (the
orchestrator, which renders cuts serially), so this limiter is a
retry-storm ceiling (30/min, well above legit render rate), not the
primary control. Also reflows motion/README.md hard-wrapped prose that
failed the markdown quality gate.

* fix(build): finish pnpm 11 migration + regen verb tables

The panel Docker image build failed on `pnpm install --frozen-lockfile`:
node:22-alpine's corepack resolved to its bundled pnpm 11, but
panel/package.json pinned packageManager to pnpm@10.25.0, and pnpm 11
refuses to run against that pin. The Dockerfiles were already written for
pnpm 11 (comments, CI=true, strictDepBuilds); the package.json pin was the
stale outlier. Finish the migration instead of working around it:

- panel/package.json: packageManager pnpm@10.25.0 -> pnpm@11.10.0; drop the
  `pnpm` field (pnpm 11 ignores it — build approval lives in
  panel/pnpm-workspace.yaml's allowBuilds). Lockfile unchanged (pnpm 11
  accepts it as-is); frozen-lockfile verified.
- remotion-renderer/package.json: pin packageManager pnpm@11.10.0 for
  determinism (was relying on corepack's implicit default); engines.node
  >=22.13 (pnpm 11 requirement).
- docker/panel.Dockerfile + docker/remotion.Dockerfile: `corepack prepare
  pnpm@11.10.0 --activate` so the build uses the pinned version explicitly
  instead of trusting corepack's bundled default (which a future
  node:22-alpine could change).
- .github/workflows/panel-ci.yml: Node 20 -> 22 (pnpm 11 requires
  Node >=22.13; Node 20 fails the engines check).

Also regenerate agents/prompts/_generated/{developer,head_marketing,verbs}.md
— the video engine added propose_video and extended propose_feature_spotlight
(wants_video, video_script) but the verb tables weren't refreshed, failing
the foundation-check quality gate.

* chore(build): approve esbuild build script in remotion pnpm-workspace.yaml

pnpm 11 generated this file with a placeholder ('set this to true or false')
during install; resolve it to true so local dev of the renderer doesn't
re-prompt. esbuild's postinstall only verifies the prebuilt platform binary
(@esbuild/<platform> is installed as an optional dep), so approving it is
safe and silences the ERR_PNPM_IGNORED_BUILDS warning.

* fix(build): copy pnpm-workspace.yaml into panel + remotion images

pnpm 11 hard-errors with [ERR_PNPM_IGNORED_BUILDS] (exit 1) when a
dependency ships a postinstall script that isn't approved in
allowBuilds. Both Dockerfiles copied only package.json + pnpm-lock.yaml,
so the build-approval map in pnpm-workspace.yaml never made it into the
image — the remotion image build died on esbuild@0.28.1's postinstall.

Copy pnpm-workspace.yaml alongside the manifests in both images. In
panel, this also drops the --config.strictDepBuilds=false workaround:
with sharp and unrs-resolver now approved, their postinstalls run and
install the platform-specific binaries (previously skipped, leaving
sharp without its @img/sharp-* binary at runtime).

Verified locally: remotion + panel `pnpm install --frozen-lockfile`
exit 0 with the workspace file present; both exit 1 without it.

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2026-07-05 13:37:17 +02:00

223 lines
7.8 KiB
Python

"""HeartbeatMutex coverage: acquire (fencing token), heartbeat renew,
compare-and-del release, fail-closed on a Redis outage, and the
run_guarded cancel-on-lock-loss dance.
No live Redis in tests (matches the project's `_no_live_redis` fixture); a
tiny in-memory fake backs the Lua compare-and-del/compare-and-expire scripts
so the fencing semantics are observable without a real Redis.
"""
from __future__ import annotations
import asyncio
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from roboco.services.heartbeat_mutex import (
HeartbeatLockUnavailable,
HeartbeatMutex,
)
_KEY = "roboco:test_mutex:abc"
_MIN_RENEWS_AFTER_RECOVERY = 2 # the failed renew, then >=1 recovered one
class _FakeRedis:
"""In-memory single-key store backing the mutex's SET NX EX + two Lua
scripts (compare-and-del release, compare-and-expire heartbeat)."""
def __init__(self) -> None:
self._store: dict[str, str] = {}
self.set_calls: list[tuple[str, str, bool, int]] = []
self.eval_calls: list[tuple[str, tuple[Any, ...]]] = []
async def set(
self, name: str, value: str, *, nx: bool = False, ex: int = 0
) -> bool:
self.set_calls.append((name, value, nx, ex))
if nx and name in self._store:
return False
self._store[name] = value
return True
async def eval(self, script: str, _numkeys: int, *args: Any) -> int:
self.eval_calls.append((script, args))
key, token = args[0], args[1]
if "expire" in script:
return 1 if self._store.get(key) == token else 0
if self._store.get(key) == token:
del self._store[key]
return 1
return 0
async def aclose(self) -> None:
return None
def _mutex(*, ttl: int = 60, heartbeat: float = 30.0) -> HeartbeatMutex:
return HeartbeatMutex(_KEY, ttl_seconds=ttl, heartbeat_seconds=heartbeat)
@pytest.mark.asyncio
async def test_acquire_sets_nx_ex_and_returns_a_fencing_token() -> None:
fake = _FakeRedis()
with patch("roboco.services.heartbeat_mutex.redis.from_url", return_value=fake):
token = await _mutex(ttl=1800).acquire()
assert token is not None
assert fake.set_calls == [(_KEY, token, True, 1800)]
@pytest.mark.asyncio
async def test_acquire_returns_none_when_already_held() -> None:
fake = _FakeRedis()
with patch("roboco.services.heartbeat_mutex.redis.from_url", return_value=fake):
first = await _mutex().acquire()
second = await _mutex().acquire()
assert first is not None
assert second is None
@pytest.mark.asyncio
async def test_release_is_compare_and_del_spares_a_usurper_lock() -> None:
fake = _FakeRedis()
with patch("roboco.services.heartbeat_mutex.redis.from_url", return_value=fake):
mutex = _mutex()
token = await mutex.acquire()
assert token is not None
# A usurper re-acquired after this token's TTL expired.
fake._store[_KEY] = "usurper-token"
await mutex.release(token)
assert fake._store.get(_KEY) == "usurper-token" # survives a stale release
await mutex.release("usurper-token")
assert _KEY not in fake._store # the owning token does clear it
@pytest.mark.asyncio
async def test_heartbeat_once_true_when_owned_false_otherwise() -> None:
fake = _FakeRedis()
with patch("roboco.services.heartbeat_mutex.redis.from_url", return_value=fake):
mutex = _mutex()
token = await mutex.acquire()
assert token is not None
assert await mutex.heartbeat_once(token) is True
assert await mutex.heartbeat_once("wrong-token") is False
@pytest.mark.asyncio
async def test_acquire_raises_lock_unavailable_on_redis_error() -> None:
broken = MagicMock()
broken.set = AsyncMock(side_effect=ConnectionError("redis down"))
broken.aclose = AsyncMock()
with (
patch("roboco.services.heartbeat_mutex.redis.from_url", return_value=broken),
pytest.raises(HeartbeatLockUnavailable),
):
await _mutex().acquire()
@pytest.mark.asyncio
async def test_run_guarded_returns_the_coroutine_result_on_success() -> None:
async def _work() -> str:
return "done"
mutex = _mutex(heartbeat=0.001)
with patch.object(HeartbeatMutex, "heartbeat_once", AsyncMock(return_value=True)):
result = await mutex.run_guarded(_work(), "tok")
assert result.lock_lost is False
assert result.value == "done"
@pytest.mark.asyncio
async def test_run_guarded_renews_the_ttl_while_the_work_runs() -> None:
calls = 0
async def _counting_heartbeat(_self: HeartbeatMutex, _token: str) -> bool:
nonlocal calls
calls += 1
return True
async def _slow_work() -> str:
await asyncio.sleep(0.02)
return "done"
mutex = _mutex(heartbeat=0.001)
with patch.object(HeartbeatMutex, "heartbeat_once", _counting_heartbeat):
result = await mutex.run_guarded(_slow_work(), "tok")
assert result.value == "done"
assert calls >= 1 # at least one renew landed while the work was in flight
@pytest.mark.asyncio
async def test_run_guarded_cancels_the_work_fail_closed_on_lock_loss() -> None:
started = asyncio.Event()
async def _blocking_work() -> str:
started.set()
await asyncio.sleep(60)
return "never"
mutex = _mutex(heartbeat=0.001)
with patch.object(HeartbeatMutex, "heartbeat_once", AsyncMock(return_value=False)):
result = await mutex.run_guarded(_blocking_work(), "tok")
assert started.is_set() # the work did start, then got cancelled
assert result.lock_lost is True
assert result.value is None
@pytest.mark.asyncio
async def test_run_guarded_tolerates_a_transient_renew_raise() -> None:
"""One raised renew error, followed by recovery, must NOT trip
lock_lost — the TTL is still alive so it's tolerated as a blip."""
calls = 0
async def _flaky_heartbeat(_self: HeartbeatMutex, _token: str) -> bool:
nonlocal calls
calls += 1
if calls == 1:
raise ConnectionError("transient redis blip")
return True
async def _work() -> str:
await asyncio.sleep(0.02)
return "done"
# ttl=60 vs. a sub-second test run: the grace window is enormous, so a
# single raise is nowhere near "unable to renew for ~the whole TTL".
mutex = _mutex(ttl=60, heartbeat=0.005)
with patch.object(HeartbeatMutex, "heartbeat_once", _flaky_heartbeat):
result = await mutex.run_guarded(_work(), "tok")
assert result.lock_lost is False
assert result.value == "done"
assert calls >= _MIN_RENEWS_AFTER_RECOVERY
@pytest.mark.asyncio
async def test_run_guarded_fails_closed_once_renew_errors_span_the_whole_ttl() -> None:
"""A `heartbeat_once` that only ever RAISES (never returns falsy) must
still fail closed once elapsed time since the last successful renew
reaches ~the whole TTL — otherwise a holder stuck erroring on every
renew would never learn its key expired server-side.
`heartbeat_seconds > ttl_seconds` collapses the grace window
(`ttl_seconds - heartbeat_seconds`) to <= 0, so the very first raise's
elapsed time (always >= 0) already exceeds it — deterministic, no
reliance on real elapsed wall-clock time (and no monkeypatching
`time.monotonic`, which is also asyncio's own scheduling clock).
"""
started = asyncio.Event()
async def _blocking_work() -> str:
started.set()
await asyncio.sleep(60)
return "never"
mock_heartbeat = AsyncMock(side_effect=ConnectionError("redis down"))
mutex = _mutex(ttl=1, heartbeat=2.0)
with patch.object(HeartbeatMutex, "heartbeat_once", mock_heartbeat):
result = await mutex.run_guarded(_blocking_work(), "tok")
assert mock_heartbeat.call_count >= 1
assert started.is_set()
assert result.lock_lost is True
assert result.value is None