From 138d9dbaa28d33bdf3d13756a4e7b17daa49c29c Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Thu, 24 Sep 2026 13:48:36 +0000 Subject: [PATCH] Backfill pre-roll videos for existing product campaigns MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds the sweep the rest of this package kept promising, a classifier for what counts as a product ad, and a backfill script that queues a render for each one. The sweep first, because without it the backfill would not work and several comments in this package were writing cheques nothing cashed. ad_video_jobs is the durable record and the BullMQ job is the transient work item, which is the right way round, but it leaves a gap: if Redis is down at save time, or the process dies between the insert and the enqueue, the row sits `queued` with nothing scheduled to look at it. processDueVideoRenders reschedules those after a minute's grace, and fails a row outright if it carries no design snapshot so a broken row cannot become a hot loop against Redis. It also means a backfill can simply insert rows from wherever it runs and let the worker schedule them. The classifier exists because "generate videos for the ads" is not the same set as "every campaign row". Of 450 campaigns, 269 point at blog posts and 4 at social links; only 177 advertise a product. Getting that wrong is expensive in one direction specifically — rendering videos for the posts that were excluded — so the shapes are taken from the live table rather than imagined. Most blog campaigns have no /blog/ path prefix at all: dev.profullstack.com/~anthony/blog/126-post.html 152 campaigns, ~user dir dev.to/chovy/ 117 campaigns, /user/slug A naive /^\/blog/ match catches neither, which is why the publishing platforms are matched on the domain. Affiliate and referral links are kept as products: they advertise something a person can buy, which is the distinction that matters. An unparseable destination is deliberately not classified as a product — something that cannot be shown to be one should not get a render. The script calls the same queueCampaignVideo the dashboard calls, so backfilled videos cannot differ from freshly saved ones with nothing to explain why. It is re-runnable: jobs dedupe on a hash of the design, so a campaign that already has a job for its current copy is skipped rather than re-encoded. It passes bumpRevision false, because a backfill is not an edit and bumping would invalidate media an earlier pass already rendered for the same unchanged design. It never rewrites campaign copy and never changes campaign status. Dry run against production: 450 campaigns, 177 product, 269 blog, 4 social, 177 queueable with no campaign lacking a usable creative. Co-Authored-By: Claude Opus 5 (1M context) --- lib/ads/video/classify.ts | 50 +++++++ lib/ads/video/sweep.ts | 99 +++++++++++++ scripts/backfill-ad-videos.ts | 171 ++++++++++++++++++++++ tests/ads-video-backfill-classify.test.ts | 101 +++++++++++++ worker/index.ts | 16 ++ 5 files changed, 437 insertions(+) create mode 100644 lib/ads/video/classify.ts create mode 100644 lib/ads/video/sweep.ts create mode 100644 scripts/backfill-ad-videos.ts create mode 100644 tests/ads-video-backfill-classify.test.ts diff --git a/lib/ads/video/classify.ts b/lib/ads/video/classify.ts new file mode 100644 index 0000000..c517f05 --- /dev/null +++ b/lib/ads/video/classify.ts @@ -0,0 +1,50 @@ +// What kind of thing a campaign advertises. +// +// ad_campaigns records no kind, so this is derived from the destination. It +// exists because "generate videos for the ads" does not mean every campaign +// row: a large share of them point at blog posts, and rendering a pre-roll for +// an article nobody asked to promote that way is the expensive mistake. +// +// The shapes below are taken from the live table rather than imagined. In +// particular, most blog campaigns do NOT have a /blog/ path prefix: +// +// dev.profullstack.com/~anthony/blog/126-post.html (under a ~user dir) +// dev.to/chovy/ (path is /user/slug) +// +// A naive /^\/blog/ match catches neither, which is why the domain families +// are matched explicitly. + +const SOCIAL_DOMAIN = + /^(twitter\.com|x\.com|.*mastodon.*|bsky\.app|threads\.net|linkedin\.com|facebook\.com|instagram\.com|reddit\.com|youtube\.com|youtu\.be|t\.me|discord\.(gg|com)|tiktok\.com)$/i; + +const BLOG_PLATFORM = + /^(dev\.to|medium\.com|.*\.medium\.com|.*\.hashnode\.dev|.*\.substack\.com|hackernoon\.com|lobste\.rs)$/i; + +// Path segments, so /newsletter-tool stays a product while /news/thing does not. +const BLOG_PATH = /(^|\/)(blog|posts?|article|articles|news)(\/|$)|-post\.html$/i; + +export type CampaignKind = "product" | "blog" | "social"; + +/** + * Classify a campaign by where it points. + * + * Affiliate and referral links count as product: they advertise something a + * person can buy, which is the distinction that matters. An unparseable + * destination is deliberately not classified as a product — something we cannot + * show to be one should not get a render. + */ +export function classifyCampaign(destinationUrl: string): CampaignKind { + let host = ""; + let path = "/"; + try { + const u = new URL(destinationUrl); + host = u.hostname.replace(/^www\./, "").toLowerCase(); + path = u.pathname || "/"; + } catch { + return "blog"; + } + if (SOCIAL_DOMAIN.test(host)) return "social"; + if (BLOG_PLATFORM.test(host)) return "blog"; + if (BLOG_PATH.test(path)) return "blog"; + return "product"; +} diff --git a/lib/ads/video/sweep.ts b/lib/ads/video/sweep.ts new file mode 100644 index 0000000..d7b362a --- /dev/null +++ b/lib/ads/video/sweep.ts @@ -0,0 +1,99 @@ +// Drain render jobs whose BullMQ job never materialised. +// +// ad_video_jobs is the durable record and the BullMQ job is the transient work +// item, which is the right way round — but it leaves a gap. If Redis is down +// when a campaign is saved, or the process dies between the insert and the +// enqueue, the row sits `queued` with nothing scheduled to look at it. Several +// comments in this package promise "a sweep can pick it up"; this is that +// sweep, and without it those rows are simply lost. +// +// It also makes a backfill trivial: a script can insert rows and let the sweep +// schedule them, instead of needing to reach Redis itself from wherever it runs. + +import type { SupabaseClient } from "@supabase/supabase-js"; +import { enqueueRender, getVideoRenderQueue, renderJobId } from "./queue"; +import type { VideoDesignSnapshot } from "./snapshot"; +import type { VideoProfileId } from "./profiles"; + +/** + * Grace period before a queued row is considered stranded. + * + * Long enough that the ordinary path — insert, then enqueue milliseconds later + * — is never second-guessed by the sweep racing it. A row that is genuinely + * only a few seconds old is almost certainly mid-save. + */ +const STRANDED_AFTER_MS = 60_000; + +/** Rows to schedule per pass. Keeps a backfill from flooding the queue at once. */ +const BATCH = Number(process.env.VIDEO_SWEEP_BATCH ?? "25"); + +export async function processDueVideoRenders( + supabase: SupabaseClient, +): Promise<{ scheduled: number; skipped: number }> { + const queue = getVideoRenderQueue(); + if (!queue) return { scheduled: 0, skipped: 0 }; + + const cutoff = new Date(Date.now() - STRANDED_AFTER_MS).toISOString(); + const { data: rows } = await supabase + .from("ad_video_jobs") + .select("id, owner_id, campaign_id, creative_id, revision, render_hash, output_profile, design, attempts") + .eq("state", "queued") + .lt("created_at", cutoff) + .order("created_at", { ascending: true }) + .limit(BATCH); + + let scheduled = 0; + let skipped = 0; + + for (const row of rows ?? []) { + const hash = row.render_hash as string; + const profile = (row.output_profile as string) ?? "default"; + + // Already scheduled: re-adding is harmless because BullMQ dedupes on the + // job id, but checking first keeps the log honest about what this did. + const existing = await queue.getJob(renderJobId(hash, profile)).catch(() => null); + if (existing) { + skipped++; + continue; + } + + // A row with no design cannot be rendered and re-queueing it forever would + // be a hot loop against Redis. Fail it so it stops being swept and shows up + // in the dashboard as something that needs attention. + const snapshot = row.design as VideoDesignSnapshot | null; + if (!snapshot || !snapshot.headline) { + await supabase + .from("ad_video_jobs") + .update({ state: "failed", error_code: "no_design_snapshot" }) + .eq("id", row.id); + skipped++; + continue; + } + + try { + const ok = await enqueueRender({ + renderHash: hash, + profile, + data: { + jobRowId: row.id as string, + ownerId: row.owner_id as string, + campaignId: (row.campaign_id as string | null) ?? null, + creativeId: (row.creative_id as string | null) ?? null, + revision: row.revision as number, + snapshot, + profile: profile as VideoProfileId, + audioSlotSupported: false, + }, + }); + if (ok) scheduled++; + else skipped++; + } catch (err) { + // Leave the row `queued`. The next pass tries again — which is the entire + // point of the row outliving the queue. + console.warn(`[worker] video sweep could not schedule ${row.id}: ${(err as Error).message}`); + skipped++; + } + } + + return { scheduled, skipped }; +} diff --git a/scripts/backfill-ad-videos.ts b/scripts/backfill-ad-videos.ts new file mode 100644 index 0000000..dd816b1 --- /dev/null +++ b/scripts/backfill-ad-videos.ts @@ -0,0 +1,171 @@ +// Queue a five-second pre-roll render for existing PRODUCT campaigns. +// +// npx tsx scripts/backfill-ad-videos.ts --env ~/crawlproof-env-backup-2026-07-28.txt --dry-run +// npx tsx scripts/backfill-ad-videos.ts --env ~/crawlproof-env-backup-2026-07-28.txt --limit 3 +// npx tsx scripts/backfill-ad-videos.ts --env ~/crawlproof-env-backup-2026-07-28.txt +// +// TypeScript through tsx, like backfill-ad-summaries, so it can call the same +// queueCampaignVideo the dashboard calls. Reimplementing the snapshot here +// would mean backfilled videos differing from freshly saved ones with nothing +// to explain why. +// +// Safe to re-run. Render jobs dedupe on a hash of the design, so a campaign +// that already has a job for its current copy is skipped rather than +// re-encoded, and an interrupted pass continues where it stopped. +// +// It never rewrites campaign copy and never changes campaign status: the only +// writes are the video creative row and its render job. + +import { readFileSync } from "node:fs"; +import { createClient } from "@supabase/supabase-js"; + +const args = process.argv.slice(2); +const flag = (name: string): string | null => { + const i = args.indexOf(`--${name}`); + return i >= 0 ? (args[i + 1] ?? "") : null; +}; +const has = (name: string) => args.includes(`--${name}`); + +const envPath = flag("env") ?? `${process.env.HOME}/crawlproof-env-backup-2026-07-28.txt`; +const dryRun = has("dry-run"); +const limit = Number(flag("limit") ?? "0") || 0; +/** Pause between queued campaigns, so a backfill does not spike the queue. */ +const DELAY_MS = Number(flag("delay") ?? "150") || 150; + +for (const [k, v] of Object.entries(readEnvFile(envPath))) { + if (!process.env[k]) process.env[k] = v; +} + +const { queueCampaignVideo } = await import("../lib/ads/video/jobs"); +const { classifyCampaign } = await import("../lib/ads/video/classify"); +type CampaignKind = "product" | "blog" | "social"; + +function readEnvFile(path: string): Record { + const out: Record = {}; + let text = ""; + try { + text = readFileSync(path, "utf8"); + } catch { + return out; + } + for (const line of text.split(/\r?\n/)) { + const m = line.match(/^\s*(?:export\s+)?([A-Z0-9_]+)\s*=\s*(.*)\s*$/); + if (!m) continue; + out[m[1]] = m[2].replace(/^["']|["']$/g, ""); + } + return out; +} + +const url = process.env.NEXT_PUBLIC_SUPABASE_URL; +const key = process.env.SUPABASE_SERVICE_ROLE_KEY; +if (!url || !key) { + console.error("Missing NEXT_PUBLIC_SUPABASE_URL / SUPABASE_SERVICE_ROLE_KEY."); + process.exit(1); +} + +const supabase = createClient(url, key, { + auth: { autoRefreshToken: false, persistSession: false }, +}); + +type CampaignRow = { + id: string; + owner_id: string; + name: string; + status: string; + destination_url: string; + destination_domain: string | null; +}; + +const { data: campaigns, error } = await supabase + .from("ad_campaigns") + .select("id, owner_id, name, status, destination_url, destination_domain") + // A rejected campaign must not gain new servable media. + .neq("status", "rejected") + .order("created_at", { ascending: true }); + +if (error) { + console.error("Could not read campaigns:", error.message); + process.exit(1); +} + +const all = (campaigns ?? []) as CampaignRow[]; +const tally: Record = { product: 0, blog: 0, social: 0 }; +const products: CampaignRow[] = []; +for (const c of all) { + const kind = classifyCampaign(c.destination_url); + tally[kind]++; + if (kind === "product") products.push(c); +} + +console.log( + `${all.length} campaigns: ${tally.product} product, ${tally.blog} blog, ${tally.social} social.`, +); + +const targets = limit > 0 ? products.slice(0, limit) : products; +console.log(`${targets.length} to queue${dryRun ? " (dry run — nothing written)" : ""}.\n`); + +let queued = 0; +let reused = 0; +let skipped = 0; +let failed = 0; + +for (const [i, c] of targets.entries()) { + const label = `${i + 1}/${targets.length} ${c.destination_url.slice(0, 70)}`; + + const { data: creatives } = await supabase + .from("ad_creatives") + .select("format, headline, cta_text, bg_color, fg_color, accent_color, font_family, logo_url, image_url") + .eq("campaign_id", c.id) + .neq("format", "video_preroll_5s") + .neq("status", "rejected"); + + const usable = (creatives ?? []).filter((r) => (r.headline ?? "").trim().length > 0); + if (usable.length === 0) { + console.log(` skip ${label} — no usable creative`); + skipped++; + continue; + } + + if (dryRun) { + console.log(` would ${label}`); + queued++; + continue; + } + + const handle = await queueCampaignVideo(supabase, { + campaignId: c.id, + ownerId: c.owner_id, + domain: c.destination_domain ?? new URL(c.destination_url).hostname.replace(/^www\./, ""), + creatives: usable.map((r) => ({ + format: r.format, + headline: r.headline ?? "", + ctaText: r.cta_text ?? "", + bgColor: r.bg_color, + fgColor: r.fg_color, + accentColor: r.accent_color, + fontFamily: r.font_family, + logoUrl: r.logo_url, + imageUrl: r.image_url, + })), + // A backfill is not an edit. Bumping the revision would invalidate media + // that a previous pass already rendered for the same unchanged design. + bumpRevision: false, + }); + + if (!handle) { + console.log(` FAIL ${label}`); + failed++; + } else if (handle.reused) { + console.log(` have ${label} — job ${handle.jobId} (${handle.state})`); + reused++; + } else { + console.log(` queue ${label} — job ${handle.jobId}${handle.enqueued ? "" : " (row only; sweep will schedule)"}`); + queued++; + } + + if (DELAY_MS > 0) await new Promise((r) => setTimeout(r, DELAY_MS)); +} + +console.log( + `\nDone. queued=${queued} already-had=${reused} skipped=${skipped} failed=${failed}`, +); diff --git a/tests/ads-video-backfill-classify.test.ts b/tests/ads-video-backfill-classify.test.ts new file mode 100644 index 0000000..5d31b64 --- /dev/null +++ b/tests/ads-video-backfill-classify.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it } from "vitest"; +import { classifyCampaign } from "@/lib/ads/video/classify"; + +// Every URL below is a real shape taken from the production ad_campaigns table. +// The classifier decides what gets a rendered video, and the expensive mistake +// is the false positive: rendering a blog post that was explicitly excluded. +describe("blog posts are excluded, however they are shaped", () => { + it("catches the own-site blog, which has no /blog/ prefix at the root", () => { + // 152 campaigns look like this. A naive /^\/blog/ match misses every one, + // because the blog lives under a tilde user directory. + expect(classifyCampaign("https://dev.profullstack.com/~anthony/blog/126-post.html")).toBe("blog"); + expect(classifyCampaign("https://dev.profullstack.com/~anthony/blog/048-post.html")).toBe("blog"); + }); + + it("catches dev.to articles, which look like ordinary deep links", () => { + // 117 campaigns. The path is // with nothing blog-shaped in it, + // so this has to be caught on the domain. + expect(classifyCampaign("https://dev.to/chovy/nice-is-half-a-fix-ep9")).toBe("blog"); + expect(classifyCampaign("https://dev.to/chovy/the-bots-now-pay-the-humans-o43")).toBe("blog"); + }); + + it("catches the other publishing platforms", () => { + for (const u of [ + "https://medium.com/@someone/a-post-abc123", + "https://someone.substack.com/p/a-post", + "https://someone.hashnode.dev/a-post", + ]) { + expect(classifyCampaign(u), u).toBe("blog"); + } + }); + + it("catches conventional blog paths", () => { + for (const u of [ + "https://example.com/blog/thing", + "https://example.com/posts/thing", + "https://example.com/news/thing", + "https://example.com/article/thing", + ]) { + expect(classifyCampaign(u), u).toBe("blog"); + } + }); +}); + +describe("social links are excluded", () => { + it("catches the platforms", () => { + for (const u of [ + "https://x.com/someone", + "https://twitter.com/someone/status/1", + "https://bsky.app/profile/someone", + "https://www.linkedin.com/in/someone", + "https://t.me/somechannel", + "https://www.reddit.com/r/something", + "https://youtube.com/watch?v=abc", + ]) { + expect(classifyCampaign(u), u).toBe("social"); + } + }); +}); + +describe("product ads are kept", () => { + it("keeps product homepages and deep links", () => { + for (const u of [ + "https://moshcoding.com/", + "https://outreachgraph.com/", + "https://bl0ggers.com/", + "https://weedforcrypto.com/", + "https://profullstack.com/pricing", + "https://openmcp.logicsrc.com/", + "https://tronbrowser.dev/store/extension.html?slug=coinpay-wallet", + "https://c0upons.com/coupons/647", + ]) { + expect(classifyCampaign(u), u).toBe("product"); + } + }); + + it("keeps affiliate and referral links — they advertise something buyable", () => { + for (const u of [ + "https://m.do.co/c/f8b2890dc9d1", + "https://www.amazon.com/Red-Bull-Energy-Drink-Pack/dp/B006O3ASKE?th=1", + "https://goldclubhosting.xyz/aff.php?aff=362", + "https://aiornot.vote/r/BAB9HME", + "https://www.netcup.com/de/", + ]) { + expect(classifyCampaign(u), u).toBe("product"); + } + }); + + it("does not mistake a product path containing 'news' inside a word", () => { + // Word-boundaried on path segments, so /newsletter-tool is a product. + expect(classifyCampaign("https://example.com/newsletter-tool")).toBe("product"); + }); +}); + +describe("an unparseable destination is skipped rather than rendered", () => { + it("treats junk as non-product", () => { + // Conservative direction: something we cannot show to be a product does not + // get a render. + expect(classifyCampaign("not a url")).not.toBe("product"); + expect(classifyCampaign("")).not.toBe("product"); + }); +}); diff --git a/worker/index.ts b/worker/index.ts index 56f701a..1f7deb4 100644 --- a/worker/index.ts +++ b/worker/index.ts @@ -49,6 +49,7 @@ import { ingestDueFeeds } from "../lib/promote/ingest"; import { refreshCookieSessions } from "../lib/sp/sessionRefresh"; import { runAutobidSweep } from "../lib/ads/bids"; import { startVideoRenderWorker } from "./video"; +import { processDueVideoRenders } from "../lib/ads/video/sweep"; const supabaseUrl = process.env.NEXT_PUBLIC_SUPABASE_URL!; const supabaseKey = process.env.SUPABASE_SERVICE_ROLE_KEY!; @@ -1559,6 +1560,20 @@ setInterval( AUTOBID_TICK_MS, ); +// Reschedule render jobs whose BullMQ job never materialised — Redis down at +// save time, or a crash between the row insert and the enqueue. Without this +// those rows sit `queued` forever, and a backfill that only inserts rows would +// never render anything. +const VIDEO_SWEEP_TICK_MS = 60_000; +async function videoSweep() { + const r = await processDueVideoRenders(supabase); + if (r.scheduled > 0) console.log(`[worker] video sweep scheduled=${r.scheduled} skipped=${r.skipped}`); +} +setInterval( + () => videoSweep().catch((e) => console.error("[worker] video sweep", e)), + VIDEO_SWEEP_TICK_MS, +); + // Bind to loopback by default so the worker isn't reachable from the public // internet when colocated with the app. Override with WORKER_BIND=0.0.0.0 to // run as a separate Railway service. @@ -1576,4 +1591,5 @@ server.listen(port, bindHost, () => { // Long-lived BullMQ consumer rather than an interval sweep: a render takes // tens of seconds and the queue, not a timer, decides when the next one runs. startVideoRenderWorker(supabase); + videoSweep().catch((e) => console.error("[worker] video sweep", e)); });