From af72a6deaeebe7b1e28a4009fe5555936f0c1d6b Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Wed, 19 Aug 2026 02:17:14 +0000 Subject: [PATCH 1/5] docs(promote): reconcile the architecture doc with #205 and #206 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The architecture doc was written in #204 and hasn't moved since, but two PRs landed under it the same day and both changed things it describes as settled. #205 added effectiveMix(): the ownership mix is narrowed to the classes a campaign can actually supply before any deficit is computed. §4.1 still described the raw configured mix, which is the version that mislabelled six production posts as via_fallback and was three posts from silencing the campaign. Someone reading §4 to extend the blend would have rebuilt the bug, so the rule and the reason it exists are now in the doc rather than only in the PR that fixed it. §4.2 gets the corresponding scope limit: covering for a class the campaign never had is not a fallback and is not capped as one. #206 added lib/sp/accountHealth.ts, which partly satisfies §14's "pause campaigns after repeated provider rejections" — at the account layer, not the campaign layer. Noted as half-built, because the Reddit provider will still need the campaign-level version: a subreddit rejection is a destination fact, not an account one. Docs only; no behaviour change. Co-Authored-By: Claude Opus 5 (1M context) --- docs/promote-engine-architecture.md | 33 +++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/docs/promote-engine-architecture.md b/docs/promote-engine-architecture.md index 0a0e456..357734a 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 @@ -212,6 +213,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 +249,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 @@ -485,6 +511,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 From 57aecd18e8f42c00288e57fd8ea28735a5fb2697 Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Wed, 19 Aug 2026 02:50:38 +0000 Subject: [PATCH 2/5] Make a Promote publication happen at most once (#AC10, #AC14) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The sweep decided and published in one pass, "claiming" a campaign 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 whenever someone clicks "Post now", so overlapping runs are designed in, not rare. Nothing downstream was idempotent either. If the process died after postViaAccount() published but before the promo_post insert, there was no record it had happened — and 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. The user sees a duplicate; the logs show nothing. promo_job makes the intended publication the unit of work, written down before anything is sent. Two mechanisms: - Plan before publishing. Jobs are keyed on sha256(list, link, account, destination, kind, slot), where the slot is the next_run_at value the sweep observed as due. A racing sweep reads the same due row, derives the same keys, loses to the unique index, and gets nothing back — so it publishes nothing. Keying on the wall clock instead would give each sweep its own key and rebuild the bug. - Claim by compare-and-swap: update ... where id = ? and state = 'queued'. Read and write are one statement, so two workers cannot both see 'queued'. The campaign-level claim keeps its place but now carries its predicate. It is an optimization — it avoids duplicated work. The guarantee is in the job. AT MOST ONCE, ON PURPOSE. A job still 'publishing' past a 10 minute lease is failed, never retried. No provider we publish through accepts an idempotency key, so an interrupted publish has genuinely unknown outcome — it may be live. Re-running it is the duplicate this exists to prevent. The reaper closes it with the outcome recorded as unknown and leaves it in history for a human; the credit is not refunded, because refunding a post that did land is the other way to be wrong. Retry stays available for failures that provably happened before the publish call, which is what 'retrying' and attempt_count are for. Stray jobs are closed rather than left queued: nothing reclaims a queued job, so one stranded by a disconnected account or a credit pause would sit there forever misreporting the campaign as backed up. MIGRATION FIRST. planJobs() logs loudly and publishes nothing if promo_job is missing, rather than failing the select and stopping every campaign in silence the way source_mix did. That is a guard, not the fix — apply 20260819120000_promote_jobs.sql before this deploys. Co-Authored-By: Claude Opus 5 (1M context) --- docs/promote-engine-architecture.md | 88 ++++- lib/promote/jobs.ts | Bin 0 -> 10269 bytes lib/promote/sweep.ts | 244 ++++++++++--- .../20260819120000_promote_jobs.sql | 151 ++++++++ tests/promote/jobs.test.ts | 338 ++++++++++++++++++ worker/index.ts | 21 +- 6 files changed, 764 insertions(+), 78 deletions(-) create mode 100644 lib/promote/jobs.ts create mode 100644 supabase/migrations/20260819120000_promote_jobs.sql create mode 100644 tests/promote/jobs.test.ts diff --git a/docs/promote-engine-architecture.md b/docs/promote-engine-architecture.md index 357734a..5ff2e01 100644 --- a/docs/promote-engine-architecture.md +++ b/docs/promote-engine-architecture.md @@ -113,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. --- @@ -375,19 +378,67 @@ 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`. +**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 its predicate +(`.eq("next_run_at", observed)`), but it is an optimization — 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. -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. +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. +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. --- @@ -533,11 +584,11 @@ rather than an account one. | 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 | --- @@ -545,8 +596,9 @@ rather than an account one. ## 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 0000000000000000000000000000000000000000..f19cf379d1a344ed891b5a26b6236534b09e2873 GIT binary patch literal 10269 zcmb_i+j85;5zVu|V!}$TNSmT$SM63p-{Q5sS$XZX%CfUnak*#^7>ck!0LDc$BdaPO zkuS`b92I_u$~b99PAvXsk)vDw=23d?@Di)HF+L+B|Jc zu`uhRIx}f)irktsZ)~2}#FXtkOXH|Xi+pUNJTYl=Yb|lSztoZsvATR2CgJ z%`(kT`CLZUOPfqEfMZ@`t)xmzQ?IQpO1nc1!1c^*#qhudj&G#v^HS4rl zn!GS&WfN}BjyYdlm}Xs=;@nnQRF>R^NLrhDf&ObqMRnED%64*?fBoa1B08T?fwnS` zqvm|lS+*+UiyR$AvDJB$$)KWb_KU^-oO_4ZW^i284Fp>c#+1fMf<17$#Man2be855 z^D`u(fK-0nV1FyKj%ou%u1cJ8-WZEL+GGEpRq>I@n02vQXNRHZm(I zM4H#OYD{Jqjcg=K^D|R7X_j$xW}~{I*lMV z#&Y;0TAx{@pRCg@Y^|QUhsKl{6orqIf%aGxYg0nE*W2Zp-9Rr?FtpUz3f7WB)|hK; zKb>;ojlqZs#HgaUJ7x^6(sM{8=f*=uCpo69sM%cfr(-v`&c~MClGsnOdF+(Z*EmEu zAnq6vhzaBW3P)=8Fr>A~p4oq3_S`u@iaj+;oL2!Guk1doV!vKT<;2h|PAe;Bu-%V3 z4jgd_Q-K+X;85(Ci$fJt4d_pkw2AeMvo@py3TY{L}(piV68 zy)C18RNI49dRhq()CYGC?|gUo$J>9pediFLnLGa(F=bMJ%*2@2OKOk%=b=?qTf#y* zyMtdYqBKJb=7gbQWgL2Cn<^z@+(==^u)hmfuZ4sDNE5y^3mB*wp|KY%ccEInN|IAGT_&x-iWCcXplRY0uY&~>0nHrvsvG}kwd z^4Mlj;x#h;1fw+f&7i2#(=?B=0cDG;qOR%R@(UKVgrV&YJ|+raR~UVPU8Knrj$c89 zyZ8mvu4erv$lqLP%FlM1qBt(ve77TV->2C{K?6bD-< zvU8iv+A0e>cI{0HllM=cg1pVL-oSa0bV~%!;ZSM{Q`1=m$JIq!#WvVmFy=gsW((ly zJc`ez2pR=oGM7PF(I7N=F{FtmBU=&UVN--~x+%ERg1v2h2*FLFEx=-r`=iCc0MV@8_r0LF%DM<6!0Ak3gs zEmS_<9)Ku1_)WcziA$EJo#XpHb1M>q(Xt+epRW+Ny$w{R?QDatQg z1sJM7Apj6QgM%iBGiP80h@a->fd?Nw@#8n&KR-J8arXS_qm!qz7dUov`w+{1$3M3L z&Radb#ztUDIZkRrzUSD;F3lq;ge0_)G5MMNa)!cFkp}K8Fq5(DMS9jTC|`@wU{5$v zip z2lnO|n+xnR5koo0dhWXW)xM5teqv`4I3B zl0jQ=m=p7+rZW1%qD$RSjCIjwi8hEx>0)NY1uoY?TgZwP3$OkH{R)?{IY$ zF|$Y#us&tU?V4Mwv2~H47sG2rb+SUyOPXNOp4l$iJZ3f#I`mI=^FvV$)4Xf}k{3W3 z8pcsC@X#LnEzTzUjkK{xJ){lH6ze?nAROQ};SKFwSh*Z^ohXYs8Ymewec!xA3s6w( zcQDo;=eONA*}+7G*Jr+k!iO!L@L^+LRP*qm8T2Nmpwij}bsY_k^)rHKl%;>~@Ctrf zsEY$n^G57??{9l@etWkL52uHRd!{ZyA8(hz8_n3!`4&G=8&9 zc~JLO!9Dbk>T9fOOIF@o&sAQNs$4fy7eue+B1g1cC1`<7BgFpjB`Om72PuS$qrXUQ zOhmXY%K~3VHbGEhUWY&+m?Q^TReFGfVLLC~jJO0~`pE2e;RN@T4uJ zuz(Q=aVQKFWeJ=|QqP%s02^iqcB(zv#2;FhHEFWSe5P2VYQbUx z{pqFxrJ%Lb$f*U8kR!^*?I&o+=}$0~#TaCBJtAptbV8WaxIHx;|+l|rZWQHY&{{NAH2h;E7NR9Bwt z0P5Cfbzc8l7Y)2u8MSr_nh`2O{3Y`GmA-_+Au;xSeH&a}SN07vmze0nzEvQA>oQHz z0YKxYpXe$nFq~E_7MIvX)MiMi8K!9tEDXVX7|hWNM^kRh#$pKlR4ft4fVmvwMYi3r zqNQX^%M*^7iHq}b8Eu4GXw598XW-*SooV8<6*NF8-Bj108|nDJ*NrwDNkrRD6XpiK zX8>kyF3lDEa}77{k@8QZgiN@)xO7b{g70V&=2Lc(OtER5=}OuaG&shh*1e@Gy7AD4 z`Amt6f3`JZq z*q~$sz`ulRN9zcne z=ov6?fSqu07Hx_6T}8{eO@ruQbJ+=81vOw4krKL*lzEjJcUeT;6eX(MD5iiwsh>qi zJV}0GSVieF9X}fjH!T9)3dJid4*yz0AsSUtwu~U|2{YG1r}9ZA_OrzNbaTdCKbuV8 z!kR1RF=vp#-Ix`F5;cL)f-^v>tG-L)Mp0rS%46{0XYf|IY-8WS_?e7wAR$E}7wQ6C zg^Euo48$i^u)M1^0XNYQt^cVq=#=_!Ga8owoq_`xdlO;gv>%rAA7zl zUvQ*gX}*^`;a?*k5-gq-^RFTwSOzo$meO$y)O*rI%8$JQ&j4H{U?i@1rgTJk#{Y;* zl2&-U_Q1cj?amstI{OvD<4VD9NIU9niiQc+bj#dUYxclqKMtvcu~=T3IJjf02!;nf zB=wPygT4TPeb;_}olk3&hmhtE7Zf#3E~oDu&R=oAzT=5gsPA3tdY&deFsL`dKjCcV zm{kZzj>{BADCP#vkrP534?*FJwLVvyW7RIaI1H*?I}BJc?^!Yb9%aU_wZv^Y$m0kL zzdhSKFVaK}uOIlX#yp>`UhJ|W&W$iO2@m=xfljx1>&}MsZU!6rZJTlTI8#B(&$tKK zv2b1yXe?v~b2ViIs1j$xk?yA9uEMcTG5_0Hf-4H0Xg66kmtRo6dVNXwngX49Wpz)} zT$Hr69{PZ}zTc`-!9;eB8n@iCs2_t&8zj2Og}@vJ5?%&~mEG;vH{bZL8X8C7zF2mW z;g)kI-La6LUw(I!N$gaeerd{}i=SAU`?FF2);6ej73FI;d!VL}?_ zwSpRyWWfy1yBLmmPvix>p;*l~ALJbb(f}tDXH!(e!;r}!e~ZkrT|ou@##gRE!NHoM z+Avge>}yPg5=){2Cl>dX7KIZkfP@bBkhqZGadXXX-^&QcOwxp`?El9e;3g&>w1z(4Jmnpba`aB9?Qi^R!O+!U_D%i7b^QaOJj?TT~V8*C+=q z^`wb1d)#qRAm!JUKlrMGeZdOF?o|ncK!!I<=ndZ<_E2c+DIo)i+oo77(7VejpFqJA z+{_F|W?x>4z)Cxr=Km|Ra>)=VKyZ=gs~SH2>QRV*+&CJ6g28am0V8OJ?#Mr-E{39! z9g<-_e>T4trj0NB@8f;HEW|LVka8$!0`($w(l%%f)WFEm4;9+!4=6j!yfFAjRTXXS zUBJ(`v8F2yxeMs-e$96=mk)v@@bS)@6Z%>MbEgogV7x1sC>#T+@LJq|!52olcVr-@ zne-R5C8UA;fof#5ZKypyJ!ii*GYeRbU%ZK<1794UGWh?|N|Z7&_sLN7%E# zxFp3rGzqgo%<7XShpVE*%5Ke)>5yJ7b(qRjwWH#WE@S5*j{iwx z&ZaL%`o`9mH2+?1+x=1lJz|L(2pDn%EEZ6OE+ { const result: PromoteSweepResult = { listsProcessed: 0, + jobsPlanned: 0, postsAttempted: 0, postsSucceeded: 0, postsFailed: 0, @@ -81,7 +101,7 @@ export async function processDuePromoteLists( 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()) @@ -90,23 +110,33 @@ export async function processDuePromoteLists( 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 not + // having moved since we read it. 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. + // + // 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) + .eq("next_run_at", l.next_run_at) + .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 +166,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 +239,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 +293,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 +320,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 +332,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 +367,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..5565ba7 --- /dev/null +++ b/supabase/migrations/20260819120000_promote_jobs.sql @@ -0,0 +1,151 @@ +-- 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. + user_id uuid not null references auth.users(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 — From ebd8d3b799b28252c91b42632885c7610f8d727d Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Wed, 19 Aug 2026 02:52:45 +0000 Subject: [PATCH 3/5] Claim a campaign on "still due" rather than an exact timestamp match The claim predicate compared next_run_at to the exact string we read back. That is correct only if a timestamptz round-trips to a byte-identical value through PostgREST, and if it ever did not, the update would match nothing, no campaign would be claimed, and Promote would stop posting entirely with nothing in the logs that looks like a failure. That is the same silent-stop shape the source_mix migration cost us a day for, and it is not worth risking to save a predicate. Re-asserting the condition the select already used gives the same atomicity with none of that exposure: the minimum cadence is 300s, so whoever wins the claim pushes next_run_at well into the future and the loser's `lte` cannot match. The scheduling slot the idempotency key is built from is still the observed next_run_at, so racing sweeps still agree on the slot. Co-Authored-By: Claude Opus 5 (1M context) --- lib/promote/sweep.ts | 22 ++++++++++++++++------ 1 file changed, 16 insertions(+), 6 deletions(-) diff --git a/lib/promote/sweep.ts b/lib/promote/sweep.ts index 9a562c4..e21e085 100644 --- a/lib/promote/sweep.ts +++ b/lib/promote/sweep.ts @@ -98,22 +98,32 @@ 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, 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, conditional on it not - // having moved since we read it. 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. + // 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 @@ -126,7 +136,7 @@ export async function processDuePromoteLists( next_run_at: new Date(Date.now() + l.cadence_seconds * 1000).toISOString(), }) .eq("id", l.id) - .eq("next_run_at", l.next_run_at) + .lte("next_run_at", dueBy) .select("id"); return (data ?? []).length > 0 ? l : null; }), From ae69fddc62200b853d69a8e5f2fa131d311c8747 Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Wed, 19 Aug 2026 02:52:59 +0000 Subject: [PATCH 4/5] =?UTF-8?q?docs(promote):=20match=20=C2=A79.2=20to=20t?= =?UTF-8?q?he=20claim=20predicate=20that=20shipped?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 5 (1M context) --- docs/promote-engine-architecture.md | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/docs/promote-engine-architecture.md b/docs/promote-engine-architecture.md index 5ff2e01..48e106f 100644 --- a/docs/promote-engine-architecture.md +++ b/docs/promote-engine-architecture.md @@ -414,9 +414,13 @@ unique index, and gets nothing back, so it publishes nothing. 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 its predicate -(`.eq("next_run_at", observed)`), but it is an optimization — it saves duplicated -work. The guarantee lives in the job. +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 From 18d9e1b7ac26fdcebdf4d0402163141b56942c37 Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Wed, 19 Aug 2026 02:53:28 +0000 Subject: [PATCH 5/5] fix(promote): key promo_job.user_id to profiles, like every other promote table promo_list.user_id references public.profiles(id). Pointing promo_job at auth.users instead would let a job outlive the profile row the rest of the feature is keyed to, and would cascade differently from its own campaign. Co-Authored-By: Claude Opus 5 (1M context) --- supabase/migrations/20260819120000_promote_jobs.sql | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/supabase/migrations/20260819120000_promote_jobs.sql b/supabase/migrations/20260819120000_promote_jobs.sql index 5565ba7..f525e12 100644 --- a/supabase/migrations/20260819120000_promote_jobs.sql +++ b/supabase/migrations/20260819120000_promote_jobs.sql @@ -53,8 +53,10 @@ 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. - user_id uuid not null references auth.users(id) on delete cascade, + -- 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,