diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8b34db2..b5f21a5 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -23,6 +23,17 @@ jobs: - name: Install (worker) run: cd worker && npm install --no-audit --no-fund + # The five-second ad renderer encodes with ffmpeg and validates with + # ffprobe. tests/ads-video-pipeline.test.ts skips itself when they are + # absent, so without this step the one test that proves a pre-roll + # actually decodes to 150 frames would quietly never run here — and a + # green check would say nothing about it. + - name: Install ffmpeg + run: sudo apt-get update && sudo apt-get install -y --no-install-recommends ffmpeg + + - name: Verify ffmpeg is present for the render tests + run: ffmpeg -version && ffprobe -version + - name: Typecheck (root) run: npm run typecheck diff --git a/lib/ads/video/compose.ts b/lib/ads/video/compose.ts new file mode 100644 index 0000000..9524965 --- /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 0000000..431056f --- /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 0000000..0bf64ff --- /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 0000000..3ccb461 --- /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 0000000..a5f39e8 --- /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 0000000..e35657a --- /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 0000000..6f3f3e2 --- /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 0000000..0728524 --- /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 1168606..49fc138 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 0000000..a4e6c6f --- /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 0000000..c507632 --- /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("