From 88ab5c6dc58a5fea72ce8018cf42d85b53974063 Mon Sep 17 00:00:00 2001 From: Backend Developer 1 Date: Fri, 10 Jul 2026 04:42:10 +0000 Subject: [PATCH] [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. --- pyproject.toml | 5 + roboco/api/routes/video.py | 111 ++++++- roboco/api/schemas/video.py | 4 +- roboco/runtime/orchestrator.py | 29 +- roboco/services/video_engine.py | 90 +++++- tests/integration/test_video_routes.py | 319 ++++++++++++++++++- tests/unit/runtime/test_video_render_loop.py | 69 ++++ tests/unit/services/test_video_engine.py | 172 ++++++++++ 8 files changed, 752 insertions(+), 47 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index c1cf544d..a1d9ad9b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -154,6 +154,11 @@ select = [ # agent_id, project_ids, route, session_id) — same >5-kwarg rationale as the # gateway verb surfaces below. "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 # verb contract (task title, description, acceptance criteria, assignee, etc.). Bundling # into a dataclass hides the field-by-field schema the LLM needs at the diff --git a/roboco/api/routes/video.py b/roboco/api/routes/video.py index 03805615..9ff5fc2e 100644 --- a/roboco/api/routes/video.py +++ b/roboco/api/routes/video.py @@ -7,7 +7,7 @@ from __future__ import annotations import asyncio from pathlib import Path -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, cast from uuid import UUID 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.security import guard_deco 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_credentials import ( TikTokCredentialsValidationError, @@ -41,6 +42,7 @@ from roboco.services.video_post_service import ( VideoCaptionTooLongError, 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_video_client import build_x_video_poster @@ -89,11 +91,15 @@ async def request_video( db: DbSession, agent: CurrentAgentContext, ) -> 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`` - when ``open_video_task`` no-ops (a duplicate occasion, the open cap, or - an unresolvable project) — neither is an error, just nothing to do. + Returns ``disabled`` when the video engine is off. 404s when + ``project_id`` doesn't resolve to a project, or that project hasn't + 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) if not settings.video_engine_enabled: @@ -101,18 +107,27 @@ async def request_video( status="disabled", 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, script=data.brief, platforms=data.platforms, brief=data.brief, + project_id=data.project_id, ) if task is None: return VideoRequestResponse( status="not_opened", detail=( - "No video task was opened (a duplicate occasion, the open-post" - " cap, or the project isn't resolvable)." + "No video task was opened (a duplicate occasion or the open-post cap)." ), ) await db.commit() @@ -196,6 +211,84 @@ async def list_video_pipeline( 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//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]: """Every ``{platform}_posted_id`` key stamped by approve, keyed by platform (e.g. ``{"x": "..", "tiktok": ".."}``).""" diff --git a/roboco/api/schemas/video.py b/roboco/api/schemas/video.py index 78ae4f75..96bdd116 100644 --- a/roboco/api/schemas/video.py +++ b/roboco/api/schemas/video.py @@ -1,6 +1,7 @@ """Schemas for the video engine's on-demand request + CEO approval surface.""" from datetime import datetime +from uuid import UUID from pydantic import BaseModel, Field @@ -10,11 +11,12 @@ from roboco.services.x_client import MAX_TWEET_CHARS 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) brief: str = Field(..., min_length=1) platforms: list[str] = Field(..., min_length=1) + project_id: UUID class VideoRequestResponse(BaseModel): diff --git a/roboco/runtime/orchestrator.py b/roboco/runtime/orchestrator.py index 2d71cffb..c82eec09 100644 --- a/roboco/runtime/orchestrator.py +++ b/roboco/runtime/orchestrator.py @@ -8106,7 +8106,7 @@ Start by: return try: 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) except Exception as exc: @@ -8151,16 +8151,29 @@ Start by: ) 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]: - """Render the vertical + square cuts from the roboco project's merged - read-clone's motion/ dir; returns {"vertical": path, "square": path}. - ``render_key`` (the source task id) scopes each cut's output path.""" + """Render the vertical + square cuts from the authoring task's OWN + project's merged read-clone's motion/ dir; returns {"vertical": 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.workspace import get_workspace_service + from roboco.services.workspace import WorkspaceError, get_workspace_service - slug = (settings.self_heal_project_slug or "roboco-api").strip() - workspace = await get_workspace_service(db).ensure_read_clone(slug) + project = await get_project_service(db).get(project_id) if project_id else None + 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") input_props = draft.get("input_props") or {} renderer = get_video_renderer() diff --git a/roboco/services/video_engine.py b/roboco/services/video_engine.py index 80960975..0e17c8ca 100644 --- a/roboco/services/video_engine.py +++ b/roboco/services/video_engine.py @@ -149,23 +149,38 @@ class VideoEngine(BaseService): service_name = "video_engine" 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() return await get_project_service(self.session).get_by_slug(slug) - async def _opted_in_project(self, occasion: str) -> ProjectTable | None: - """The RoboCo project if resolvable AND opted into video, else None. + async def resolve_authoring_project( + 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 - config 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 this repo into - authoring against its ``motion/`` dir (mirrors ``ci_watch_enabled``). + ``project_id`` (the on-demand ``/video/request`` caller, and the + render loop's own per-task resolution) resolves by id; omitted (the + release/spotlight hooks), it falls back to the fixed RoboCo project. + Both paths share the same two skip reasons, both logged: unresolvable + 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: self.log.warning( - "video-engine: RoboCo project not resolvable; skipping video task", + "video-engine: project not resolvable; skipping video task", occasion=occasion, + project_id=str(project_id) if project_id else None, ) return None if not getattr(project, "video_engine_enabled", False): @@ -231,16 +246,22 @@ class VideoEngine(BaseService): platforms: list[str], brief: str, suggested_input_props: dict[str, Any] | None = None, + project_id: UUID | None = None, ) -> TaskTable | 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 - (``video_engine_enabled``), a task for this occasion is already open - (authoring or held draft), the open cap is reached, or the RoboCo - project isn't resolvable. The opened task is a normal, ASSIGNED - delivery task (``source=VIDEO_SOURCE``, ``confirmed_by_human=True``) - — NOT held — so it dispatches straight to the assigned ux-dev like any - other pre-assigned code task. + No-ops when the global flag is off, the target project hasn't opted + in (``video_engine_enabled``), a task for this occasion is already + open (authoring or held draft), the open cap is reached, or the + target project isn't resolvable. The opened task is a normal, + ASSIGNED delivery task (``source=VIDEO_SOURCE``, + ``confirmed_by_human=True``) — NOT held — so it dispatches straight to + 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 appended) before becoming the task description and the marker's @@ -263,7 +284,9 @@ class VideoEngine(BaseService): occasion=occasion, ) 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: return None from sqlalchemy.exc import SQLAlchemyError @@ -425,6 +448,39 @@ class VideoEngine(BaseService): ) 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: """Build a VideoEngine for ``session``.""" diff --git a/tests/integration/test_video_routes.py b/tests/integration/test_video_routes.py index 695e9482..d256cfdb 100644 --- a/tests/integration/test_video_routes.py +++ b/tests/integration/test_video_routes.py @@ -260,14 +260,17 @@ async def test_request_video_opens_authoring_task( ) -> 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) + project = ( + await db_session.execute(select(ProjectTable).where(ProjectTable.slug == SLUG)) + ).scalar_one() 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"], + "project_id": str(project.id), }, ) 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.source == VIDEO_SOURCE assert task.status == TaskStatus.PENDING + assert task.project_id == project.id finally: # The route's commit durably persists this task past this test's own # 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()) resp = await ceo_client.post( "/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 body = resp.json() @@ -315,27 +324,91 @@ async def test_request_video_disabled_returns_clear_response( @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 ) -> None: - """An unresolvable project makes open_video_task no-op — a clear - ``not_opened`` response, not a 500 or a fabricated task.""" + """A project_id that doesn't resolve to any project 404s — no fabricated + task, no silent not_opened.""" 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"]}, + json={ + "occasion": "occ-unresolvable", + "brief": "brief", + "platforms": ["x"], + "project_id": str(uuid4()), + }, ) - assert resp.status_code == HTTPStatus.OK - body = resp.json() - assert body["status"] == "not_opened" - assert body["task_id"] is None + assert resp.status_code == HTTPStatus.NOT_FOUND 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_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 async def test_list_posts_returns_open_draft( 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: request_resp = await client.post( "/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") 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.content == b"local-file-bytes" # FileResponse fallback, not MinIO 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("") + 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("intro") + 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 == "intro" + 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() diff --git a/tests/unit/runtime/test_video_render_loop.py b/tests/unit/runtime/test_video_render_loop.py index c29d65e0..b19d5837 100644 --- a/tests/unit/runtime/test_video_render_loop.py +++ b/tests/unit/runtime/test_video_render_loop.py @@ -42,6 +42,7 @@ UX_DEV_2_UUID = _foundation.AGENTS["ux-dev-2"].uuid SLUG = "roboco" ONE = 1 TWO = 2 +FOUR = 4 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" +@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 async def test_render_video_task_second_call_is_idempotent( db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch @@ -356,6 +386,45 @@ async def test_render_video_task_second_call_is_idempotent( 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 async def test_render_video_task_skips_task_without_composition_id( db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch diff --git a/tests/unit/services/test_video_engine.py b/tests/unit/services/test_video_engine.py index a73ab9ed..1b0a4009 100644 --- a/tests/unit/services/test_video_engine.py +++ b/tests/unit/services/test_video_engine.py @@ -10,6 +10,7 @@ from __future__ import annotations from typing import TYPE_CHECKING from unittest.mock import AsyncMock +from uuid import uuid4 import pytest 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 +# --------------------------------------------------------------------------- # +# 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 # --------------------------------------------------------------------------- # @@ -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) assert first is not 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() assert len(open_tasks) == ONE