[7f2c881a] feat(video): scope on-demand video requests + render loop to project_id

Require project_id on VideoRequestBody (404 when unresolvable or not
opted into the video engine), thread it through VideoEngine.open_video_task
via a shared resolve_authoring_project helper, and resolve the render
loop's motion/ workspace from the authoring task's own project_id instead
of the hardcoded self_heal_project_slug.
This commit is contained in:
Backend Developer 1
2026-07-10 04:42:10 +00:00
parent f601e32788
commit 88ab5c6dc5
8 changed files with 752 additions and 47 deletions
+5
View File
@@ -154,6 +154,11 @@ select = [
# agent_id, project_ids, route, session_id) — same >5-kwarg rationale as the # agent_id, project_ids, route, session_id) — same >5-kwarg rationale as the
# gateway verb surfaces below. # gateway verb surfaces below.
"roboco/services/prompter.py" = ["PLR0913"] "roboco/services/prompter.py" = ["PLR0913"]
# open_video_task's kwargs (occasion, script, platforms, brief,
# suggested_input_props, project_id) are the authoring-task contract shared
# by the release/spotlight/on-demand callers — same "bundling would just
# relocate the same named fields behind one hop" rationale as prompter.py.
"roboco/services/video_engine.py" = ["PLR0913"]
# Gateway methods are typed verb surfaces — agent-facing kwargs reflect the # Gateway methods are typed verb surfaces — agent-facing kwargs reflect the
# verb contract (task title, description, acceptance criteria, assignee, etc.). Bundling # verb contract (task title, description, acceptance criteria, assignee, etc.). Bundling
# into a dataclass hides the field-by-field schema the LLM needs at the # into a dataclass hides the field-by-field schema the LLM needs at the
+102 -9
View File
@@ -7,7 +7,7 @@ from __future__ import annotations
import asyncio import asyncio
from pathlib import Path from pathlib import Path
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any, cast
from uuid import UUID from uuid import UUID
from fastapi import APIRouter, HTTPException, Query, status from fastapi import APIRouter, HTTPException, Query, status
@@ -30,7 +30,8 @@ from roboco.config import settings
from roboco.foundation.policy.content import markers from roboco.foundation.policy.content import markers
from roboco.security import guard_deco from roboco.security import guard_deco
from roboco.services import minio_client from roboco.services import minio_client
from roboco.services.task import VIDEO_POST_SOURCE, get_task_service from roboco.services.project import get_project_service
from roboco.services.task import VIDEO_POST_SOURCE, VIDEO_SOURCE, get_task_service
from roboco.services.tiktok_client import build_tiktok_poster from roboco.services.tiktok_client import build_tiktok_poster
from roboco.services.tiktok_credentials import ( from roboco.services.tiktok_credentials import (
TikTokCredentialsValidationError, TikTokCredentialsValidationError,
@@ -41,6 +42,7 @@ from roboco.services.video_post_service import (
VideoCaptionTooLongError, VideoCaptionTooLongError,
get_video_post_service, get_video_post_service,
) )
from roboco.services.workspace import WorkspaceError, get_workspace_service
from roboco.services.x_credentials import get_x_credentials_service from roboco.services.x_credentials import get_x_credentials_service
from roboco.services.x_video_client import build_x_video_poster from roboco.services.x_video_client import build_x_video_poster
@@ -89,11 +91,15 @@ async def request_video(
db: DbSession, db: DbSession,
agent: CurrentAgentContext, agent: CurrentAgentContext,
) -> VideoRequestResponse: ) -> VideoRequestResponse:
"""Open a UX/UI video-authoring task for the CEO's on-demand brief. """Open a UX/UI video-authoring task for the CEO's on-demand brief,
scoped to ``data.project_id``.
Returns ``disabled`` when the video engine is off, and ``not_opened`` Returns ``disabled`` when the video engine is off. 404s when
when ``open_video_task`` no-ops (a duplicate occasion, the open cap, or ``project_id`` doesn't resolve to a project, or that project hasn't
an unresolvable project) — neither is an error, just nothing to do. opted into the video engine (``video_engine_enabled`` is off on it).
Returns ``not_opened`` when ``open_video_task`` no-ops for any other
reason (a duplicate occasion or the open-post cap) — not an error, just
nothing to do.
""" """
_require_ceo(agent) _require_ceo(agent)
if not settings.video_engine_enabled: if not settings.video_engine_enabled:
@@ -101,18 +107,27 @@ async def request_video(
status="disabled", status="disabled",
detail="The video engine is disabled (video_engine_enabled is off).", detail="The video engine is disabled (video_engine_enabled is off).",
) )
task = await get_video_engine(db).open_video_task( engine = get_video_engine(db)
project = await engine.resolve_authoring_project(
project_id=data.project_id, occasion=data.occasion
)
if project is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="project not found or not opted into the video engine",
)
task = await engine.open_video_task(
occasion=data.occasion, occasion=data.occasion,
script=data.brief, script=data.brief,
platforms=data.platforms, platforms=data.platforms,
brief=data.brief, brief=data.brief,
project_id=data.project_id,
) )
if task is None: if task is None:
return VideoRequestResponse( return VideoRequestResponse(
status="not_opened", status="not_opened",
detail=( detail=(
"No video task was opened (a duplicate occasion, the open-post" "No video task was opened (a duplicate occasion or the open-post cap)."
" cap, or the project isn't resolvable)."
), ),
) )
await db.commit() await db.commit()
@@ -196,6 +211,84 @@ async def list_video_pipeline(
return [_to_pipeline_item(t) for t in tasks] return [_to_pipeline_item(t) for t in tasks]
@router.post("/pipeline/{task_id}/rerender", response_model=VideoPipelineItemResponse)
@guard_deco.rate_limit(requests=20, window=60)
@guard_deco.block_clouds()
async def rerender_video_task(
task_id: UUID, db: DbSession, agent: CurrentAgentContext
) -> VideoPipelineItemResponse:
"""Clear a completed video-authoring task's render idempotency keys
(``render_status``/``render_attempts``/``render_error``) so the next
render cycle re-picks it up and re-renders it. 404s when there's no such
completed authoring task with a proposed composition."""
_require_ceo(agent)
task = await get_video_engine(db).rerender(task_id)
if task is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="No such completed video task with a proposed composition",
)
await db.commit()
return _to_pipeline_item(task)
def _resolve_preview_path(root: Path, file_path: str) -> Path | None:
"""Resolve ``file_path`` against the workspace ``root``, refusing
anything that escapes it. A leading ``/`` is stripped before joining —
pathlib's ``/`` operator otherwise lets an absolute right operand
discard ``root`` entirely — then the joined path must resolve to an
existing file still under ``root``. The sole confinement check for the
CEO preview proxy."""
candidate = (root / file_path.lstrip("/")).resolve()
if not candidate.is_relative_to(root) or not candidate.is_file():
return None
return candidate
@router.get("/preview/{task_id}/{file_path:path}", response_model=None)
async def get_video_preview(
task_id: UUID,
file_path: str,
db: DbSession,
agent: CurrentAgentContext,
) -> FileResponse:
"""Serve a video-authoring task's composition HTML + sibling assets
(kit/public/etc.) straight off its project's merged read-clone — the
panel's live preview iframe. ``file_path`` is relative to the resolved
workspace root (e.g. ``motion/compositions/<id>/vertical.html``);
confined there so it can't traverse out, per ``_resolve_preview_path``.
CEO-only; the response carries explicit iframe-permitting headers so the
panel can embed it.
"""
_require_ceo(agent)
task = await get_task_service(db).get(task_id)
if task is None or task.source != VIDEO_SOURCE or task.project_id is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="No such video task"
)
project = await get_project_service(db).get(cast("UUID", task.project_id))
if project is None or not project.slug:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="Project not found"
)
try:
workspace = await get_workspace_service(db).ensure_read_clone(project.slug)
except WorkspaceError as e:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) from e
resolved = _resolve_preview_path(workspace.resolve(), file_path)
if resolved is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="No such preview file"
)
return FileResponse(
resolved,
headers={
"X-Frame-Options": "SAMEORIGIN",
"Content-Security-Policy": "frame-ancestors 'self'",
},
)
def _posted_ids(draft: dict[str, Any]) -> dict[str, str]: def _posted_ids(draft: dict[str, Any]) -> dict[str, str]:
"""Every ``{platform}_posted_id`` key stamped by approve, keyed by """Every ``{platform}_posted_id`` key stamped by approve, keyed by
platform (e.g. ``{"x": "..", "tiktok": ".."}``).""" platform (e.g. ``{"x": "..", "tiktok": ".."}``)."""
+3 -1
View File
@@ -1,6 +1,7 @@
"""Schemas for the video engine's on-demand request + CEO approval surface.""" """Schemas for the video engine's on-demand request + CEO approval surface."""
from datetime import datetime from datetime import datetime
from uuid import UUID
from pydantic import BaseModel, Field from pydantic import BaseModel, Field
@@ -10,11 +11,12 @@ from roboco.services.x_client import MAX_TWEET_CHARS
class VideoRequestBody(BaseModel): class VideoRequestBody(BaseModel):
"""The CEO's on-demand video brief.""" """The CEO's on-demand video brief, scoped to a specific project."""
occasion: str = Field(..., min_length=1) occasion: str = Field(..., min_length=1)
brief: str = Field(..., min_length=1) brief: str = Field(..., min_length=1)
platforms: list[str] = Field(..., min_length=1) platforms: list[str] = Field(..., min_length=1)
project_id: UUID
class VideoRequestResponse(BaseModel): class VideoRequestResponse(BaseModel):
+21 -8
View File
@@ -8106,7 +8106,7 @@ Start by:
return return
try: try:
mp4_paths = await self._render_both_cuts( mp4_paths = await self._render_both_cuts(
db, draft, composition_id, str(task.id) db, draft, composition_id, str(task.id), task.project_id
) )
await self._materialize_video_post(db, task, draft, mp4_paths) await self._materialize_video_post(db, task, draft, mp4_paths)
except Exception as exc: except Exception as exc:
@@ -8151,16 +8151,29 @@ Start by:
) )
async def _render_both_cuts( async def _render_both_cuts(
self, db: Any, draft: dict[str, Any], composition_id: str, render_key: str self,
db: Any,
draft: dict[str, Any],
composition_id: str,
render_key: str,
project_id: Any,
) -> dict[str, str]: ) -> dict[str, str]:
"""Render the vertical + square cuts from the roboco project's merged """Render the vertical + square cuts from the authoring task's OWN
read-clone's motion/ dir; returns {"vertical": path, "square": path}. project's merged read-clone's motion/ dir; returns {"vertical": path,
``render_key`` (the source task id) scopes each cut's output path.""" "square": path}. ``render_key`` (the source task id) scopes each
cut's output path. ``project_id`` is the task's own ``project_id``
never a fixed slug so a video task authored against any opted-in
project renders from that project's ``motion/`` dir, not RoboCo's."""
from roboco.services.project import get_project_service
from roboco.services.video_renderer_client import get_video_renderer from roboco.services.video_renderer_client import get_video_renderer
from roboco.services.workspace import get_workspace_service from roboco.services.workspace import WorkspaceError, get_workspace_service
slug = (settings.self_heal_project_slug or "roboco-api").strip() project = await get_project_service(db).get(project_id) if project_id else None
workspace = await get_workspace_service(db).ensure_read_clone(slug) if project is None or not project.slug:
raise WorkspaceError(
f"video-render: task's project not resolvable ({project_id})"
)
workspace = await get_workspace_service(db).ensure_read_clone(project.slug)
motion_dir = str(workspace / "motion") motion_dir = str(workspace / "motion")
input_props = draft.get("input_props") or {} input_props = draft.get("input_props") or {}
renderer = get_video_renderer() renderer = get_video_renderer()
+73 -17
View File
@@ -149,23 +149,38 @@ class VideoEngine(BaseService):
service_name = "video_engine" service_name = "video_engine"
async def _roboco_project(self) -> ProjectTable | None: async def _roboco_project(self) -> ProjectTable | None:
"""The fixed RoboCo project the release/spotlight hooks author
against (``self_heal_project_slug``). The on-demand caller instead
supplies its own ``project_id`` — see ``resolve_authoring_project``."""
slug = (settings.self_heal_project_slug or "roboco-api").strip() slug = (settings.self_heal_project_slug or "roboco-api").strip()
return await get_project_service(self.session).get_by_slug(slug) return await get_project_service(self.session).get_by_slug(slug)
async def _opted_in_project(self, occasion: str) -> ProjectTable | None: async def resolve_authoring_project(
"""The RoboCo project if resolvable AND opted into video, else None. self, *, project_id: UUID | None, occasion: str
) -> ProjectTable | None:
"""The project to author this video against, or None when
unresolvable or not opted into the video engine.
Two skip reasons, both logged: unresolvable project (warning — a ``project_id`` (the on-demand ``/video/request`` caller, and the
config gap) vs. project not opted in (info — the operator hasn't render loop's own per-task resolution) resolves by id; omitted (the
flipped the per-project ``video_engine_enabled`` toggle). The global release/spotlight hooks), it falls back to the fixed RoboCo project.
flag arms the subsystem; the project's flag opts this repo into Both paths share the same two skip reasons, both logged: unresolvable
authoring against its ``motion/`` dir (mirrors ``ci_watch_enabled``). project (warning — a config/data gap) vs. project not opted in (info
— the operator hasn't flipped the per-project ``video_engine_enabled``
toggle). The global flag arms the subsystem; the project's flag opts
that repo into authoring against its own ``motion/`` dir (mirrors
``ci_watch_enabled``).
""" """
project = await self._roboco_project() project = (
await get_project_service(self.session).get(project_id)
if project_id is not None
else await self._roboco_project()
)
if project is None or project.id is None: if project is None or project.id is None:
self.log.warning( self.log.warning(
"video-engine: RoboCo project not resolvable; skipping video task", "video-engine: project not resolvable; skipping video task",
occasion=occasion, occasion=occasion,
project_id=str(project_id) if project_id else None,
) )
return None return None
if not getattr(project, "video_engine_enabled", False): if not getattr(project, "video_engine_enabled", False):
@@ -231,16 +246,22 @@ class VideoEngine(BaseService):
platforms: list[str], platforms: list[str],
brief: str, brief: str,
suggested_input_props: dict[str, Any] | None = None, suggested_input_props: dict[str, Any] | None = None,
project_id: UUID | None = None,
) -> TaskTable | None: ) -> TaskTable | None:
"""Originate ONE UX/UI authoring task for a bespoke video, or None. """Originate ONE UX/UI authoring task for a bespoke video, or None.
No-ops when the global flag is off, the RoboCo project hasn't opted in No-ops when the global flag is off, the target project hasn't opted
(``video_engine_enabled``), a task for this occasion is already open in (``video_engine_enabled``), a task for this occasion is already
(authoring or held draft), the open cap is reached, or the RoboCo open (authoring or held draft), the open cap is reached, or the
project isn't resolvable. The opened task is a normal, ASSIGNED target project isn't resolvable. The opened task is a normal,
delivery task (``source=VIDEO_SOURCE``, ``confirmed_by_human=True``) ASSIGNED delivery task (``source=VIDEO_SOURCE``,
— NOT held — so it dispatches straight to the assigned ux-dev like any ``confirmed_by_human=True``) — NOT held — so it dispatches straight to
other pre-assigned code task. the assigned ux-dev like any other pre-assigned code task.
``project_id`` scopes authoring to a specific project (the on-demand
``/video/request`` caller); omitted, the release/spotlight hooks
default to the fixed RoboCo project — see
``resolve_authoring_project``.
``brief`` is enriched (brand-voice + motion design-bar pointer ``brief`` is enriched (brand-voice + motion design-bar pointer
appended) before becoming the task description and the marker's appended) before becoming the task description and the marker's
@@ -263,7 +284,9 @@ class VideoEngine(BaseService):
occasion=occasion, occasion=occasion,
) )
return None return None
project = await self._opted_in_project(occasion) project = await self.resolve_authoring_project(
project_id=project_id, occasion=occasion
)
if project is None: if project is None:
return None return None
from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.exc import SQLAlchemyError
@@ -425,6 +448,39 @@ class VideoEngine(BaseService):
) )
return task return task
# ---- re-render (CEO-triggered retry) -----------------------------------
async def rerender(self, task_id: UUID) -> TaskTable | None:
"""Clear ``render_status``/``render_attempts``/``render_error`` off a
completed video-authoring task's ``video_draft`` marker, so the next
render cycle's scan (``render_status`` unset) re-picks it up and
re-renders it — e.g. after the CEO fixes something and wants a fresh
pass, or wants to retry past a terminal ``failed`` state.
None (a 404 to the route) when there is no such completed authoring
task, or the dev hasn't called ``propose_video`` yet (no
``composition_id`` — nothing to render).
"""
task = await get_task_service(self.session).get(task_id)
if (
task is None
or task.source != VIDEO_SOURCE
or task.status != TaskStatus.COMPLETED
):
return None
draft = markers.get_video_draft(task) or {}
if not draft.get("composition_id"):
return None
cleared = {
k: v
for k, v in draft.items()
if k not in ("render_status", "render_attempts", "render_error")
}
markers.set_video_draft(task, cleared)
await self.session.flush()
self.log.info("video-engine: re-render requested", task_id=str(task.id))
return task
def get_video_engine(session: AsyncSession) -> VideoEngine: def get_video_engine(session: AsyncSession) -> VideoEngine:
"""Build a VideoEngine for ``session``.""" """Build a VideoEngine for ``session``."""
+307 -12
View File
@@ -260,14 +260,17 @@ async def test_request_video_opens_authoring_task(
) -> None: ) -> None:
await _seed(db_session) await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True) monkeypatch.setattr(cfg, "video_engine_enabled", True)
monkeypatch.setattr(cfg, "self_heal_project_slug", SLUG)
monkeypatch.setattr(cfg, "video_max_open_posts", 5) monkeypatch.setattr(cfg, "video_max_open_posts", 5)
project = (
await db_session.execute(select(ProjectTable).where(ProjectTable.slug == SLUG))
).scalar_one()
resp = await ceo_client.post( resp = await ceo_client.post(
"/api/video/request", "/api/video/request",
json={ json={
"occasion": "CEO on-demand: launch teaser", "occasion": "CEO on-demand: launch teaser",
"brief": "A short teaser for the new dashboard", "brief": "A short teaser for the new dashboard",
"platforms": ["x", "tiktok"], "platforms": ["x", "tiktok"],
"project_id": str(project.id),
}, },
) )
assert resp.status_code == HTTPStatus.OK assert resp.status_code == HTTPStatus.OK
@@ -283,6 +286,7 @@ async def test_request_video_opens_authoring_task(
assert task is not None assert task is not None
assert task.source == VIDEO_SOURCE assert task.source == VIDEO_SOURCE
assert task.status == TaskStatus.PENDING assert task.status == TaskStatus.PENDING
assert task.project_id == project.id
finally: finally:
# The route's commit durably persists this task past this test's own # The route's commit durably persists this task past this test's own
# rollback teardown — a non-terminal source=video row left behind # rollback teardown — a non-terminal source=video row left behind
@@ -304,7 +308,12 @@ async def test_request_video_disabled_returns_clear_response(
before = len(await get_task_service(db_session).list_open_video_posts()) before = len(await get_task_service(db_session).list_open_video_posts())
resp = await ceo_client.post( resp = await ceo_client.post(
"/api/video/request", "/api/video/request",
json={"occasion": "occ-disabled", "brief": "brief", "platforms": ["x"]}, json={
"occasion": "occ-disabled",
"brief": "brief",
"platforms": ["x"],
"project_id": str(uuid4()),
},
) )
assert resp.status_code == HTTPStatus.OK assert resp.status_code == HTTPStatus.OK
body = resp.json() body = resp.json()
@@ -315,27 +324,91 @@ async def test_request_video_disabled_returns_clear_response(
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_request_video_not_opened_when_project_unresolvable( async def test_request_video_404s_on_unresolvable_project_id(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None: ) -> None:
"""An unresolvable project makes open_video_task no-op — a clear """A project_id that doesn't resolve to any project 404s — no fabricated
``not_opened`` response, not a 500 or a fabricated task.""" task, no silent not_opened."""
await _seed(db_session) await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True) 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()) before = len(await get_task_service(db_session).list_open_video_posts())
resp = await ceo_client.post( resp = await ceo_client.post(
"/api/video/request", "/api/video/request",
json={"occasion": "occ-unresolvable", "brief": "brief", "platforms": ["x"]}, json={
"occasion": "occ-unresolvable",
"brief": "brief",
"platforms": ["x"],
"project_id": str(uuid4()),
},
) )
assert resp.status_code == HTTPStatus.OK assert resp.status_code == HTTPStatus.NOT_FOUND
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()) after = len(await get_task_service(db_session).list_open_video_posts())
assert after == before # nothing new was opened assert after == before # nothing new was opened
@pytest.mark.asyncio
async def test_request_video_404s_on_non_opted_in_project_id(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A project that resolves but hasn't flipped video_engine_enabled also
404s, distinct from the unresolvable-project case above."""
await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True)
project = ProjectTable(
id=uuid4(),
name="Not Opted In",
slug=f"not-opted-{uuid4().hex[:6]}",
git_url="https://example.com/notopted.git",
assigned_cell=Team.BACKEND,
created_by=SYSTEM_UUID,
video_engine_enabled=False,
)
db_session.add(project)
await db_session.flush()
resp = await ceo_client.post(
"/api/video/request",
json={
"occasion": "occ-not-opted",
"brief": "brief",
"platforms": ["x"],
"project_id": str(project.id),
},
)
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_request_video_not_opened_on_duplicate_occasion(
db_session: AsyncSession, ceo_client: AsyncClient, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A valid, opted-in project_id still no-ops (not a 404) for a duplicate
occasion — the open-cap/dedup reasons stay a clear 200 not_opened."""
await _seed(db_session)
monkeypatch.setattr(cfg, "video_engine_enabled", True)
monkeypatch.setattr(cfg, "video_max_open_posts", 5)
project = (
await db_session.execute(select(ProjectTable).where(ProjectTable.slug == SLUG))
).scalar_one()
payload = {
"occasion": "occ-dupe",
"brief": "brief",
"platforms": ["x"],
"project_id": str(project.id),
}
first = await ceo_client.post("/api/video/request", json=payload)
assert first.status_code == HTTPStatus.OK
assert first.json()["status"] == "opened"
task_id = first.json()["task_id"]
try:
second = await ceo_client.post("/api/video/request", json=payload)
assert second.status_code == HTTPStatus.OK
assert second.json()["status"] == "not_opened"
assert second.json()["task_id"] is None
finally:
await db_session.execute(delete(TaskTable).where(TaskTable.id == UUID(task_id)))
await db_session.commit()
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_list_posts_returns_open_draft( async def test_list_posts_returns_open_draft(
db_session: AsyncSession, ceo_client: AsyncClient db_session: AsyncSession, ceo_client: AsyncClient
@@ -820,7 +893,12 @@ async def test_non_ceo_is_forbidden(db_session: AsyncSession) -> None:
async with AsyncClient(transport=transport, base_url="http://test") as client: async with AsyncClient(transport=transport, base_url="http://test") as client:
request_resp = await client.post( request_resp = await client.post(
"/api/video/request", "/api/video/request",
json={"occasion": "occ", "brief": "brief", "platforms": ["x"]}, json={
"occasion": "occ",
"brief": "brief",
"platforms": ["x"],
"project_id": str(uuid4()),
},
) )
list_resp = await client.get("/api/video/posts") list_resp = await client.get("/api/video/posts")
media_resp = await client.get(f"/api/video/posts/{task.id}/media?cut=vertical") media_resp = await client.get(f"/api/video/posts/{task.id}/media?cut=vertical")
@@ -977,3 +1055,220 @@ async def test_media_falls_back_to_local_file_when_minio_missing(
assert resp.headers["content-type"] == "video/mp4" assert resp.headers["content-type"] == "video/mp4"
assert resp.content == b"local-file-bytes" # FileResponse fallback, not MinIO assert resp.content == b"local-file-bytes" # FileResponse fallback, not MinIO
app.dependency_overrides.clear() app.dependency_overrides.clear()
# --- re-render (task 3, 2026-07-10) -------------------------------------------
@pytest.mark.asyncio
async def test_rerender_clears_state_and_returns_pipeline_item(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_authoring_task(
db_session,
status=TaskStatus.COMPLETED,
draft_extra={
"composition_id": "Intro",
"render_status": "failed",
"render_attempts": markers.MAX_VIDEO_RENDER_ATTEMPTS,
"render_error": "sidecar timeout",
},
)
resp = await ceo_client.post(f"/api/video/pipeline/{task.id}/rerender")
assert resp.status_code == HTTPStatus.OK
body = resp.json()
assert body["task_id"] == str(task.id)
assert body["render_status"] is None
assert body["render_attempts"] == 0
assert body["render_error"] is None
draft = markers.get_video_draft(task)
assert draft is not None
assert "render_status" not in draft
assert "render_attempts" not in draft
assert "render_error" not in draft
assert draft["composition_id"] == "Intro" # everything else preserved
@pytest.mark.asyncio
async def test_rerender_missing_task_is_404(ceo_client: AsyncClient) -> None:
resp = await ceo_client.post(f"/api/video/pipeline/{uuid4()}/rerender")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_rerender_non_completed_task_is_404(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_authoring_task(
db_session,
status=TaskStatus.IN_PROGRESS,
draft_extra={"composition_id": "Intro"},
)
resp = await ceo_client.post(f"/api/video/pipeline/{task.id}/rerender")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_rerender_without_composition_id_is_404(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_authoring_task(db_session, status=TaskStatus.COMPLETED)
resp = await ceo_client.post(f"/api/video/pipeline/{task.id}/rerender")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_rerender_non_ceo_is_forbidden(db_session: AsyncSession) -> None:
task = await _seed_authoring_task(
db_session,
status=TaskStatus.COMPLETED,
draft_extra={"composition_id": "Intro"},
)
app = _build_app(db_session, AgentRole.DEVELOPER, uuid4())
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.post(f"/api/video/pipeline/{task.id}/rerender")
assert resp.status_code == HTTPStatus.FORBIDDEN
app.dependency_overrides.clear()
# --- preview proxy (task 3, 2026-07-10) ---------------------------------------
def _fake_workspace_service(root: Path) -> object:
"""A ``get_workspace_service``-shaped stub whose ``ensure_read_clone``
returns a fixed local dir — no real git clone touched."""
class _Svc:
async def ensure_read_clone(self, _slug: str, *, force: bool = False) -> Path:
_ = force
return root
return _Svc()
def test_resolve_preview_path_serves_file_inside_root(tmp_path: Path) -> None:
root = (tmp_path / "clone").resolve()
(root / "motion" / "compositions" / "Intro").mkdir(parents=True)
target = root / "motion" / "compositions" / "Intro" / "vertical.html"
target.write_text("<html></html>")
resolved = video_module._resolve_preview_path(
root, "motion/compositions/Intro/vertical.html"
)
assert resolved == target.resolve()
def test_resolve_preview_path_blocks_dot_dot_traversal(tmp_path: Path) -> None:
root = (tmp_path / "clone").resolve()
(root / "motion").mkdir(parents=True)
secret = tmp_path / "secret.txt"
secret.write_text("nope")
assert video_module._resolve_preview_path(root, "../secret.txt") is None
assert video_module._resolve_preview_path(root, "motion/../../secret.txt") is None
def test_resolve_preview_path_blocks_absolute_path_override(tmp_path: Path) -> None:
"""A leading '/' in file_path would otherwise let pathlib's ``/``
operator discard ``root`` entirely and resolve straight to the absolute
path — the ``lstrip("/")`` guard neutralizes that."""
root = (tmp_path / "clone").resolve()
root.mkdir()
outside = tmp_path / "outside.txt"
outside.write_text("nope")
assert video_module._resolve_preview_path(root, str(outside)) is None
def test_resolve_preview_path_missing_file_is_none(tmp_path: Path) -> None:
root = (tmp_path / "clone").resolve()
root.mkdir()
assert video_module._resolve_preview_path(root, "motion/nope.html") is None
@pytest.mark.asyncio
async def test_preview_serves_composition_html_and_sibling_kit_asset(
db_session: AsyncSession,
ceo_client: AsyncClient,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
comp_dir = tmp_path / "motion" / "compositions" / "Intro"
comp_dir.mkdir(parents=True)
(comp_dir / "vertical.html").write_text("<html>intro</html>")
kit_dir = tmp_path / "motion" / "kit"
kit_dir.mkdir(parents=True)
(kit_dir / "kit.css").write_text("body{}")
monkeypatch.setattr(
video_module,
"get_workspace_service",
lambda _db: _fake_workspace_service(tmp_path),
)
task = await _seed_authoring_task(
db_session,
status=TaskStatus.IN_PROGRESS,
draft_extra={"composition_id": "Intro"},
)
resp = await ceo_client.get(
f"/api/video/preview/{task.id}/motion/compositions/Intro/vertical.html"
)
assert resp.status_code == HTTPStatus.OK
assert resp.text == "<html>intro</html>"
assert resp.headers.get("x-frame-options") == "SAMEORIGIN"
assert "frame-ancestors" in resp.headers.get("content-security-policy", "")
kit_resp = await ceo_client.get(f"/api/video/preview/{task.id}/motion/kit/kit.css")
assert kit_resp.status_code == HTTPStatus.OK
assert kit_resp.text == "body{}"
@pytest.mark.asyncio
async def test_preview_missing_file_is_404(
db_session: AsyncSession,
ceo_client: AsyncClient,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
(tmp_path / "motion").mkdir()
monkeypatch.setattr(
video_module,
"get_workspace_service",
lambda _db: _fake_workspace_service(tmp_path),
)
task = await _seed_authoring_task(
db_session,
status=TaskStatus.IN_PROGRESS,
draft_extra={"composition_id": "Intro"},
)
resp = await ceo_client.get(
f"/api/video/preview/{task.id}/motion/compositions/Intro/vertical.html"
)
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_preview_missing_task_is_404(ceo_client: AsyncClient) -> None:
resp = await ceo_client.get(f"/api/video/preview/{uuid4()}/vertical.html")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_preview_non_video_task_is_404(
db_session: AsyncSession, ceo_client: AsyncClient
) -> None:
task = await _seed_draft(db_session) # source=video_post, not video
resp = await ceo_client.get(f"/api/video/preview/{task.id}/vertical.html")
assert resp.status_code == HTTPStatus.NOT_FOUND
@pytest.mark.asyncio
async def test_preview_non_ceo_is_forbidden(db_session: AsyncSession) -> None:
task = await _seed_authoring_task(
db_session,
status=TaskStatus.IN_PROGRESS,
draft_extra={"composition_id": "Intro"},
)
app = _build_app(db_session, AgentRole.DEVELOPER, uuid4())
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get(f"/api/video/preview/{task.id}/vertical.html")
assert resp.status_code == HTTPStatus.FORBIDDEN
app.dependency_overrides.clear()
@@ -42,6 +42,7 @@ UX_DEV_2_UUID = _foundation.AGENTS["ux-dev-2"].uuid
SLUG = "roboco" SLUG = "roboco"
ONE = 1 ONE = 1
TWO = 2 TWO = 2
FOUR = 4
def _orch() -> Any: def _orch() -> Any:
@@ -333,6 +334,35 @@ async def test_render_video_task_renders_both_cuts_and_materializes_post(
assert source_draft["render_status"] == "rendered" assert source_draft["render_status"] == "rendered"
@pytest.mark.asyncio
async def test_render_video_task_resolves_workspace_from_task_project_not_settings(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
"""The render loop resolves the read-clone from the authoring task's OWN
project_id flipping self_heal_project_slug to a bogus value AFTER the
task was authored must not affect the render, proving the loop no longer
reads that setting live."""
await _seed(db_session)
_enable(monkeypatch)
task = await _make_completed_video_task(
db_session, occasion="own-project-not-settings", composition_id="Intro"
)
monkeypatch.setattr(cfg, "self_heal_project_slug", "no-such-project-anymore")
renderer = _FakeRenderer()
workspace = _fake_workspace()
orch = _orch()
p1, p2 = _render_patches(renderer, workspace)
with p1, p2:
await orch._render_video_task(db_session, task)
# Still resolved via the task's own project_id -> slug "roboco", not the
# now-bogus self_heal_project_slug.
workspace.ensure_read_clone.assert_awaited_once_with(SLUG)
posts = await get_task_service(db_session).list_open_video_posts()
assert len(posts) == ONE
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_render_video_task_second_call_is_idempotent( async def test_render_video_task_second_call_is_idempotent(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
@@ -356,6 +386,45 @@ async def test_render_video_task_second_call_is_idempotent(
assert len(posts) == ONE assert len(posts) == ONE
@pytest.mark.asyncio
async def test_rerender_clears_state_so_next_cycle_re_renders(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
"""CEO re-render flow end to end: a rendered task is a no-op on a second
render pass (idempotent); clearing render_status/render_attempts via
VideoEngine.rerender makes the NEXT pass pick it up and render it again."""
await _seed(db_session)
_enable(monkeypatch)
task = await _make_completed_video_task(
db_session, occasion="rerender-cycle", composition_id="Intro"
)
renderer = _FakeRenderer()
workspace = _fake_workspace()
orch = _orch()
p1, p2 = _render_patches(renderer, workspace)
with p1, p2:
await orch._render_video_task(db_session, task)
assert len(renderer.calls) == TWO
draft = markers.get_video_draft(task)
assert draft is not None
assert draft["render_status"] == "rendered"
rerendered = await VideoEngine(db_session).rerender(task.id)
assert rerendered is not None
draft = markers.get_video_draft(task)
assert draft is not None
assert "render_status" not in draft
assert "render_attempts" not in draft
with p1, p2:
await orch._render_video_task(db_session, task) # re-picked up
assert len(renderer.calls) == FOUR # rendered a second time, not skipped
draft = markers.get_video_draft(task)
assert draft is not None
assert draft["render_status"] == "rendered"
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_render_video_task_skips_task_without_composition_id( async def test_render_video_task_skips_task_without_composition_id(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
+172
View File
@@ -10,6 +10,7 @@ from __future__ import annotations
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from unittest.mock import AsyncMock from unittest.mock import AsyncMock
from uuid import uuid4
import pytest import pytest
from roboco.config import settings as cfg from roboco.config import settings as cfg
@@ -307,6 +308,93 @@ async def test_open_video_task_insert_error_returns_none_session_usable(
assert check.scalar_one_or_none() is not None assert check.scalar_one_or_none() is not None
# --------------------------------------------------------------------------- #
# open_video_task(project_id=...) — the on-demand caller's own project scope
# --------------------------------------------------------------------------- #
@pytest.mark.asyncio
async def test_open_video_task_with_explicit_project_id_ignores_self_heal_slug(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A caller-supplied project_id resolves independent of
self_heal_project_slug the authoring task lands on THAT project, not
the fixed RoboCo one."""
await _seed(db_session)
_enable(monkeypatch)
other = ProjectTable(
name="Other Repo",
slug="other-repo",
git_url="https://github.com/x/other.git",
default_branch="master",
protected_branches=["master"],
assigned_cell=Team.BACKEND,
created_by=SYSTEM_UUID,
is_active=True,
video_engine_enabled=True,
)
db_session.add(other)
await db_session.flush()
monkeypatch.setattr(cfg, "self_heal_project_slug", "no-such-slug")
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="on-demand other repo",
script="s",
platforms=["x"],
brief="b",
project_id=other.id,
)
assert task is not None
assert task.project_id == other.id
@pytest.mark.asyncio
async def test_open_video_task_explicit_project_id_not_opted_in_opens_nothing(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
_enable(monkeypatch)
other = ProjectTable(
name="Not Opted In",
slug="not-opted-in",
git_url="https://github.com/x/notopted.git",
default_branch="master",
protected_branches=["master"],
assigned_cell=Team.BACKEND,
created_by=SYSTEM_UUID,
is_active=True,
video_engine_enabled=False,
)
db_session.add(other)
await db_session.flush()
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="on-demand not opted",
script="s",
platforms=["x"],
brief="b",
project_id=other.id,
)
assert task is None
@pytest.mark.asyncio
async def test_open_video_task_explicit_project_id_unresolvable_opens_nothing(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
_enable(monkeypatch)
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="on-demand missing project",
script="s",
platforms=["x"],
brief="b",
project_id=uuid4(),
)
assert task is None
# --------------------------------------------------------------------------- # # --------------------------------------------------------------------------- #
# _originate_video_post # _originate_video_post
# --------------------------------------------------------------------------- # # --------------------------------------------------------------------------- #
@@ -644,5 +732,89 @@ async def test_draft_release_video_dedupes_same_version(
second = await engine.draft_release_video(version="1.0.0", changelog=_CHANGELOG) second = await engine.draft_release_video(version="1.0.0", changelog=_CHANGELOG)
assert first is not None assert first is not None
assert second is None assert second is None
# --------------------------------------------------------------------------- #
# rerender — CEO-triggered clear of the render idempotency keys
# --------------------------------------------------------------------------- #
@pytest.mark.asyncio
async def test_rerender_clears_render_state_keeps_the_rest(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
_enable(monkeypatch)
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="rerender-me", script="s", platforms=["x"], brief="b"
)
assert task is not None
draft = markers.get_video_draft(task) or {}
markers.set_video_draft(
task,
{
**draft,
"composition_id": "Intro",
"render_status": "failed",
"render_attempts": THREE,
"render_error": "sidecar timeout",
},
)
task.status = TS.COMPLETED
await db_session.flush()
result = await engine.rerender(task.id)
assert result is not None
cleared = markers.get_video_draft(result)
assert cleared is not None
assert "render_status" not in cleared
assert "render_attempts" not in cleared
assert "render_error" not in cleared
assert cleared["composition_id"] == "Intro" # everything else preserved
@pytest.mark.asyncio
async def test_rerender_none_for_non_completed_task(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
_enable(monkeypatch)
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="not-yet-done", script="s", platforms=["x"], brief="b"
)
assert task is not None
draft = markers.get_video_draft(task) or {}
markers.set_video_draft(task, {**draft, "composition_id": "Intro"})
await db_session.flush() # still PENDING, not COMPLETED
result = await engine.rerender(task.id)
assert result is None
@pytest.mark.asyncio
async def test_rerender_none_without_composition_id(
db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
await _seed(db_session)
_enable(monkeypatch)
engine = video_engine_module.VideoEngine(db_session)
task = await engine.open_video_task(
occasion="no-composition-yet", script="s", platforms=["x"], brief="b"
)
assert task is not None
task.status = TS.COMPLETED
await db_session.flush() # no propose_video call yet -> no composition_id
result = await engine.rerender(task.id)
assert result is None
@pytest.mark.asyncio
async def test_rerender_none_for_missing_task(db_session: AsyncSession) -> None:
engine = video_engine_module.VideoEngine(db_session)
result = await engine.rerender(uuid4())
assert result is None
open_tasks = await get_task_service(db_session).list_open_video_posts() open_tasks = await get_task_service(db_session).list_open_video_posts()
assert len(open_tasks) == ONE assert len(open_tasks) == ONE