diff --git a/docs/promote-engine-architecture.md b/docs/promote-engine-architecture.md index 0a0e456..48e106f 100644 --- a/docs/promote-engine-architecture.md +++ b/docs/promote-engine-architecture.md @@ -1,6 +1,7 @@ # CrawlProof Promote — Multi-Channel Promotion Engine **Status:** architecture baseline. Phase 1 content sources are built; the rest is design. +**Last reconciled with `master`:** 2026-08-19, at `a1a7d30` (PR #206). **Primary interface:** `/dashboard/promote` **Surfaces:** PWA/Web, CLI, HTTP API, MCP **First provider:** Reddit @@ -112,9 +113,12 @@ section; `app/actions/promote.ts` gains `addKeywordSources`, `addFeedSource`, ### 2.3 Not built yet -Reddit destinations and subreddit discovery, the durable job model with idempotency -keys, approval modes, relevance scoring, crossposts, per-destination cooldowns, -account groups, the HTTP API surface, and the CLI. §5 onward describes these. +Reddit destinations and subreddit discovery, approval modes, relevance scoring, +crossposts, per-destination cooldowns, account groups, the HTTP API surface, and the +CLI. §5 onward describes these. + +The durable job model *is* built — see §9. It is listed under §16 phase 1 as the +piece the rest of the scheduler hangs off. --- @@ -212,6 +216,27 @@ least-recently-promoted first. The window is the last 50 posts, long enough to be a ratio and short enough that changing the mix takes effect within a day. +**The mix is narrowed to the classes the campaign can actually supply** before any +deficit is computed — `effectiveMix()` in `lib/promote/blend.ts`. A class survives the +narrowing if the campaign has a link ready in it *or* an enabled source feeding it; +if nothing survives, the configured mix is used unchanged. So an all-owned campaign +runs on a 100%-owned effective mix and posts on target forever, and shared re-enters +the mix the moment a keyword source is added. + +This is not a refinement, it is what keeps the feature alive. The sources migration +backfilled every pre-existing list with the 70/30 default. Without the narrowing, a +campaign holding only owned links reads that as "30% short on shared", finds no +shared inventory, covers with owned content, and marks each post `via_fallback` — +and then the daily fallback cap below stops it posting at all. In production that +mislabelled six posts within minutes of the migration and was three posts from +silencing the campaign (PR #205). + +The general rule, which outlives this feature: **a class with no inventory and no +source is not starved, it is not part of that campaign's mix.** Backfilling a policy +default onto rows that predate the policy makes those rows look permanently in +violation of it, and a quota on the violation path turns that into silent death +rather than a visible error. + ### 4.2 Fallback ```json @@ -227,6 +252,10 @@ stops it becoming an uncontrolled shared-content firehose. Fallback posts are ma `via_fallback` and counted over a rolling 24 hours — rolling rather than calendar, so the cap cannot be gamed at a midnight boundary. +Fallback only applies **within the effective mix of §4.1**. Covering for a class the +campaign never had is not a fallback and must not be marked or capped as one; only a +class that is genuinely in the mix and genuinely ran dry counts against the cap. + --- ## 5. Provider adapter contract @@ -349,19 +378,71 @@ UNIQUE (campaign_id, connection_id, destination_key, normalized_url) ## 9. Scheduling and jobs -Not built. Today the sweep publishes inline on the tick. The target is durable -publication jobs carrying immutable resolved inputs — resolved title, body and URL, -`scheduledAt`, an `idempotencyKey`, an attempt count, and a state machine of -`queued → preflighting → blocked → publishing → published | retrying | failed | -cancelled`. - -Two properties matter most and neither exists yet: **no publication without an -idempotency key**, and **a failed worker cannot double-publish after retry or -failover**. The existing claim-then-work pattern gives the second property for list -scheduling but not for individual publications. - -Schedule model: interval / times-of-day / cron, timezone, days of week, quiet hours, -jitter, and per-day caps overall and per provider. +**Built.** `promo_job` (migration `20260819120000_promote_jobs.sql`) and +`lib/promote/jobs.ts`. A job is one intended publication: one link, to one account, +at one destination, for one scheduling slot. It carries the resolved URL and title +frozen at plan time, the body once a worker has written it, the ownership and source +that selected it, an attempt count, a state, and an idempotency key. + +The state machine is `queued → preflighting → blocked → publishing → published | +retrying | failed | cancelled`. Only `queued`, `publishing`, `published`, `failed` +and `cancelled` are reached today; the other three are in the CHECK constraint +already so the Reddit provider does not need a migration to use them. + +### 9.1 Why the old claim did not hold + +The sweep read every due campaign and "claimed" each by pushing `next_run_at` +forward. That UPDATE carried no predicate on `next_run_at`, so it was a +read-then-write: two sweeps that both read the row both won it. And the worker runs +the sweep on a 60s interval *and* out-of-band for "Post now" +(`POST /dashboard/promote/sweep`), so overlapping runs are designed in, not rare. + +Downstream, nothing was idempotent. A crash between `postViaAccount()` returning and +the `promo_post` insert left no record; `last_promoted_at` is only stamped at the end +of the campaign, so the same link was still least-recently-promoted on the next tick +and went out again. + +### 9.2 The two mechanisms + +**Plan before publishing.** Jobs are inserted before anything is sent, keyed on +`sha256(list, link, account, destination, kind, slot)`. The slot is the `next_run_at` +value the sweep *observed as due* — not the wall clock, which would give each sweep +its own key and rebuild the bug. A racing sweep derives the same keys, loses to the +unique index, and gets nothing back, so it publishes nothing. + +**Claim by compare-and-swap.** `update ... where id = ? and state = 'queued'`. The +read and the write are one statement, so two workers cannot both observe `queued`. +Postgres decides ownership; a prior SELECT does not. + +The campaign-level claim is still there and now carries a predicate — it re-asserts +the same "still due" condition the select used (`.lte("next_run_at", dueBy)`), so the +loser's update matches nothing once the winner has pushed the campaign forward. It +re-asserts the condition rather than matching the exact timestamp read back, because +a predicate that silently never matched would stop every campaign posting with +nothing in the logs. It is an optimization either way — it saves duplicated work. The +guarantee lives in the job. + +### 9.3 At most once, deliberately + +A job still `publishing` past its lease (`PUBLISH_LEASE_MS`, 10 minutes) is **failed, +never retried**. No provider we publish through accepts an idempotency key, so an +interrupted publish has genuinely unknown outcome — the post may be live. Re-running +it is the duplicate this exists to prevent. The reaper closes it with the outcome +recorded as unknown and surfaces it in history for a human. The credit is not +refunded, because refunding a post that was in fact delivered is the other way to be +wrong; support can refund from history. + +Retry is therefore reserved for failures that provably happened *before* the publish +call. Bounded retry with backoff for those is still open, and is what `retrying` and +`attempt_count` are for. + +A pending cookie-auth post counts as published for the job: it has been handed to the +Playwright worker, so the job must not run again. `reconcilePromo` settles the +`promo_post` later, as before. + +Schedule model — interval / times-of-day / cron, timezone, days of week, quiet hours, +jitter, and per-day caps overall and per provider — is still the existing single +`cadence_seconds` plus quiet hours. Not built. --- @@ -485,6 +566,13 @@ campaigns; retain an auditable record of requesting user and exact publication. Of these, provenance for shared content and duplicate prevention are built. The `maxFallbackItemsPerDay` cap is a mass-posting control as much as an editorial one. +"Pause after repeated provider rejections" is now half-built, at the connected-account +layer rather than the campaign layer: `lib/sp/accountHealth.ts` (PR #206) stops +retrying an account whose consecutive failures have run away — the case that prompted +it had logged 2,953. Campaign-level pausing on destination rejection is still open, and +the Reddit provider will need it, since a subreddit rejection is a destination fact +rather than an account one. + --- ## 15. Acceptance criteria @@ -500,11 +588,11 @@ Of these, provenance for shared content and duplicate prevention are built. The | 7 | One or more authorized connected accounts can be selected | built (pre-existing) | | 8 | Web, CLI, API and MCP use the same campaign and job records | partial — web + MCP; no CLI or API | | 9 | A shared topic feed is fetched once and reused across campaigns | **built** | -| 10 | No publication runs without an idempotency key | not built | +| 10 | No publication runs without an idempotency key | **built** | | 11 | Reddit checks activity, links, crossposts, rules, requirements | not built | | 12 | Reddit supports one original link and up to two delayed crossposts | not built | | 13 | Every attempted publication appears in history and audit log | built (pre-existing) | -| 14 | A failed worker cannot double-publish after retry or failover | partial — per list, not per publication | +| 14 | A failed worker cannot double-publish after retry or failover | **built** — per publication | | 15 | Preview and approve the exact item, copy, account and destination | not built | --- @@ -512,8 +600,9 @@ Of these, provenance for shared content and duplicate prevention are built. The ## 16. Delivery phases 1. **Shared Promote core** — source normalization, keyword and custom-feed ingestion, - campaigns and blending, connected accounts, dashboard. *Sources, blending and - ingestion are built; the durable scheduler and queue are not.* + campaigns and blending, connected accounts, dashboard. *Sources, blending, + ingestion and the durable job model are built. The richer schedule model + (times-of-day, cron, per-provider caps) and the review queue are not.* 2. **Reddit provider** — server-side OAuth, subreddit discovery, activity/link/ crosspost filtering, rules and requirements, original link posting, delayed crossposts, Reddit-specific preflight and metrics. diff --git a/lib/promote/jobs.ts b/lib/promote/jobs.ts new file mode 100644 index 0000000..f19cf37 Binary files /dev/null and b/lib/promote/jobs.ts differ diff --git a/lib/promote/sweep.ts b/lib/promote/sweep.ts index e18099d..e21e085 100644 --- a/lib/promote/sweep.ts +++ b/lib/promote/sweep.ts @@ -1,6 +1,13 @@ -// Promote sweep: processes due promo_lists, picks the next link, -// generates a fresh pitch via LLM, publishes via the sp platform layer, -// debits 1 credit per post, and advances the scheduler. +// Promote sweep: processes due promo_lists, picks the next link, plans a +// durable job per intended publication, generates a fresh pitch via LLM, +// publishes via the sp platform layer, debits 1 credit per post, and advances +// the scheduler. +// +// The job layer (lib/promote/jobs.ts) is what makes a publication happen at +// most once. This sweep runs on a 60s interval *and* out-of-band whenever a +// user clicks "Post now", so two runs racing the same campaign is normal. They +// now converge: both derive the same idempotency key for the same slot, one +// insert wins, and only the worker that wins the claim publishes. import type { SupabaseClient } from "@supabase/supabase-js"; import type Anthropic from "@anthropic-ai/sdk"; @@ -8,6 +15,13 @@ import type OpenAI from "openai"; import { generatePitch } from "@/lib/promote/generatePitch"; import { selectNextLink } from "@/lib/promote/selectLink"; import type { Ownership } from "@/lib/promote/blend"; +import { + claimJob, + planJobs, + recordJobBody, + settleJob, + type PlanJobInput, +} from "@/lib/promote/jobs"; import { postViaAccount, type PostResult } from "@/lib/sp/post"; export type PromoteSweepClients = { @@ -17,6 +31,7 @@ export type PromoteSweepClients = { export type PromoteSweepResult = { listsProcessed: number; + jobsPlanned: number; postsAttempted: number; postsSucceeded: number; postsFailed: number; @@ -37,6 +52,9 @@ type PromoList = { quiet_start: number | null; quiet_end: number | null; timezone: string | null; + // The due time this sweep observed. It is the scheduling slot every job + // planned this tick is keyed on, so a racing sweep keys on the same one. + next_run_at: string; source_mix?: unknown; fallback_policy?: unknown; }; @@ -48,6 +66,7 @@ type PromoLink = { angle: string | null; summary?: string | null; source_name?: string | null; + source_id?: string | null; ownership?: Ownership | null; times_promoted?: number | null; }; @@ -71,6 +90,7 @@ export async function processDuePromoteLists( ): Promise { const result: PromoteSweepResult = { listsProcessed: 0, + jobsPlanned: 0, postsAttempted: 0, postsSucceeded: 0, postsFailed: 0, @@ -78,35 +98,55 @@ export async function processDuePromoteLists( listsPaused: 0, }; + const dueBy = new Date().toISOString(); + const { data: lists, error } = await supabase .from("promo_list") .select( - "id, user_id, name, cadence_seconds, post_mode, target_account_ids, brand_voice, quiet_start, quiet_end, timezone, source_mix, fallback_policy", + "id, user_id, name, cadence_seconds, post_mode, target_account_ids, brand_voice, quiet_start, quiet_end, timezone, next_run_at, source_mix, fallback_policy", ) .eq("status", "running") - .lte("next_run_at", new Date().toISOString()) + .lte("next_run_at", dueBy) .order("next_run_at", { ascending: true }) .limit(limit); if (error || !lists || lists.length === 0) return result; - // Claim each due list by pushing next_run_at forward before processing, so an - // overlapping sweep (e.g. the periodic tick racing a manual "Post now" - // trigger) can't pick up the same list and double-post. advanceScheduler - // re-stamps next_run_at at the end of a successful run. - await Promise.all( - (lists as PromoList[]).map((l) => - supabase + // Claim each due list by pushing next_run_at forward, conditional on it still + // being due. The predicate is the point: without it this is a read-then-write + // and an overlapping sweep wins the same list. Only the lists whose update + // matched a row are ours to process. + // + // "Still due" rather than "unchanged since we read it": the minimum cadence is + // 300s, so a claim always pushes next_run_at well into the future and the + // loser's predicate cannot match. Re-asserting the same condition the select + // used avoids depending on a timestamptz round-tripping to a byte-identical + // string — and a claim that silently never matched would stop every campaign + // posting with nothing in the logs, which is the failure mode this codebase + // has already paid for once. + // + // This is an optimization, not the safety property — it saves duplicated + // work. The guarantee that nothing publishes twice lives in the job's + // idempotency key and claim, one level down. + const claimed = await Promise.all( + (lists as PromoList[]).map(async (l) => { + const { data } = await supabase .from("promo_list") - .update({ next_run_at: new Date(Date.now() + l.cadence_seconds * 1000).toISOString() }) - .eq("id", l.id), - ), + .update({ + next_run_at: new Date(Date.now() + l.cadence_seconds * 1000).toISOString(), + }) + .eq("id", l.id) + .lte("next_run_at", dueBy) + .select("id"); + return (data ?? []).length > 0 ? l : null; + }), ); - for (const list of lists as PromoList[]) { + for (const list of claimed.filter((l): l is PromoList => l !== null)) { try { const r = await processOneList(supabase, list, clients); result.listsProcessed++; + result.jobsPlanned += r.planned; result.postsAttempted += r.attempted; result.postsSucceeded += r.succeeded; result.postsFailed += r.failed; @@ -136,8 +176,22 @@ async function processOneList( supabase: SupabaseClient, list: PromoList, clients: PromoteSweepClients, -): Promise<{ attempted: number; succeeded: number; failed: number; pending: number; paused: boolean }> { - const out = { attempted: 0, succeeded: 0, failed: 0, pending: 0, paused: false }; +): Promise<{ + planned: number; + attempted: number; + succeeded: number; + failed: number; + pending: number; + paused: boolean; +}> { + const out = { + planned: 0, + attempted: 0, + succeeded: 0, + failed: 0, + pending: 0, + paused: false, + }; // Resolve accounts const accounts = await resolveAccounts(supabase, list); @@ -195,7 +249,53 @@ async function processOneList( postTargets.push({ link, account: nextAccount }); } - for (const { link: targetLink, account } of postTargets) { + // Write down what we intend to publish, before publishing any of it. A + // racing sweep that reached the same decision derives the same idempotency + // keys and gets nothing back here, so it publishes nothing. + const plans: PlanJobInput[] = postTargets.map(({ link: targetLink, account }) => ({ + userId: list.user_id, + listId: list.id, + linkId: targetLink.id, + accountId: account.id, + platform: account.platform, + resolvedUrl: targetLink.url, + resolvedTitle: targetLink.title, + ownership, + sourceId: targetLink.source_id ?? null, + viaFallback, + slotAt: list.next_run_at, + })); + + const jobs = await planJobs(supabase, plans); + out.planned = jobs.length; + + if (jobs.length === 0) { + // Either another sweep already owns this slot, or the job table could not + // be written (planJobs has already said which, loudly). Nothing to do, and + // in particular nothing to publish. + await advanceScheduler(supabase, list); + return out; + } + + const accountsById = new Map(accounts.map((a) => [a.id, a])); + + for (const [index, job] of jobs.entries()) { + const account = accountsById.get(job.account_id); + if (!account) { + // The account was disconnected between resolve and plan. Close the job + // rather than leaving it queued: nothing reclaims a queued job, so it + // would sit there forever misrepresenting the campaign as backed up. + await settleJob(supabase, job.id, { + state: "cancelled", + error: "connected account is no longer available", + }); + continue; + } + + // Take the job. Nothing below this line runs twice for the same job, and + // nothing above it published anything. + if (!(await claimJob(supabase, job))) continue; + // Check credits before each post const { data: hasCredit } = await supabase.rpc("consume_credit", { p_owner: list.user_id, @@ -203,7 +303,15 @@ async function processOneList( }); if (!hasCredit) { - // Insufficient credits — auto-pause + // Insufficient credits — auto-pause. Cancel this job and every one still + // queued behind it, so a paused campaign does not leave a tick's worth of + // jobs stranded in 'queued' with nothing that will ever claim them. + for (const stranded of jobs.slice(index)) { + await settleJob(supabase, stranded.id, { + state: "cancelled", + error: "insufficient_credits", + }); + } await supabase .from("promo_list") .update({ @@ -222,7 +330,7 @@ async function processOneList( const { data: recentPitches } = await supabase .from("promo_post") .select("body") - .eq("link_id", targetLink.id) + .eq("link_id", job.link_id) .eq("platform", account.platform) .eq("status", "posted") .order("created_at", { ascending: false }) @@ -234,19 +342,23 @@ async function processOneList( // Generate a fresh pitch const pitch = await generatePitch({ - url: targetLink.url, - title: targetLink.title, - angle: targetLink.angle, + url: job.resolved_url, + title: job.resolved_title, + angle: link.angle, platform: account.platform, brandVoice: list.brand_voice, recentBodies, anthropic: clients.anthropic, openai: clients.openai, - summary: targetLink.summary ?? null, - sourceName: targetLink.source_name ?? null, - ownership, + summary: link.summary ?? null, + sourceName: link.source_name ?? null, + ownership: job.ownership, }); + // Freeze the copy on the job before it goes out, so a job interrupted + // during publish can be read back and shows exactly what was sent. + await recordJobBody(supabase, job.id, pitch.body); + // Publish via the existing sp platform layer const postResult: PostResult = await postViaAccount({ supabase, @@ -265,56 +377,80 @@ async function processOneList( // the real URL/status (and refund on failure). Only synchronous API posts // (bluesky/telegram/discord + OAuth reddit/mastodon) are 'posted' now. const isPending = postResult.ok && postResult.pending === true; - await supabase.from("promo_post").insert({ - list_id: list.id, - link_id: targetLink.id, - account_id: account.id, - platform: account.platform, - ownership, - source_id: (targetLink as { source_id?: string | null }).source_id ?? null, - via_fallback: viaFallback, - body: pitch.body, - provider: pitch.provider, - model: pitch.model, - status: !postResult.ok ? "failed" : isPending ? "pending" : "posted", - external_post_id: postResult.ok && !isPending ? postResult.platformPostId : null, - post_url: postResult.ok && !isPending ? postResult.webUrl || null : null, - error: postResult.ok ? null : postResult.error, - credits_spent: 1, - posted_at: postResult.ok && !isPending ? new Date().toISOString() : null, - sp_post_id: postResult.ok && isPending ? postResult.postId : null, - }); + const { data: inserted } = await supabase + .from("promo_post") + .insert({ + list_id: list.id, + link_id: job.link_id, + account_id: account.id, + platform: account.platform, + ownership: job.ownership, + source_id: job.source_id, + via_fallback: job.via_fallback, + body: pitch.body, + provider: pitch.provider, + model: pitch.model, + status: !postResult.ok ? "failed" : isPending ? "pending" : "posted", + external_post_id: postResult.ok && !isPending ? postResult.platformPostId : null, + post_url: postResult.ok && !isPending ? postResult.webUrl || null : null, + error: postResult.ok ? null : postResult.error, + credits_spent: 1, + posted_at: postResult.ok && !isPending ? new Date().toISOString() : null, + sp_post_id: postResult.ok && isPending ? postResult.postId : null, + }) + .select("id") + .maybeSingle(); + + const promoPostId = (inserted as { id?: string } | null)?.id ?? null; if (!postResult.ok) { out.failed++; + await settleJob(supabase, job.id, { + state: "failed", + error: postResult.error ?? "publish failed", + promoPostId, + }); // Synchronous failure — refund the credit now. (Async cookie failures // are refunded later by reconcilePromo in the worker.) await supabase.rpc("consume_credit", { p_owner: list.user_id, p_count: -1, }); - } else if (isPending) { - out.pending++; } else { - out.succeeded++; + // A pending cookie-auth post has been handed to the browser worker, so + // as far as this job is concerned the publish happened: the job must + // not be re-run. reconcilePromo settles the promo_post later. + await settleJob(supabase, job.id, { state: "published", promoPostId }); + if (isPending) out.pending++; + else out.succeeded++; } } catch (err) { out.failed++; const message = err instanceof Error ? err.message : "Unknown error"; // Record the failed post - await supabase.from("promo_post").insert({ - list_id: list.id, - link_id: targetLink.id, - account_id: account.id, - platform: account.platform, - ownership, - source_id: (targetLink as { source_id?: string | null }).source_id ?? null, - via_fallback: viaFallback, - body: `[generation failed: ${message}]`, - status: "failed", + const { data: inserted } = await supabase + .from("promo_post") + .insert({ + list_id: list.id, + link_id: job.link_id, + account_id: account.id, + platform: account.platform, + ownership: job.ownership, + source_id: job.source_id, + via_fallback: job.via_fallback, + body: `[generation failed: ${message}]`, + status: "failed", + error: message, + credits_spent: 0, + }) + .select("id") + .maybeSingle(); + + await settleJob(supabase, job.id, { + state: "failed", error: message, - credits_spent: 0, + promoPostId: (inserted as { id?: string } | null)?.id ?? null, }); // Refund credit diff --git a/supabase/migrations/20260819120000_promote_jobs.sql b/supabase/migrations/20260819120000_promote_jobs.sql new file mode 100644 index 0000000..f525e12 --- /dev/null +++ b/supabase/migrations/20260819120000_promote_jobs.sql @@ -0,0 +1,153 @@ +-- Promote durable jobs: make a publication the unit of work, and make it +-- impossible to publish the same thing twice. +-- +-- WHAT WAS WRONG +-- +-- The drip sweep decided and published in one pass. It selected every +-- promo_list whose next_run_at was due, "claimed" each one by pushing +-- next_run_at forward, then generated a pitch and posted it. +-- +-- That claim does not hold. The update carries no predicate on next_run_at, so +-- two sweeps that both read the same due row both win it. The worker runs the +-- sweep on a 60s interval *and* out-of-band whenever a user clicks "Post now" +-- (worker/index.ts, POST /dashboard/promote/sweep), so overlapping runs are a +-- designed-in feature, not a rare race. +-- +-- Worse, nothing downstream is idempotent. If the process dies after +-- postViaAccount() has published but before the promo_post insert, there is no +-- record that it happened. last_promoted_at is only stamped at the end of the +-- list, so the same link is still least-recently-promoted on the next tick and +-- goes out again. The user sees a duplicate; we see nothing. +-- +-- WHAT THIS CHANGES +-- +-- A promo_job is one intended publication: one link, to one account, at one +-- destination, for one scheduling slot. It is written *before* anything is +-- published and carries a deterministic idempotency_key over exactly that +-- intent, with a unique index behind it. +-- +-- Two sweeps racing the same due list derive the same key from the same slot, +-- so the second insert loses to the index and plans nothing. Then a worker +-- takes a job by compare-and-swap -- update ... where state = 'queued' -- so +-- only one worker can move a job to 'publishing'. Postgres decides, not +-- read-then-write. +-- +-- AT MOST ONCE, DELIBERATELY +-- +-- A job found stuck in 'publishing' past its lease is NOT retried. None of the +-- social providers accept an idempotency key, so a publish that was +-- interrupted has genuinely unknown outcome: it may be live on the platform. +-- Retrying it is the duplicate we are here to prevent. Such a job is failed +-- with the outcome recorded as unknown and surfaced in history, and a human +-- decides. Retry is reserved for failures that provably happened *before* the +-- publish call. +-- +-- ORDERING: this table must exist before the code that selects it ships. A +-- sweep selecting a missing column makes PostgREST error the whole select, +-- which returns null rows and stops every campaign silently -- the failure +-- mode that cost us a day when promo_list.source_mix shipped ahead of its +-- migration. planJobs() logs loudly rather than quietly skipping if this table +-- is missing, but the ordering is still the actual fix. + +create table if not exists public.promo_job ( + id uuid primary key default gen_random_uuid(), + + -- Denormalized owner, so the worker can bill and the user can read their own + -- jobs without a join through promo_list. References profiles, matching + -- promo_list.user_id -- not auth.users, which would let a job outlive the + -- profile row every other promote table is keyed to. + user_id uuid not null references public.profiles(id) on delete cascade, + + list_id uuid not null references public.promo_list(id) on delete cascade, + link_id uuid not null references public.promo_link(id) on delete cascade, + account_id uuid not null references public.sp_account(id) on delete cascade, + + platform text not null, + + -- Where inside the platform this goes. Empty string for platforms with a + -- single destination (a Bluesky account has one timeline); 'r/bitcoin' once + -- the Reddit provider lands. Part of the idempotency key, so it is not null + -- -- a null would make every key distinct under SQL comparison. + destination_key text not null default '', + + kind text not null default 'original' + check (kind in ('original', 'crosspost', 'reshare')), + + -- A crosspost is scheduled off the original it follows. + parent_job_id uuid references public.promo_job(id) on delete set null, + + -- Resolved inputs. Frozen at plan time so a job publishes what it was + -- planned to publish, even if the link or the campaign changes underneath. + -- resolved_body is null until the job is claimed, because writing the pitch + -- costs an LLM call and only the worker that wins the claim should make it. + resolved_url text not null, + resolved_title text, + resolved_body text, + + -- Denormalized selection context, so history explains itself without + -- re-deriving the blend that produced it. + ownership text not null default 'owned' + check (ownership in ('owned', 'partner', 'shared')), + source_id uuid references public.promo_source(id) on delete set null, + via_fallback boolean not null default false, + + -- The tick this job belongs to: the promo_list.next_run_at value the sweep + -- observed as due. Two sweeps reading the same due row see the same slot. + slot_at timestamptz not null, + scheduled_at timestamptz not null default now(), + + -- 'preflighting', 'blocked' and 'retrying' are unused today. They are in the + -- constraint now so the Reddit provider, which needs all three, does not + -- need a migration to change a check. + state text not null default 'queued' + check (state in ( + 'queued', 'preflighting', 'blocked', 'publishing', + 'published', 'retrying', 'failed', 'cancelled' + )), + + attempt_count int not null default 0, + last_error text, + + -- Lease bookkeeping. locked_at is stamped when a worker wins the claim; a + -- job still 'publishing' long after it is a crashed worker, not a slow one. + locked_at timestamptz, + + idempotency_key text not null, + + -- The attempt this job produced, once it has produced one. + promo_post_id uuid references public.promo_post(id) on delete set null, + + created_at timestamptz not null default now(), + updated_at timestamptz not null default now() +); + +-- The guarantee. Everything above is bookkeeping; this is the part that makes +-- a double publication impossible rather than unlikely. +create unique index if not exists promo_job_idempotency_key_uidx + on public.promo_job (idempotency_key); + +-- The claim scan: queued work, oldest slot first. +create index if not exists promo_job_claimable_idx + on public.promo_job (state, scheduled_at) + where state in ('queued', 'publishing'); + +-- History, and the campaign detail page. +create index if not exists promo_job_list_idx + on public.promo_job (list_id, created_at desc); + +-- ---------- RLS ---------- +-- Owner-readable. Jobs are written by the worker under the service role, which +-- bypasses this; a user may read their own and cancel a queued one, which is +-- what the approval modes in the architecture doc will need. + +alter table public.promo_job enable row level security; + +create policy "promo_job owner all" + on public.promo_job for all + using (auth.uid() = user_id) + with check (auth.uid() = user_id); + +-- ---------- updated_at ---------- +create trigger promo_job_updated_at + before update on public.promo_job + for each row execute function public.promo_set_updated_at(); diff --git a/tests/promote/jobs.test.ts b/tests/promote/jobs.test.ts new file mode 100644 index 0000000..5a2f829 --- /dev/null +++ b/tests/promote/jobs.test.ts @@ -0,0 +1,338 @@ +import { describe, it, expect, beforeEach, vi, afterEach } from "vitest"; +import { + claimJob, + idempotencyKeyFor, + planJobs, + reapStalePublishingJobs, + settleJob, + PUBLISH_LEASE_MS, + type PlanJobInput, +} from "@/lib/promote/jobs"; +import { makeFakeSupabase, resetIds, type FakeDb, type UniqueConstraint } from "./fake-supabase"; + +// The unique index that carries the whole guarantee. +const CONSTRAINTS: UniqueConstraint[] = [ + { table: "promo_job", columns: ["idempotency_key"] }, +]; + +const SLOT = "2026-08-19T12:00:00.000Z"; + +function plan(over: Partial = {}): PlanJobInput { + return { + userId: "user-1", + listId: "list-a", + linkId: "link-1", + accountId: "acct-1", + platform: "bluesky", + resolvedUrl: "https://example.com/a", + resolvedTitle: "A", + ownership: "owned", + sourceId: null, + viaFallback: false, + slotAt: SLOT, + ...over, + }; +} + +function db(over: Partial = {}): FakeDb { + return { promo_job: [], ...over }; +} + +beforeEach(() => { + resetIds(); +}); + +describe("idempotencyKeyFor", () => { + it("is stable for the same intended publication", () => { + const a = idempotencyKeyFor({ listId: "l", linkId: "k", accountId: "a", slotAt: SLOT }); + const b = idempotencyKeyFor({ listId: "l", linkId: "k", accountId: "a", slotAt: SLOT }); + expect(a).toBe(b); + }); + + it("separates every axis of the intent", () => { + const base = { listId: "l", linkId: "k", accountId: "a", slotAt: SLOT }; + const key = idempotencyKeyFor(base); + + expect(idempotencyKeyFor({ ...base, listId: "l2" })).not.toBe(key); + expect(idempotencyKeyFor({ ...base, linkId: "k2" })).not.toBe(key); + expect(idempotencyKeyFor({ ...base, accountId: "a2" })).not.toBe(key); + expect(idempotencyKeyFor({ ...base, slotAt: "2026-08-19T13:00:00.000Z" })).not.toBe(key); + expect(idempotencyKeyFor({ ...base, destinationKey: "r/bitcoin" })).not.toBe(key); + expect(idempotencyKeyFor({ ...base, kind: "crosspost" })).not.toBe(key); + }); + + it("reads two spellings of the same instant as one slot", () => { + // Postgres and JS disagree about how to render a timestamptz. If the key + // took the raw string, the same due time read back differently would look + // like a different slot and the second sweep would publish again. + const base = { listId: "l", linkId: "k", accountId: "a" }; + expect(idempotencyKeyFor({ ...base, slotAt: "2026-08-19T12:00:00.000Z" })).toBe( + idempotencyKeyFor({ ...base, slotAt: "2026-08-19T12:00:00+00:00" }), + ); + }); + + it("treats a missing destination as the empty destination, not as absent", () => { + const base = { listId: "l", linkId: "k", accountId: "a", slotAt: SLOT }; + expect(idempotencyKeyFor(base)).toBe(idempotencyKeyFor({ ...base, destinationKey: "" })); + expect(idempotencyKeyFor(base)).toBe(idempotencyKeyFor({ ...base, destinationKey: null })); + }); +}); + +describe("planJobs", () => { + it("writes one queued job per intended publication", async () => { + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + + const jobs = await planJobs(client, [ + plan({ accountId: "acct-1", platform: "bluesky" }), + plan({ accountId: "acct-2", platform: "mastodon" }), + ]); + + expect(jobs).toHaveLength(2); + expect(store.promo_job).toHaveLength(2); + expect(jobs.every((j) => j.state === "queued")).toBe(true); + // The copy is not written until a worker wins the claim and pays for it. + expect(jobs.every((j) => j.resolved_body === null)).toBe(true); + }); + + it("plans nothing for a slot another sweep already planned", async () => { + // The race the whole change exists for: the 60s tick and a "Post now" + // trigger both reach the same due campaign and reach the same decision. + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + + const first = await planJobs(client, [plan()]); + const second = await planJobs(client, [plan()]); + + expect(first).toHaveLength(1); + expect(second).toHaveLength(0); + expect(store.promo_job).toHaveLength(1); + }); + + it("plans the same publication again on the next slot", async () => { + // Deduping must be per tick, not forever — a drip campaign is supposed to + // post the same link again later. + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + + await planJobs(client, [plan({ slotAt: SLOT })]); + const next = await planJobs(client, [plan({ slotAt: "2026-08-19T12:30:00.000Z" })]); + + expect(next).toHaveLength(1); + expect(store.promo_job).toHaveLength(2); + }); + + it("returns only the jobs this caller created when a slot is half planned", async () => { + const { client } = makeFakeSupabase(db(), CONSTRAINTS); + + await planJobs(client, [plan({ accountId: "acct-1" })]); + const rest = await planJobs(client, [ + plan({ accountId: "acct-1" }), + plan({ accountId: "acct-2" }), + ]); + + expect(rest.map((j) => j.account_id)).toEqual(["acct-2"]); + }); + + it("publishes nothing, loudly, when the job table cannot be written", async () => { + // The migration-not-applied case. Returning [] means the sweep publishes + // nothing; the console.error is what stops it being a silent stop. + const spy = vi.spyOn(console, "error").mockImplementation(() => {}); + const client = { + from: () => ({ + upsert: () => ({ + select: async () => ({ + data: null, + error: { message: 'relation "promo_job" does not exist' }, + }), + }), + }), + } as any; + + const jobs = await planJobs(client, [plan()]); + + expect(jobs).toEqual([]); + expect(spy).toHaveBeenCalled(); + expect(String(spy.mock.calls[0]?.[0])).toContain("publishing nothing this tick"); + spy.mockRestore(); + }); + + it("is a no-op for an empty plan", async () => { + const { client } = makeFakeSupabase(db(), CONSTRAINTS); + expect(await planJobs(client, [])).toEqual([]); + }); +}); + +describe("claimJob", () => { + it("lets exactly one worker take a job", async () => { + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + const [job] = await planJobs(client, [plan()]); + + expect(await claimJob(client, job)).toBe(true); + // The second worker read 'queued' before the first wrote — its update + // matches nothing, which is the point of the predicate. + expect(await claimJob(client, job)).toBe(false); + + const stored = store.promo_job[0]; + expect(stored.state).toBe("publishing"); + expect(stored.attempt_count).toBe(1); + expect(stored.locked_at).toBeTruthy(); + }); + + it("will not take a job that has already been settled", async () => { + const { client } = makeFakeSupabase(db(), CONSTRAINTS); + const [job] = await planJobs(client, [plan()]); + + await claimJob(client, job); + await settleJob(client, job.id, { state: "published", promoPostId: "post-1" }); + + expect(await claimJob(client, job)).toBe(false); + }); + + it("does not take a job when the claim errors", async () => { + const spy = vi.spyOn(console, "error").mockImplementation(() => {}); + const client = { + from: () => ({ + update: () => ({ + eq: () => ({ + eq: () => ({ + select: async () => ({ data: null, error: { message: "boom" } }), + }), + }), + }), + }), + } as any; + + expect(await claimJob(client, { id: "job-1", attempt_count: 0 })).toBe(false); + spy.mockRestore(); + }); +}); + +describe("settleJob", () => { + it("records the attempt the job produced", async () => { + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + const [job] = await planJobs(client, [plan()]); + + await settleJob(client, job.id, { state: "published", promoPostId: "post-9" }); + + expect(store.promo_job[0].state).toBe("published"); + expect(store.promo_job[0].promo_post_id).toBe("post-9"); + expect(store.promo_job[0].locked_at).toBeNull(); + }); + + it("keeps the reason a job failed", async () => { + const { client, db: store } = makeFakeSupabase(db(), CONSTRAINTS); + const [job] = await planJobs(client, [plan()]); + + await settleJob(client, job.id, { state: "failed", error: "rate limited" }); + + expect(store.promo_job[0].state).toBe("failed"); + expect(store.promo_job[0].last_error).toBe("rate limited"); + }); +}); + +describe("reapStalePublishingJobs", () => { + const NOW = new Date("2026-08-19T12:00:00.000Z"); + + beforeEach(() => { + vi.useFakeTimers(); + vi.setSystemTime(NOW); + }); + afterEach(() => { + vi.useRealTimers(); + }); + + function publishingJob(over: Record = {}) { + return { + id: "job-1", + list_id: "list-a", + platform: "bluesky", + state: "publishing", + locked_at: new Date(NOW.getTime() - PUBLISH_LEASE_MS - 1000).toISOString(), + attempt_count: 1, + ...over, + }; + } + + it("fails an interrupted job instead of retrying it", async () => { + // The publish may well have landed. Re-running it would be the duplicate + // this whole model exists to prevent, so the job is closed, not requeued. + const spy = vi.spyOn(console, "warn").mockImplementation(() => {}); + const { client, db: store } = makeFakeSupabase( + db({ promo_job: [publishingJob()] }), + CONSTRAINTS, + ); + + const result = await reapStalePublishingJobs(client); + + expect(result.reaped).toBe(1); + expect(store.promo_job[0].state).toBe("failed"); + expect(store.promo_job[0].state).not.toBe("queued"); + expect(store.promo_job[0].last_error).toContain("outcome unknown"); + expect(store.promo_job[0].locked_at).toBeNull(); + spy.mockRestore(); + }); + + it("leaves a job that is merely slow alone", async () => { + const { client, db: store } = makeFakeSupabase( + db({ + promo_job: [ + publishingJob({ locked_at: new Date(NOW.getTime() - 30_000).toISOString() }), + ], + }), + CONSTRAINTS, + ); + + expect((await reapStalePublishingJobs(client)).reaped).toBe(0); + expect(store.promo_job[0].state).toBe("publishing"); + }); + + it("ignores jobs that are not publishing", async () => { + const stale = new Date(NOW.getTime() - PUBLISH_LEASE_MS - 1000).toISOString(); + const { client, db: store } = makeFakeSupabase( + db({ + promo_job: [ + publishingJob({ id: "job-q", state: "queued", locked_at: null }), + publishingJob({ id: "job-p", state: "published", locked_at: stale }), + ], + }), + CONSTRAINTS, + ); + + expect((await reapStalePublishingJobs(client)).reaped).toBe(0); + expect(store.promo_job.map((j) => j.state)).toEqual(["queued", "published"]); + }); + + it("does not overwrite a slow worker that finished first", async () => { + // The reaper re-checks state in the update, so a worker that settled the + // job between the select and the write keeps its result. + const { client, db: store } = makeFakeSupabase( + db({ promo_job: [publishingJob()] }), + CONSTRAINTS, + ); + + // Simulate the finish landing after the reaper's select. + const original = client.from.bind(client); + let selected = false; + client.from = (table: string) => { + const builder = original(table); + if (table === "promo_job" && !selected) { + const select = builder.select.bind(builder); + builder.select = (...args: unknown[]) => { + const chained = select(...args); + const then = chained.then.bind(chained); + chained.then = (onOk: any, onErr: any) => + then((value: any) => { + if (!selected) { + selected = true; + store.promo_job[0].state = "published"; + } + return onOk ? onOk(value) : value; + }, onErr); + return chained; + }; + } + return builder; + }; + + expect((await reapStalePublishingJobs(client)).reaped).toBe(0); + expect(store.promo_job[0].state).toBe("published"); + }); +}); diff --git a/worker/index.ts b/worker/index.ts index fa8bded..f695caa 100644 --- a/worker/index.ts +++ b/worker/index.ts @@ -42,6 +42,7 @@ import { processUserAlerts } from "../lib/alerts/worker"; import { processDuePortScans } from "../lib/prober-queue"; import { processDueMonitors } from "../lib/uptime"; import { processDuePromoteLists } from "../lib/promote/sweep"; +import { reapStalePublishingJobs } from "../lib/promote/jobs"; import { ingestDueFeeds } from "../lib/promote/ingest"; import { refreshCookieSessions } from "../lib/sp/sessionRefresh"; @@ -1377,7 +1378,7 @@ async function promoteSweep() { const r = await processDuePromoteLists(supabase, { anthropic, openai }); if (r.postsAttempted > 0) { console.log( - `[worker] promote sweep lists=${r.listsProcessed} attempted=${r.postsAttempted} ok=${r.postsSucceeded} pending=${r.postsPending} fail=${r.postsFailed} paused=${r.listsPaused}`, + `[worker] promote sweep lists=${r.listsProcessed} planned=${r.jobsPlanned} attempted=${r.postsAttempted} ok=${r.postsSucceeded} pending=${r.postsPending} fail=${r.postsFailed} paused=${r.listsPaused}`, ); } } @@ -1386,6 +1387,24 @@ setInterval( 60_000, ); +// Promote job reaper: close out jobs whose worker died mid-publish. +// +// These are failed, never retried. No provider we publish through takes an +// idempotency key, so an interrupted publish may already be live on the +// platform and re-running it is the duplicate the job model exists to prevent. +// Runs on its own timer because it is about workers that are no longer running +// a sweep at all. +async function promoteReapSweep() { + const r = await reapStalePublishingJobs(supabase); + if (r.reaped > 0) { + console.warn(`[worker] promote reaped ${r.reaped} interrupted job(s)`); + } +} +setInterval( + () => promoteReapSweep().catch((e) => console.error("[worker] promote reap", e)), + 5 * 60_000, +); + // Promote content sources: refresh each subscribed feed once and fan its new // items out to every list that subscribes to it. Runs on its own timer rather // than inside the sweep because the feed registry is shared across users —