diff --git a/lib/ads/fraud.ts b/lib/ads/fraud.ts index 3b0e1cd..a8077ba 100644 --- a/lib/ads/fraud.ts +++ b/lib/ads/fraud.ts @@ -21,6 +21,22 @@ function safeId(v: string | null | undefined): string | null { return /^[\w-]{1,128}$/.test(s) ? s : null; } +// Impressions get a much shorter window than clicks, and are keyed on the slot +// rather than the campaign. +// +// Both choices come from what the inflation actually looks like. A scheduled +// pool refresher fetches one slot N times back to back to fill N cache entries; +// each fetch picks a *different* campaign at random, so campaign-keyed dedupe +// (the shape used for clicks) collapses none of it. Slot-keyed dedupe collapses +// the whole burst to one counted impression. +// +// 60s rather than the click path's 6h because a repeat view is a real thing: a +// human reloading an article an hour later genuinely saw the ad twice, and an +// hours-long window would erase legitimate delivery. 60s is long enough to +// swallow a machine-driven burst (the observed one fires 12 fetches in ~3s) and +// short enough that no human pattern lands inside it twice by accident. +export const IMPRESSION_DEDUPE_WINDOW_MS = 60 * 1000; + export type ClickValidity = { valid: boolean; reason?: string }; export async function assessClickValidity(input: { @@ -83,3 +99,58 @@ export async function assessClickValidity(input: { return { valid: true }; } + +/** + * Has this viewer already been counted on this slot inside the dedupe window? + * + * Deliberately does NOT stop the ad being served or the row being written — the + * caller still inserts an impression, just flagged. Two reasons the row has to + * exist either way: + * + * 1. Click attribution. A terminal click URL is /a/, which + * resolves the campaign and creative back off the impression row. Skipping + * the insert would serve a real advertiser's creative with a click link + * that resolves to nothing — the click would go unbilled and the publisher + * unpaid, which is strictly worse than an inflated count. + * 2. Each fetch in a burst renders a *different* campaign, so there is no one + * earlier row that could stand in for the rest without misattributing + * every later click to the first campaign. + * + * Best-effort: a failed lookup returns false (count it) rather than throwing. + * Losing an impression is worse than counting one twice. + */ +export async function isDuplicateImpression(input: { + slotId: string; + visitorId?: string | null; + ipHashes?: string[] | null; +}): Promise { + const visitor = safeId(input.visitorId); + const ipHashes = (input.ipHashes ?? []) + .map((h) => safeId(h)) + .filter((h): h is string => h !== null); + if (!visitor && ipHashes.length === 0) return false; // nothing to dedupe on + + const terms = [ + ...(visitor ? [`visitor_id.eq.${visitor}`] : []), + ...ipHashes.map((h) => `ip_hash.eq.${h}`), + ]; + + const since = new Date(Date.now() - IMPRESSION_DEDUPE_WINDOW_MS).toISOString(); + try { + const { data, error } = await serviceClient() + .from("ad_impressions") + .select("id") + .eq("slot_id", input.slotId) + .gte("ts", since) + .limit(1) + .or(terms.join(",")); + + if (error) return false; + return !!data && data.length > 0; + } catch { + // Never let the dedupe probe take serving down with it. This runs on the + // hot path of every fill; if it throws, the right answer is "count it" and + // carry on, not to lose the impression and the click that may follow. + return false; + } +} diff --git a/lib/ads/serve.ts b/lib/ads/serve.ts index e7aa62b..c83bae2 100644 --- a/lib/ads/serve.ts +++ b/lib/ads/serve.ts @@ -12,7 +12,13 @@ import { import { TERMINAL_FORMAT_ID } from "./formats"; import { houseFill, HOUSE_AD_ROTATION_RATE } from "./house"; import { CREDIT_CENTS, DEFAULT_BID_CREDITS, PLATFORM_RATE } from "./pricing"; -import { assessClickValidity, isBotDevice, CLICK_DEDUPE_WINDOW_MS } from "./fraud"; +import { + assessClickValidity, + isBotDevice, + isDuplicateImpression, + CLICK_DEDUPE_WINDOW_MS, + IMPRESSION_DEDUPE_WINDOW_MS, +} from "./fraud"; import { hashIpRotating, rotatingIpHashCandidates } from "@/lib/ipHash"; import { runAuction } from "./auction"; import { generateShortCode } from "./shortcode"; @@ -235,18 +241,27 @@ export async function serveAd( if (!campaign) return null; // Record the impression first so we have an id to bind the click to. + const ipHash = hashIpRotating(ctx.ip ?? null); const base = { slot_id: slotId, campaign_id: campaign.id, creative_id: pick.id, visitor_id: ctx.visitorId ?? null, - ip_hash: hashIpRotating(ctx.ip ?? null), + ip_hash: ipHash, geo_country: ctx.country ?? null, device: ctx.device ?? null, billable: false, tier, }; + // Flag — never skip. The row still has to exist so /a/ can resolve + // this exact campaign and creative; reporting excludes flagged rows instead. + const duplicate = await isDuplicateImpression({ + slotId, + visitorId: ctx.visitorId, + ipHashes: rotatingIpHashCandidates(ctx.ip ?? null, IMPRESSION_DEDUPE_WINDOW_MS), + }); + // The short code is what lets a terminal click URL fit inside the box, and // ctx.src records the publisher's surface tag on the row instead of in the // printed URL. Both live behind `add column if not exists`, and migrations @@ -257,7 +272,7 @@ export async function serveAd( const shortCode = generateShortCode(); let { data: imp } = await sb .from("ad_impressions") - .insert({ ...base, short_code: shortCode, src: ctx.src ?? null }) + .insert({ ...base, short_code: shortCode, src: ctx.src ?? null, duplicate }) .select("id, short_code") .single(); diff --git a/supabase/migrations/20260810120000_ad_impression_dedupe.sql b/supabase/migrations/20260810120000_ad_impression_dedupe.sql new file mode 100644 index 0000000..3fae3eb --- /dev/null +++ b/supabase/migrations/20260810120000_ad_impression_dedupe.sql @@ -0,0 +1,254 @@ +-- Ad network: mark repeat impressions so machine-driven prefetch stops reading +-- as advertiser delivery. +-- +-- Impressions are metered server-side in serveAd at *fill* time, on all three +-- serving paths (/api/ads/serve, /api/ads/frame, /api/ads/motd). Nothing about +-- that requires a browser, and the terminal path deliberately treats curl as a +-- real client. A scheduled pool refresher that fetches one slot N times back to +-- back therefore books N impressions per run, every run, whether or not a human +-- ever loads the page those fills land on. The observed case fires 12 fetches +-- in ~3 seconds every 10 minutes: ~1,700 impressions/day from one machine. +-- +-- Clicks already had a dedupe window (lib/ads/fraud.ts); impressions had none +-- at all. This closes that, with two deliberate differences from the click +-- rules, both explained in fraud.ts: +-- +-- * keyed on the SLOT, not the campaign — each fetch in a burst draws a +-- different campaign at random, so campaign-keyed dedupe collapses nothing; +-- * 60 seconds, not 6 hours — a repeat view hours later is real delivery and +-- must keep counting. +-- +-- Flag rather than drop. The row is what /a/ resolves a terminal +-- click back to, so skipping the insert would hand a real advertiser's creative +-- a click link pointing at nothing: unbilled click, unpaid publisher. Strictly +-- worse than an inflated count. Reporting excludes flagged rows instead. +-- +-- Existing rows default to false, so no historical figure moves. +-- +-- Apply via psql over the pooler / MCP (prod history diverged), not `db push`. + +alter table public.ad_impressions + add column if not exists duplicate boolean not null default false; + +comment on column public.ad_impressions.duplicate is + 'True when this viewer was already counted on this slot within the impression dedupe window. The row still exists for click attribution; reporting excludes it.'; + +-- Serves the dedupe lookup itself: (slot, ts) filtered, then visitor/ip matched. +create index if not exists ad_impressions_slot_ts_idx + on public.ad_impressions (slot_id, ts desc); + +-- Reporting: exclude flagged rows from both series functions. +-- Bodies are otherwise unchanged from 20260731140000_ad_range_series.sql. + +create or replace function public.ad_account_series( + p_since timestamptz default null, + p_bucket_seconds integer default 86400 +) +returns table ( + bucket timestamptz, + impressions bigint, + free_impressions bigint, + clicks bigint, + free_clicks bigint, + spent_cents bigint +) +language sql +stable +security invoker +set search_path = public +as $$ + with b as ( + select make_interval(secs => greatest(coalesce(p_bucket_seconds, 86400), 60)) as step + ), + owned as ( + select id from public.ad_campaigns where owner_id = auth.uid() + ), + ev as ( + select date_bin((select step from b), i.ts, timestamptz 'epoch') as bucket, + case when i.tier = 'free' then 0 else 1 end as imp, + case when i.tier = 'free' then 1 else 0 end as free_imp, + 0 as clk, + 0 as free_clk, + 0 as spent + from public.ad_impressions i + where i.campaign_id in (select id from owned) + and not i.duplicate + and (p_since is null or i.ts >= p_since) + union all + select date_bin((select step from b), cl.ts, timestamptz 'epoch'), + 0, + 0, + case when cl.valid then 1 else 0 end, + case when not cl.valid and cl.tier = 'free' then 1 else 0 end, + case when cl.valid then cl.charged_cents else 0 end + from public.ad_clicks cl + where cl.campaign_id in (select id from owned) + and (p_since is null or cl.ts >= p_since) + ) + select bucket, + sum(imp)::bigint, + sum(free_imp)::bigint, + sum(clk)::bigint, + sum(free_clk)::bigint, + sum(spent)::bigint + from ev + group by bucket + order by bucket; +$$; + +grant execute on function public.ad_account_series(timestamptz, integer) to authenticated, service_role; + +create or replace function public.ad_campaign_totals( + p_since timestamptz default null +) +returns table ( + campaign_id uuid, + impressions bigint, + free_impressions bigint, + clicks bigint, + free_clicks bigint, + spent_cents bigint +) +language sql +stable +security invoker +set search_path = public +as $$ + with owned as ( + select id from public.ad_campaigns where owner_id = auth.uid() + ), + ev as ( + select i.campaign_id, + case when i.tier = 'free' then 0 else 1 end as imp, + case when i.tier = 'free' then 1 else 0 end as free_imp, + 0 as clk, + 0 as free_clk, + 0 as spent + from public.ad_impressions i + where i.campaign_id in (select id from owned) + and not i.duplicate + and (p_since is null or i.ts >= p_since) + union all + select cl.campaign_id, + 0, + 0, + case when cl.valid then 1 else 0 end, + case when not cl.valid and cl.tier = 'free' then 1 else 0 end, + case when cl.valid then cl.charged_cents else 0 end + from public.ad_clicks cl + where cl.campaign_id in (select id from owned) + and (p_since is null or cl.ts >= p_since) + ) + select campaign_id, + sum(imp)::bigint, + sum(free_imp)::bigint, + sum(clk)::bigint, + sum(free_clk)::bigint, + sum(spent)::bigint + from ev + group by campaign_id; +$$; + +grant execute on function public.ad_campaign_totals(timestamptz) to authenticated, service_role; + +-- The dashboard reads impressions through three more surfaces. All of them get +-- the same exclusion, or the spike simply reappears on a different screen. + +-- Per-campaign / per-slot totals. Column lists are unchanged, so `create or +-- replace view` is enough here (the free-tier migration had to drop first only +-- because it inserted columns mid-list). +create or replace view public.ad_campaign_stats + with (security_invoker = true) as +select + c.id as campaign_id, + (select count(*) from public.ad_impressions i + where i.campaign_id = c.id and i.tier = 'paid' and not i.duplicate) as impressions, + (select count(*) from public.ad_impressions i + where i.campaign_id = c.id and i.tier = 'free' and not i.duplicate) as free_impressions, + (select count(*) from public.ad_clicks cl + where cl.campaign_id = c.id and cl.valid) as clicks, + (select count(*) from public.ad_clicks cl + where cl.campaign_id = c.id and not cl.valid and cl.tier = 'free') as free_clicks, + (select coalesce(sum(cl.charged_cents), 0) + from public.ad_clicks cl where cl.campaign_id = c.id and cl.valid) as spent_cents, + c.spend_today_cents, + c.total_spent_cents +from public.ad_campaigns c; + +grant select on public.ad_campaign_stats to authenticated, service_role; + +create or replace view public.ad_slot_stats + with (security_invoker = true) as +select + s.id as slot_id, + (select count(*) from public.ad_impressions i + where i.slot_id = s.id and i.tier = 'paid' and not i.duplicate) as impressions, + (select count(*) from public.ad_impressions i + where i.slot_id = s.id and i.tier = 'free' and not i.duplicate) as free_impressions, + (select count(*) from public.ad_clicks cl + where cl.slot_id = s.id and cl.valid) as clicks, + (select count(*) from public.ad_clicks cl + where cl.slot_id = s.id and not cl.valid and cl.tier = 'free') as free_clicks, + (select coalesce(sum(cl.publisher_earn_cents), 0) + from public.ad_clicks cl where cl.slot_id = s.id and cl.valid) as earned_cents +from public.ad_slots s; + +grant select on public.ad_slot_stats to authenticated, service_role; + +-- Daily series. Body otherwise unchanged from +-- 20260717032002_ad_campaign_daily_series_rpc.sql. +create or replace function public.ad_campaign_daily_series(days integer default 30) +returns table ( + campaign_id uuid, + day date, + impressions bigint, + clicks bigint, + spent_cents bigint +) +language sql +stable +security invoker +set search_path = public +as $$ + with bounds as ( + select greatest(coalesce(days, 30), 1) as n + ), + since as ( + select ((now() at time zone 'UTC')::date - (n - 1))::timestamptz as from_ts + from bounds + ), + owned as ( + select id from public.ad_campaigns where owner_id = auth.uid() + ), + imps as ( + select i.campaign_id, + (i.ts at time zone 'UTC')::date as day, + count(*)::bigint as impressions + from public.ad_impressions i + where i.campaign_id in (select id from owned) + and not i.duplicate + and i.ts >= (select from_ts from since) + group by 1, 2 + ), + clk as ( + select cl.campaign_id, + (cl.ts at time zone 'UTC')::date as day, + count(*)::bigint as clicks, + coalesce(sum(cl.charged_cents), 0)::bigint as spent_cents + from public.ad_clicks cl + where cl.valid + and cl.campaign_id in (select id from owned) + and cl.ts >= (select from_ts from since) + group by 1, 2 + ) + select coalesce(i.campaign_id, c.campaign_id) as campaign_id, + coalesce(i.day, c.day) as day, + coalesce(i.impressions, 0) as impressions, + coalesce(c.clicks, 0) as clicks, + coalesce(c.spent_cents, 0) as spent_cents + from imps i + full outer join clk c + on c.campaign_id = i.campaign_id and c.day = i.day; +$$; + +grant execute on function public.ad_campaign_daily_series(integer) to authenticated, service_role; diff --git a/tests/contract/ads-impression-dedupe.test.ts b/tests/contract/ads-impression-dedupe.test.ts new file mode 100644 index 0000000..d3bea29 --- /dev/null +++ b/tests/contract/ads-impression-dedupe.test.ts @@ -0,0 +1,138 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +// Impressions had no dedupe at all while clicks had a 6h window. A scheduled +// pool refresher fetching one slot 12 times in ~3 seconds therefore booked 12 +// advertiser impressions per run, every 10 minutes, with no human involved. + +const H = vi.hoisted(() => ({ + state: { + // Rows the dedupe probe will "find" for the slot inside the window. + existing: [] as Array<{ id: string }>, + lastQuery: null as null | Record, + throwOnSelect: false, + }, +})); +const { state } = H; + +vi.mock("@/lib/supabase/service", () => ({ + serviceClient: () => ({ + from(table: string) { + if (table !== "ad_impressions") throw new Error(`unexpected table ${table}`); + if (state.throwOnSelect) { + // Mirrors a client whose shape doesn't match — the probe must swallow it. + return {} as never; + } + const q: Record = {}; + const rec = (k: string) => (v: unknown) => { + (state.lastQuery ??= {})[k] = v; + return q; + }; + q.select = rec("select"); + q.eq = (col: string, v: unknown) => { + (state.lastQuery ??= {})[col] = v; + return q; + }; + q.gte = (col: string, v: unknown) => { + (state.lastQuery ??= {})[`gte_${col}`] = v; + return q; + }; + q.limit = rec("limit"); + q.or = (terms: string) => { + (state.lastQuery ??= {}).or = terms; + return Promise.resolve({ data: state.existing, error: null }); + }; + return q; + }, + }), +})); + +import { isDuplicateImpression, IMPRESSION_DEDUPE_WINDOW_MS } from "@/lib/ads/fraud"; + +beforeEach(() => { + state.existing = []; + state.lastQuery = null; + state.throwOnSelect = false; +}); + +describe("isDuplicateImpression", () => { + it("flags a repeat view of the same slot inside the window", async () => { + state.existing = [{ id: "imp-earlier" }]; + const dup = await isDuplicateImpression({ + slotId: "slot-1", + visitorId: "v-abc", + ipHashes: ["a".repeat(32)], + }); + expect(dup).toBe(true); + }); + + it("does not flag the first view", async () => { + state.existing = []; + const dup = await isDuplicateImpression({ + slotId: "slot-1", + visitorId: "v-abc", + ipHashes: ["a".repeat(32)], + }); + expect(dup).toBe(false); + }); + + it("keys on the slot, not the campaign", async () => { + // The burst draws a different campaign per fetch, so campaign-keyed dedupe + // would collapse none of it. The probe must never constrain campaign_id. + state.existing = [{ id: "imp-earlier" }]; + await isDuplicateImpression({ slotId: "slot-1", ipHashes: ["b".repeat(32)] }); + expect(state.lastQuery).toBeTruthy(); + expect(state.lastQuery).toHaveProperty("slot_id", "slot-1"); + expect(state.lastQuery).not.toHaveProperty("campaign_id"); + }); + + it("matches on visitor id or any rotating ip hash", async () => { + state.existing = [{ id: "imp-earlier" }]; + await isDuplicateImpression({ + slotId: "slot-1", + visitorId: "v-abc", + ipHashes: ["c".repeat(32), "d".repeat(32)], + }); + const or = String((state.lastQuery ?? {}).or ?? ""); + expect(or).toContain("visitor_id.eq.v-abc"); + expect(or).toContain(`ip_hash.eq.${"c".repeat(32)}`); + // Yesterday's salt window still has to match after a rotation. + expect(or).toContain(`ip_hash.eq.${"d".repeat(32)}`); + }); + + it("returns false when there is nothing to dedupe on", async () => { + // No visitor and no IP: a terminal with no ?v= behind an unknown address. + // Flagging those together would collapse unrelated viewers into one. + const dup = await isDuplicateImpression({ slotId: "slot-1" }); + expect(dup).toBe(false); + expect(state.lastQuery).toBeNull(); + }); + + it("rejects identifiers that could inject PostgREST filter syntax", async () => { + state.existing = [{ id: "imp-earlier" }]; + await isDuplicateImpression({ + slotId: "slot-1", + visitorId: "v-abc,ip_hash.not.is.null", + ipHashes: ["e".repeat(32)], + }); + const or = String((state.lastQuery ?? {}).or ?? ""); + expect(or).not.toContain("ip_hash.not.is.null"); + }); + + it("counts the impression rather than throwing when the probe fails", async () => { + // Hot path on every fill: losing an impression (and the click that may + // follow) is worse than counting one twice. + state.throwOnSelect = true; + const dup = await isDuplicateImpression({ + slotId: "slot-1", + visitorId: "v-abc", + ipHashes: ["f".repeat(32)], + }); + expect(dup).toBe(false); + }); + + it("uses a much shorter window than the click path", async () => { + // A repeat view hours later is real delivery; only a machine-driven burst + // lands twice inside a minute. + expect(IMPRESSION_DEDUPE_WINDOW_MS).toBe(60_000); + }); +});