Files
roboco/tests/integration/test_video_routes.py
T
4923ee3ff3 MinIO video storage (chunk 1: config+deps+compose) + event-loop perf fix (#308)
* 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.

* fix(perf): offload conventions + release-readiness blocking I/O off the event loop

The orchestrator runs uvicorn and the orchestration background loops on a
single shared event loop, so any sync I/O anywhere — even inside a background
loop — blocks API responsiveness for its duration. Two call sites were missing
asyncio.to_thread wrappers:

- ConventionsService.get_map/health/restore called the sync _resolve
  (`git rev-parse`), _read_committed_standard (file read + yaml parse), and
  _derive (filesystem walk via derive_from_scan) inline. Reachable from
  GET /api/projects/{id}/conventions and from the agent spawn-prepare path.
- ReleaseManagerEngine._production_assess called gather_snapshot inline —
  multiple `subprocess.run` git calls + a filesystem walk, running inside the
  release-manager background loop.

Wrap each blocking call in asyncio.to_thread at the async boundary. No
signature changes; helpers stay sync. Verified: targeted tests pass
(196 passed, 36 DB-skipped), ruff + format clean.

These were the only responsiveness gaps surfaced by the concurrency audit —
the rest of the heavy paths (agent spawn via `docker run -d`, video render
loop, git ops via the 16-worker ThreadPoolExecutor, workspace subprocess
calls) already offload correctly. No API/worker container split needed.

* feat(storage): add MinIO config + dep + compose (no-op, default-off)

Chunk 1 of the MinIO video-storage plan (§1, §2, §6). No behavior change:
minio_endpoint defaults to empty = disabled, the existing FileResponse serve
path is untouched (chunk 4 wires the serve path; chunk 2 adds the client).

- pyproject.toml: add `minio` (minio-py) to dependencies; regenerate uv.lock
  (resolves minio v7.2.20 + pycryptodome transitive).
- roboco/config.py: add 5 settings fields after video_output_dir
  (minio_endpoint/_access_key/_secret_key/_bucket/_region). Plain str Fields
  matching the existing ROBOCO_ENCRYPTION_KEY style; no SecretStr, no
  presign_ttl_seconds (YAGNI — we don't presign in phase 1).
- docker-compose.yml: add `minio` service (data network only, named
  minio-data volume, host ports 19000/19001 for debugging, mc healthcheck)
  and a one-shot `minio-init` service mirroring the ollama-init pattern
  (mc alias set + mb -p, idempotent via || true). Add ROBOCO_MINIO_* env to
  the orchestrator env block (endpoint, access/secret key, bucket, region).
- docker-compose.registry.yml: intentionally omit the minio/minio-init
  services and leave ROBOCO_MINIO_* unset (NAS default-on, registry
  default-off — the established pattern); comment added to the orchestrator
  env block noting the omission.

* docs(storage): 0.19.0 CHANGELOG + RAG + map reference for MinIO chunk 1

Backfills the release-polish docs for MinIO chunk 1 (§10 of the plan):
- docker-compose.yaml synced to docker-compose.yml (the two NAS compose files
  must stay byte-identical; .yml was edited in chunk 1, .yaml was stale).
- CHANGELOG [0.19.0]: Added (MinIO scaffolding) + Fixed (event-loop I/O offload).
- docs/rag/architecture/minio-storage.md: RAG doc mirroring video-engine.md.
- docs/map/deployment-tooling.md: one-line storage reference.

* MinIO chunk 2: minio_client module (singleton + unconfigured guard) (#309)

* feat(storage): minio_client module (singleton + unconfigured guard)

Chunk 2 of the MinIO plan (§3). roboco/services/minio_client.py adds:
- get_client(): singleton minio-py Minio from settings; returns None when
  minio_endpoint is empty (the disabled path used by the chunk 3/4 guards).
  Parses http://... endpoint into host:port + secure flag.
- put_object(bytes, key): no-ops when unconfigured; otherwise PUTs to
  settings.minio_bucket with ContentType video/mp4.
- get_object_stream(key): yields object bytes for StreamingResponse; lets
  S3Error propagate so the serve route (chunk 4) can fall back to disk.

Sync calls — every call site wraps in asyncio.to_thread (chunks 3/4). One
unit test covers the unconfigured guard + endpoint scheme parsing (mocks,
no real MinIO). Not yet wired into remotion_client._save or the media route.

* MinIO chunk 3: wire write path (remotion_client._save PUT) (#310)

* feat(storage): wire MinIO write path in remotion_client._save

Chunk 3 of the MinIO plan (§3). After the local mp4 write, _save PUTs the bytes
to MinIO under key = Path(mp4_path).name (already {render_key}-{orientation}.mp4),
guarded by minio_client.get_client() (None when minio_endpoint empty) and
wrapped in asyncio.to_thread. Local disk stays the source of truth for the
poster publish path (x_video_client/tiktok_client read mp4_path from disk);
the PUT is additive. _save still returns the local path str — mp4_paths,
marker, and schema unchanged. Disabled (local-only) when MinIO unconfigured.

One test: asserts put_object is called with the basename key when configured
and the local file is still written; existing test stays green via the
unconfigured-default path. Mocks only.

* fix(storage): make MinIO PUT non-fatal in remotion_client._save

A configured-but-down MinIO made put_object raise inside the worker thread,
failing the render and retry-looping a task whose local file was already
written. Local disk is the source of truth and the serve route falls back to
FileResponse on S3Error, so a failed durable-copy PUT must never fail the
render — log and continue; the next render re-attempts the PUT.

Adds test_save_swallows_minio_put_failure (PUT raises -> _save still returns
the local path and the local file is written). Extends the CHANGELOG write-
path bullet with the non-fatal guarantee.

* MinIO chunk 4: serve path (StreamingResponse + FileResponse fallback) (#311)

* feat(storage): serve MinIO via the media route (StreamingResponse + FileResponse fallback)

Chunk 4 of the MinIO plan (§4 — the crux). GET /api/video/posts/{id}/media
derives key = Path(mp4_path).name and, when minio_endpoint is set, returns a
StreamingResponse over minio_client.get_object_stream(key), keeping
_require_ceo so auth stays end-to-end (no presigned URLs). Falls back to
FileResponse on S3Error (old render not in MinIO) or when MinIO is
unconfigured — the panel's axios-blob flow is unchanged (same URL, headers,
body, just chunked). The confinement check is kept as defense-in-depth (the
key is a basename so traversal is impossible, but the check is cheap and
protects the poster path).

Two integration tests: configured serve path streams from a stubbed
get_object_stream (CEO 200, non-CEO 403); unconfigured fallback serves the
local file via FileResponse. Mocks only — no real MinIO.

* fix(storage): eager stat_object probe so the MinIO serve fallback actually fires

The chunk-4 route wrapped StreamingResponse(get_object_stream(key), ...) in a
try/except, but get_object_stream is a lazy generator — its client.get_object
call runs on the first next(), i.e. AFTER the route returned and Starlette
started streaming. An S3Error (NoSuchKey / MinIO down) there is uncatchable;
the try/except caught nothing and the FileResponse fallback never triggered.

Add minio_client.stat_object(key): an eager existence/readiness probe that
runs INSIDE the route's try/except, so a missing object or down MinIO raises
before the StreamingResponse starts and the fallback serves the local file.
stat-then-get is two round trips; a mid-stream failure after a successful stat
is a rare race the CEO can retry (documented ceiling).

Tests: the configured test now stubs stat_object; a new test asserts the
S3Error fallback serves the local file via FileResponse and that
get_object_stream is never called. RAG doc updated to record the eager-probe
correctness detail + the non-fatal PUT.

* docs(rag): mark MinIO deployment note landed (chunk 5) (#312)

Co-authored-by: Renn F <rennf93@users.noreply.github.com>

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>

* fix(video): offload minio stat_object off the event loop

stat_object was called inline in the async media route, blocking the
shared event loop for one sync urllib3 round-trip per preview request —
contradicting minio_client's own 'every call site wraps in to_thread'
docstring and this PR's perf-fix theme. Wrap in asyncio.to_thread; the
try/except still catches S3Error (to_thread re-raises) so the
FileResponse fallback is unchanged. Also add the trailing newline to
the minio-storage RAG doc.

* Fix red CI

* Make CI green

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2026-07-05 16:12:44 +02:00

696 lines
27 KiB
Python

"""Video engine route coverage — the on-demand request trigger, the held
video_post draft list/approve/reject queue, and the TikTok credentials
sub-router. CEO-only throughout."""
from __future__ import annotations
from http import HTTPStatus
from types import SimpleNamespace
from typing import TYPE_CHECKING
from unittest.mock import AsyncMock, patch
from uuid import UUID, uuid4
import pytest
import pytest_asyncio
from fastapi import FastAPI
from httpx import ASGITransport, AsyncClient
from roboco.api.deps import get_agent_context, get_db
from roboco.api.routes import video as video_module
from roboco.api.routes.video import router as video_router
from roboco.api.routes.video import tiktok_router
from roboco.config import settings as cfg
from roboco.db.tables import AgentTable, ProjectTable, TaskTable
from roboco.foundation import identity as _foundation
from roboco.foundation.policy.content import markers
from roboco.models import AgentRole, AgentStatus, Team
from roboco.models.base import Complexity, TaskNature, TaskStatus, TaskType
from roboco.models.permissions import AgentContext
from roboco.services import minio_client
from roboco.services.heartbeat_mutex import HeartbeatMutex
from roboco.services.task import VIDEO_POST_SOURCE, VIDEO_SOURCE, get_task_service
from roboco.services.tiktok_credentials import get_tiktok_credentials_service
from roboco.services.video_post_service import XVideoPostResult
from roboco.services.x_credentials import get_x_credentials_service
from roboco.services.x_video_client import LiveXVideoPoster
from sqlalchemy import delete, select
if TYPE_CHECKING:
from collections.abc import AsyncIterator
from pathlib import Path
from sqlalchemy.ext.asyncio import AsyncSession
SLUG = "roboco-video-route-test"
SYSTEM_UUID = _foundation.AGENTS["system"].uuid
UX_DEV_1_UUID = _foundation.AGENTS["ux-dev-1"].uuid
UX_DEV_2_UUID = _foundation.AGENTS["ux-dev-2"].uuid
async def _seed(session: AsyncSession) -> None:
for uuid_, slug, role, team in (
(SYSTEM_UUID, "system", AgentRole.SYSTEM, None),
(UX_DEV_1_UUID, "ux-dev-1", AgentRole.DEVELOPER, Team.UX_UI),
(UX_DEV_2_UUID, "ux-dev-2", AgentRole.DEVELOPER, Team.UX_UI),
):
if await session.get(AgentTable, uuid_) is None:
session.add(
AgentTable(
id=uuid_,
name=slug,
slug=slug,
role=role,
team=team,
status=AgentStatus.ACTIVE,
model_config={},
system_prompt="x",
capabilities=[],
permissions={},
metrics={},
)
)
await session.flush()
existing = await session.execute(
select(ProjectTable).where(ProjectTable.slug == SLUG)
)
if existing.scalar_one_or_none() is None:
session.add(
ProjectTable(
name="RoboCo",
slug=SLUG,
git_url="https://github.com/x/roboco.git",
default_branch="master",
protected_branches=["master"],
assigned_cell=Team.BACKEND,
created_by=SYSTEM_UUID,
is_active=True,
)
)
await session.flush()
async def _seed_agent(session: AsyncSession, role: AgentRole, slug: str) -> AgentTable:
agent = AgentTable(
id=uuid4(),
name=slug,
slug=f"{slug}-{uuid4().hex[:6]}",
role=role,
team=None,
status=AgentStatus.ACTIVE,
model_config={},
system_prompt="x",
capabilities=[],
permissions={},
metrics={},
)
session.add(agent)
await session.flush()
return agent
async def _seed_draft(
session: AsyncSession,
*,
platforms: list[str] | None = None,
mp4_paths: dict[str, str] | None = None,
) -> TaskTable:
"""A held ``video_post`` draft — the approve/reject/list queue basis."""
system = await _seed_agent(session, AgentRole.SYSTEM, "system")
secretary = await _seed_agent(session, AgentRole.SECRETARY, "secretary")
project = ProjectTable(
id=uuid4(),
name="RoboCo",
slug=f"roboco-{uuid4().hex[:6]}",
git_url="https://example.com/roboco.git",
assigned_cell=Team.BACKEND,
created_by=system.id,
)
session.add(project)
await session.flush()
task = TaskTable(
id=uuid4(),
title="Video post: release 1.0",
description="script",
acceptance_criteria=["CEO approves or rejects the draft"],
status=TaskStatus.PENDING,
priority=2,
task_type=TaskType.ADMINISTRATIVE,
nature=TaskNature.NON_TECHNICAL,
estimated_complexity=Complexity.LOW,
project_id=project.id,
created_by=system.id,
assigned_to=secretary.id,
team=Team.MAIN_PM,
source=VIDEO_POST_SOURCE,
confirmed_by_human=False,
)
session.add(task)
await session.flush()
markers.set_video_draft(
task,
{
"occasion": "release 1.0",
"script": "script",
"platforms": platforms if platforms is not None else ["x"],
"mp4_paths": mp4_paths
if mp4_paths is not None
else {
"square": "/render/out/1-square.mp4",
"vertical": "/render/out/1-vertical.mp4",
},
"x_caption": "Check out this clip",
"tiktok_caption": "Check out this clip on TikTok",
"render_status": "rendered",
},
)
await session.flush()
return task
def _build_app(db_session: AsyncSession, role: AgentRole, agent_id: UUID) -> FastAPI:
app = FastAPI()
app.include_router(video_router, prefix="/api/video")
app.include_router(tiktok_router, prefix="/api/tiktok")
async def _override_db() -> AsyncIterator[AsyncSession]:
yield db_session
async def _override_agent() -> AgentContext:
return AgentContext(agent_id=agent_id, role=role, team=None)
app.dependency_overrides[get_db] = _override_db
app.dependency_overrides[get_agent_context] = _override_agent
return app
@pytest_asyncio.fixture
async def ceo_client(db_session: AsyncSession) -> AsyncIterator[AsyncClient]:
app = _build_app(db_session, AgentRole.CEO, uuid4())
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client
app.dependency_overrides.clear()
_LOCKED = (
patch.object(HeartbeatMutex, "acquire", AsyncMock(return_value="tok")),
patch.object(HeartbeatMutex, "release", AsyncMock(return_value=None)),
)
@pytest.mark.asyncio
async def test_request_video_opens_authoring_task(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True)
monkeypatch.setattr(cfg, "self_heal_project_slug", SLUG)
monkeypatch.setattr(cfg, "video_max_open_posts", 5)
resp = await ceo_client.post(
"/api/video/request",
json={
"occasion": "CEO on-demand: launch teaser",
"brief": "A short teaser for the new dashboard",
"platforms": ["x", "tiktok"],
},
)
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["status"] == "opened"
assert body["task_id"] is not None
try:
# The route commits (mirrors the X route), so identity/field checks on
# the specific created row — not a global open-list count — keep this
# robust against other committed rows in the shared session-scoped
# test DB.
task = await db_session.get(TaskTable, UUID(body["task_id"]))
assert task is not None
assert task.source == VIDEO_SOURCE
assert task.status == TaskStatus.PENDING
finally:
# The route's commit durably persists this task past this test's own
# rollback teardown — a non-terminal source=video row left behind
# pollutes every later test in this session that counts open video
# tasks (test_video_engine.py / test_video_render_loop.py), so it
# must be deleted explicitly, not just rolled back.
await db_session.execute(
delete(TaskTable).where(TaskTable.id == UUID(body["task_id"]))
)
await db_session.commit()
@pytest.mark.asyncio
async def test_request_video_disabled_returns_clear_response(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", False)
before = len(await get_task_service(db_session).list_open_video_posts())
resp = await ceo_client.post(
"/api/video/request",
json={"occasion": "occ-disabled", "brief": "brief", "platforms": ["x"]},
)
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["status"] == "disabled"
assert body["task_id"] is None
after = len(await get_task_service(db_session).list_open_video_posts())
assert after == before # nothing new was opened
@pytest.mark.asyncio
async def test_request_video_not_opened_when_project_unresolvable(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None:
"""An unresolvable project makes open_video_task no-op — a clear
``not_opened`` response, not a 500 or a fabricated task."""
await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True)
monkeypatch.setattr(cfg, "self_heal_project_slug", "no-such-project")
before = len(await get_task_service(db_session).list_open_video_posts())
resp = await ceo_client.post(
"/api/video/request",
json={"occasion": "occ-unresolvable", "brief": "brief", "platforms": ["x"]},
)
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["status"] == "not_opened"
assert body["task_id"] is None
after = len(await get_task_service(db_session).list_open_video_posts())
assert after == before # nothing new was opened
@pytest.mark.asyncio
async def test_list_posts_returns_open_draft(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_draft(db_session)
resp = await ceo_client.get("/api/video/posts")
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert len(body) == 1
assert body[0]["task_id"] == str(task.id)
assert body[0]["occasion"] == "release 1.0"
assert body[0]["platforms"] == ["x"]
assert body[0]["mp4_paths"] == {
"square": "/render/out/1-square.mp4",
"vertical": "/render/out/1-vertical.mp4",
}
@pytest.mark.asyncio
async def test_media_returns_the_rendered_cut(
db_session: AsyncSession,
ceo_client: AsyncClient,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(cfg, "video_output_dir", str(tmp_path))
vertical = tmp_path / "clip-vertical.mp4"
vertical.write_bytes(b"fake-mp4-bytes-vertical")
task = await _seed_draft(
db_session,
mp4_paths={"vertical": str(vertical), "square": str(tmp_path / "missing.mp4")},
)
resp = await ceo_client.get(f"/api/video/posts/{task.id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.OK
assert resp.headers["content-type"] == "video/mp4"
assert resp.content == b"fake-mp4-bytes-vertical"
@pytest.mark.asyncio
async def test_media_outside_output_dir_is_404(
db_session: AsyncSession,
ceo_client: AsyncClient,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""A mp4_paths entry that resolves outside video_output_dir is refused
even though the file exists on disk — defense-in-depth against any
future writer of mp4_paths."""
outside = tmp_path / "outside" / "clip-vertical.mp4"
outside.parent.mkdir(parents=True)
outside.write_bytes(b"fake-mp4-bytes")
monkeypatch.setattr(cfg, "video_output_dir", str(tmp_path / "confined"))
task = await _seed_draft(db_session, mp4_paths={"vertical": str(outside)})
resp = await ceo_client.get(f"/api/video/posts/{task.id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_media_bad_cut_is_400(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_draft(db_session)
resp = await ceo_client.get(f"/api/video/posts/{task.id}/media?cut=diagonal")
assert resp.status_code == HTTPStatus.BAD_REQUEST
@pytest.mark.asyncio
async def test_media_missing_task_is_404(ceo_client: AsyncClient) -> None:
resp = await ceo_client.get(f"/api/video/posts/{uuid4()}/media?cut=vertical")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_media_unrendered_cut_is_404(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
"""The seeded draft's paths never exist on disk — a 404, not a crash."""
task = await _seed_draft(db_session)
resp = await ceo_client.get(f"/api/video/posts/{task.id}/media?cut=square")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_approve_without_credentials_fails_gracefully(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
"""No X/TikTok credentials configured in this test DB: the route still
builds real (Null) posters and the approve completes without raising —
just with nothing posted."""
task = await _seed_draft(db_session, platforms=["x", "tiktok"])
try:
with _LOCKED[0], _LOCKED[1]:
resp = await ceo_client.post(f"/api/video/posts/{task.id}/approve", json={})
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["status"] == "post_failed"
assert body["posted"] == {}
await db_session.refresh(task)
assert task.status == TaskStatus.PENDING # never advanced without a real post
finally:
# The approve route commits durably even on a post_failed outcome, so
# this non-terminal source=video_post row survives this test's own
# rollback teardown — left behind, it pollutes every later test in
# this session that counts open video tasks.
await db_session.execute(delete(TaskTable).where(TaskTable.id == task.id))
await db_session.commit()
@pytest.mark.asyncio
async def test_approve_with_credentials_posts_via_the_real_poster_wiring(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
"""Once X credentials are configured, the route builds a LiveXVideoPoster
(not the Null default) — the network call itself is mocked here; the
real HTTP sequence is covered by test_x_video_client.py.
The approve route commits durably (mirrors the X-post pattern), so the
x_credentials singleton row must be cleared afterward — left behind, it
leaks into any later test in this shared session-scoped test DB that
asserts a fresh "unset" state (e.g. test_x_credentials_service.py)."""
task = await _seed_draft(db_session, platforms=["x"])
creds_svc = get_x_credentials_service(db_session)
await creds_svc.set_credentials(
api_key="ak", api_secret="as", access_token="at", access_token_secret="ats"
)
try:
with (
_LOCKED[0],
_LOCKED[1],
patch.object(
LiveXVideoPoster,
"post_video",
AsyncMock(
return_value=XVideoPostResult(
posted=True, video_id="xid1", detail="posted"
)
),
),
):
resp = await ceo_client.post(f"/api/video/posts/{task.id}/approve", json={})
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["status"] == "posted"
assert body["posted"] == {"x": "xid1"}
await db_session.refresh(task)
assert task.status == TaskStatus.COMPLETED
finally:
await creds_svc.set_credentials(
api_key="", api_secret="", access_token="", access_token_secret=""
)
await db_session.commit()
@pytest.mark.asyncio
async def test_approve_edited_x_caption_over_limit_is_422(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_draft(db_session)
resp = await ceo_client.post(
f"/api/video/posts/{task.id}/approve", json={"x_caption": "x" * 281}
)
assert resp.status_code == HTTPStatus.UNPROCESSABLE_ENTITY
@pytest.mark.asyncio
async def test_approve_missing_task_is_404(ceo_client: AsyncClient) -> None:
resp = await ceo_client.post(f"/api/video/posts/{uuid4()}/approve", json={})
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_reject_cancels_and_records_reason(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_draft(db_session)
resp = await ceo_client.post(
f"/api/video/posts/{task.id}/reject", json={"reason": "Not our voice"}
)
assert resp.status_code == HTTPStatus.OK
assert resp.json()["reject_reason"] == "Not our voice"
refreshed = await db_session.get(TaskTable, task.id)
assert refreshed is not None
assert refreshed.status == TaskStatus.CANCELLED
@pytest.mark.asyncio
async def test_reject_missing_task_is_404(ceo_client: AsyncClient) -> None:
resp = await ceo_client.post(
f"/api/video/posts/{uuid4()}/reject", json={"reason": "not relevant here"}
)
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_tiktok_credentials_default_is_unset(ceo_client: AsyncClient) -> None:
resp = await ceo_client.get("/api/tiktok/credentials")
assert resp.status_code == HTTPStatus.OK
assert resp.json()["has_credentials"] is False
@pytest.mark.asyncio
async def test_set_tiktok_credentials_reports_status_never_plaintext(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
"""The route commits durably, so the tiktok_credentials singleton row is
cleared afterward — left behind, it leaks into any later test in this
shared session-scoped test DB (e.g. test_tiktok_credentials_service.py's
"unset" assertions)."""
try:
resp = await ceo_client.post(
"/api/tiktok/credentials",
json={
"client_key": "secret-key-value",
"client_secret": "secret-clientsecret-value",
"access_token": "secret-token-value",
"refresh_token": "secret-refresh-value",
},
)
assert resp.status_code == HTTPStatus.OK
assert resp.json() == {"has_credentials": True}
assert "secret-key-value" not in resp.text
assert "secret-clientsecret-value" not in resp.text
assert "secret-token-value" not in resp.text
assert "secret-refresh-value" not in resp.text
status_resp = await ceo_client.get("/api/tiktok/credentials")
assert status_resp.json()["has_credentials"] is True
finally:
await get_tiktok_credentials_service(db_session).set_credentials(
client_key="", client_secret="", access_token="", refresh_token=""
)
await db_session.commit()
@pytest.mark.asyncio
async def test_set_tiktok_credentials_partial_is_400(ceo_client: AsyncClient) -> None:
resp = await ceo_client.post(
"/api/tiktok/credentials",
json={
"client_key": "only-one",
"client_secret": "",
"access_token": "",
"refresh_token": "",
},
)
assert resp.status_code == HTTPStatus.BAD_REQUEST
@pytest.mark.asyncio
async def test_non_ceo_is_forbidden(db_session: AsyncSession) -> None:
await _seed(db_session)
task = await _seed_draft(db_session)
app = _build_app(db_session, AgentRole.DEVELOPER, uuid4())
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
request_resp = await client.post(
"/api/video/request",
json={"occasion": "occ", "brief": "brief", "platforms": ["x"]},
)
list_resp = await client.get("/api/video/posts")
media_resp = await client.get(f"/api/video/posts/{task.id}/media?cut=vertical")
creds_resp = await client.get("/api/tiktok/credentials")
assert request_resp.status_code == HTTPStatus.FORBIDDEN
assert list_resp.status_code == HTTPStatus.FORBIDDEN
assert media_resp.status_code == HTTPStatus.FORBIDDEN
assert creds_resp.status_code == HTTPStatus.FORBIDDEN
app.dependency_overrides.clear()
# --- MinIO serve path (chunk 4) — unit-style, no DB / no real MinIO ------------
# These two tests monkeypatch ``get_task_service`` in the video routes module
# so they run without postgres (the ``db_session``-based tests above are
# skipped when Postgres is unreachable). Mocks only — no testcontainers.
def _stub_task_service_factory(task: object) -> object:
"""A ``get_task_service``-shaped stub (the real one is a sync factory
returning a service with an async ``.get``). Patched in place of
``video_module.get_task_service`` so the route runs without postgres."""
class _Svc:
async def get(self, _task_id: UUID) -> object:
return task
return _Svc()
def _make_task(mp4_path: str, task_id: UUID) -> SimpleNamespace:
"""A minimal task-shaped stub carrying the video_draft marker the route
reads — enough for the media route, no DB row needed."""
return SimpleNamespace(
id=task_id,
source=VIDEO_POST_SOURCE,
orchestration_markers={"video_draft": {"mp4_paths": {"vertical": mp4_path}}},
)
def _patch_task_service(monkeypatch: pytest.MonkeyPatch, task: object) -> None:
monkeypatch.setattr(
video_module,
"get_task_service",
lambda _db: _stub_task_service_factory(task),
)
def _minio_stream(_key: str) -> object:
"""Stub ``get_object_stream`` yielding fixed bytes for ``StreamingResponse``."""
return iter([b"minio-stream-bytes"])
@pytest.mark.asyncio
async def test_media_serves_from_minio_when_configured(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""Configured serve path: when MinIO is configured, the media route streams
the object via ``minio_client.get_object_stream`` (key = basename) and the
panel-preview URL/headers stay identical. ``_require_ceo`` still 403s a
non-CEO agent. No DB / no real MinIO — ``get_task_service`` is stubbed so
the route runs without postgres."""
# A real local file so the route's is_file() + confinement checks pass.
# The served bytes come from the stubbed MinIO stream below, NOT this
# file — that's what proves the MinIO path was taken rather than the
# FileResponse fallback.
monkeypatch.setattr(cfg, "video_output_dir", str(tmp_path))
vertical = tmp_path / "clip-vertical.mp4"
vertical.write_bytes(b"local-file-bytes")
task_id = uuid4()
_patch_task_service(monkeypatch, _make_task(str(vertical), task_id))
# non-None sentinel so the route takes the MinIO branch.
monkeypatch.setattr(minio_client, "get_client", lambda: True)
# The route probes stat_object eagerly before streaming; stub it to pass.
monkeypatch.setattr(minio_client, "stat_object", lambda _key: None)
monkeypatch.setattr(minio_client, "get_object_stream", _minio_stream)
# CEO 200 — streamed from MinIO.
app = _build_app(None, AgentRole.CEO, uuid4()) # type: ignore[arg-type]
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get(f"/api/video/posts/{task_id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.OK
assert resp.headers["content-type"] == "video/mp4"
assert resp.content == b"minio-stream-bytes"
app.dependency_overrides.clear()
# Non-CEO 403 — _require_ceo still gates end-to-end (no presigned URL).
app = _build_app(None, AgentRole.DEVELOPER, uuid4()) # type: ignore[arg-type]
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get(f"/api/video/posts/{task_id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.FORBIDDEN
app.dependency_overrides.clear()
@pytest.mark.asyncio
async def test_media_falls_back_to_local_file_when_minio_unconfigured(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""Unconfigured fallback: with ``get_client`` returning None
(``minio_endpoint`` empty), the media route serves the local file via
``FileResponse`` — the body equals the local file's bytes. No DB / no
real MinIO."""
monkeypatch.setattr(cfg, "video_output_dir", str(tmp_path))
vertical = tmp_path / "clip-vertical.mp4"
vertical.write_bytes(b"local-file-bytes")
task_id = uuid4()
_patch_task_service(monkeypatch, _make_task(str(vertical), task_id))
monkeypatch.setattr(minio_client, "get_client", lambda: None)
app = _build_app(None, AgentRole.CEO, uuid4()) # type: ignore[arg-type]
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get(f"/api/video/posts/{task_id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.OK
assert resp.headers["content-type"] == "video/mp4"
assert resp.content == b"local-file-bytes"
app.dependency_overrides.clear()
@pytest.mark.asyncio
async def test_media_falls_back_to_local_file_when_minio_missing(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""S3Error fallback: when MinIO is configured but the object is missing
(NoSuchKey — an old render not yet in MinIO) or MinIO is down, the route's
eager ``stat_object`` probe raises, the ``try/except`` catches it, and the
route serves the local file via ``FileResponse``. ``get_object_stream`` is
never called. No DB / no real MinIO."""
monkeypatch.setattr(cfg, "video_output_dir", str(tmp_path))
vertical = tmp_path / "clip-vertical.mp4"
vertical.write_bytes(b"local-file-bytes")
task_id = uuid4()
_patch_task_service(monkeypatch, _make_task(str(vertical), task_id))
monkeypatch.setattr(minio_client, "get_client", lambda: True)
def _stat_raises(_key: str) -> None:
raise RuntimeError("minio NoSuchKey / down")
monkeypatch.setattr(minio_client, "stat_object", _stat_raises)
# If the route wrongly takes the MinIO stream branch, this would be called
# and the assertion below would fail — guard against a regression.
def _stream_must_not_be_called(_key: str) -> object:
pytest.fail("get_object_stream must not be called when stat_object raises")
monkeypatch.setattr(minio_client, "get_object_stream", _stream_must_not_be_called)
app = _build_app(None, AgentRole.CEO, uuid4()) # type: ignore[arg-type]
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get(f"/api/video/posts/{task_id}/media?cut=vertical")
assert resp.status_code == HTTPStatus.OK
assert resp.headers["content-type"] == "video/mp4"
assert resp.content == b"local-file-bytes" # FileResponse fallback, not MinIO
app.dependency_overrides.clear()