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)); });