[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
# 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
+102 -9
View File
@@ -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/<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]:
"""Every ``{platform}_posted_id`` key stamped by approve, keyed by
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."""
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):
+21 -8
View File
@@ -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()
+73 -17
View File
@@ -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``."""
+307 -12
View File
@@ -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("<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"
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
+172
View File
@@ -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