From 3bedc26d033ca19902d855005caede65975bab52 Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Thu, 24 Sep 2026 12:29:47 +0000 Subject: [PATCH 1/2] Render the five-second pre-roll: compositor, encoder, validation (package B) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A design snapshot now produces the media a pre-roll actually needs: three MP4 renditions, an fMP4 HLS package, a poster, and captions plus an audible companion where there is narration to carry. Package B of the streaming-ads spec. Nothing selects or serves this yet — that is package D. Two properties the pipeline is built around. It is deterministic. The composition exposes window.__seek(frame) and every animated value is a pure function of that index; nothing reads a clock, and there is no CSS animation or rAF loop whose output would depend on how quickly Playwright got round to the screenshot. That is what lets a render be cached by a hash of its inputs without the hash being a lie. RENDERER_VERSION is part of that hash, so changing the look invalidates every cached encode rather than serving the old bytes forever. It is closed. Artwork reaches the compositor only as a data: URI the caller already fetched and hashed, advertiser copy is escaped before it enters a document a browser executes, and the capture context aborts every request that is not data:/blob:/about:. A renderer that fetched a URL for itself would be a request-forgery primitive running on our own network, and it would also make renders depend on what that URL served today. The snapshot stores artwork content hashes rather than URLs for the same reason: the same URL can serve different bytes, and a cache keyed on the URL would reuse a render of artwork the advertiser has replaced. Validation asks the decoded output, never the job that produced it. ffprobe runs with -count_frames, so "150 frames" means 150 frames decoded rather than whatever the muxer wrote in the header; 29.97 fails where 30 passes; audio gets exactly one AAC frame (21.333ms at 48kHz) of endpoint rounding and no more, and the ad is never padded to six seconds to land on a round segment size. The GOP is pinned to 60 frames with scene-change detection off so the HLS boundaries land on 0/2/4s regardless of artwork, and HLS is packaged by stream copy so the segments carry the media the advertiser approved. EXT-X-INDEPENDENT-SEGMENTS is deliberately never emitted: our segments open on a keyframe but are not independently decodable in the sense that tag asserts, and a player acts on the claim. Publishing a revision is a compare-and-swap against requested_revision. Two quick edits queue revisions 4 and 5; if 4 finishes second, a naive write would point the creative back at media for a design already replaced. Conditioned on requested_revision still being 4, the stale write matches no rows and is discarded. Render jobs dedupe on a hash of immutable inputs only, and the BullMQ job id is colon-free — bullmq parses a colon-bearing custom id as a structured key and rejects anything that is not exactly three parts, and an id derived from state the job itself resets would either stop deduping or collide with a finished job and drop the retry silently. The Redis connection parser moved out of lib/prober-queue into lib/redis-connection now that there are two queues; prober behaviour is unchanged. Frame capture is injected rather than imported so Playwright stays in the worker image, and so the pipeline can be driven in a test without Chromium. tests/ads-video-pipeline.test.ts does exactly that and then runs real ffmpeg: 150 decoded frames at 30fps in every rendition, each inside its byte budget, segments summing to five seconds, init map and ENDLIST present. The spec is explicit that a filename or a five-second timer is not evidence a pre-roll works, so the one test that could have been faked is the one that is not. Co-Authored-By: Claude Opus 5 (1M context) --- lib/ads/video/compose.ts | 220 ++++++++++++++ lib/ads/video/encode.ts | 288 ++++++++++++++++++ lib/ads/video/profiles.ts | 206 +++++++++++++ lib/ads/video/queue.ts | 85 ++++++ lib/ads/video/render.ts | 354 +++++++++++++++++++++++ lib/ads/video/snapshot.ts | 155 ++++++++++ lib/ads/video/storage.ts | 144 +++++++++ lib/ads/video/validate.ts | 255 ++++++++++++++++ lib/prober-queue.ts | 21 +- lib/redis-connection.ts | 22 ++ tests/ads-video-compose-validate.test.ts | 340 ++++++++++++++++++++++ tests/ads-video-contract.test.ts | 199 +++++++++++++ tests/ads-video-pipeline.test.ts | 213 ++++++++++++++ worker/Dockerfile | 9 +- worker/frames.ts | 87 ++++++ worker/index.ts | 4 + worker/video.ts | 162 +++++++++++ 17 files changed, 2743 insertions(+), 21 deletions(-) create mode 100644 lib/ads/video/compose.ts create mode 100644 lib/ads/video/encode.ts create mode 100644 lib/ads/video/profiles.ts create mode 100644 lib/ads/video/queue.ts create mode 100644 lib/ads/video/render.ts create mode 100644 lib/ads/video/snapshot.ts create mode 100644 lib/ads/video/storage.ts create mode 100644 lib/ads/video/validate.ts create mode 100644 lib/redis-connection.ts create mode 100644 tests/ads-video-compose-validate.test.ts create mode 100644 tests/ads-video-contract.test.ts create mode 100644 tests/ads-video-pipeline.test.ts create mode 100644 worker/frames.ts create mode 100644 worker/video.ts diff --git a/lib/ads/video/compose.ts b/lib/ads/video/compose.ts new file mode 100644 index 00000000..9524965d --- /dev/null +++ b/lib/ads/video/compose.ts @@ -0,0 +1,220 @@ +// The five-second composition, as a self-contained HTML document the worker +// screenshots frame by frame. +// +// Two properties matter more than how it looks: +// +// * It is deterministic. Nothing reads a clock. The document exposes +// `window.__seek(frame)` and every animated value is a pure function of +// that frame index, so frame 73 is identical whether it was rendered first +// or last, on a fast machine or a loaded one. A CSS animation or a +// requestAnimationFrame loop would make the output depend on how quickly +// Playwright got round to the screenshot, which is exactly the kind of +// nondeterminism that turns a cached render hash into a lie. +// +// * It is closed. Artwork arrives as a data: URI the caller has already +// fetched and hashed; the document references no network origin, no local +// file and no font service. A compositor that fetched a URL for itself +// would be a request-forgery primitive pointed at our own infrastructure, +// and it would also make renders depend on whatever that URL served today. + +import { PREROLL_FRAMES, PREROLL_FPS, TIMELINE, frameTimeMs } from "./profiles"; +import type { VideoDesignSnapshot } from "./snapshot"; + +/** Artwork the caller resolved, already fetched and inlined. */ +export type ComposeAssets = { + /** `data:image/...;base64,...` or null. */ + logo: string | null; + hero: string | null; +}; + +/** + * Copy is advertiser-controlled and lands inside a document a browser executes. + * Escaping it is not defence-in-depth, it is the boundary: an advertiser who + * can inject markup here runs script in the render container. + */ +export function escapeHtml(s: string): string { + return s + .replace(/&/g, "&") + .replace(//g, ">") + .replace(/"/g, """) + .replace(/'/g, "'"); +} + +/** + * A data: URI, or nothing. + * + * The compositor accepts only data: URIs for artwork. Anything else — an http + * URL, a file path, a javascript: scheme — is dropped rather than sanitised, + * because there is no legitimate caller that needs one and a partial sanitiser + * is how the interesting ones get through. + */ +export function safeDataUri(uri: string | null): string | null { + if (!uri) return null; + return /^data:image\/(png|jpeg|webp|gif|svg\+xml);base64,[A-Za-z0-9+/=]+$/.test(uri) ? uri : null; +} + +/** Cubic ease-out. Motion that decelerates reads as deliberate rather than mechanical. */ +function easeOut(t: number): number { + const c = Math.min(1, Math.max(0, t)); + return 1 - Math.pow(1 - c, 3); +} + +/** Progress through a window, clamped to [0,1]. */ +function phase(ms: number, startMs: number, endMs: number): number { + if (endMs <= startMs) return ms >= endMs ? 1 : 0; + return Math.min(1, Math.max(0, (ms - startMs) / (endMs - startMs))); +} + +export type FrameState = { + frame: number; + timeMs: number; + /** Headline/brand entrance: a short rise, already legible at frame 0. */ + entrance: number; + /** Slow artwork drift across the hold. */ + drift: number; + /** CTA emphasis over the final beat. */ + cta: number; +}; + +/** + * The animated state at a frame. Exported because the test suite asserts the + * timeline's shape directly rather than by decoding pixels, and because the + * document below embeds this same function verbatim. + */ +export function frameState(frame: number, reducedMotion: boolean): FrameState { + const timeMs = frameTimeMs(frame); + if (reducedMotion) { + // Same five-second timeline, no movement: every beat is already at rest. + // The ad still *ends* on its CTA, it just never animates toward it. + return { frame, timeMs, entrance: 1, drift: 0, cta: 1 }; + } + return { + frame, + timeMs, + entrance: easeOut(phase(timeMs, 0, TIMELINE.entranceEndMs)), + drift: phase(timeMs, TIMELINE.entranceEndMs, TIMELINE.holdEndMs), + cta: easeOut(phase(timeMs, TIMELINE.holdEndMs, TIMELINE.endMs)), + }; +} + +function initial(s: string): string { + const m = s.match(/[a-z0-9]/i); + return m ? m[0].toUpperCase() : "★"; +} + +/** + * The composition document. + * + * Laid out at 1920x1080 and scaled by the caller's viewport rather than + * re-laid-out per rendition, so the 480p encode is a downscale of the same + * pixels the 1080p master shows. Copy that fits the master therefore fits every + * delivery rendition, which is the only way a legibility check at one size + * means anything at the others. + */ +export function composeDocument( + snapshot: VideoDesignSnapshot, + assets: ComposeAssets = { logo: null, hero: null }, +): string { + const headline = escapeHtml(snapshot.headline); + const cta = escapeHtml(snapshot.ctaText); + const domain = escapeHtml(snapshot.domain); + const logo = safeDataUri(assets.logo); + const hero = safeDataUri(assets.hero); + const font = escapeHtml(snapshot.fontFamily || "system-ui, -apple-system, Segoe UI, Roboto, sans-serif"); + + // Colours are validated as hex by validateSnapshot() before a render is + // queued; escaping them too means a bad one degrades to an inert attribute + // rather than closing the style block. + const bg = escapeHtml(snapshot.bgColor); + const fg = escapeHtml(snapshot.fgColor); + const accent = escapeHtml(snapshot.accentColor); + + const heroLayer = hero + ? `
` + : `
`; + + const mark = logo + ? `` + : ``; + + return ` +ad + +
+ ${heroLayer} +
+
${mark}${domain}
+

${headline}

+
${cta}
+
+
+ +`; +} + +/** + * Every frame's state, without a browser. + * + * The encoder needs the frame count and the tests need the timeline; neither + * should have to launch Chromium to get them. + */ +export function timeline(reducedMotion: boolean): FrameState[] { + return Array.from({ length: PREROLL_FRAMES }, (_, i) => frameState(i, reducedMotion)); +} diff --git a/lib/ads/video/encode.ts b/lib/ads/video/encode.ts new file mode 100644 index 00000000..431056fa --- /dev/null +++ b/lib/ads/video/encode.ts @@ -0,0 +1,288 @@ +// ffmpeg invocations for the five-second pre-roll. +// +// Argument construction is separated from execution on purpose: the arguments +// are where the media contract actually lives (GOP length is what makes the +// HLS segment boundaries land on 0/2/4s; +faststart is what makes the MP4 +// playable before it has finished downloading), and a pure builder is something +// a test can assert without a 90-second encode. + +import { spawn } from "node:child_process"; +import path from "node:path"; +import { + AAC_SAMPLE_RATE, + HLS_KEYFRAME_SECONDS, + PREROLL_FPS, + PREROLL_FRAMES, + videoProfile, + type VideoProfileId, +} from "./profiles"; + +export const FFMPEG_BIN = process.env.FFMPEG_PATH || "ffmpeg"; +export const FFPROBE_BIN = process.env.FFPROBE_PATH || "ffprobe"; + +/** + * Frames between keyframes. + * + * HLS_KEYFRAME_SECONDS is [0, 2, 4] and the ad runs at 30 fps, so a keyframe + * every 60 frames puts one exactly at each boundary. `-sc_threshold 0` stops + * x264 inserting extra keyframes on scene changes, which would be harmless for + * playback but would make segment durations depend on the artwork. + */ +export const GOP_FRAMES = Math.round( + (HLS_KEYFRAME_SECONDS[1] - HLS_KEYFRAME_SECONDS[0]) * PREROLL_FPS, +); + +export type Mp4EncodeOptions = { + /** printf-style pattern of the PNG frame sequence, e.g. `/tmp/x/f-%04d.png`. */ + framePattern: string; + /** Optional narration track. AAC-LC at 48 kHz is produced regardless of input. */ + audioPath: string | null; + outPath: string; + profile: VideoProfileId; + /** Video bitrate in kbps. The caller lowers it and retries if over budget. */ + videoKbps: number; +}; + +/** + * Frames in, one MP4 out. + * + * `-frames:v 150` rather than `-t 5`: the contract is a frame count, and a + * duration flag lets a timebase rounding error produce 149 or 151 frames that + * still measure "5.0 seconds". Counting frames is the check that actually + * catches a dropped one. + */ +export function mp4Args(o: Mp4EncodeOptions): string[] { + const p = videoProfile(o.profile); + const args = [ + "-y", + "-nostdin", + "-framerate", + String(PREROLL_FPS), + "-i", + o.framePattern, + ]; + + if (o.audioPath) args.push("-i", o.audioPath); + + args.push( + "-frames:v", + String(PREROLL_FRAMES), + "-r", + String(PREROLL_FPS), + "-c:v", + "libx264", + // High profile, 8-bit 4:2:0. The spec's compatibility floor, and what every + // browser and receiver in the inventory can decode in hardware. + "-profile:v", + "high", + "-pix_fmt", + "yuv420p", + "-b:v", + `${o.videoKbps}k`, + "-maxrate", + `${Math.round(o.videoKbps * 1.3)}k`, + "-bufsize", + `${o.videoKbps * 2}k`, + "-g", + String(GOP_FRAMES), + "-keyint_min", + String(GOP_FRAMES), + "-sc_threshold", + "0", + ); + + if (p.width && p.height) { + // Explicit scale rather than relying on the source size, so a compositor + // viewport change cannot silently ship a differently-sized rendition. + args.push("-vf", `scale=${p.width}:${p.height}:flags=lanczos`); + } + + if (o.audioPath) { + args.push("-c:a", "aac", "-profile:a", "aac_low", "-ar", String(AAC_SAMPLE_RATE), "-b:a", "128k", "-ac", "2", "-shortest"); + } else { + args.push("-an"); + } + + // Move the moov atom to the front. Without it a player must fetch the tail of + // the file before it can start, which on a five-second pre-roll is most of + // the budget spent before the first frame. + args.push("-movflags", "+faststart", o.outPath); + return args; +} + +export type HlsOptions = { + /** Source MP4 for this rendition — already encoded with aligned keyframes. */ + inPath: string; + outDir: string; + /** Basename, e.g. "720p": yields 720p.m3u8, 720p-init.mp4, 720p-1.m4s… */ + name: string; +}; + +/** + * Package one rendition as fMP4 HLS. + * + * `-c copy` is load-bearing. Re-encoding here would move the keyframes and + * break the 0/2/4s boundaries the MP4 was encoded to hit, and it would also + * mean the HLS a viewer streams is not the same media the advertiser + * downloaded and approved. + */ +export function hlsArgs(o: HlsOptions): string[] { + return [ + "-y", + "-nostdin", + "-i", + o.inPath, + "-c", + "copy", + "-f", + "hls", + "-hls_time", + String(HLS_KEYFRAME_SECONDS[1] - HLS_KEYFRAME_SECONDS[0]), + "-hls_playlist_type", + "vod", // writes EXT-X-ENDLIST; this is a finite VOD package, not a live window + "-hls_segment_type", + "fmp4", + "-hls_fmp4_init_filename", + `${o.name}-init.mp4`, + "-hls_list_size", + "0", + "-hls_segment_filename", + path.join(o.outDir, `${o.name}-%d.m4s`), + path.join(o.outDir, `${o.name}.m3u8`), + ]; +} + +/** + * The multivariant playlist. + * + * Written by hand rather than by ffmpeg's var_stream_map because the bandwidth + * and codec declarations have to describe the renditions we actually produced, + * and because EXT-X-INDEPENDENT-SEGMENTS is deliberately absent: our segments + * open on a keyframe but are not independently decodable in the sense that tag + * asserts, and claiming it would be a lie a player acts on. + */ +export function multivariantPlaylist( + renditions: { name: string; width: number; height: number; bandwidth: number; codecs: string }[], +): string { + const lines = ["#EXTM3U", "#EXT-X-VERSION:7"]; + for (const r of renditions) { + lines.push( + `#EXT-X-STREAM-INF:BANDWIDTH=${r.bandwidth},RESOLUTION=${r.width}x${r.height},CODECS="${r.codecs}"`, + ); + lines.push(`${r.name}.m3u8`); + } + return lines.join("\n") + "\n"; +} + +/** + * Poster frame. + * + * Taken from the hold, not from frame 0: the first frame is mid-entrance, and a + * poster is the still a viewer stares at while the player buffers. 2 seconds in + * the composition is settled and the headline is fully up. + */ +export function posterArgs(inPath: string, outPath: string): string[] { + return [ + "-y", + "-nostdin", + "-ss", + "2", + "-i", + inPath, + "-frames:v", + "1", + "-c:v", + "libwebp", + "-quality", + "82", + "-vf", + "scale=1280:720:flags=lanczos", + outPath, + ]; +} + +/** + * The audible companion, as AAC in an MP4 container. + * + * Produced from the same narration as the video's own track so the two cannot + * drift into describing different offers. A silent creative has no companion: + * five seconds of silence would be recorded as an audio ad having played, and + * the spec is explicit that silence never satisfies an audio slot. + */ +export function audioCompanionArgs(inPath: string, outPath: string): string[] { + return [ + "-y", + "-nostdin", + "-i", + inPath, + "-vn", + "-c:a", + "aac", + "-profile:a", + "aac_low", + "-ar", + String(AAC_SAMPLE_RATE), + "-b:a", + "128k", + "-ac", + "2", + "-movflags", + "+faststart", + outPath, + ]; +} + +export class FfmpegError extends Error { + constructor( + message: string, + readonly code: number | null, + readonly stderr: string, + ) { + super(message); + this.name = "FfmpegError"; + } +} + +/** + * Run a binary and collect stderr. + * + * stderr is captured rather than inherited because ffmpeg's failure message is + * the only useful diagnostic when a render fails on a worker nobody is watching, + * and it has to reach the job row. Only the tail is kept: a progress-spammed + * stderr can run to megabytes and none of the early lines say why it failed. + */ +export function run(bin: string, args: string[], timeoutMs = 120_000): Promise { + return new Promise((resolve, reject) => { + const child = spawn(bin, args, { stdio: ["ignore", "pipe", "pipe"] }); + let stdout = ""; + let stderr = ""; + const timer = setTimeout(() => { + child.kill("SIGKILL"); + reject(new FfmpegError(`${bin} timed out after ${timeoutMs}ms`, null, stderr.slice(-4000))); + }, timeoutMs); + + child.stdout.on("data", (d) => { + stdout += String(d); + }); + child.stderr.on("data", (d) => { + stderr += String(d); + if (stderr.length > 200_000) stderr = stderr.slice(-100_000); + }); + child.on("error", (e) => { + clearTimeout(timer); + reject(new FfmpegError(`${bin} failed to start: ${e.message}`, null, stderr.slice(-4000))); + }); + child.on("close", (code) => { + clearTimeout(timer); + if (code === 0) resolve(stdout); + else reject(new FfmpegError(`${bin} exited ${code}`, code, stderr.slice(-4000))); + }); + }); +} + +export const encodeMp4 = (o: Mp4EncodeOptions) => run(FFMPEG_BIN, mp4Args(o)); +export const packageHls = (o: HlsOptions) => run(FFMPEG_BIN, hlsArgs(o)); +export const extractPoster = (inPath: string, outPath: string) => + run(FFMPEG_BIN, posterArgs(inPath, outPath)); +export const encodeAudioCompanion = (inPath: string, outPath: string) => + run(FFMPEG_BIN, audioCompanionArgs(inPath, outPath)); diff --git a/lib/ads/video/profiles.ts b/lib/ads/video/profiles.ts new file mode 100644 index 00000000..0bf64ffd --- /dev/null +++ b/lib/ads/video/profiles.ts @@ -0,0 +1,206 @@ +// The media contract for a five-second streaming pre-roll: durations, frame +// counts, output profiles and the budgets validation enforces. +// +// Pure constants and arithmetic, kept free of node built-ins for the same +// reason ../formats is — the dashboard renders render-job state and asset sizes +// in a client component, and importing this must not drag ffmpeg or the +// Playwright compositor into the browser bundle. + +/** + * Five seconds, as frames rather than milliseconds. + * + * The spec's duration requirement is "exactly 150 frames at 30 fps", and that + * is deliberately not the same statement as "5000 ms". A frame count is what + * ffprobe can count and what a compositor can seek to; a millisecond duration + * is a float that container timebases round. Everything downstream measures + * frames and derives time from them, never the other way round. + */ +export const PREROLL_FPS = 30; +export const PREROLL_FRAMES = 150; +export const PREROLL_MS = (PREROLL_FRAMES / PREROLL_FPS) * 1000; // 5000 + +/** Presentation time of a frame index, in milliseconds. */ +export function frameTimeMs(frame: number): number { + return (frame / PREROLL_FPS) * 1000; +} + +/** + * The composition timeline. Three beats, from the spec's §4.2 table. + * + * Brand and headline are already legible at frame 0 — the entrance animates + * around copy that is readable from the first frame rather than revealing it. + * A five-second ad that spends its first half second assembling itself has + * spent a tenth of its life saying nothing. + */ +export const TIMELINE = { + entranceEndMs: 350, + holdEndMs: 3500, + endMs: PREROLL_MS, +} as const; + +/** + * Keyframe positions, in seconds, for the HLS packaging. + * + * Segments come out ~2s/2s/1s. The third is deliberately short rather than + * padding the ad to six seconds to land on a round segment size: the viewer's + * ad is five seconds, and a nominal segment duration is not a reason to make + * somebody watch a sixth. + */ +export const HLS_KEYFRAME_SECONDS = [0, 2, 4] as const; + +/** + * Audio endpoint tolerance: one AAC frame at 48 kHz. + * + * AAC is framed at 1024 samples, so an encoder cannot land an audio track on an + * arbitrary boundary — it pads to the next frame and signals the remainder as + * priming/padding. 1024/48000 = 21.333 ms of legitimate rounding. Validation + * allows exactly one frame of it and no more, which catches a genuinely wrong + * duration while not failing a correct encode for being AAC. + */ +export const AAC_FRAME_SAMPLES = 1024; +export const AAC_SAMPLE_RATE = 48_000; +export const AUDIO_TOLERANCE_MS = (AAC_FRAME_SAMPLES / AAC_SAMPLE_RATE) * 1000; // 21.333… + +export type VideoProfileId = + | "master_1080p" + | "mp4_720p" + | "mp4_480p" + | "hls" + | "poster" + | "captions" + | "audio"; + +export type VideoProfile = { + id: VideoProfileId; + /** What the object is for, in the advertiser-facing dashboard. */ + label: string; + width: number | null; + height: number | null; + /** Hard ceiling. A profile that exceeds it re-encodes or fails validation. */ + maxBytes: number | null; + contentType: string; + /** False for the derived/optional outputs that are not always produced. */ + required: boolean; +}; + +const KB = 1024; +const MB = 1024 * KB; + +/** + * The outputs one ready revision consists of. + * + * The byte ceilings are product targets, not HLS requirements — they exist + * because a pre-roll that buffers has already failed, and a viewer on a phone + * pays for these bytes. The 1080p master is the one an advertiser downloads, + * so it gets the loosest budget; the 480p rendition is what a constrained + * connection actually receives, so it gets the tightest. + */ +export const VIDEO_PROFILES: VideoProfile[] = [ + { + id: "master_1080p", + label: "Downloadable master (1080p)", + width: 1920, + height: 1080, + maxBytes: 3 * MB, + contentType: "video/mp4", + required: true, + }, + { + id: "mp4_720p", + label: "Delivery MP4 (720p)", + width: 1280, + height: 720, + maxBytes: Math.round(1.5 * MB), + contentType: "video/mp4", + required: true, + }, + { + id: "mp4_480p", + label: "Delivery MP4 (480p)", + width: 854, + height: 480, + maxBytes: 750 * KB, + contentType: "video/mp4", + required: true, + }, + { + id: "hls", + label: "HLS package", + width: null, + height: null, + maxBytes: null, + contentType: "application/vnd.apple.mpegurl", + required: true, + }, + { + id: "poster", + label: "Poster frame", + width: 1280, + height: 720, + maxBytes: 300 * KB, + contentType: "image/webp", + required: true, + }, + { + id: "captions", + label: "Captions (WebVTT)", + width: null, + height: null, + maxBytes: 64 * KB, + contentType: "text/vtt", + // Only when the creative has narration. A silent ad has nothing to caption, + // and an empty VTT would be worse than none. + required: false, + }, + { + id: "audio", + label: "Audible companion", + width: null, + height: null, + maxBytes: 200 * KB, + contentType: "audio/mp4", + // Required only for a property with an audio-only slot. A silent gap does + // not satisfy the ad requirement there, so when it is required it is + // genuinely required — see audioRequiredFor() below. + required: false, + }, +]; + +export function videoProfile(id: VideoProfileId): VideoProfile { + const p = VIDEO_PROFILES.find((x) => x.id === id); + if (!p) throw new Error(`unknown video profile: ${id}`); + return p; +} + +/** The MP4 renditions, largest first — the download master and both deliveries. */ +export const MP4_PROFILE_IDS: VideoProfileId[] = ["master_1080p", "mp4_720p", "mp4_480p"]; + +/** + * Which profiles a given creative must produce before its revision can publish. + * + * `audioMode` comes from the design snapshot: a creative with narration always + * carries captions and an audio companion, and a silent creative carries an + * audio companion only if some property it may serve has an audio slot. The + * caller passes that in rather than this module guessing, because it is a + * publisher-inventory question, not a media one. + */ +export function requiredProfiles(opts: { + narrated: boolean; + audioSlotSupported: boolean; +}): VideoProfileId[] { + const ids = VIDEO_PROFILES.filter((p) => p.required).map((p) => p.id); + if (opts.narrated) ids.push("captions"); + if (opts.narrated || opts.audioSlotSupported) ids.push("audio"); + return ids; +} + +/** + * Whether a rendition's byte size is within its budget. + * + * Separated from validation so the encoder can check a result and re-encode at + * a lower bitrate before the whole job is failed for being 40 KB over. + */ +export function withinBudget(id: VideoProfileId, byteSize: number): boolean { + const max = videoProfile(id).maxBytes; + return max === null || byteSize <= max; +} diff --git a/lib/ads/video/queue.ts b/lib/ads/video/queue.ts new file mode 100644 index 00000000..3ccb4611 --- /dev/null +++ b/lib/ads/video/queue.ts @@ -0,0 +1,85 @@ +// The render queue. +// +// Enqueue is driven from ad_video_jobs rows the same way port scans are driven +// from port_scans (see ../../prober-queue): the row is the durable record and +// the BullMQ job is the transient work item, so losing Redis loses throughput +// rather than losing an advertiser's render. + +import { Queue } from "bullmq"; +import { redisConnectionOptions } from "@/lib/redis-connection"; +import type { VideoProfileId } from "./profiles"; +import type { VideoDesignSnapshot } from "./snapshot"; + +export const VIDEO_RENDER_QUEUE = "ad-video-render"; + +export type VideoRenderJobData = { + jobRowId: string; + ownerId: string; + campaignId: string | null; + creativeId: string | null; + revision: number; + snapshot: VideoDesignSnapshot; + profile: VideoProfileId | "default"; + audioSlotSupported: boolean; +}; + +let queue: Queue | null = null; + +export function getVideoRenderQueue(): Queue | null { + const connection = redisConnectionOptions(); + if (!connection) return null; + if (!queue) queue = new Queue(VIDEO_RENDER_QUEUE, { connection }); + return queue; +} + +/** + * BullMQ job id for a render. + * + * Two constraints, both learned the hard way and neither obvious from the API: + * + * * No colons. BullMQ parses a custom jobId containing ":" as a structured + * key and rejects anything that does not split into exactly three parts. + * A render hash never contains one, but a caller passing a composed + * "owner:campaign:rev" id would work by accident and then break the moment + * a fourth component was added. + * + * * Derived only from immutable input. The id is the dedupe key, so if it + * were derived from anything the job itself changes — the job row's state, + * an attempt counter, the creative's published revision — then a retry + * would compute a different id and the "already queued" check would stop + * working, or worse, the id would collide with a finished job and the retry + * would be silently dropped as a duplicate. The render hash is a hash of + * inputs that by construction do not change. + */ +export function renderJobId(renderHash: string, profile: string): string { + const id = `vr-${renderHash}-${profile}`; + if (id.includes(":")) { + throw new Error(`render job id must not contain ':' (got ${id})`); + } + return id; +} + +/** + * Enqueue a render, or no-op if one is already queued for these exact inputs. + * + * Returns false when Redis is not configured — the caller leaves the job row + * `queued` and a later sweep picks it up, rather than failing an advertiser's + * campaign save because a worker dependency is down. + */ +export async function enqueueRender(args: { + renderHash: string; + profile: string; + data: VideoRenderJobData; +}): Promise { + const q = getVideoRenderQueue(); + if (!q) return false; + + await q.add("render", args.data, { + jobId: renderJobId(args.renderHash, args.profile), + attempts: 3, + backoff: { type: "exponential", delay: 10_000 }, + removeOnComplete: { age: 3600 }, + removeOnFail: { age: 86_400 }, + }); + return true; +} diff --git a/lib/ads/video/render.ts b/lib/ads/video/render.ts new file mode 100644 index 00000000..a5f39e8b --- /dev/null +++ b/lib/ads/video/render.ts @@ -0,0 +1,354 @@ +// Snapshot in, validated assets out. +// +// The orchestration only. Frame capture is injected rather than imported so +// that Playwright stays in the worker image where it belongs, and so the +// pipeline can be driven in a test by a capturer that writes synthetic frames +// instead of launching Chromium for 150 screenshots. + +import { mkdir, readFile, readdir, writeFile } from "node:fs/promises"; +import path from "node:path"; +import { + MP4_PROFILE_IDS, + PREROLL_FRAMES, + requiredProfiles, + videoProfile, + withinBudget, + type VideoProfileId, +} from "./profiles"; +import { composeDocument, type ComposeAssets } from "./compose"; +import { validateSnapshot, type VideoDesignSnapshot } from "./snapshot"; +import { + encodeAudioCompanion, + encodeMp4, + extractPoster, + multivariantPlaylist, + packageHls, +} from "./encode"; +import { evaluateProbe, probeMedia, validateMediaPlaylist, type ValidationProblem } from "./validate"; +import { statSync } from "node:fs"; +import { createHash } from "node:crypto"; + +/** + * Writes `frames` PNGs into `outDir` named `f-0000.png` … and resolves. + * + * The contract the worker's Playwright implementation satisfies, and the seam + * the tests substitute at. + */ +export type FrameCapturer = (args: { + html: string; + outDir: string; + frames: number; + width: number; + height: number; +}) => Promise; + +export const FRAME_PATTERN = "f-%04d.png"; +export function framePath(dir: string, i: number): string { + return path.join(dir, `f-${String(i).padStart(4, "0")}.png`); +} + +/** + * Starting bitrates, chosen to land comfortably inside each profile's byte + * budget rather than at its edge. + * + * A five-second 1080p file at 4.8 Mbps is exactly 3 MB, so encoding at 4.8 + * would put every render one rounding error from failing validation. 4.0 leaves + * room for the container overhead and the moov atom. + */ +const START_KBPS: Record = { + master_1080p: 4000, + mp4_720p: 2000, + mp4_480p: 1000, +}; + +/** How many times a rendition may be re-encoded lower before the job fails. */ +const BUDGET_RETRIES = 2; + +export type RenderedAsset = { + profile: VideoProfileId; + /** Path on disk, relative to the job's working directory. */ + filePath: string; + contentType: string; + byteSize: number; + sha256: string; + width: number | null; + height: number | null; + durationMs: number | null; + codecs: string | null; + validation: { ok: boolean; problems: ValidationProblem[] }; + /** For the HLS profile: the segment and init files that travel with it. */ + extraFiles?: string[]; +}; + +export type RenderResult = { + assets: RenderedAsset[]; + problems: ValidationProblem[]; +}; + +export class RenderError extends Error { + constructor( + message: string, + readonly code: string, + ) { + super(message); + this.name = "RenderError"; + } +} + +async function fileFacts(filePath: string) { + const bytes = await readFile(filePath); + return { + byteSize: bytes.byteLength, + sha256: createHash("sha256").update(bytes).digest("hex"), + }; +} + +/** + * Encode one MP4 rendition, lowering the bitrate if it lands over budget. + * + * Retrying is worth the extra encode: the alternative is failing an otherwise + * correct render because a dense hero image pushed a rendition 5% over, which + * an advertiser cannot act on and which a lower bitrate fixes invisibly. + */ +async function encodeWithinBudget(args: { + profile: VideoProfileId; + framePattern: string; + audioPath: string | null; + outPath: string; +}): Promise { + let kbps = START_KBPS[args.profile] ?? 2000; + + for (let attempt = 0; attempt <= BUDGET_RETRIES; attempt++) { + await encodeMp4({ + framePattern: args.framePattern, + audioPath: args.audioPath, + outPath: args.outPath, + profile: args.profile, + videoKbps: kbps, + }); + const size = statSync(args.outPath).size; + if (withinBudget(args.profile, size)) return size; + + const max = videoProfile(args.profile).maxBytes!; + // Aim at 90% of budget rather than exactly at it, so the next attempt has + // margin instead of landing on the boundary again. + kbps = Math.max(200, Math.floor((kbps * max * 0.9) / size)); + } + + throw new RenderError( + `${args.profile} exceeds its byte budget after ${BUDGET_RETRIES + 1} attempts`, + "budget_exceeded", + ); +} + +/** + * Render one design snapshot into every required output. + * + * `audioPath` is a narration track the caller has already synthesised; this + * module does not do text-to-speech. Keeping that out means a render is pure + * media work with no model call in it, which is what lets the same snapshot be + * re-rendered years later and produce the same bytes. + */ +export async function renderPreroll(args: { + snapshot: VideoDesignSnapshot; + assets?: ComposeAssets; + workDir: string; + captureFrames: FrameCapturer; + audioPath?: string | null; + audioSlotSupported?: boolean; +}): Promise { + const { snapshot, workDir, captureFrames } = args; + + const snapshotProblems = validateSnapshot(snapshot); + if (snapshotProblems.length > 0) { + throw new RenderError( + `snapshot rejected: ${snapshotProblems.map((p) => `${p.field} (${p.reason})`).join(", ")}`, + "invalid_snapshot", + ); + } + + const narrated = snapshot.audioMode === "narrated"; + const audioPath = narrated ? (args.audioPath ?? null) : null; + if (narrated && !audioPath) { + // The snapshot says narrated and no track arrived. Encoding anyway would + // produce a silent "narrated" ad, and a silent audio companion is the one + // outcome the spec singles out as never acceptable. + throw new RenderError("narrated snapshot with no audio track", "missing_narration_audio"); + } + + const framesDir = path.join(workDir, "frames"); + const outDir = path.join(workDir, "out"); + await mkdir(framesDir, { recursive: true }); + await mkdir(outDir, { recursive: true }); + + // 1. Compose and capture. Always at the master's dimensions; the renditions + // are downscales of these pixels, not separate layouts. + const html = composeDocument(snapshot, args.assets ?? { logo: null, hero: null }); + await captureFrames({ html, outDir: framesDir, frames: PREROLL_FRAMES, width: 1920, height: 1080 }); + + const captured = (await readdir(framesDir)).filter((f) => f.endsWith(".png")); + if (captured.length !== PREROLL_FRAMES) { + throw new RenderError( + `compositor produced ${captured.length} frames, expected ${PREROLL_FRAMES}`, + "frame_count", + ); + } + + const framePattern = path.join(framesDir, FRAME_PATTERN); + const assets: RenderedAsset[] = []; + const problems: ValidationProblem[] = []; + + // 2. The MP4 renditions. + for (const profile of MP4_PROFILE_IDS) { + const outPath = path.join(outDir, `${profile}.mp4`); + await encodeWithinBudget({ profile, framePattern, audioPath, outPath }); + + const probe = await probeMedia(outPath); + const facts = await fileFacts(outPath); + const result = evaluateProbe(probe, profile, facts.byteSize, { expectAudio: !!audioPath }); + if (!result.ok) problems.push(...result.problems); + + const spec = videoProfile(profile); + assets.push({ + profile, + filePath: outPath, + contentType: spec.contentType, + byteSize: facts.byteSize, + sha256: facts.sha256, + width: result.measured.width, + height: result.measured.height, + durationMs: result.measured.frames !== null ? (result.measured.frames / 30) * 1000 : null, + codecs: result.measured.videoCodec, + validation: { ok: result.ok, problems: result.problems }, + }); + } + + // 3. HLS, packaged by stream-copying the delivery renditions so the segments + // carry exactly the media the advertiser approved. + const hlsDir = path.join(outDir, "hls"); + await mkdir(hlsDir, { recursive: true }); + const renditions: { name: string; width: number; height: number; bandwidth: number; codecs: string }[] = []; + + for (const profile of ["mp4_720p", "mp4_480p"] as const) { + const spec = videoProfile(profile); + const name = `${spec.height}p`; + await packageHls({ inPath: path.join(outDir, `${profile}.mp4`), outDir: hlsDir, name }); + + const playlist = await readFile(path.join(hlsDir, `${name}.m3u8`), "utf8"); + problems.push(...validateMediaPlaylist(playlist)); + + const mp4Size = statSync(path.join(outDir, `${profile}.mp4`)).size; + renditions.push({ + name, + width: spec.width!, + height: spec.height!, + // Peak bandwidth, derived from the rendition we actually produced. + bandwidth: Math.round((mp4Size * 8) / 5), + codecs: "avc1.640028,mp4a.40.2", + }); + } + + const master = multivariantPlaylist(renditions); + const masterPath = path.join(hlsDir, "master.m3u8"); + await writeFile(masterPath, master, "utf8"); + + const hlsFiles = (await readdir(hlsDir)).map((f) => path.join(hlsDir, f)); + const hlsFacts = await fileFacts(masterPath); + assets.push({ + profile: "hls", + filePath: masterPath, + contentType: videoProfile("hls").contentType, + byteSize: hlsFacts.byteSize, + sha256: hlsFacts.sha256, + width: null, + height: null, + durationMs: null, + codecs: "avc1.640028,mp4a.40.2", + validation: { ok: true, problems: [] }, + extraFiles: hlsFiles.filter((f) => f !== masterPath), + }); + + // 4. Poster, from the settled hold rather than the entrance. + const posterPath = path.join(outDir, "poster.webp"); + await extractPoster(path.join(outDir, "master_1080p.mp4"), posterPath); + const posterFacts = await fileFacts(posterPath); + assets.push({ + profile: "poster", + filePath: posterPath, + contentType: videoProfile("poster").contentType, + byteSize: posterFacts.byteSize, + sha256: posterFacts.sha256, + width: 1280, + height: 720, + durationMs: null, + codecs: null, + validation: { ok: true, problems: [] }, + }); + + // 5. Captions and the audible companion, when there is narration to carry. + if (narrated && snapshot.narration) { + const vttPath = path.join(outDir, "captions.vtt"); + await writeFile(vttPath, narrationVtt(snapshot.narration), "utf8"); + const vttFacts = await fileFacts(vttPath); + assets.push({ + profile: "captions", + filePath: vttPath, + contentType: videoProfile("captions").contentType, + byteSize: vttFacts.byteSize, + sha256: vttFacts.sha256, + width: null, + height: null, + durationMs: null, + codecs: null, + validation: { ok: true, problems: [] }, + }); + } + + if (audioPath && (narrated || args.audioSlotSupported)) { + const companionPath = path.join(outDir, "audio.m4a"); + await encodeAudioCompanion(audioPath, companionPath); + const probe = await probeMedia(companionPath); + const facts = await fileFacts(companionPath); + const audioStream = probe.streams.find((s) => s.codec_type === "audio"); + assets.push({ + profile: "audio", + filePath: companionPath, + contentType: videoProfile("audio").contentType, + byteSize: facts.byteSize, + sha256: facts.sha256, + width: null, + height: null, + durationMs: audioStream?.duration ? Number(audioStream.duration) * 1000 : null, + codecs: audioStream?.codec_name ?? null, + validation: { ok: true, problems: [] }, + }); + } + + // 6. Everything the caller said this creative owes must be present. + const required = requiredProfiles({ + narrated, + audioSlotSupported: !!args.audioSlotSupported, + }); + for (const id of required) { + if (!assets.some((a) => a.profile === id)) { + problems.push({ check: "required profile", expected: id, actual: "absent" }); + } + } + + return { assets, problems }; +} + +/** + * A single caption cue spanning the ad. + * + * Five seconds of narration is one sentence; splitting it into timed cues would + * be inventing timings we did not measure. One cue over the whole timeline is + * honest and is what a five-second read actually looks like. + */ +export function narrationVtt(narration: string): string { + // Collapse every run of whitespace, not just newlines: a script arrives with + // the indentation of wherever it was authored, and "NicheDB today" is what + // a viewer would otherwise see rendered as a caption. + const text = narration.trim().replace(/\s+/g, " "); + return `WEBVTT\n\n00:00:00.000 --> 00:00:05.000\n${text}\n`; +} diff --git a/lib/ads/video/snapshot.ts b/lib/ads/video/snapshot.ts new file mode 100644 index 00000000..e35657a3 --- /dev/null +++ b/lib/ads/video/snapshot.ts @@ -0,0 +1,155 @@ +// The immutable input to a render, and the hash that dedupes it. +// +// A render job is expensive and perfectly deterministic: the same snapshot, the +// same renderer, the same profile always produce the same bytes. So the job is +// keyed by a hash of everything that can change the output, and an unchanged +// design reuses the encode that already exists instead of burning a worker on +// it again. + +import { createHash } from "node:crypto"; +import type { VideoProfileId } from "./profiles"; + +/** + * Bump when the compositor's output changes for unchanged input. + * + * This is part of the hash, which is the whole point: a change to the timeline, + * the type scale, or the safe area must invalidate every cached render, and + * without a version in the key a fixed compositor would keep serving the old + * bytes forever. Bumping it is the deliberate cost of changing the look. + */ +export const RENDERER_VERSION = "1"; + +export type AudioMode = "silent" | "narrated"; + +/** + * Everything the compositor is allowed to read. + * + * Deliberately flat and primitive. A snapshot is stored as JSON, hashed, and + * replayed weeks later by a worker that cannot re-fetch anything — so it holds + * resolved values, not references. `logoSha256`/`heroSha256` are the *content* + * hashes of the source artwork rather than its URLs, because the same URL can + * serve different bytes and a cache keyed on the URL would then reuse a render + * of artwork the advertiser has since replaced. + */ +export type VideoDesignSnapshot = { + headline: string; + ctaText: string; + /** Bare host, e.g. "nichedb.dev" — shown as the destination, never a full URL. */ + domain: string; + bgColor: string; + fgColor: string; + accentColor: string; + fontFamily: string; + logoUrl: string | null; + logoSha256: string | null; + heroUrl: string | null; + heroSha256: string | null; + audioMode: AudioMode; + /** Narration script, when audioMode is "narrated". Drives the captions too. */ + narration: string | null; + locale: string; + /** Static composition over the same five-second timeline. */ + reducedMotion: boolean; +}; + +/** + * Canonical JSON: keys sorted at every level, no incidental whitespace. + * + * JSON.stringify preserves insertion order, so two snapshots that differ only + * in the order their fields were assigned would otherwise hash differently and + * re-render identical bytes. Sorting makes the hash a function of the content. + */ +export function canonicalJson(value: unknown): string { + if (value === null || typeof value !== "object") return JSON.stringify(value) ?? "null"; + if (Array.isArray(value)) return `[${value.map(canonicalJson).join(",")}]`; + const entries = Object.entries(value as Record) + // Undefined members are absent, not null: an explicitly-undefined field and + // a missing one describe the same design and must hash the same. + .filter(([, v]) => v !== undefined) + .sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)); + return `{${entries.map(([k, v]) => `${JSON.stringify(k)}:${canonicalJson(v)}`).join(",")}}`; +} + +/** + * The dedupe key for one output of one design. + * + * Profile is part of the key rather than a separate column lookup because the + * 720p and 480p renditions of the same snapshot are different bytes with + * different budgets, and a single key across both would let one satisfy the + * other. + */ +export function renderHash(snapshot: VideoDesignSnapshot, profile: VideoProfileId): string { + return createHash("sha256") + .update( + canonicalJson({ + v: RENDERER_VERSION, + profile, + snapshot, + }), + ) + .digest("hex"); +} + +/** Content hash of a source asset, for the snapshot's *Sha256 fields. */ +export function assetHash(bytes: Uint8Array): string { + return createHash("sha256").update(bytes).digest("hex"); +} + +/** + * Headline length guard. + * + * Eight words is the spec's target, but the real constraint is that the + * headline has to be legible at the 480p rendition's height for the roughly + * three seconds it holds. Measured layout happens in the compositor; this is + * the cheap check that runs before a worker is spent on it. + */ +export const MAX_HEADLINE_WORDS = 8; + +export function headlineWords(headline: string): number { + return headline.trim().split(/\s+/).filter(Boolean).length; +} + +export type SnapshotProblem = { field: string; reason: string }; + +/** + * Reject a snapshot that cannot produce a legible ad, before it costs a render. + * + * Returns every problem rather than the first, so the dashboard can show an + * advertiser all of what needs fixing in one pass instead of one per attempt. + */ +export function validateSnapshot(s: VideoDesignSnapshot): SnapshotProblem[] { + const problems: SnapshotProblem[] = []; + const hex = /^#([0-9a-fA-F]{6}|[0-9a-fA-F]{8})$/; + + if (!s.headline.trim()) { + problems.push({ field: "headline", reason: "empty" }); + } else if (headlineWords(s.headline) > MAX_HEADLINE_WORDS) { + problems.push({ + field: "headline", + reason: `${headlineWords(s.headline)} words, max ${MAX_HEADLINE_WORDS}`, + }); + } + + if (!s.ctaText.trim()) problems.push({ field: "ctaText", reason: "empty" }); + if (!s.domain.trim()) problems.push({ field: "domain", reason: "empty" }); + + for (const [field, value] of [ + ["bgColor", s.bgColor], + ["fgColor", s.fgColor], + ["accentColor", s.accentColor], + ] as const) { + if (!hex.test(value)) problems.push({ field, reason: "not a hex colour" }); + } + + // A narrated ad with no script would render five seconds of silence and be + // recorded as an audible companion, which is the failure the spec calls out + // by name: silence is never an audio ad. + if (s.audioMode === "narrated" && !s.narration?.trim()) { + problems.push({ field: "narration", reason: "narrated but no script" }); + } + if (s.audioMode === "silent" && s.narration?.trim()) { + problems.push({ field: "narration", reason: "script present but audioMode is silent" }); + } + + return problems; +} diff --git a/lib/ads/video/storage.ts b/lib/ads/video/storage.ts new file mode 100644 index 00000000..6f3f3e25 --- /dev/null +++ b/lib/ads/video/storage.ts @@ -0,0 +1,144 @@ +// Where a rendered revision's bytes live. +// +// Reuses the existing public `ad-assets` bucket rather than provisioning new +// storage: it is already how server-generated ad art (see ../heroImage) is +// hosted, its RLS already confines authenticated writes to a user's own prefix, +// and service-role uploads already bypass that for generated content. + +import { serviceClient } from "@/lib/supabase/service"; +import { readFile } from "node:fs/promises"; +import path from "node:path"; +import type { VideoProfileId } from "./profiles"; +import type { RenderedAsset } from "./render"; + +export const ASSET_BUCKET = "ad-assets"; + +/** + * Immutable object prefix for one revision. + * + * Revision is in the path, so publishing a new one adds objects rather than + * replacing them. That is what lets a decision already issued keep resolving to + * the exact bytes the viewer was promised while a newer revision serves every + * session that starts afterwards — and it is why these objects can be cached + * immutably at the edge. + */ +export function revisionPrefix(args: { + ownerId: string; + campaignId: string; + creativeId: string; + revision: number; +}): string { + return `video/${args.ownerId}/${args.campaignId}/${args.creativeId}/r${args.revision}`; +} + +/** Object key for one profile's primary file. */ +export function objectKey( + prefix: string, + profile: VideoProfileId, + filename: string, +): string { + // HLS is a directory of files, so it keeps its own subtree; every other + // profile is a single object named for the profile. + return profile === "hls" ? `${prefix}/hls/${filename}` : `${prefix}/${filename}`; +} + +export type UploadedAsset = { + profile: VideoProfileId; + objectKey: string; + publicUrl: string; + byteSize: number; + sha256: string; + contentType: string; + width: number | null; + height: number | null; + durationMs: number | null; + codecs: string | null; +}; + +/** + * Upload one render's outputs. + * + * Draft renders are uploaded to the same bucket as published ones and are + * distinguished by the `published` flag on ad_video_assets rather than by + * location, so publishing a revision is a database transaction rather than a + * byte copy. The flip side — draft bytes are reachable by URL to anyone who + * guesses a uuid quadruple — is acceptable for ad creative and is the same + * posture the hero images already have. Anything genuinely private would need a + * separate private bucket and signed delivery. + */ +export async function uploadRenderedAssets(args: { + assets: RenderedAsset[]; + ownerId: string; + campaignId: string; + creativeId: string; + revision: number; +}): Promise { + const svc = serviceClient(); + const prefix = revisionPrefix(args); + const uploaded: UploadedAsset[] = []; + + for (const asset of args.assets) { + // HLS travels as a set: the playlist is useless without its init segment + // and media segments, so they upload together or the profile does not + // count as stored at all. + const files = [asset.filePath, ...(asset.extraFiles ?? [])]; + let primaryKey = ""; + + for (const filePath of files) { + const bytes = await readFile(filePath); + const key = objectKey(prefix, asset.profile, path.basename(filePath)); + const { error } = await svc.storage.from(ASSET_BUCKET).upload(key, bytes, { + contentType: contentTypeFor(filePath, asset.contentType), + upsert: false, + // A revision's objects never change, so the edge may keep them for a + // year. A new revision is a new path. + cacheControl: "public, max-age=31536000, immutable", + }); + if (error) { + throw new Error(`upload failed for ${key}: ${error.message}`); + } + if (filePath === asset.filePath) primaryKey = key; + } + + uploaded.push({ + profile: asset.profile, + objectKey: primaryKey, + publicUrl: svc.storage.from(ASSET_BUCKET).getPublicUrl(primaryKey).data.publicUrl, + byteSize: asset.byteSize, + sha256: asset.sha256, + contentType: asset.contentType, + width: asset.width, + height: asset.height, + durationMs: asset.durationMs, + codecs: asset.codecs, + }); + } + + return uploaded; +} + +/** + * Per-file content type. + * + * The HLS profile's declared type describes its playlist; the segments and init + * fragment beside it are MP4. Serving an .m4s as application/vnd.apple.mpegurl + * makes some players refuse it outright. + */ +export function contentTypeFor(filePath: string, profileContentType: string): string { + const ext = path.extname(filePath).toLowerCase(); + switch (ext) { + case ".m3u8": + return "application/vnd.apple.mpegurl"; + case ".m4s": + case ".mp4": + return "video/mp4"; + case ".m4a": + return "audio/mp4"; + case ".webp": + return "image/webp"; + case ".vtt": + return "text/vtt"; + default: + return profileContentType; + } +} diff --git a/lib/ads/video/validate.ts b/lib/ads/video/validate.ts new file mode 100644 index 00000000..07285245 --- /dev/null +++ b/lib/ads/video/validate.ts @@ -0,0 +1,255 @@ +// Does this encode actually satisfy the media contract? +// +// The question is asked of the decoded output, never of the job that produced +// it. A render worker reporting success proves it did not crash; it does not +// prove the file has 150 frames, that they are 30 fps, that the audio ends +// where the video does, or that the thing is under its byte budget. Those are +// facts about bytes on disk, so they are measured from bytes on disk. + +import { statSync } from "node:fs"; +import { + AUDIO_TOLERANCE_MS, + PREROLL_FPS, + PREROLL_FRAMES, + PREROLL_MS, + videoProfile, + withinBudget, + type VideoProfileId, +} from "./profiles"; +import { FFPROBE_BIN, run } from "./encode"; + +export type ProbeStream = { + codec_type?: string; + codec_name?: string; + profile?: string; + pix_fmt?: string; + width?: number; + height?: number; + /** "30/1" — a rational, because 29.97 must not silently pass as 30. */ + r_frame_rate?: string; + avg_frame_rate?: string; + /** Present only with -count_frames; it is a decode, not a header read. */ + nb_read_frames?: string; + duration?: string; + sample_rate?: string; + channels?: number; +}; + +export type ProbeResult = { + streams: ProbeStream[]; + format?: { duration?: string; format_name?: string; size?: string }; +}; + +/** + * -count_frames is the point of this invocation. + * + * Without it ffprobe reports `nb_frames` from the container header, which is + * whatever the muxer wrote down — including on a file whose frames were + * truncated. Counting means decoding every frame, which is slower and is the + * only way the number means anything. + */ +export function probeArgs(filePath: string): string[] { + return [ + "-v", + "error", + "-count_frames", + "-show_streams", + "-show_format", + "-of", + "json", + filePath, + ]; +} + +export async function probeMedia(filePath: string): Promise { + const out = await run(FFPROBE_BIN, probeArgs(filePath), 180_000); + return JSON.parse(out) as ProbeResult; +} + +/** "30/1" → 30. Returns null for a malformed or zero-denominator rational. */ +export function parseRational(r: string | undefined): number | null { + if (!r) return null; + const [n, d] = r.split("/"); + const num = Number(n); + const den = d === undefined ? 1 : Number(d); + if (!Number.isFinite(num) || !Number.isFinite(den) || den === 0) return null; + return num / den; +} + +export type ValidationProblem = { check: string; expected: string; actual: string }; + +export type ValidationResult = { + ok: boolean; + problems: ValidationProblem[]; + measured: { + frames: number | null; + fps: number | null; + width: number | null; + height: number | null; + videoCodec: string | null; + pixFmt: string | null; + audioCodec: string | null; + audioDurationMs: number | null; + byteSize: number | null; + }; +}; + +/** + * Evaluate a probe against a profile's requirements. + * + * Pure: takes the probe and the byte size, returns problems. Kept separate from + * probeMedia so the awkward cases — 29.97 fps, a 149-frame encode, an audio + * track two frames long — can be tested as data instead of by manufacturing a + * broken file for each one. + */ +export function evaluateProbe( + probe: ProbeResult, + profileId: VideoProfileId, + byteSize: number | null, + opts: { expectAudio: boolean } = { expectAudio: false }, +): ValidationResult { + const problems: ValidationProblem[] = []; + const profile = videoProfile(profileId); + const video = probe.streams.find((s) => s.codec_type === "video"); + const audio = probe.streams.find((s) => s.codec_type === "audio"); + + const frames = video?.nb_read_frames ? Number(video.nb_read_frames) : null; + const fps = parseRational(video?.r_frame_rate); + const audioDurationMs = audio?.duration ? Number(audio.duration) * 1000 : null; + + const measured = { + frames, + fps, + width: video?.width ?? null, + height: video?.height ?? null, + videoCodec: video?.codec_name ?? null, + pixFmt: video?.pix_fmt ?? null, + audioCodec: audio?.codec_name ?? null, + audioDurationMs, + byteSize, + }; + + const fail = (check: string, expected: string, actual: unknown) => + problems.push({ check, expected, actual: String(actual) }); + + if (!video) { + fail("video stream", "present", "absent"); + return { ok: false, problems, measured }; + } + + // The duration requirement, stated the only way it can be checked. + if (frames !== PREROLL_FRAMES) fail("frame count", String(PREROLL_FRAMES), frames); + + // Exactly 30, not 29.97. A 30000/1001 encode drifts ~5ms across five seconds, + // which is small — but it means the frame at index 149 is not at 4966.67ms, + // and the quartile boundaries measured against played media time stop lining + // up with the frames they name. + if (fps !== PREROLL_FPS) fail("frame rate", `${PREROLL_FPS}/1`, video.r_frame_rate ?? fps); + + if (video.codec_name !== "h264") fail("video codec", "h264", video.codec_name); + if (video.pix_fmt !== "yuv420p") fail("pixel format", "yuv420p (8-bit)", video.pix_fmt); + + if (profile.width && video.width !== profile.width) { + fail("width", String(profile.width), video.width); + } + if (profile.height && video.height !== profile.height) { + fail("height", String(profile.height), video.height); + } + + if (opts.expectAudio) { + if (!audio) { + fail("audio stream", "present", "absent"); + } else { + if (audio.codec_name !== "aac") fail("audio codec", "aac", audio.codec_name); + if (Number(audio.sample_rate) !== 48_000) fail("audio sample rate", "48000", audio.sample_rate); + if (audioDurationMs !== null) { + // One AAC frame of slack, no more. AAC cannot land on an arbitrary + // boundary, so some rounding here is correct rather than a defect — + // but two frames of it means the track is genuinely the wrong length. + const drift = Math.abs(audioDurationMs - PREROLL_MS); + if (drift > AUDIO_TOLERANCE_MS) { + fail( + "audio duration", + `${PREROLL_MS}ms ±${AUDIO_TOLERANCE_MS.toFixed(3)}ms (one AAC frame)`, + `${audioDurationMs.toFixed(3)}ms (drift ${drift.toFixed(3)}ms)`, + ); + } + } + } + } else if (audio) { + // A silent profile that somehow carries a track is not a harmless extra: + // it is how an ad ends up making noise over a stream nobody expected it to. + fail("audio stream", "absent", audio.codec_name ?? "present"); + } + + if (byteSize !== null && !withinBudget(profileId, byteSize)) { + fail("byte size", `<= ${profile.maxBytes} bytes`, `${byteSize} bytes`); + } + + return { ok: problems.length === 0, problems, measured }; +} + +/** Probe a file on disk and evaluate it. */ +export async function validateRendition( + filePath: string, + profileId: VideoProfileId, + opts: { expectAudio: boolean } = { expectAudio: false }, +): Promise { + const probe = await probeMedia(filePath); + let byteSize: number | null = null; + try { + byteSize = statSync(filePath).size; + } catch { + byteSize = null; + } + return evaluateProbe(probe, profileId, byteSize, opts); +} + +/** + * Structural checks on a generated HLS media playlist. + * + * Text, not media — this asks whether the packaging says what we require, which + * is a different question from whether the segments decode. Both are checked; + * this is the cheap half. + */ +export function validateMediaPlaylist(playlist: string): ValidationProblem[] { + const problems: ValidationProblem[] = []; + const has = (tag: string) => playlist.includes(tag); + + if (!has("#EXTM3U")) problems.push({ check: "playlist", expected: "#EXTM3U", actual: "absent" }); + // A finite VOD package. Without ENDLIST a player treats it as a live window + // and keeps reloading a playlist that will never change. + if (!has("#EXT-X-ENDLIST")) { + problems.push({ check: "playlist", expected: "#EXT-X-ENDLIST", actual: "absent" }); + } + if (!has("#EXT-X-MAP")) { + problems.push({ check: "playlist", expected: "#EXT-X-MAP (fMP4 init)", actual: "absent" }); + } + if (has("#EXT-X-INDEPENDENT-SEGMENTS")) { + problems.push({ + check: "playlist", + expected: "no EXT-X-INDEPENDENT-SEGMENTS (we do not guarantee it)", + actual: "present", + }); + } + + const durations = [...playlist.matchAll(/#EXTINF:([\d.]+)/g)].map((m) => Number(m[1])); + if (durations.length === 0) { + problems.push({ check: "segments", expected: ">= 1", actual: "0" }); + } else { + const total = durations.reduce((a, b) => a + b, 0) * 1000; + // Generous: segment durations are written to 6dp from a timebase, and the + // sum of three of them accumulates rounding. Half a frame is plenty tight + // to catch a missing or duplicated segment, which is what this is for. + const tolerance = 1000 / PREROLL_FPS / 2; + if (Math.abs(total - PREROLL_MS) > tolerance) { + problems.push({ + check: "segment total duration", + expected: `${PREROLL_MS}ms ±${tolerance.toFixed(2)}ms`, + actual: `${total.toFixed(3)}ms`, + }); + } + } + + return problems; +} diff --git a/lib/prober-queue.ts b/lib/prober-queue.ts index 11686060..49fc1381 100644 --- a/lib/prober-queue.ts +++ b/lib/prober-queue.ts @@ -6,9 +6,10 @@ // The droplet writes no DB — it returns the result as the job's return value, // which we read here and persist (scan row + findings). Keeps Supabase creds // off the droplet. -import { Queue, type ConnectionOptions, type Job } from "bullmq"; +import { Queue, type Job } from "bullmq"; import type { SupabaseClient } from "@supabase/supabase-js"; import { PROBER_QUEUE, type PortScanResult } from "./prober"; +import { redisConnectionOptions } from "./redis-connection"; // Ports that are alarming when publicly exposed (databases, caches, admin, etc.). const HIGH_RISK_PORTS = new Set([ @@ -23,24 +24,8 @@ const RUNNING_TIMEOUT_MS = 40 * 60 * 1000; let queue: Queue | null = null; -// Parse REDIS_URL into bullmq connection options so bullmq builds its own -// ioredis client (avoids a version clash between our ioredis and bullmq's). -function connectionOptions(): ConnectionOptions | null { - const url = process.env.REDIS_URL; - if (!url) return null; - const u = new URL(url); - return { - host: u.hostname, - port: Number(u.port || "6379"), - username: u.username ? decodeURIComponent(u.username) : undefined, - password: u.password ? decodeURIComponent(u.password) : undefined, - tls: u.protocol === "rediss:" ? {} : undefined, - maxRetriesPerRequest: null, - }; -} - export function getProberQueue(): Queue | null { - const connection = connectionOptions(); + const connection = redisConnectionOptions(); if (!connection) return null; if (!queue) queue = new Queue(PROBER_QUEUE, { connection }); return queue; diff --git a/lib/redis-connection.ts b/lib/redis-connection.ts new file mode 100644 index 00000000..a4e6c6f4 --- /dev/null +++ b/lib/redis-connection.ts @@ -0,0 +1,22 @@ +// REDIS_URL -> bullmq connection options. +// +// Lifted out of ./prober-queue when the video render queue became the second +// caller. Kept as options rather than a shared ioredis instance on purpose: +// bullmq then builds its own client, which avoids a version clash between the +// ioredis we depend on and the one bullmq bundles. + +import type { ConnectionOptions } from "bullmq"; + +export function redisConnectionOptions(): ConnectionOptions | null { + const url = process.env.REDIS_URL; + if (!url) return null; + const u = new URL(url); + return { + host: u.hostname, + port: Number(u.port || "6379"), + username: u.username ? decodeURIComponent(u.username) : undefined, + password: u.password ? decodeURIComponent(u.password) : undefined, + tls: u.protocol === "rediss:" ? {} : undefined, + maxRetriesPerRequest: null, + }; +} diff --git a/tests/ads-video-compose-validate.test.ts b/tests/ads-video-compose-validate.test.ts new file mode 100644 index 00000000..c5076328 --- /dev/null +++ b/tests/ads-video-compose-validate.test.ts @@ -0,0 +1,340 @@ +import { describe, expect, it } from "vitest"; +import { + composeDocument, + escapeHtml, + frameState, + safeDataUri, + timeline, +} from "@/lib/ads/video/compose"; +import type { VideoDesignSnapshot } from "@/lib/ads/video/snapshot"; +import { GOP_FRAMES, hlsArgs, mp4Args, multivariantPlaylist, posterArgs } from "@/lib/ads/video/encode"; +import { + evaluateProbe, + parseRational, + probeArgs, + validateMediaPlaylist, + type ProbeResult, +} from "@/lib/ads/video/validate"; +import { PREROLL_FRAMES, TIMELINE } from "@/lib/ads/video/profiles"; + +const snapshot: VideoDesignSnapshot = { + headline: "Sources in, feeds out", + ctaText: "Start free", + domain: "nichedb.dev", + bgColor: "#12161f", + fgColor: "#e7e9ee", + accentColor: "#6ee7b7", + fontFamily: "system-ui, sans-serif", + logoUrl: null, + logoSha256: null, + heroUrl: null, + heroSha256: null, + audioMode: "silent", + narration: null, + locale: "en", + reducedMotion: false, +}; + +describe("the composition is deterministic", () => { + it("is a pure function of the frame index", () => { + expect(frameState(73, false)).toEqual(frameState(73, false)); + expect(timeline(false)).toHaveLength(PREROLL_FRAMES); + }); + + it("has the headline legible from the very first frame", () => { + // The entrance animates around copy that is already readable. A five-second + // ad that spends its first beat assembling itself says nothing for a tenth + // of its life. + const f0 = frameState(0, false); + expect(f0.entrance).toBe(0); + // ...and the document's opacity floor for the copy block is 0.55, not 0. + expect(composeDocument(snapshot)).toContain("0.55 + 0.45 * s.entrance"); + }); + + it("runs its three beats in order and finishes settled", () => { + // Frame boundaries do not land on the beat boundaries (350ms is frame + // 10.5), so assert on the frames either side rather than on a rounded one. + const lastEntranceFrame = Math.floor((TIMELINE.entranceEndMs / 1000) * 30); // 10 + const beforeHold = frameState(lastEntranceFrame, false); + expect(beforeHold.entrance).toBeLessThan(1); + expect(beforeHold.drift).toBe(0); + + const afterEntrance = frameState(lastEntranceFrame + 1, false); + expect(afterEntrance.entrance).toBe(1); + expect(afterEntrance.drift).toBeGreaterThan(0); + + // The CTA beat has not started before the hold ends. + const lastHoldFrame = Math.floor((TIMELINE.holdEndMs / 1000) * 30); // 105 + expect(frameState(lastHoldFrame, false).cta).toBe(0); + expect(frameState(lastHoldFrame + 1, false).cta).toBeGreaterThan(0); + + // Each beat is monotonic across the whole timeline. + const frames = timeline(false); + for (let i = 1; i < frames.length; i++) { + expect(frames[i].entrance).toBeGreaterThanOrEqual(frames[i - 1].entrance); + expect(frames[i].drift).toBeGreaterThanOrEqual(frames[i - 1].drift); + expect(frames[i].cta).toBeGreaterThanOrEqual(frames[i - 1].cta); + } + + expect(frames[PREROLL_FRAMES - 1].cta).toBeGreaterThan(0.9); + }); + + it("holds everything still for a reduced-motion viewer, on the same timeline", () => { + for (const frame of [0, 40, 149]) { + const s = frameState(frame, true); + expect(s.entrance).toBe(1); + expect(s.drift).toBe(0); + expect(s.cta).toBe(1); + } + // Same number of frames — it is a static composition over five seconds, not + // a shorter ad. + expect(timeline(true)).toHaveLength(PREROLL_FRAMES); + }); +}); + +describe("the composition is closed against its own input", () => { + it("escapes advertiser copy", () => { + const evil = { + ...snapshot, + headline: ``, + ctaText: `" onerror="alert(1)`, + }; + const html = composeDocument(evil); + expect(html).not.toContain("