Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 50 additions & 0 deletions lib/ads/video/classify.ts
Original file line number Diff line number Diff line change
@@ -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/<slug> (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";
}
99 changes: 99 additions & 0 deletions lib/ads/video/sweep.ts
Original file line number Diff line number Diff line change
@@ -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 };
}
171 changes: 171 additions & 0 deletions scripts/backfill-ad-videos.ts
Original file line number Diff line number Diff line change
@@ -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<string, string> {
const out: Record<string, string> = {};
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<CampaignKind, number> = { 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}`,
);
Loading