diff --git a/app/(app)/dashboard/promote/[id]/page.tsx b/app/(app)/dashboard/promote/[id]/page.tsx index 3e608425..484f512c 100644 --- a/app/(app)/dashboard/promote/[id]/page.tsx +++ b/app/(app)/dashboard/promote/[id]/page.tsx @@ -7,6 +7,9 @@ import { PromoteRealtime } from "@/components/promote/promote-realtime"; import { PromoteEditForm } from "@/components/promote/promote-edit-form"; import { AddLinksForm } from "@/components/promote/add-links-form"; import { LinkList } from "@/components/promote/link-list"; +import { SourceList, type PromoteSourceRow } from "@/components/promote/source-list"; +import { BlendForm } from "@/components/promote/blend-form"; +import { parseMix } from "@/lib/promote/blend"; export const metadata = { title: "Promote list" }; @@ -28,10 +31,13 @@ export default async function PromoteDetailPage({ params }: Props) { .maybeSingle(); if (!list) notFound(); - const [{ data: links }, { data: posts }, { data: accounts }] = await Promise.all([ + const [{ data: links }, { data: posts }, { data: accounts }, { data: sources }] = + await Promise.all([ supabase .from("promo_link") - .select("id, url, title, angle, enabled, times_promoted, last_promoted_at, created_at") + .select( + "id, url, title, angle, enabled, times_promoted, last_promoted_at, created_at, ownership, source_name", + ) .eq("list_id", id) .order("created_at", { ascending: true }), supabase @@ -46,6 +52,13 @@ export default async function PromoteDetailPage({ params }: Props) { .eq("user_id", user.id) .eq("status", "active") .order("platform"), + supabase + .from("promo_source") + .select( + "id, type, ownership, label, keyword, enabled, items_imported, last_ingested_at, promo_feed(feed_url, topic_slug, consecutive_failures, last_error)", + ) + .eq("list_id", id) + .order("created_at", { ascending: true }), ]); const linkRows = (links ?? []) as Array<{ @@ -57,6 +70,8 @@ export default async function PromoteDetailPage({ params }: Props) { times_promoted: number; last_promoted_at: string | null; created_at: string; + ownership: string | null; + source_name: string | null; }>; const postRows = (posts ?? []) as Array<{ @@ -79,6 +94,27 @@ export default async function PromoteDetailPage({ params }: Props) { handle: string; }>; + // The embedded promo_feed comes back as an object (or null) per row; flatten + // it so the client component gets one shape. + const sourceRows: PromoteSourceRow[] = ( + (sources ?? []) as Array> + ).map((s) => ({ + id: s.id, + type: s.type, + ownership: s.ownership, + label: s.label, + keyword: s.keyword ?? null, + enabled: s.enabled, + items_imported: s.items_imported ?? 0, + last_ingested_at: s.last_ingested_at ?? null, + feed_url: s.promo_feed?.feed_url ?? null, + topic_slug: s.promo_feed?.topic_slug ?? null, + consecutive_failures: s.promo_feed?.consecutive_failures ?? 0, + last_error: s.promo_feed?.last_error ?? null, + })); + + const mix = parseMix(list.source_mix); + const totalPosts = postRows.filter((p) => p.status === "posted").length; const totalCredits = postRows.reduce((sum, p) => sum + (p.credits_spent ?? 0), 0); @@ -134,6 +170,27 @@ export default async function PromoteDetailPage({ params }: Props) { /> + {/* Content sources */} +
+

Content sources

+

+ Standing subscriptions that keep this campaign supplied with fresh links. +

+
+ +
+
+ + {/* Content mix */} + {sourceRows.length > 0 && ( +
+

Content mix

+
+ +
+
+ )} + {/* Links */}
diff --git a/app/actions/promote.ts b/app/actions/promote.ts index a39269fd..30e59b77 100644 --- a/app/actions/promote.ts +++ b/app/actions/promote.ts @@ -4,6 +4,16 @@ import { revalidatePath } from "next/cache"; import { createClient } from "@/lib/supabase/server"; import { serviceClient } from "@/lib/supabase/service"; import { parseLinks, fetchLinkTitle } from "@/lib/promote/generatePitch"; +import { parseKeywords } from "@/lib/promote/keywords"; +import { + addKeywordSources as addKeywordSourcesToList, + ensureFeed, + normalizeFeedUrl, + validateFeedUrl, + type KeywordSourceOutcome, +} from "@/lib/promote/sources"; +import { fanOutToSubscribers, ingestFeedNow } from "@/lib/promote/ingest"; +import { parseFallback, parseMix, type Ownership } from "@/lib/promote/blend"; import { env } from "@/lib/env"; type Ok> = { ok: true } & T; @@ -179,6 +189,200 @@ export async function addLinksToList(input: { return { ok: true, added: count ?? urls.length }; } +// ============================================================ +// Content sources +// ============================================================ + +async function ownedList(listId: string, userId: string): Promise { + const svc = serviceClient(); + const { data } = await svc + .from("promo_list") + .select("id") + .eq("id", listId) + .eq("user_id", userId) + .maybeSingle(); + return Boolean(data); +} + +/** + * Add one source per keyword. "bitcoin, blockchain" becomes two RSS Amplifier + * topic sources, never one combined feed. + */ +export async function addKeywordSources(input: { + listId: string; + keywords: string; + ownership?: Ownership; +}): Promise | Err> { + const auth = await requireUser(); + if (!auth.ok) return auth; + if (!(await ownedList(input.listId, auth.userId))) { + return { ok: false, error: "List not found." }; + } + + const keywords = parseKeywords(input.keywords); + if (keywords.length === 0) return { ok: false, error: "Enter at least one keyword." }; + if (keywords.length > 25) { + return { ok: false, error: "Add at most 25 keywords at a time." }; + } + + const svc = serviceClient(); + const results = await addKeywordSourcesToList(svc, { + listId: input.listId, + keywords, + // Keyword sources are other people's writing by default. + ownership: input.ownership ?? "shared", + }); + + // Backfill each new source now, so the campaign has something to post + // without waiting for the next ingestion tick. + for (const result of results) { + if (!result.ok || !result.feedId) continue; + try { + await ingestFeedNow(svc, result.feedId); + await fanOutToSubscribers(svc, result.feedId, new Date(), result.sourceId); + } catch { + // The worker retries on its own schedule; a slow feed is not a failure + // to add the source. + } + } + + revalidatePath(`/dashboard/promote/${input.listId}`); + return { ok: true, results, added: results.filter((r) => r.ok).length }; +} + +/** Add a single RSS or Atom feed the user owns or follows. */ +export async function addFeedSource(input: { + listId: string; + feedUrl: string; + ownership?: Ownership; + label?: string; +}): Promise | Err> { + const auth = await requireUser(); + if (!auth.ok) return auth; + if (!(await ownedList(input.listId, auth.userId))) { + return { ok: false, error: "List not found." }; + } + + const normalized = normalizeFeedUrl(input.feedUrl); + if (!normalized.ok) return { ok: false, error: normalized.error }; + + const validation = await validateFeedUrl(normalized.url); + if (!validation.ok) return { ok: false, error: validation.error }; + + const svc = serviceClient(); + const feed = await ensureFeed(svc, { + feedUrl: validation.feedUrl, + kind: "custom_feed", + title: validation.title, + }); + if (!feed) return { ok: false, error: "Could not register that feed." }; + + const { data: source, error } = await svc + .from("promo_source") + .insert({ + list_id: input.listId, + feed_id: feed.id, + type: "custom_feed", + // A feed the user went out of their way to add is usually their own. + ownership: input.ownership ?? "owned", + label: (input.label ?? "").trim() || validation.title || validation.feedUrl, + }) + .select("id") + .single(); + + if (error || !source) { + return { ok: false, error: "This campaign already tracks that feed." }; + } + + try { + await ingestFeedNow(svc, feed.id); + await fanOutToSubscribers(svc, feed.id, new Date(), source.id as string); + } catch { + // Ingestion retries on the worker's schedule. + } + + revalidatePath(`/dashboard/promote/${input.listId}`); + return { ok: true, sourceId: source.id as string, title: validation.title }; +} + +/** Turn a source off without losing the links it has already contributed. */ +export async function toggleSource(sourceId: string, enabled: boolean): Promise { + const auth = await requireUser(); + if (!auth.ok) return auth; + const svc = serviceClient(); + + const { data: source } = await svc + .from("promo_source") + .select("id, list_id, promo_list!inner(user_id)") + .eq("id", sourceId) + .maybeSingle(); + if (!source) return { ok: false, error: "Source not found." }; + if ((source as any).promo_list?.user_id !== auth.userId) { + return { ok: false, error: "Not authorized." }; + } + + const { error } = await svc.from("promo_source").update({ enabled }).eq("id", sourceId); + if (error) return { ok: false, error: error.message }; + + revalidatePath(`/dashboard/promote/${source.list_id}`); + return { ok: true }; +} + +/** + * Remove a source. Links it already imported stay: they are in the rotation, + * and silently deleting posts a user has seen queued would be a surprise. + * promo_link.source_id is ON DELETE SET NULL for exactly this reason. + */ +export async function removeSource(sourceId: string): Promise { + const auth = await requireUser(); + if (!auth.ok) return auth; + const svc = serviceClient(); + + const { data: source } = await svc + .from("promo_source") + .select("id, list_id, promo_list!inner(user_id)") + .eq("id", sourceId) + .maybeSingle(); + if (!source) return { ok: false, error: "Source not found." }; + if ((source as any).promo_list?.user_id !== auth.userId) { + return { ok: false, error: "Not authorized." }; + } + + const { error } = await svc.from("promo_source").delete().eq("id", sourceId); + if (error) return { ok: false, error: error.message }; + + revalidatePath(`/dashboard/promote/${source.list_id}`); + return { ok: true }; +} + +/** Set the owned/partner/shared publishing ratio and the fallback policy. */ +export async function updateBlend(input: { + listId: string; + mix: Partial>; + fallback?: Record; +}): Promise { + const auth = await requireUser(); + if (!auth.ok) return auth; + + const mix = parseMix(input.mix); + const total = mix.owned + mix.partner + mix.shared; + if (total <= 0) return { ok: false, error: "At least one source group needs a weight." }; + + const svc = serviceClient(); + const { error } = await svc + .from("promo_list") + .update({ + source_mix: mix, + ...(input.fallback ? { fallback_policy: parseFallback(input.fallback) } : {}), + }) + .eq("id", input.listId) + .eq("user_id", auth.userId); + if (error) return { ok: false, error: error.message }; + + revalidatePath(`/dashboard/promote/${input.listId}`); + return { ok: true }; +} + // -------- Remove a link -------- export async function removeLink(linkId: string): Promise { const auth = await requireUser(); diff --git a/components/promote/blend-form.tsx b/components/promote/blend-form.tsx new file mode 100644 index 00000000..9b56936f --- /dev/null +++ b/components/promote/blend-form.tsx @@ -0,0 +1,80 @@ +"use client"; + +import { useState, useTransition } from "react"; +import { useRouter } from "next/navigation"; +import { updateBlend } from "@/app/actions/promote"; + +type Mix = { owned: number; partner: number; shared: number }; + +export function BlendForm({ listId, mix }: { listId: string; mix: Mix }) { + const router = useRouter(); + const [pending, start] = useTransition(); + const [values, setValues] = useState(mix); + const [error, setError] = useState(""); + const [notice, setNotice] = useState(""); + + const total = values.owned + values.partner + values.shared; + // Weights are relative, so what the user cares about is the resulting share. + const share = (value: number) => (total > 0 ? Math.round((value / total) * 100) : 0); + + const submit = (e: React.FormEvent) => { + e.preventDefault(); + setError(""); + setNotice(""); + start(async () => { + const result = await updateBlend({ listId, mix: values }); + if (!result.ok) { + setError(result.error); + return; + } + setNotice("Content mix saved."); + router.refresh(); + }); + }; + + return ( +
+

+ How the campaign splits its posts between your own content and everybody + else’s. Weights are relative — 70 and 30 gives the same result as 7 and 3. +

+ {error &&

{error}

} + {notice &&

{notice}

} + +
+ {( + [ + ["owned", "Our content"], + ["partner", "Partner"], + ["shared", "Industry"], + ] as const + ).map(([key, label]) => ( + + ))} +
+ + {total === 0 && ( +

+ Give at least one group a weight, or the campaign has nothing to post. +

+ )} + + +
+ ); +} diff --git a/components/promote/link-list.tsx b/components/promote/link-list.tsx index 5c9bd88b..c5e76b0c 100644 --- a/components/promote/link-list.tsx +++ b/components/promote/link-list.tsx @@ -13,6 +13,8 @@ type LinkRow = { times_promoted: number; last_promoted_at: string | null; created_at: string; + ownership?: string | null; + source_name?: string | null; }; export function LinkList({ links }: { links: LinkRow[] }) { @@ -61,6 +63,14 @@ function LinkItem({ link }: { link: LinkRow }) { {link.title && (
{link.title}
)} + {/* Where this link came from: a link imported from somebody else's feed + reads very differently from one the user pasted. */} + {link.ownership && link.ownership !== "owned" && ( +
+ {link.ownership === "shared" ? "Industry" : "Partner"} content + {link.source_name ? ` · ${link.source_name}` : ""} +
+ )} {link.angle && (
Angle: {link.angle} diff --git a/components/promote/source-list.tsx b/components/promote/source-list.tsx new file mode 100644 index 00000000..7b6cb537 --- /dev/null +++ b/components/promote/source-list.tsx @@ -0,0 +1,225 @@ +"use client"; + +import { useState, useTransition } from "react"; +import { useRouter } from "next/navigation"; +import { + addFeedSource, + addKeywordSources, + removeSource, + toggleSource, +} from "@/app/actions/promote"; + +export type PromoteSourceRow = { + id: string; + type: string; + ownership: string; + label: string; + keyword: string | null; + enabled: boolean; + items_imported: number; + last_ingested_at: string | null; + feed_url: string | null; + topic_slug: string | null; + consecutive_failures: number; + last_error: string | null; +}; + +const OWNERSHIP_LABELS: Record = { + owned: "Our content", + partner: "Partner", + shared: "Industry", +}; + +export function SourceList({ + listId, + sources, +}: { + listId: string; + sources: PromoteSourceRow[]; +}) { + const router = useRouter(); + const [pending, start] = useTransition(); + const [keywords, setKeywords] = useState(""); + const [feedUrl, setFeedUrl] = useState(""); + const [feedOwnership, setFeedOwnership] = useState<"owned" | "partner" | "shared">("owned"); + const [error, setError] = useState(""); + const [notice, setNotice] = useState(""); + + const reset = () => { + setError(""); + setNotice(""); + }; + + const submitKeywords = (e: React.FormEvent) => { + e.preventDefault(); + reset(); + start(async () => { + const result = await addKeywordSources({ listId, keywords }); + if (!result.ok) { + setError(result.error); + return; + } + // Report per keyword: "bitcoin added, zzz is not a topic yet" beats one + // blanket success or failure. + const failed = result.results.filter((r) => !r.ok); + setNotice( + `Added ${result.added} source${result.added !== 1 ? "s" : ""}.` + + (failed.length + ? ` Skipped: ${failed.map((f) => `${f.keyword} (${f.error})`).join(", ")}` + : ""), + ); + setKeywords(""); + router.refresh(); + }); + }; + + const submitFeed = (e: React.FormEvent) => { + e.preventDefault(); + reset(); + start(async () => { + const result = await addFeedSource({ listId, feedUrl, ownership: feedOwnership }); + if (!result.ok) { + setError(result.error); + return; + } + setNotice(`Added ${result.title ?? "the feed"}.`); + setFeedUrl(""); + router.refresh(); + }); + }; + + return ( +
+ {error &&

{error}

} + {notice &&

{notice}

} + + {sources.length === 0 ? ( +

+ No content sources yet. Add a keyword and this campaign keeps finding fresh + links to post on its own. +

+ ) : ( +
    + {sources.map((source) => ( +
  • +
    +
    +
    + {source.label} + + {OWNERSHIP_LABELS[source.ownership] ?? source.ownership} + + {!source.enabled && paused} + {source.consecutive_failures > 0 && ( + not responding + )} +
    + {source.feed_url && ( + + {source.feed_url} + + )} +

    + {source.items_imported} link{source.items_imported !== 1 ? "s" : ""} imported + {source.last_ingested_at + ? ` · last checked ${new Date(source.last_ingested_at).toLocaleString()}` + : " · not checked yet"} +

    + {source.last_error && source.consecutive_failures > 0 && ( +

    {source.last_error}

    + )} +
    +
    + + +
    +
    +
  • + ))} +
+ )} + +
+ +

+ One source per keyword, from the RSS Amplifier directory. Separate several with + commas — a phrase like “artificial intelligence” counts as one keyword. +

+ setKeywords(e.target.value)} + placeholder="bitcoin, blockchain, ethereum" + className="input w-full text-sm" + /> + +
+ +
+ +

+ Any RSS or Atom URL — your own blog, or a publication worth resharing. +

+
+ setFeedUrl(e.target.value)} + placeholder="https://example.com/feed.xml" + className="input min-w-0 flex-1 font-mono text-sm" + /> + +
+ +
+
+ ); +} diff --git a/docs/promote-engine-architecture.md b/docs/promote-engine-architecture.md new file mode 100644 index 00000000..0a0e4561 --- /dev/null +++ b/docs/promote-engine-architecture.md @@ -0,0 +1,523 @@ +# CrawlProof Promote — Multi-Channel Promotion Engine + +**Status:** architecture baseline. Phase 1 content sources are built; the rest is design. +**Primary interface:** `/dashboard/promote` +**Surfaces:** PWA/Web, CLI, HTTP API, MCP +**First provider:** Reddit + +This document is the target architecture for Promote. `promote-prd.md` describes the +drip engine that shipped in July 2026 and is still accurate about that engine; this +one describes what Promote becomes. Where they disagree, this document wins. + +--- + +## 0. Runtime baseline — read this before building anything here + +The architecture brief this document is derived from specified *"Bun services on +Railway, Turso for control-plane data, Redis/BullMQ-compatible durable queues and +distributed rate limiting."* **That is not the stack CrawlProof runs on**, and building +to it would fork the product. + +What is actually true: + +| Concern | Brief said | CrawlProof actually uses | +|---|---|---| +| App runtime | Bun services | **Next.js 16** on Node, one app (`app/`) | +| Control-plane data | Turso | **Supabase Postgres** (`ywcizjsgrcmhgyplldac`), with RLS | +| Background work | Separate Bun services | **One worker** (`worker/index.ts`) of `setInterval` sweeps on Railway | +| Queues | Redis/BullMQ everywhere | BullMQ exists but is used **only** for the port-scan prober; every other sweep claims rows in Postgres | +| Rate limiting | Distributed limiter | Per-sweep throttles in Postgres (`lib/sp/feedAutopost.ts` `POST_THROTTLE_MS`) | + +The service boundaries in §11 of the brief (`promote-api`, `feed-ingestor`, +`content-selector`, `promote-scheduler`, `provider-workers`, `metrics-worker`) are +**module boundaries here, not deployments**. They map onto `lib/promote/*` and sweeps +in the single worker. Keep the seams — they are good seams, and they are what would +make a later split cheap — but do not stand up six services for a feature whose whole +job is to post a few links an hour. + +The claim-a-row-then-work pattern the existing sweeps use (`next_run_at` pushed +forward before processing) is the house substitute for a queue lease, and it is what +the new ingestion sweep uses too. + +--- + +## 1. Product definition + +Promote is a multi-channel publishing and content-discovery engine. A user connects +one or more social accounts, configures one or more content sources, and creates +campaigns that publish relevant content through selected connected accounts. + +The Reddit integration is the first provider, not a standalone subsystem. Every +surface — PWA/Web, CLI, API, MCP — calls the same service layer and uses the same +campaign, authorization, scheduling, dedupe and audit model. + +A campaign may publish the user's own content, content from custom RSS/Atom feeds, +shared topical content from RSS Amplifier, a configurable blend of owned and shared, +manually submitted URLs, and later CrawlProof-generated content such as Autoblog items. + +--- + +## 2. What exists today + +### 2.1 Built before this document + +The drip engine (`promote-prd.md`), shipped July 2026: + +- `promo_list` — a campaign: cadence, post mode, target accounts, brand voice. +- `promo_link` — the rotation unit. Hand-pasted URLs. +- `promo_post` — one publication attempt, with credits and error history. +- `lib/promote/sweep.ts` — the 60s sweep: claim due lists, pick a link, write a + fresh pitch, publish, debit a credit. +- `lib/promote/generatePitch.ts` — dual-provider (Anthropic + OpenAI) copywriter + with per-platform voice profiles and anti-repeat. +- Connected accounts reuse `sp_account`, and publishing reuses `lib/sp/post.ts` and + `lib/sp/platforms/*` (bluesky, discord, facebook, linkedin, mastodon, reddit, + telegram, threads, x, plus a Playwright browser path). + +### 2.2 Built by this change — content sources + +Campaigns can now feed themselves. Migration +`supabase/migrations/20260818190000_promote_sources.sql`. + +- **`promo_feed`** — the shared fetch registry, keyed on feed URL, with no `user_id`. + Two hundred users tracking "bitcoin" poll RSS Amplifier **once between them**. + Carries ETag/Last-Modified, a fetch interval, and geometric failure backoff. +- **`promo_feed_item`** — normalized entries of a feed, also shared, unique on + `(feed_id, url_hash)`. +- **`promo_source`** — one campaign's subscription to one feed, carrying the + ownership classification and a per-ingest cap. +- **`promo_link`** gains provenance (`source_id`, `ownership`, `summary`, + `image_url`, `author_name`, `source_name`, `normalized_url`, `url_hash`, + `published_at`) and a partial unique index on `(list_id, url_hash)`. +- **`promo_post`** gains `ownership`, `source_id` and `via_fallback`, denormalized + the way `platform` already is, so the blend can read its own history. +- **`promo_list`** gains `source_mix` and `fallback_policy`. + +Modules: + +| Module | Job | +|---|---| +| `lib/promote/keywords.ts` | keyword → RSS Amplifier topic slug and feed URL | +| `lib/promote/normalizeUrl.ts` | canonical (publishable) vs identity (dedupe) URL forms | +| `lib/promote/feedParse.ts` | RSS/Atom → normalized items with summary, image, author, `` attribution | +| `lib/promote/ingest.ts` | conditional fetch, store, fan out to subscribers | +| `lib/promote/sources.ts` | validate and register sources | +| `lib/promote/blend.ts` | deficit-based ownership selection (pure) | +| `lib/promote/selectLink.ts` | the database side of selection | + +Surfaces: `/dashboard/promote/[id]` gains a **Content sources** and a **Content mix** +section; `app/actions/promote.ts` gains `addKeywordSources`, `addFeedSource`, +`toggleSource`, `removeSource`, `updateBlend`; MCP gains `promote_list_campaigns`, +`promote_add_keyword_source`, `promote_add_feed_source`, `promote_list_sources`. + +### 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. + +--- + +## 3. Content sources + +### 3.1 Types + +```ts +type PromoteSourceType = + | "project_feed" + | "custom_feed" + | "rssamplifier_topic" + | "manual_url" + | "crawlproof_autoblog"; +``` + +### 3.2 Keyword sources + +A keyword becomes exactly one RSS Amplifier topic feed: + +``` +bitcoin → https://rssamplifier.com/topics/bitcoin.rss +``` + +Several keywords become several independent sources, never one combined URL: + +``` +bitcoin, blockchain, ethereum + → /topics/bitcoin.rss + → /topics/blockchain.rss + → /topics/ethereum.rss +``` + +**Slugging matters and is verified against the live directory.** RSS Amplifier serves +`/topics/artificial-intelligence.rss`, *not* `/topics/artificial%20intelligence.rss`, +so keywords are hyphenated, not percent-encoded. An unknown topic returns **404**, +which is what lets a keyword be validated the moment the user adds it rather than +silently producing a campaign that never posts. + +Normalization: trim, collapse whitespace, lowercase, strip diacritics and +punctuation, hyphenate, deduplicate by slug, keep a display label. Keyword lists +split on commas and newlines **but never on spaces** — "artificial intelligence" is +one keyword. + +### 3.3 Custom feeds + +Any RSS or Atom URL. Fetched and parsed before it is saved, so "that address is +reachable but is not a feed" is caught in the form. Bare hostnames are accepted and +assumed `https`. Private and link-local addresses are refused, the same guard +`lib/audit/engine.ts` applies to user-supplied targets. + +### 3.4 Ownership + +```ts +type ContentOwnership = "owned" | "partner" | "shared"; +``` + +Ownership drives blend ratios, attribution, and fallback. Keyword sources default to +`shared`; a feed the user deliberately added defaults to `owned`; hand-pasted links +are `owned`. + +### 3.5 URL identity + +Two forms, deliberately kept apart: + +- **canonical** — what we publish. Tracking parameters removed; host case, `www.`, + scheme and trailing slash untouched, so the link resolves as the publisher meant. +- **normalized** — dedupe identity only. Scheme folded to `https`, `www.` dropped, + trailing slash removed, query sorted. `url_hash` is its sha256. + +A bare `ref` parameter is **not** stripped: it is tracking on some sites and routing +on others, and losing an attribution tag is cheaper than publishing a link that 404s. + +--- + +## 4. Campaigns, blend and fallback + +A campaign joins content sources to connected accounts and destinations. + +### 4.1 Blend + +`promo_list.source_mix` weights each ownership class: + +```json +{ "owned": 70, "partner": 0, "shared": 30 } +``` + +**Selection is deficit-based, not weighted-random.** Weighted random gives streaks, +and a streak of shared content is exactly what makes an automated account read as a +content farm. On every tick the class furthest below its target share posts. Over 100 +ticks a 70/30 campaign posts exactly 70 owned and 30 shared, and never runs the same +class more than three times consecutively. Within a class the rotation is unchanged: +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. + +### 4.2 Fallback + +```json +{ + "whenOwnedQueueEmpty": "use_shared", + "whenSharedQueueEmpty": "use_owned", + "maxFallbackItemsPerDay": 3 +} +``` + +This is what lets a user with no original content still run a campaign, while the cap +stops it becoming an uncontrolled shared-content firehose. Fallback posts are marked +`via_fallback` and counted over a rolling 24 hours — rolling rather than calendar, so +the cap cannot be gamed at a midnight boundary. + +--- + +## 5. Provider adapter contract + +All providers implement one interface. The scheduler must never contain +provider-specific posting logic; it creates jobs and invokes the selected adapter. + +```ts +interface PromoteProviderAdapter { + provider: string; + + getAuthorizationUrl(input: AuthorizationInput): Promise; + exchangeAuthorizationCode(input: AuthorizationCodeInput): Promise; + refreshAccessToken(connection: PromoteConnection): Promise; + + listAccounts(connection: PromoteConnection): Promise; + discoverDestinations(input: DestinationDiscoveryInput): Promise; + inspectDestination(input: DestinationInspectionInput): Promise; + + preflight(input: PromotePreflightInput): Promise; + publish(input: PromotePublishInput): Promise; + + reshare?(input: PromoteReshareInput): Promise; + deletePublication?(input: DeletePublicationInput): Promise; + readMetrics?(input: ReadMetricsInput): Promise; + + getRateLimitState(connection: PromoteConnection): Promise; +} +``` + +`lib/sp/platforms/*` is the existing informal version of this. Formalizing it is a +refactor of what is already there, not a rewrite — the publish path in +`lib/sp/post.ts` already dispatches per platform. + +--- + +## 6. Reddit provider + +The adapter supports OAuth with refresh, keyword-based subreddit discovery, filtering +by recent activity / link-submission support / crosspost support, rule and +post-requirement preflight, original link submissions, one primary subreddit plus +zero to two delayed crossposts, required flair, per-destination copy templates, +duplicate and cooldown checks, and publication URL capture. + +```ts +interface RedditPromoteDestination { + provider: "reddit"; + connectionId: string; + subreddit: string; + + enabled: boolean; + approvedByUser: boolean; + + rulesReviewedAt: string | null; + rulesHash: string | null; + + allowOriginalLinks: boolean; + allowCrossposts: boolean; + + flairTemplateId: string | null; + + minDestinationGapHours: number; + minDomainGapHours: number; + + crosspostRole: "primary" | "secondary" | "either"; +} +``` + +Publication pattern: primary subreddit → original link submission → delay → secondary +subreddit 1 → delay → secondary subreddit 2. Every secondary publication must be +destination-approved, rules-reviewed, relevant, and individually preflighted **at +execution time**, not at scheduling time. + +Existing groundwork: `lib/sp/platforms/reddit.ts`, `lib/sp/platforms/redditOutreach.ts`, +`lib/sp/redditSubreddit.ts`, and OAuth at `app/api/sp/oauth/reddit/*`. + +--- + +## 7. Relevance and selection + +Before an item reaches the queue, score it against the campaign and destination. + +```ts +interface PromoteRelevanceScore { + total: number; + keywordScore: number; + destinationScore: number; + freshnessScore: number; + sourceQualityScore: number; + ownershipPriorityScore: number; + duplicationPenalty: number; +} +``` + +Sequence: reject expired or malformed → canonicalize → drop already-published +duplicates → apply blocked/competitor-domain rules → match campaign keywords → match +destination context → freshness → source quality → **blend selection (built)** → +generate copy → provider preflight → schedule or send to review. + +--- + +## 8. Deduplication + +Levels: global canonical URL, per account, per connected provider account, per +destination, per campaign, same-story similarity across different URLs, domain +cooldown, destination cooldown. + +Built today: per-campaign dedupe via the partial unique index on +`promo_link (list_id, url_hash)`, and per-feed dedupe via +`promo_feed_item (feed_id, url_hash)`. `normalizedTitleHash` exists for same-story +detection but is not yet stored. + +The eventual publication key: + +```sql +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. + +--- + +## 10. Ingestion and fan-out + +``` +normalized keyword: bitcoin + ↓ +one promo_feed row (shared, no user_id) + ↓ +one scheduled fetch, conditional on ETag/Last-Modified + ↓ +promo_feed_item rows, normalized and fingerprinted + ↓ +fan out into promo_link for every subscribing campaign +``` + +Rules, all of which the built ingestion follows: fetch each normalized shared feed +once per interval; cache conditional request metadata; cap items per source per pass; +back off geometrically on failure and surface the error on the source; never let one +bad feed take down a sweep. + +Rules still to come: partition provider queues by provider and connection; rate-limit +by provider client, connected account and endpoint; idempotency keys on every +publication; bounded retries with jitter; dead-letter permanent failures. Never log +raw OAuth tokens (already true — tokens live in the `sp_account` vault). + +--- + +## 11. Surfaces + +`/dashboard/promote` should carry overview (active campaigns, queue depth, posts +today, failed/blocked jobs, account health, engagement, source freshness warnings), +campaigns, sources, connected accounts, destinations, queue and history. + +Built: campaigns, sources (add keywords / add a feed / pause / remove / health and +import counts), content mix, links with provenance, recent posts. + +Not built: queue preview and approval, destination management, history filtering by +destination, engagement metrics. + +### HTTP API (not built) + +``` +POST /api/v1/promote/sources +GET /api/v1/promote/sources +GET /api/v1/promote/sources/:sourceId/items +PATCH /api/v1/promote/sources/:sourceId +DELETE /api/v1/promote/sources/:sourceId +... +``` + +Adding several keywords returns **one source per normalized keyword**: + +```json +{ "type": "rssamplifier_topic", "keywords": ["bitcoin", "blockchain"], "ownership": "shared" } +``` + +### CLI (not built) + +`cli/index.ts` currently has no promote commands. Target surface per the brief: +`crawlproof promote source add --keywords bitcoin,blockchain`, `campaign create +--mix owned=70,shared=30`, `queue approve`, `history`, all with +`--format table|json|jsonl`. + +### MCP (partly built) + +Built: `promote_list_campaigns`, `promote_add_keyword_source`, +`promote_add_feed_source`, `promote_list_sources`, alongside the existing +`list_accounts`, `generate_promo_post`, `post_to_socials`, `promote_url`. + +Any MCP publication tool must require an explicit connected account and destination, +and its response should include the preflight summary and idempotency key. + +--- + +## 12. Copy generation + +```ts +interface PromoteCopyPolicy { + mode: "feed_title" | "template" | "ai_rewrite"; + template: string | null; + includeSummary: boolean; + includeSourceName: boolean; + includeHashtags: boolean; + maxHashtags: number; + preserveOriginalTitle: boolean; + prohibitedPhrases: string[]; +} +``` + +Not built as a policy object, but the substance of the attribution requirement is: +`generatePitch` now takes `ownership`, `summary` and `sourceName`, and shared content +gets an explicit instruction never to imply we wrote it and to credit the source by +name. Without that the model happily announces somebody else's blog post as though +the account shipped it. + +--- + +## 13. Review and automation modes (not built) + +`manual`, `review_first`, `automatic`. Recommended defaults: new connected accounts +and newly discovered destinations start at `review_first`; shared-content-only +campaigns start at `review_first`; `automatic` unlocks after a campaign has valid +destinations, healthy sources and successful reviewed publications. + +--- + +## 14. Policy and abuse controls + +Promote provides publishing automation, not indiscriminate mass posting. + +Required: explicit authorization per connected account; explicit destination +selection; per-account and per-destination limits; domain and destination cooldowns; +duplicate prevention; destination-rule preflight; **no** automated voting, liking, +following, joining or account creation; **no** proxy rotation to evade limits; **no** +hidden account fan-out; clear provenance for shared content; pause campaigns after +repeated provider rejections; allow providers and administrators to disable abusive +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. + +--- + +## 15. Acceptance criteria + +| # | Criterion | State | +|---|---|---| +| 1 | One keyword creates one RSS Amplifier topic source | **built** | +| 2 | Comma-separated keywords create deduplicated sources | **built** | +| 3 | Multiple custom RSS/Atom feeds can be added | **built** | +| 4 | A campaign can contain owned and shared source groups | **built** | +| 5 | A campaign maintains a configured owned/shared ratio | **built** | +| 6 | Shared content acts as fallback when the owned queue is empty | **built** | +| 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 | +| 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 | +| 15 | Preview and approve the exact item, copy, account and destination | not built | + +--- + +## 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.* +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. +3. **MCP** — read tools, campaign preview, explicit schedule and publish tools, + audit-log integration. *Source tools built.* +4. **Additional providers** — added through the adapter contract without modifying + campaign or source architecture. diff --git a/docs/promote-prd.md b/docs/promote-prd.md index 6fa3b332..1be2a2de 100644 --- a/docs/promote-prd.md +++ b/docs/promote-prd.md @@ -1,5 +1,13 @@ # Crawlproof Promote — PRD +> **Superseded in part.** This PRD describes the link-drip engine that shipped in +> July 2026, and is still accurate about it. The target architecture — content +> sources, campaigns with blend ratios, the provider adapter contract and the Reddit +> provider — lives in [promote-engine-architecture.md](./promote-engine-architecture.md). +> Where the two disagree, that document wins. Content sources (§3 there) are built; +> the paste-a-list-of-links flow below still works unchanged alongside them. + + > Goal: a new **global, top-level** feature (peer of Ads and Alerts, *not* scoped to a single project) where the customer pastes a list of links, and AI writes a fresh, custom marketing pitch for **each link × each platform**, then drip-publishes them across **all connected social accounts on a recurring cadence** (default every 30 minutes). Every send is uniquely generated — never the same copy twice — using OpenAI and/or Anthropic. > > Topnav label: **Promote**. Route: `/promote`. diff --git a/lib/mcp/promote.ts b/lib/mcp/promote.ts index d276bfac..6bddb579 100644 --- a/lib/mcp/promote.ts +++ b/lib/mcp/promote.ts @@ -12,6 +12,14 @@ import { env } from "@/lib/env"; import { serviceClient } from "@/lib/supabase/service"; import { generatePitch, fetchLinkTitle } from "@/lib/promote/generatePitch"; import { postViaAccount, type PostResult } from "@/lib/sp/post"; +import { parseKeywords, topicPageUrl } from "@/lib/promote/keywords"; +import { + addKeywordSources, + ensureFeed, + normalizeFeedUrl, + validateFeedUrl, +} from "@/lib/promote/sources"; +import { fanOutToSubscribers, ingestFeedNow } from "@/lib/promote/ingest"; type Account = { id: string; platform: string; handle: string; status: string }; @@ -194,4 +202,204 @@ export function registerPromoteTools(server: McpServer): void { return textResult(lines.join("\n")); }, ); + + // ---- Content sources ------------------------------------------------- + // Same service layer the dashboard uses, so a campaign built by an agent is + // indistinguishable from one built in the UI. + + server.registerTool( + "promote_list_campaigns", + { + description: + "List the caller's Promote campaigns with their ids, status and content mix. Use an id with the source tools.", + inputSchema: {}, + }, + async (_args, extra) => { + const userId = getUserId(extra); + const sb = serviceClient(); + const { data } = await sb + .from("promo_list") + .select("id, name, status, cadence_seconds, source_mix") + .eq("user_id", userId) + .order("created_at", { ascending: false }); + const rows = (data ?? []) as Array>; + if (!rows.length) return textResult("No Promote campaigns yet."); + return textResult( + rows + .map((r) => { + const mix = r.source_mix ?? {}; + return `${r.id} ${r.name} [${r.status}] every ${Math.round((r.cadence_seconds ?? 0) / 60)}m mix owned=${mix.owned ?? 0}/partner=${mix.partner ?? 0}/shared=${mix.shared ?? 0}`; + }) + .join("\n"), + ); + }, + ); + + server.registerTool( + "promote_add_keyword_source", + { + description: + "Add one RSS Amplifier topic source per keyword to a Promote campaign. 'bitcoin, ethereum' creates two independent sources, never one combined feed. Reports each keyword individually.", + inputSchema: { + campaign_id: z.string().describe("Promote campaign id (from promote_list_campaigns)."), + keywords: z + .string() + .describe("One keyword, or several separated by commas. Phrases are one keyword."), + ownership: z + .enum(["owned", "partner", "shared"]) + .optional() + .describe("Defaults to 'shared' — topic feeds are other people's writing."), + }, + }, + async (args, extra) => { + const userId = getUserId(extra); + const sb = serviceClient(); + const { data: list } = await sb + .from("promo_list") + .select("id") + .eq("id", args.campaign_id) + .eq("user_id", userId) + .maybeSingle(); + if (!list) return errorResult("Campaign not found."); + + const keywords = parseKeywords(args.keywords); + if (!keywords.length) return errorResult("No usable keywords in that input."); + + const results = await addKeywordSources(sb, { + listId: args.campaign_id, + keywords, + ownership: args.ownership ?? "shared", + }); + + for (const r of results) { + if (!r.ok || !r.feedId) continue; + try { + await ingestFeedNow(sb, r.feedId); + await fanOutToSubscribers(sb, r.feedId, new Date(), r.sourceId); + } catch { + // The worker retries; adding the source still succeeded. + } + } + + return textResult( + results + .map((r) => + r.ok + ? `✓ ${r.keyword} → ${topicPageUrl(r.slug)}` + : `✗ ${r.keyword}: ${r.error}`, + ) + .join("\n"), + ); + }, + ); + + server.registerTool( + "promote_add_feed_source", + { + description: + "Add an RSS or Atom feed as a content source for a Promote campaign. The feed is fetched and validated before it is saved.", + inputSchema: { + campaign_id: z.string().describe("Promote campaign id."), + feed_url: z.string().describe("The RSS or Atom feed URL."), + ownership: z + .enum(["owned", "partner", "shared"]) + .optional() + .describe("Defaults to 'owned' — a feed you added deliberately is usually yours."), + label: z.string().optional().describe("Display name; defaults to the feed's own title."), + }, + }, + async (args, extra) => { + const userId = getUserId(extra); + const sb = serviceClient(); + const { data: list } = await sb + .from("promo_list") + .select("id") + .eq("id", args.campaign_id) + .eq("user_id", userId) + .maybeSingle(); + if (!list) return errorResult("Campaign not found."); + + const normalized = normalizeFeedUrl(args.feed_url); + if (!normalized.ok) return errorResult(normalized.error); + const validation = await validateFeedUrl(normalized.url); + if (!validation.ok) return errorResult(validation.error); + + const feed = await ensureFeed(sb, { + feedUrl: validation.feedUrl, + kind: "custom_feed", + title: validation.title, + }); + if (!feed) return errorResult("Could not register that feed."); + + const { data: source, error } = await sb + .from("promo_source") + .insert({ + list_id: args.campaign_id, + feed_id: feed.id, + type: "custom_feed", + ownership: args.ownership ?? "owned", + label: (args.label ?? "").trim() || validation.title || validation.feedUrl, + }) + .select("id") + .single(); + if (error || !source) + return errorResult("This campaign already tracks that feed."); + + try { + await ingestFeedNow(sb, feed.id); + await fanOutToSubscribers(sb, feed.id, new Date(), source.id as string); + } catch { + // Ingestion retries on the worker's schedule. + } + + return textResult( + `✓ Added ${validation.title ?? validation.feedUrl} (${validation.itemCount} entries) as a ${args.ownership ?? "owned"} source.`, + ); + }, + ); + + server.registerTool( + "promote_list_sources", + { + description: + "List the content sources of a Promote campaign, with ownership, health and how many links each has contributed.", + inputSchema: { campaign_id: z.string().describe("Promote campaign id.") }, + }, + async (args, extra) => { + const userId = getUserId(extra); + const sb = serviceClient(); + const { data: list } = await sb + .from("promo_list") + .select("id") + .eq("id", args.campaign_id) + .eq("user_id", userId) + .maybeSingle(); + if (!list) return errorResult("Campaign not found."); + + const { data } = await sb + .from("promo_source") + .select( + "id, type, ownership, label, enabled, items_imported, last_ingested_at, promo_feed(feed_url, last_success_at, consecutive_failures, last_error)", + ) + .eq("list_id", args.campaign_id) + .order("created_at", { ascending: true }); + + const rows = (data ?? []) as Array>; + if (!rows.length) return textResult("This campaign has no content sources yet."); + + return textResult( + rows + .map((r) => { + const feed = r.promo_feed ?? {}; + const health = + (feed.consecutive_failures ?? 0) > 0 + ? `failing (${feed.last_error ?? "unknown error"})` + : "ok"; + return `${r.id} ${r.label} [${r.type}/${r.ownership}]${r.enabled ? "" : " (paused)"} imported=${r.items_imported ?? 0} ${health}`; + }) + .join("\n"), + ); + }, + ); + } diff --git a/lib/promote/blend.ts b/lib/promote/blend.ts new file mode 100644 index 00000000..76af6de6 --- /dev/null +++ b/lib/promote/blend.ts @@ -0,0 +1,177 @@ +// Blend selection: deciding whether the next post draws from the user's own +// content or from shared content. +// +// Random selection weighted 70/30 does not give you 70/30. Over the handful of +// posts a drip campaign makes in a day it gives you streaks, and a streak of +// shared content is exactly what makes an automated account look like a +// content farm. So selection is deficit-based instead: on every tick the class +// furthest *below* its target share is the one that posts. Over a rolling +// window that converges on the configured ratio, and it never produces a run +// of one class while the other sits starved. +// +// This module is pure. The database side lives in lib/promote/selectLink.ts. + +export type Ownership = "owned" | "partner" | "shared"; + +export const OWNERSHIPS: readonly Ownership[] = ["owned", "partner", "shared"] as const; + +export type BlendMix = Record; + +export type FallbackAction = "pause" | "use_shared" | "use_owned" | "use_any_available"; + +export type FallbackPolicy = { + whenOwnedQueueEmpty: FallbackAction; + whenSharedQueueEmpty: FallbackAction; + maxFallbackItemsPerDay: number | null; +}; + +export const DEFAULT_MIX: BlendMix = { owned: 70, partner: 0, shared: 30 }; + +export const DEFAULT_FALLBACK: FallbackPolicy = { + whenOwnedQueueEmpty: "use_shared", + whenSharedQueueEmpty: "use_owned", + maxFallbackItemsPerDay: 3, +}; + +/** Read a stored source_mix, falling back to the default on anything unusable. */ +export function parseMix(raw: unknown): BlendMix { + if (!raw || typeof raw !== "object") return { ...DEFAULT_MIX }; + const source = raw as Record; + const mix: BlendMix = { owned: 0, partner: 0, shared: 0 }; + let total = 0; + for (const key of OWNERSHIPS) { + const value = Number(source[key]); + const weight = Number.isFinite(value) && value > 0 ? value : 0; + mix[key] = weight; + total += weight; + } + // A mix that weights nothing would post nothing. + return total > 0 ? mix : { ...DEFAULT_MIX }; +} + +/** Read a stored fallback_policy, falling back to the default per field. */ +export function parseFallback(raw: unknown): FallbackPolicy { + const source = (raw ?? {}) as Record; + const action = (value: unknown, fallback: FallbackAction): FallbackAction => + value === "pause" || + value === "use_shared" || + value === "use_owned" || + value === "use_any_available" + ? value + : fallback; + + const cap = source.maxFallbackItemsPerDay; + return { + whenOwnedQueueEmpty: action( + source.whenOwnedQueueEmpty, + DEFAULT_FALLBACK.whenOwnedQueueEmpty, + ), + whenSharedQueueEmpty: action( + source.whenSharedQueueEmpty, + DEFAULT_FALLBACK.whenSharedQueueEmpty, + ), + maxFallbackItemsPerDay: + cap === null ? null : Number.isFinite(Number(cap)) ? Number(cap) : DEFAULT_FALLBACK.maxFallbackItemsPerDay, + }; +} + +export type BlendDecision = { + /** The class to draw the next link from, or null when nothing may post. */ + ownership: Ownership | null; + /** True when the target class was empty and another one covered for it. */ + viaFallback: boolean; + /** Why, for the history view and for debugging a stalled list. */ + reason: + | "on_target" + | "fallback" + | "no_inventory" + | "fallback_disabled" + | "fallback_cap_reached"; +}; + +export type BlendInput = { + mix: BlendMix; + /** Posts per ownership class over the rolling window. */ + posted: Partial>; + /** Which classes have a link ready to post right now. */ + available: Partial>; + fallback: FallbackPolicy; + /** Fallback posts already made today, against maxFallbackItemsPerDay. */ + fallbackUsedToday?: number; +}; + +/** + * Rank the ownership classes by how far each is below its target share. + * Exported for the preview UI, which shows why a given item is next. + */ +export function rankByDeficit(mix: BlendMix, posted: Partial>): Ownership[] { + const totalWeight = OWNERSHIPS.reduce((sum, key) => sum + mix[key], 0); + const totalPosted = OWNERSHIPS.reduce((sum, key) => sum + (posted[key] ?? 0), 0); + + return OWNERSHIPS.filter((key) => mix[key] > 0) + .map((key) => { + const target = mix[key] / totalWeight; + const actual = totalPosted === 0 ? 0 : (posted[key] ?? 0) / totalPosted; + return { key, deficit: target - actual, weight: mix[key] }; + }) + // Largest deficit first; ties break toward the heavier weight so a fresh + // list opens with its dominant class rather than by object key order. + .sort((a, b) => b.deficit - a.deficit || b.weight - a.weight) + .map((entry) => entry.key); +} + +/** + * Choose which ownership class the next post draws from. + */ +export function chooseOwnership(input: BlendInput): BlendDecision { + const { mix, posted, available, fallback } = input; + const ranked = rankByDeficit(mix, posted); + + // The class that is furthest behind and actually has something to post. + const target = ranked[0]; + if (!target) return { ownership: null, viaFallback: false, reason: "no_inventory" }; + if (available[target]) { + return { ownership: target, viaFallback: false, reason: "on_target" }; + } + + // The target is starved. What the campaign is allowed to do about it depends + // on which side ran dry. + const action: FallbackAction = + target === "owned" + ? fallback.whenOwnedQueueEmpty + : target === "shared" + ? fallback.whenSharedQueueEmpty + : "use_any_available"; + + if (action === "pause") { + return { ownership: null, viaFallback: false, reason: "fallback_disabled" }; + } + + const preferred: Ownership[] = + action === "use_shared" + ? ["shared"] + : action === "use_owned" + ? ["owned"] + : ranked.filter((key) => key !== target); + + // "use_any_available" also considers classes with zero weight — a list mixed + // 100/0 still has partner links it may fall back to. + const candidates = + action === "use_any_available" + ? [...preferred, ...OWNERSHIPS.filter((key) => key !== target && !preferred.includes(key))] + : preferred; + + const replacement = candidates.find((key) => available[key]); + if (!replacement) { + return { ownership: null, viaFallback: false, reason: "no_inventory" }; + } + + // The cap is what stops a user with no original content from turning into an + // uncontrolled shared-content firehose. + const cap = fallback.maxFallbackItemsPerDay; + if (cap !== null && (input.fallbackUsedToday ?? 0) >= cap) { + return { ownership: null, viaFallback: false, reason: "fallback_cap_reached" }; + } + + return { ownership: replacement, viaFallback: true, reason: "fallback" }; +} diff --git a/lib/promote/feedParse.ts b/lib/promote/feedParse.ts new file mode 100644 index 00000000..14e2de53 --- /dev/null +++ b/lib/promote/feedParse.ts @@ -0,0 +1,192 @@ +// RSS/Atom parsing for Promote content sources. +// +// The social-feed engine already parses feeds (lib/sp/feedAutopost.ts), but it +// only needs a URL, a title and a date. Promote needs enough to write a pitch +// without re-fetching the page — summary, image, author, and a stable guid for +// dedupe — so it gets its own parser rather than widening that one. +// +// cheerio in xmlMode is the house pattern for feed XML here. + +import * as cheerio from "cheerio"; + +const MAX_SUMMARY_LENGTH = 600; + +export type ParsedFeedItem = { + url: string; + title: string | null; + summary: string | null; + /** The publisher's own id for the entry, when it gives one. */ + guid: string | null; + imageUrl: string | null; + author: string | null; + /** + * The publication the entry came from. Aggregator feeds (RSS Amplifier topic + * feeds among them) name the original publisher in ; shared content + * is attributed with this rather than with the aggregator's own name. + */ + sourceName: string | null; + publishedAt: string | null; +}; + +export type ParsedFeed = { + /** The feed's own title, used to attribute shared content. */ + title: string | null; + items: ParsedFeedItem[]; +}; + +/** + * Parse an RSS or Atom document. Unknown or malformed documents yield an empty + * item list rather than throwing — a bad feed should pause one source, not + * take down an ingestion sweep. + */ +export function parseFeed(xml: string, feedUrl: string, maxItems = 50): ParsedFeed { + let $: cheerio.CheerioAPI; + try { + $ = cheerio.load(xml, { xmlMode: true }); + } catch { + return { title: null, items: [] }; + } + + const feedTitle = + cleanText($("channel > title").first().text()) ?? + cleanText($("feed > title").first().text()); + + const items: ParsedFeedItem[] = []; + + $("item").each((_, el) => { + if (items.length >= maxItems) return false; + const node = $(el); + const link = node.children("link").first().text().trim(); + const guid = node.children("guid").first().text().trim() || null; + const url = absolutize(link || guid || "", feedUrl); + if (!url) return; + items.push({ + url, + title: cleanText(node.children("title").first().text()), + summary: summarize( + node.children("content\\:encoded").first().text() || + node.children("description").first().text(), + ), + guid, + imageUrl: itemImage($, node, feedUrl), + author: + cleanText(node.children("dc\\:creator").first().text()) ?? + cleanText(node.children("author").first().text()), + sourceName: cleanText(node.children("source").first().text()), + publishedAt: parseDate( + node.children("pubDate").first().text() || + node.children("dc\\:date").first().text() || + node.children("date").first().text(), + ), + }); + return; + }); + + $("entry").each((_, el) => { + if (items.length >= maxItems) return false; + const node = $(el); + const link = + node.children("link[rel='alternate']").first().attr("href") ?? + node.children("link").first().attr("href") ?? + node.children("id").first().text().trim(); + const url = absolutize(link ?? "", feedUrl); + if (!url) return; + items.push({ + url, + title: cleanText(node.children("title").first().text()), + summary: summarize( + node.children("summary").first().text() || + node.children("content").first().text(), + ), + guid: node.children("id").first().text().trim() || null, + imageUrl: itemImage($, node, feedUrl), + author: cleanText(node.children("author").first().children("name").first().text()), + sourceName: cleanText(node.children("source").first().children("title").first().text()), + publishedAt: parseDate( + node.children("published").first().text() || + node.children("updated").first().text(), + ), + }); + return; + }); + + return { title: feedTitle, items }; +} + +function itemImage( + $: cheerio.CheerioAPI, + node: cheerio.Cheerio | ReturnType, + feedUrl: string, +): string | null { + const candidates = [ + node.children("media\\:content").first().attr("url"), + node.children("media\\:thumbnail").first().attr("url"), + node.children("enclosure[type^='image']").first().attr("url"), + node.children("image").first().text().trim(), + ]; + for (const candidate of candidates) { + const url = absolutize((candidate ?? "").trim(), feedUrl); + if (url) return url; + } + // Last resort: the first inside the rendered body. + const body = + node.children("content\\:encoded").first().text() || + node.children("description").first().text() || + node.children("content").first().text(); + if (body) { + const match = body.match(/]+src=["']([^"']+)["']/i); + if (match) return absolutize(match[1], feedUrl); + } + return null; +} + +/** Strip markup and entities out of feed prose, then trim it to a usable length. */ +export function summarize(raw: string | null | undefined): string | null { + if (!raw) return null; + let text = raw + .replace(//gi, " ") + .replace(//gi, " ") + .replace(/<[^>]+>/g, " "); + text = decodeEntities(text).replace(/\s+/g, " ").trim(); + if (!text) return null; + if (text.length <= MAX_SUMMARY_LENGTH) return text; + // Prefer cutting at a word boundary so the summary does not end mid-word. + const clipped = text.slice(0, MAX_SUMMARY_LENGTH); + const lastSpace = clipped.lastIndexOf(" "); + return (lastSpace > MAX_SUMMARY_LENGTH * 0.6 ? clipped.slice(0, lastSpace) : clipped).trim() + "…"; +} + +function decodeEntities(text: string): string { + return text + .replace(/&(?:#3[49]|quot|apos);/g, (m) => (m === """ || m === """ ? '"' : "'")) + .replace(/ /g, " ") + .replace(/</g, "<") + .replace(/>/g, ">") + .replace(/&#(\d+);/g, (_, code) => String.fromCharCode(Number(code))) + .replace(/&/g, "&"); +} + +function cleanText(text: string | null | undefined): string | null { + const trimmed = decodeEntities(text ?? "").replace(/\s+/g, " ").trim(); + return trimmed || null; +} + +function absolutize(raw: string, base: string): string | null { + const value = (raw ?? "").trim(); + if (!value) return null; + try { + const url = new URL(value, base); + if (url.protocol !== "http:" && url.protocol !== "https:") return null; + return url.toString(); + } catch { + return null; + } +} + +function parseDate(value: string | null | undefined): string | null { + const raw = (value ?? "").trim(); + if (!raw) return null; + const ms = Date.parse(raw); + if (Number.isNaN(ms)) return null; + return new Date(ms).toISOString(); +} diff --git a/lib/promote/generatePitch.ts b/lib/promote/generatePitch.ts index 146ee3c3..7eeb773b 100644 --- a/lib/promote/generatePitch.ts +++ b/lib/promote/generatePitch.ts @@ -115,6 +115,15 @@ export type GeneratePitchArgs = { recentBodies: string[]; anthropic: Anthropic | null; openai: OpenAI | null; + /** Feed summary, when the link came from a content source. */ + summary?: string | null; + /** The originating publication, for attributing shared content. */ + sourceName?: string | null; + /** + * Whose content this is. Shared content is written about rather than as — + * "here's a good piece by X", not "we just shipped". + */ + ownership?: "owned" | "partner" | "shared" | null; }; export type PitchResult = { @@ -125,6 +134,25 @@ export type PitchResult = { model: string; }; +/** + * Shared content has to be written *about*, not *as*. Without this the model + * happily announces somebody else's blog post as though the account shipped + * it, which is both misleading and the fastest way to get an account banned. + */ +export function ownershipGuidance( + ownership: GeneratePitchArgs["ownership"], + sourceName: string | null | undefined, +): string { + if (ownership !== "shared" && ownership !== "partner") return ""; + const credit = sourceName ? ` Credit the source by name: ${sourceName}.` : ""; + return ( + "This link is somebody else's content, not ours. Share it the way a reader " + + "recommends a good article — never imply we wrote, built or launched it, and " + + "never use first-person product language like \"we just shipped\"." + + credit + ); +} + export async function generatePitch( args: GeneratePitchArgs, ): Promise { @@ -162,7 +190,9 @@ export async function generatePitch( `${HASHTAG_GUIDANCE[profile.usesHashtags]}`, brandVoice ? `Brand voice / instructions: ${brandVoice}` : "", angle ? `Marketing angle to emphasize: ${angle}` : "", + ownershipGuidance(args.ownership, args.sourceName), `Page title: ${title ?? "(no title — derive from URL)"}`, + args.summary ? `What the page says: ${args.summary}` : "", `URL to promote (do NOT include in the text field; the renderer appends it): ${url}`, avoidSection, ] diff --git a/lib/promote/ingest.ts b/lib/promote/ingest.ts new file mode 100644 index 00000000..253d8275 --- /dev/null +++ b/lib/promote/ingest.ts @@ -0,0 +1,394 @@ +// Promote ingestion: fetch each source feed once, fan the results out to every +// list that subscribes to it. +// +// The fan-out is the point. Two hundred users tracking "bitcoin" all read the +// same RSS Amplifier topic feed, so the registry is keyed on the feed URL and +// polled once per interval; each subscribing list then gets promo_link rows by +// reference. Nothing here is per-user until the fan-out step. +// +// Conditional requests (ETag / Last-Modified) mean an unchanged feed costs a +// 304 rather than a parse, which is most polls of most feeds. + +import type { SupabaseClient } from "@supabase/supabase-js"; +import { parseFeed, type ParsedFeedItem } from "@/lib/promote/feedParse"; +import { canonicalizeUrl, normalizeUrlForIdentity, urlHash } from "@/lib/promote/normalizeUrl"; + +const USER_AGENT = "CrawlProofPromote/1.0 (+https://crawlproof.com)"; +const FETCH_TIMEOUT_MS = 15_000; +const MAX_ITEMS_PER_FEED = 50; +const MAX_BODY_BYTES = 5_000_000; + +// A failing feed backs off geometrically instead of being polled every 15 +// minutes forever. Capped so a feed that comes back is noticed within a day. +const MAX_BACKOFF_MULTIPLIER = 32; +const MAX_BACKOFF_SECONDS = 86_400; + +export type FetchLike = ( + url: string, + init?: { headers?: Record; signal?: AbortSignal }, +) => Promise<{ + ok: boolean; + status: number; + headers: { get(name: string): string | null }; + text(): Promise; +}>; + +export type IngestResult = { + feedsChecked: number; + feedsUnchanged: number; + feedsFailed: number; + itemsStored: number; + linksCreated: number; +}; + +export type IngestOptions = { + /** How many due feeds to process in one pass. */ + limit?: number; + /** Injected for tests; defaults to global fetch. */ + fetchImpl?: FetchLike; + /** Injected for tests so scheduling maths is deterministic. */ + now?: Date; +}; + +type FeedRow = { + id: string; + feed_url: string; + etag: string | null; + last_modified: string | null; + fetch_interval_seconds: number; + consecutive_failures: number; +}; + +type SourceRow = { + id: string; + list_id: string; + ownership: string; + max_items_per_ingest: number; + items_imported: number; +}; + +const emptyResult = (): IngestResult => ({ + feedsChecked: 0, + feedsUnchanged: 0, + feedsFailed: 0, + itemsStored: 0, + linksCreated: 0, +}); + +/** + * Main ingestion entry point, called by the worker on a timer. Claims due + * feeds, refreshes them, and fans new items out to subscribers. + */ +export async function ingestDueFeeds( + supabase: SupabaseClient, + options: IngestOptions = {}, +): Promise { + const now = options.now ?? new Date(); + const result = emptyResult(); + + const { data: feeds, error } = await supabase + .from("promo_feed") + .select("id, feed_url, etag, last_modified, fetch_interval_seconds, consecutive_failures") + .lte("next_fetch_at", now.toISOString()) + .order("next_fetch_at", { ascending: true }) + .limit(options.limit ?? 20); + + if (error || !feeds || feeds.length === 0) return result; + + // Claim every due feed up front by pushing next_fetch_at forward, so an + // overlapping pass cannot fetch the same feed twice. + await Promise.all( + (feeds as FeedRow[]).map((feed) => + supabase + .from("promo_feed") + .update({ next_fetch_at: nextFetchAt(feed, now, false) }) + .eq("id", feed.id), + ), + ); + + for (const feed of feeds as FeedRow[]) { + try { + const one = await ingestOneFeed(supabase, feed, options, now); + result.feedsChecked++; + if (one.unchanged) result.feedsUnchanged++; + if (one.failed) result.feedsFailed++; + result.itemsStored += one.itemsStored; + result.linksCreated += one.linksCreated; + } catch (err) { + result.feedsChecked++; + result.feedsFailed++; + await recordFailure(supabase, feed, now, err); + } + } + + return result; +} + +/** + * Refresh one feed immediately, whatever its schedule says. Used when a user + * has just added a source and expects to see items without waiting for the + * next tick. + */ +export async function ingestFeedNow( + supabase: SupabaseClient, + feedId: string, + options: IngestOptions = {}, +): Promise { + const now = options.now ?? new Date(); + const result = emptyResult(); + + const { data: feed } = await supabase + .from("promo_feed") + .select("id, feed_url, etag, last_modified, fetch_interval_seconds, consecutive_failures") + .eq("id", feedId) + .maybeSingle(); + if (!feed) return result; + + try { + // A first fetch must not be answered with 304, or a brand new source shows + // up empty: ignore any stored validators. + const one = await ingestOneFeed( + supabase, + { ...(feed as FeedRow), etag: null, last_modified: null }, + options, + now, + ); + result.feedsChecked++; + if (one.unchanged) result.feedsUnchanged++; + if (one.failed) result.feedsFailed++; + result.itemsStored += one.itemsStored; + result.linksCreated += one.linksCreated; + } catch (err) { + result.feedsChecked++; + result.feedsFailed++; + await recordFailure(supabase, feed as FeedRow, now, err); + } + + return result; +} + +async function ingestOneFeed( + supabase: SupabaseClient, + feed: FeedRow, + options: IngestOptions, + now: Date, +): Promise<{ unchanged: boolean; failed: boolean; itemsStored: number; linksCreated: number }> { + const doFetch = options.fetchImpl ?? (globalThis.fetch as unknown as FetchLike); + + const headers: Record = { + accept: + "application/rss+xml,application/atom+xml,application/xml,text/xml,*/*;q=0.5", + "user-agent": USER_AGENT, + }; + if (feed.etag) headers["if-none-match"] = feed.etag; + if (feed.last_modified) headers["if-modified-since"] = feed.last_modified; + + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), FETCH_TIMEOUT_MS); + let response: Awaited>; + try { + response = await doFetch(feed.feed_url, { headers, signal: controller.signal }); + } finally { + clearTimeout(timer); + } + + // Unchanged since last time: the common case, and the cheap one. + if (response.status === 304) { + await supabase + .from("promo_feed") + .update({ + last_fetched_at: now.toISOString(), + last_success_at: now.toISOString(), + consecutive_failures: 0, + last_error: null, + next_fetch_at: nextFetchAt({ ...feed, consecutive_failures: 0 }, now, false), + }) + .eq("id", feed.id); + return { unchanged: true, failed: false, itemsStored: 0, linksCreated: 0 }; + } + + if (!response.ok) { + await recordFailure(supabase, feed, now, new Error(`Feed returned HTTP ${response.status}`)); + return { unchanged: false, failed: true, itemsStored: 0, linksCreated: 0 }; + } + + const body = await response.text(); + if (body.length > MAX_BODY_BYTES) { + await recordFailure(supabase, feed, now, new Error("Feed body too large")); + return { unchanged: false, failed: true, itemsStored: 0, linksCreated: 0 }; + } + + const parsed = parseFeed(body, feed.feed_url, MAX_ITEMS_PER_FEED); + if (parsed.items.length === 0) { + await recordFailure(supabase, feed, now, new Error("Feed contained no usable entries")); + return { unchanged: false, failed: true, itemsStored: 0, linksCreated: 0 }; + } + + const rows = parsed.items.map((item) => toItemRow(feed.id, item, now)).filter(isPresent); + + // ignoreDuplicates: a feed re-lists the same entries every poll by design, so + // a conflict on (feed_id, url_hash) is the normal case, not an error. + if (rows.length > 0) { + await supabase + .from("promo_feed_item") + .upsert(rows, { onConflict: "feed_id,url_hash", ignoreDuplicates: true }); + } + + await supabase + .from("promo_feed") + .update({ + title: parsed.title, + etag: response.headers.get("etag"), + last_modified: response.headers.get("last-modified"), + last_fetched_at: now.toISOString(), + last_success_at: now.toISOString(), + consecutive_failures: 0, + last_error: null, + next_fetch_at: nextFetchAt({ ...feed, consecutive_failures: 0 }, now, false), + }) + .eq("id", feed.id); + + const linksCreated = await fanOutToSubscribers(supabase, feed.id, now); + return { unchanged: false, failed: false, itemsStored: rows.length, linksCreated }; +} + +/** + * Copy a feed's newest items into every subscribing list as promo_link rows. + * + * Exported so a freshly added source can be backfilled without re-fetching the + * feed it shares with somebody else. + */ +export async function fanOutToSubscribers( + supabase: SupabaseClient, + feedId: string, + now: Date = new Date(), + onlySourceId?: string, +): Promise { + let sourceQuery = supabase + .from("promo_source") + .select("id, list_id, ownership, max_items_per_ingest, items_imported") + .eq("feed_id", feedId) + .eq("enabled", true); + if (onlySourceId) sourceQuery = sourceQuery.eq("id", onlySourceId); + + const { data: sources } = await sourceQuery; + if (!sources || sources.length === 0) return 0; + + // Newest first: a list that can only take ten links should take the ten most + // recent ones, not ten arbitrary ones. + const widestCap = Math.max( + ...(sources as SourceRow[]).map((s) => s.max_items_per_ingest ?? 10), + ); + const { data: items } = await supabase + .from("promo_feed_item") + .select( + "id, url, normalized_url, url_hash, title, summary, image_url, author_name, source_name, published_at", + ) + .eq("feed_id", feedId) + .order("published_at", { ascending: false, nullsFirst: false }) + .limit(Math.min(widestCap, MAX_ITEMS_PER_FEED)); + + if (!items || items.length === 0) return 0; + + let created = 0; + for (const source of sources as SourceRow[]) { + const slice = items.slice(0, source.max_items_per_ingest ?? 10); + const linkRows = slice.map((item: any) => ({ + list_id: source.list_id, + source_id: source.id, + ownership: source.ownership, + url: item.url, + normalized_url: item.normalized_url, + url_hash: item.url_hash, + title: item.title, + summary: item.summary, + image_url: item.image_url, + author_name: item.author_name, + source_name: item.source_name, + published_at: item.published_at, + discovered_at: now.toISOString(), + enabled: true, + })); + + // The unique index on (list_id, url_hash) does the deduplication: a story + // this list has already seen is silently skipped, including one that + // arrived earlier from a different source. + const { count, error } = await supabase + .from("promo_link") + .upsert(linkRows, { + onConflict: "list_id,url_hash", + ignoreDuplicates: true, + count: "exact", + }); + + const added = error ? 0 : (count ?? 0); + created += added; + + await supabase + .from("promo_source") + .update({ + last_ingested_at: now.toISOString(), + items_imported: (source.items_imported ?? 0) + added, + }) + .eq("id", source.id); + } + + return created; +} + +function toItemRow(feedId: string, item: ParsedFeedItem, now: Date) { + const url = canonicalizeUrl(item.url); + const normalized = normalizeUrlForIdentity(item.url); + const hash = urlHash(item.url); + // A feed entry we cannot identify is one we could neither dedupe nor publish. + if (!url || !normalized || !hash) return null; + return { + feed_id: feedId, + url, + normalized_url: normalized, + url_hash: hash, + title: item.title, + summary: item.summary, + image_url: item.imageUrl, + author_name: item.author, + source_name: item.sourceName, + guid: item.guid, + published_at: item.publishedAt, + discovered_at: now.toISOString(), + }; +} + +function isPresent(value: T | null): value is T { + return value !== null; +} + +async function recordFailure( + supabase: SupabaseClient, + feed: FeedRow, + now: Date, + err: unknown, +): Promise { + const failures = (feed.consecutive_failures ?? 0) + 1; + await supabase + .from("promo_feed") + .update({ + last_fetched_at: now.toISOString(), + consecutive_failures: failures, + last_error: err instanceof Error ? err.message : String(err), + next_fetch_at: nextFetchAt({ ...feed, consecutive_failures: failures }, now, true), + }) + .eq("id", feed.id); +} + +/** Geometric backoff while a feed is failing, plain interval while it is healthy. */ +export function nextFetchAt( + feed: Pick, + now: Date, + failing: boolean, +): string { + const base = feed.fetch_interval_seconds || 900; + const failures = failing ? Math.max(1, feed.consecutive_failures ?? 1) : 0; + const multiplier = Math.min(2 ** Math.max(0, failures - 1), MAX_BACKOFF_MULTIPLIER); + const seconds = Math.min(base * (failures > 0 ? multiplier : 1), MAX_BACKOFF_SECONDS); + return new Date(now.getTime() + seconds * 1000).toISOString(); +} diff --git a/lib/promote/keywords.ts b/lib/promote/keywords.ts new file mode 100644 index 00000000..4770d053 --- /dev/null +++ b/lib/promote/keywords.ts @@ -0,0 +1,68 @@ +// Keyword sources: a user types "bitcoin" and gets the RSS Amplifier topic +// feed for it. One keyword is one source — several keywords never collapse +// into a single ambiguous URL. + +export const RSSAMPLIFIER_BASE_URL = ( + process.env.RSSAMPLIFIER_BASE_URL ?? "https://rssamplifier.com" +).replace(/\/+$/, ""); + +// RSS Amplifier slugs its topics: "artificial intelligence" is served at +// /topics/artificial-intelligence, not /topics/artificial%20intelligence. +const MAX_SLUG_LENGTH = 80; + +export type KeywordInput = { + /** Identity: lowercase, hyphenated, safe to put in a path segment. */ + slug: string; + /** What the user typed, tidied up. Shown in the UI. */ + label: string; +}; + +/** + * Fold a raw keyword into an RSS Amplifier topic slug. Returns null when + * nothing usable survives (empty input, or punctuation only). + */ +export function slugifyKeyword(raw: string): string | null { + const slug = (raw ?? "") + .normalize("NFKD") + .replace(/[\u0300-\u036f]/g, "") + .toLowerCase() + .replace(/[^a-z0-9]+/g, "-") + .replace(/^-+|-+$/g, "") + .slice(0, MAX_SLUG_LENGTH) + .replace(/-+$/g, ""); + return slug || null; +} + +/** Collapse runs of whitespace so the stored label is tidy. */ +export function keywordLabel(raw: string): string { + return (raw ?? "").trim().replace(/\s+/g, " "); +} + +/** + * Parse a user-supplied keyword list. Splits on commas and newlines only — + * never on spaces, because "artificial intelligence" is one keyword. + * Deduplicates by slug, keeping the first label seen. + */ +export function parseKeywords(raw: string): KeywordInput[] { + const out: KeywordInput[] = []; + const seen = new Set(); + for (const chunk of (raw ?? "").split(/[,\n\r]+/)) { + const label = keywordLabel(chunk); + if (!label) continue; + const slug = slugifyKeyword(label); + if (!slug || seen.has(slug)) continue; + seen.add(slug); + out.push({ slug, label }); + } + return out; +} + +/** The topic feed URL a keyword source polls. */ +export function topicFeedUrl(slug: string, base: string = RSSAMPLIFIER_BASE_URL): string { + return `${base.replace(/\/+$/, "")}/topics/${encodeURIComponent(slug)}.rss`; +} + +/** The human-facing topic page, for "view source" links in the UI. */ +export function topicPageUrl(slug: string, base: string = RSSAMPLIFIER_BASE_URL): string { + return `${base.replace(/\/+$/, "")}/topics/${encodeURIComponent(slug)}`; +} diff --git a/lib/promote/normalizeUrl.ts b/lib/promote/normalizeUrl.ts new file mode 100644 index 00000000..0eef8383 --- /dev/null +++ b/lib/promote/normalizeUrl.ts @@ -0,0 +1,153 @@ +// URL canonicalization and identity hashing for Promote content items. +// +// Two different jobs, deliberately kept apart: +// +// canonicalUrl what we actually publish. Tracking junk is removed, but the +// resource itself is untouched (host case, `www.`, scheme and +// trailing slash all preserved) so the link resolves exactly +// the way the publisher meant it to. +// +// normalizedUrl an identity key used only for dedupe, never published. More +// aggressive: scheme folded to https, `www.` dropped, trailing +// slash removed, remaining query sorted. +// +// Keeping them apart matters: folding `www.` is right for "have I posted this +// story before?" and wrong for "which URL do I hand to Reddit?". + +import { createHash } from "node:crypto"; + +// Whole families of analytics parameters, matched by prefix. +const TRACKING_PREFIXES = [ + "utm_", // Google/analytics standard + "pk_", // Matomo (legacy Piwik) + "mtm_", // Matomo + "hsa_", // HubSpot ads + "_hs", // HubSpot email (_hsenc, _hsmi) + "at_", // AT Internet + "wt_", // Webtrekk +]; + +// Individually named click/campaign identifiers. +// +// Deliberately NOT stripped: a bare `ref`. It is a tracking parameter on some +// sites and a routing parameter on others (CrawlProof's own short links use +// `ref_slug`), so removing it can change which page loads. Losing one +// attribution tag is cheaper than publishing a link that 404s. +const TRACKING_PARAMS = new Set([ + "fbclid", + "gclid", + "gbraid", + "wbraid", + "dclid", + "msclkid", + "yclid", + "twclid", + "ttclid", + "igshid", + "igsh", + "mc_cid", + "mc_eid", + "mkt_tok", + "_openstat", + "oly_anon_id", + "oly_enc_id", + "ref_src", + "s_cid", + "cmpid", + "vero_id", + "vero_conv", + "ck_subscriber_id", + "hsctatracking", +]); + +export function isTrackingParam(name: string): boolean { + const key = name.toLowerCase(); + if (TRACKING_PARAMS.has(key)) return true; + return TRACKING_PREFIXES.some((prefix) => key.startsWith(prefix)); +} + +function parse(raw: string): URL | null { + const trimmed = (raw ?? "").trim(); + if (!trimmed) return null; + let url: URL; + try { + url = new URL(trimmed); + } catch { + return null; + } + // Only ever publish or dedupe web links. + if (url.protocol !== "http:" && url.protocol !== "https:") return null; + if (!url.hostname) return null; + return url; +} + +function dropTracking(url: URL): void { + for (const key of [...url.searchParams.keys()]) { + if (isTrackingParam(key)) url.searchParams.delete(key); + } +} + +/** + * The publishable form of a URL: tracking parameters removed, fragment + * dropped, everything else left exactly as the publisher wrote it. + * Returns null when the input is not a usable http(s) URL. + */ +export function canonicalizeUrl(raw: string): string | null { + const url = parse(raw); + if (!url) return null; + dropTracking(url); + url.hash = ""; + // A bare "?" left behind after stripping every parameter is noise. + if ([...url.searchParams.keys()].length === 0) url.search = ""; + return url.toString(); +} + +/** + * The dedupe identity of a URL. Not for publishing — this intentionally + * rewrites the URL into a shape that compares well. + */ +export function normalizeUrlForIdentity(raw: string): string | null { + const url = parse(raw); + if (!url) return null; + dropTracking(url); + url.hash = ""; + url.protocol = "https:"; + url.hostname = url.hostname.toLowerCase().replace(/^www\./, ""); + // Default ports carry no meaning once the scheme is fixed. + if (url.port === "80" || url.port === "443") url.port = ""; + if (url.pathname.length > 1 && url.pathname.endsWith("/")) { + url.pathname = url.pathname.replace(/\/+$/, ""); + } + if ([...url.searchParams.keys()].length === 0) { + url.search = ""; + } else { + url.searchParams.sort(); + } + return url.toString(); +} + +/** + * Stable dedupe key: sha256 of the identity form. Returns null for input that + * is not a usable http(s) URL, so callers can reject rather than store a hash + * of garbage. + */ +export function urlHash(raw: string): string | null { + const identity = normalizeUrlForIdentity(raw); + if (!identity) return null; + return createHash("sha256").update(identity).digest("hex"); +} + +/** + * Titles drift ("Foo — Bar" vs "Foo - Bar"), so same-story detection across + * different URLs compares a flattened title instead of the raw one. + */ +export function normalizedTitleHash(title: string | null | undefined): string | null { + const flat = (title ?? "") + .toLowerCase() + .normalize("NFKD") + .replace(/[\u0300-\u036f]/g, "") + .replace(/[^a-z0-9]+/g, " ") + .trim(); + if (!flat) return null; + return createHash("sha256").update(flat).digest("hex"); +} diff --git a/lib/promote/selectLink.ts b/lib/promote/selectLink.ts new file mode 100644 index 00000000..c16eb3f8 --- /dev/null +++ b/lib/promote/selectLink.ts @@ -0,0 +1,129 @@ +// Pick the next link a Promote list should post, honouring its blend. +// +// The old rule was one line: the least recently promoted enabled link. That is +// still the rule *within* an ownership class — it is what keeps a single link +// from dominating — but which class to draw from is now a blend decision. + +import type { SupabaseClient } from "@supabase/supabase-js"; +import { + chooseOwnership, + parseFallback, + parseMix, + type BlendDecision, + type Ownership, + OWNERSHIPS, +} from "@/lib/promote/blend"; + +// How many recent posts the ratio is measured over. Long enough to be a ratio, +// short enough that changing the mix takes effect within a day of posting. +const BLEND_WINDOW = 50; + +export type SelectableLink = { + id: string; + url: string; + title: string | null; + angle: string | null; + summary: string | null; + source_name: string | null; + ownership: Ownership; + source_id: string | null; + times_promoted: number | null; +}; + +export type LinkSelection = { + link: SelectableLink | null; + decision: BlendDecision; +}; + +type ListLike = { + id: string; + source_mix?: unknown; + fallback_policy?: unknown; +}; + +/** + * Choose the next link for a list. Returns the blend decision alongside it so + * the caller can record *why* this link was picked. + */ +export async function selectNextLink( + supabase: SupabaseClient, + list: ListLike, + now: Date = new Date(), +): Promise { + const mix = parseMix(list.source_mix); + const fallback = parseFallback(list.fallback_policy); + + // What the list has actually posted lately, per class. + const { data: recent } = await supabase + .from("promo_post") + .select("ownership") + .eq("list_id", list.id) + .in("status", ["posted", "pending"]) + .order("created_at", { ascending: false }) + .limit(BLEND_WINDOW); + + const posted: Partial> = {}; + for (const row of (recent ?? []) as Array<{ ownership: string | null }>) { + // Posts made before sources existed carry no ownership; they were all the + // user's own hand-pasted links, so they count as owned. + const key = (row.ownership ?? "owned") as Ownership; + if (OWNERSHIPS.includes(key)) posted[key] = (posted[key] ?? 0) + 1; + } + + // The best candidate in each class: least recently promoted first, so the + // rotation stays fair inside the class. + const candidates: Partial> = {}; + const available: Partial> = {}; + await Promise.all( + OWNERSHIPS.map(async (ownership) => { + const { data } = await supabase + .from("promo_link") + .select( + "id, url, title, angle, summary, source_name, ownership, source_id, times_promoted", + ) + .eq("list_id", list.id) + .eq("enabled", true) + .eq("ownership", ownership) + .order("last_promoted_at", { ascending: true, nullsFirst: true }) + .limit(1); + const row = (data ?? [])[0] as SelectableLink | undefined; + if (row) { + candidates[ownership] = row; + available[ownership] = true; + } + }), + ); + + const fallbackUsedToday = await countFallbackToday(supabase, list.id, now); + + const decision = chooseOwnership({ + mix, + posted, + available, + fallback, + fallbackUsedToday, + }); + + return { + link: decision.ownership ? (candidates[decision.ownership] ?? null) : null, + decision, + }; +} + +async function countFallbackToday( + supabase: SupabaseClient, + listId: string, + now: Date, +): Promise { + // A rolling 24 hours rather than a calendar day: the list has a timezone but + // the cap is about pacing, and a rolling window cannot be gamed by a + // midnight boundary. + const since = new Date(now.getTime() - 24 * 60 * 60 * 1000).toISOString(); + const { count } = await supabase + .from("promo_post") + .select("id", { count: "exact", head: true }) + .eq("list_id", listId) + .eq("via_fallback", true) + .gte("created_at", since); + return count ?? 0; +} diff --git a/lib/promote/sources.ts b/lib/promote/sources.ts new file mode 100644 index 00000000..1b492e58 --- /dev/null +++ b/lib/promote/sources.ts @@ -0,0 +1,244 @@ +// Creating and validating Promote content sources. +// +// The shared registry is what makes this cheap at scale: adding "bitcoin" to a +// list does not create a feed, it *joins* one. The second user to track +// bitcoin reuses the first user's promo_feed row and starts from the items it +// has already collected. + +import type { SupabaseClient } from "@supabase/supabase-js"; +import { parseFeed } from "@/lib/promote/feedParse"; +import { topicFeedUrl, type KeywordInput } from "@/lib/promote/keywords"; +import type { FetchLike } from "@/lib/promote/ingest"; + +// Same guard the audit engine applies to user-supplied targets: a feed URL is +// fetched by our server, so it must not be able to point at our own network. +// Wider than the audit copy — it also covers the 172.16/12 private range. +const PRIVATE_HOSTS = + /^(localhost|127\.|10\.|192\.168\.|169\.254\.|172\.(1[6-9]|2\d|3[01])\.|::1|fc00:|fd00:|.*\.local)$/i; + +const VALIDATE_TIMEOUT_MS = 15_000; + +export type FeedKind = "rssamplifier_topic" | "custom_feed" | "project_feed"; + +export type SourceValidation = + | { ok: true; feedUrl: string; title: string | null; itemCount: number } + | { ok: false; error: string }; + +/** + * Normalize a user-supplied feed URL, or explain why it is unusable. + * Bare hostnames are accepted and assumed https, the way users paste them. + */ +export function normalizeFeedUrl(raw: string): { ok: true; url: string } | { ok: false; error: string } { + const trimmed = (raw ?? "").trim(); + if (!trimmed) return { ok: false, error: "Enter a feed URL." }; + + let url: URL; + try { + url = new URL(/^https?:\/\//i.test(trimmed) ? trimmed : `https://${trimmed}`); + } catch { + return { ok: false, error: "That is not a valid URL." }; + } + if (url.protocol !== "http:" && url.protocol !== "https:") { + return { ok: false, error: "Only http and https feeds are supported." }; + } + if (PRIVATE_HOSTS.test(url.hostname)) { + return { ok: false, error: "That address is not reachable from the public internet." }; + } + url.hash = ""; + return { ok: true, url: url.toString() }; +} + +/** + * Fetch a candidate feed and confirm it parses into entries. + * + * Adding a source that turns out to be an HTML page is a mistake worth + * catching while the user is still looking at the form, rather than leaving + * them to wonder why a campaign never posts. + */ +export async function validateFeedUrl( + raw: string, + fetchImpl?: FetchLike, +): Promise { + const normalized = normalizeFeedUrl(raw); + if (!normalized.ok) return normalized; + + const doFetch = fetchImpl ?? (globalThis.fetch as unknown as FetchLike); + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), VALIDATE_TIMEOUT_MS); + try { + const response = await doFetch(normalized.url, { + headers: { + accept: + "application/rss+xml,application/atom+xml,application/xml,text/xml,*/*;q=0.5", + "user-agent": "CrawlProofPromote/1.0 (+https://crawlproof.com)", + }, + signal: controller.signal, + }); + if (!response.ok) { + return { + ok: false, + error: + response.status === 404 + ? "No feed at that address (404)." + : `That address returned HTTP ${response.status}.`, + }; + } + const body = await response.text(); + const parsed = parseFeed(body, normalized.url); + if (parsed.items.length === 0) { + return { ok: false, error: "That address is reachable but is not an RSS or Atom feed." }; + } + return { + ok: true, + feedUrl: normalized.url, + title: parsed.title, + itemCount: parsed.items.length, + }; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + return { + ok: false, + error: message.includes("abort") ? "That feed took too long to respond." : message, + }; + } finally { + clearTimeout(timer); + } +} + +/** + * Find or create the shared registry row for a feed URL. Concurrent callers + * are safe: feed_url is unique, so a lost race re-reads the winner's row. + */ +export async function ensureFeed( + supabase: SupabaseClient, + input: { feedUrl: string; kind: FeedKind; topicSlug?: string | null; title?: string | null }, +): Promise<{ id: string; created: boolean } | null> { + const { data: existing } = await supabase + .from("promo_feed") + .select("id") + .eq("feed_url", input.feedUrl) + .maybeSingle(); + if (existing) return { id: existing.id as string, created: false }; + + const { data: inserted, error } = await supabase + .from("promo_feed") + .insert({ + feed_url: input.feedUrl, + kind: input.kind, + topic_slug: input.topicSlug ?? null, + title: input.title ?? null, + // Due immediately: a new feed should have items before the user has + // finished reading the confirmation. + next_fetch_at: new Date().toISOString(), + }) + .select("id") + .single(); + + if (inserted) return { id: inserted.id as string, created: true }; + + // Unique-violation: somebody else created it between our read and write. + if (error) { + const { data: raced } = await supabase + .from("promo_feed") + .select("id") + .eq("feed_url", input.feedUrl) + .maybeSingle(); + if (raced) return { id: raced.id as string, created: false }; + } + return null; +} + +export type KeywordSourceOutcome = { + keyword: string; + slug: string; + ok: boolean; + /** Set when the topic could not be added. */ + error?: string; + sourceId?: string; + feedId?: string; +}; + +/** + * Turn a parsed keyword list into one source per keyword. + * + * Every keyword is reported on individually: "bitcoin and ethereum were added, + * zzz is not a topic yet" is far more useful than one blanket failure. + */ +export async function addKeywordSources( + supabase: SupabaseClient, + input: { + listId: string; + keywords: KeywordInput[]; + ownership: "owned" | "partner" | "shared"; + }, + fetchImpl?: FetchLike, +): Promise { + const outcomes: KeywordSourceOutcome[] = []; + + for (const keyword of input.keywords) { + const feedUrl = topicFeedUrl(keyword.slug); + const validation = await validateFeedUrl(feedUrl, fetchImpl); + if (!validation.ok) { + outcomes.push({ + keyword: keyword.label, + slug: keyword.slug, + ok: false, + error: + validation.error.includes("404") + ? `No RSS Amplifier topic for "${keyword.label}" yet.` + : validation.error, + }); + continue; + } + + const feed = await ensureFeed(supabase, { + feedUrl, + kind: "rssamplifier_topic", + topicSlug: keyword.slug, + title: validation.title, + }); + if (!feed) { + outcomes.push({ + keyword: keyword.label, + slug: keyword.slug, + ok: false, + error: "Could not register that topic feed.", + }); + continue; + } + + const { data: source, error } = await supabase + .from("promo_source") + .insert({ + list_id: input.listId, + feed_id: feed.id, + type: "rssamplifier_topic", + ownership: input.ownership, + label: keyword.label, + keyword: keyword.label, + }) + .select("id") + .single(); + + if (error || !source) { + // unique (list_id, feed_id): the list already tracks this topic. + outcomes.push({ + keyword: keyword.label, + slug: keyword.slug, + ok: false, + error: "Already tracked by this campaign.", + }); + continue; + } + + outcomes.push({ + keyword: keyword.label, + slug: keyword.slug, + ok: true, + sourceId: source.id as string, + feedId: feed.id, + }); + } + + return outcomes; +} diff --git a/lib/promote/sweep.ts b/lib/promote/sweep.ts index 0a5173f1..e18099d5 100644 --- a/lib/promote/sweep.ts +++ b/lib/promote/sweep.ts @@ -6,6 +6,8 @@ import type { SupabaseClient } from "@supabase/supabase-js"; import type Anthropic from "@anthropic-ai/sdk"; 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 { postViaAccount, type PostResult } from "@/lib/sp/post"; export type PromoteSweepClients = { @@ -35,6 +37,8 @@ type PromoList = { quiet_start: number | null; quiet_end: number | null; timezone: string | null; + source_mix?: unknown; + fallback_policy?: unknown; }; type PromoLink = { @@ -42,6 +46,10 @@ type PromoLink = { url: string; title: string | null; angle: string | null; + summary?: string | null; + source_name?: string | null; + ownership?: Ownership | null; + times_promoted?: number | null; }; type SpAccount = { @@ -73,7 +81,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", + "id, user_id, name, cadence_seconds, post_mode, target_account_ids, brand_voice, quiet_start, quiet_end, timezone, source_mix, fallback_policy", ) .eq("status", "running") .lte("next_run_at", new Date().toISOString()) @@ -139,21 +147,26 @@ async function processOneList( return out; } - // Pick the next link(s) using round-robin (least-recently-promoted first) - const { data: links } = await supabase - .from("promo_link") - .select("id, url, title, angle") - .eq("list_id", list.id) - .eq("enabled", true) - .order("last_promoted_at", { ascending: true, nullsFirst: true }) - .limit(1); - - if (!links || links.length === 0) { + // Pick the next link. Within an ownership class this is still + // least-recently-promoted-first round-robin; which class to draw from is the + // list's blend decision (70% our content, 30% industry content, and so on). + const selection = await selectNextLink(supabase, list); + if (!selection.link) { + // Nothing eligible: an empty list, or a blend whose target class is starved + // and whose fallback policy says not to cover for it. Either way this is a + // quiet no-op, not a failure — the next tick tries again. + if (selection.decision.reason !== "no_inventory") { + console.log( + `[promote] list ${list.id} skipped a tick: ${selection.decision.reason}`, + ); + } await advanceScheduler(supabase, list); return out; } - const link = links[0] as PromoLink; + const link = selection.link as PromoLink; + const ownership: Ownership = (link.ownership ?? "owned") as Ownership; + const viaFallback = selection.decision.viaFallback; // Determine which (link, account) pairs to post this tick const postTargets: Array<{ link: PromoLink; account: SpAccount }> = []; @@ -229,6 +242,9 @@ async function processOneList( recentBodies, anthropic: clients.anthropic, openai: clients.openai, + summary: targetLink.summary ?? null, + sourceName: targetLink.source_name ?? null, + ownership, }); // Publish via the existing sp platform layer @@ -254,6 +270,9 @@ async function processOneList( 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, @@ -289,6 +308,9 @@ async function processOneList( 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", error: message, @@ -303,14 +325,14 @@ async function processOneList( } } - // Update the link's round-robin cursor + // Update the link's round-robin cursor. times_promoted is a real counter now + // that selection reads it back — it used to be re-stamped to 1 every tick, + // because the old query never selected the column it was incrementing. await supabase .from("promo_link") .update({ last_promoted_at: new Date().toISOString(), - times_promoted: (link as any).times_promoted - ? (link as any).times_promoted + 1 - : 1, + times_promoted: (link.times_promoted ?? 0) + 1, }) .eq("id", link.id); diff --git a/supabase/migrations/20260818190000_promote_sources.sql b/supabase/migrations/20260818190000_promote_sources.sql new file mode 100644 index 00000000..48193cce --- /dev/null +++ b/supabase/migrations/20260818190000_promote_sources.sql @@ -0,0 +1,273 @@ +-- Promote content sources: campaigns that feed themselves. +-- +-- Until now a Promote list was a hand-pasted set of links. The user typed 20 +-- URLs and the drip engine rotated through them forever. That works for a fixed +-- set of product pages and nothing else: it cannot promote what the user +-- published this morning, and it gives a user with no back catalogue nothing to +-- post at all. +-- +-- A *source* is a standing subscription that keeps supplying links: +-- +-- rssamplifier_topic a keyword. "bitcoin" becomes the RSS Amplifier topic +-- feed https://rssamplifier.com/topics/bitcoin.rss. +-- One keyword is one source — several keywords never +-- collapse into a single ambiguous URL. +-- custom_feed any RSS or Atom URL the user owns or follows. +-- manual_url the existing hand-pasted behaviour, now named. +-- +-- Three tables, because the fetch and the subscription are different things: +-- +-- promo_feed one row per feed URL, shared by every list that subscribes +-- to it. Two hundred users tracking "bitcoin" poll RSS +-- Amplifier once between them, not two hundred times. This +-- is the whole reason the registry is keyed on the URL and +-- carries no user_id. +-- promo_feed_item the normalized entries of a feed, also shared. Fetched +-- once, then fanned out by reference. +-- promo_source one list's subscription to one feed, with the ownership +-- classification that drives blend ratios. +-- +-- Items still land in promo_link, which stays the unit the drip engine rotates +-- through. A source-fed link is a promo_link with source_id set; a hand-pasted +-- one has source_id null. Nothing about the existing engine changes shape. + +-- ---------- promo_feed: the shared fetch registry ---------- +create table if not exists public.promo_feed ( + id uuid primary key default gen_random_uuid(), + + -- The identity of the registry. One row per feed, globally. + feed_url text not null unique, + + kind text not null default 'custom_feed' + check (kind in ('rssamplifier_topic', 'custom_feed', 'project_feed')), + + -- Set for rssamplifier_topic feeds: the topic slug the URL was built from. + topic_slug text, + + -- The feed's own , used to attribute shared content. + title text, + + -- Conditional-request state, so a feed that has not changed costs us a 304 + -- rather than a parse. + etag text, + last_modified text, + + -- Scheduler bookkeeping. next_fetch_at is the claim point. + fetch_interval_seconds int not null default 900 + check (fetch_interval_seconds between 300 and 86400), + next_fetch_at timestamptz not null default now(), + last_fetched_at timestamptz, + last_success_at timestamptz, + + -- A feed that keeps failing backs off and eventually stops being polled; + -- the sources that subscribe to it surface the error. + consecutive_failures int not null default 0, + last_error text, + + created_at timestamptz not null default now(), + updated_at timestamptz not null default now() +); + +-- The claim query: due feeds, oldest first. +create index if not exists promo_feed_due_idx + on public.promo_feed (next_fetch_at); + +-- ---------- promo_feed_item: normalized entries, fetched once ---------- +create table if not exists public.promo_feed_item ( + id uuid primary key default gen_random_uuid(), + feed_id uuid not null references public.promo_feed(id) on delete cascade, + + -- url is what we publish; normalized_url and url_hash are dedupe identity + -- only (scheme folded, www dropped, tracking parameters removed). They are + -- deliberately different values — see lib/promote/normalizeUrl.ts. + url text not null, + normalized_url text not null, + url_hash text not null, + + title text, + summary text, + image_url text, + author_name text, + + -- The originating publication. Aggregator feeds name it in <source>, which + -- is what shared content gets attributed to. + source_name text, + + -- The publisher's own id for the entry, when it gives one. + guid text, + + published_at timestamptz, + discovered_at timestamptz not null default now(), + + -- One entry per feed. The hash, not the raw URL, so a publisher re-emitting + -- the same story with a fresh campaign tag does not create a second item. + unique (feed_id, url_hash) +); + +create index if not exists promo_feed_item_feed_recent_idx + on public.promo_feed_item (feed_id, published_at desc nulls last); + +-- ---------- promo_source: a list's subscription ---------- +create table if not exists public.promo_source ( + id uuid primary key default gen_random_uuid(), + list_id uuid not null references public.promo_list(id) on delete cascade, + + -- Null for manual_url sources, which have no feed behind them. + feed_id uuid references public.promo_feed(id) on delete cascade, + + type text not null + check (type in ('rssamplifier_topic', 'custom_feed', 'manual_url', 'project_feed')), + + -- Drives blend ratios, attribution and fallback. Keyword sources default to + -- 'shared' because the content belongs to somebody else. + ownership text not null default 'shared' + check (ownership in ('owned', 'partner', 'shared')), + + -- What the user typed: the keyword, or a name for the feed. + label text not null, + + -- The display form of the keyword for topic sources ("Artificial + -- Intelligence"); promo_feed.topic_slug holds the normalized form. + keyword text, + + enabled boolean not null default true, + + -- Ceiling on how many new links one ingestion pass may import from this + -- source, so a feed with a 500-entry backlog cannot flood a list. + max_items_per_ingest int not null default 10 + check (max_items_per_ingest between 1 and 100), + + last_ingested_at timestamptz, + items_imported int not null default 0, + + created_at timestamptz not null default now(), + + -- A list subscribes to a given feed once. + unique (list_id, feed_id) +); + +create index if not exists promo_source_list_idx + on public.promo_source (list_id, enabled); +create index if not exists promo_source_feed_idx + on public.promo_source (feed_id, enabled); + +-- ---------- promo_link: provenance for source-fed links ---------- +alter table public.promo_link + add column if not exists source_id uuid references public.promo_source(id) on delete set null, + add column if not exists ownership text not null default 'owned' + check (ownership in ('owned', 'partner', 'shared')), + add column if not exists summary text, + add column if not exists image_url text, + add column if not exists author_name text, + add column if not exists source_name text, + add column if not exists normalized_url text, + add column if not exists url_hash text, + add column if not exists published_at timestamptz, + add column if not exists discovered_at timestamptz not null default now(); + +-- Hand-pasted links predate sources and are the user's own: 'owned' is the +-- right default for them, which is why the column defaults that way rather +-- than to 'shared'. + +-- Dedupe on identity as well as on the raw URL. Partial, because links that +-- predate this migration have no hash and must not collide with each other. +create unique index if not exists promo_link_list_hash_idx + on public.promo_link (list_id, url_hash) + where url_hash is not null; + +create index if not exists promo_link_source_idx + on public.promo_link (source_id); + +-- Selection reads "enabled links of this list by ownership, least recently +-- promoted first" on every tick. +create index if not exists promo_link_blend_idx + on public.promo_link (list_id, enabled, ownership, last_promoted_at); + +-- ---------- promo_post: what the blend actually did ---------- +-- Ownership is denormalized onto the post the same way platform already is, +-- because the selector reads "what have I posted lately" on every tick and +-- must not join back through promo_link — a link can be deleted, and the +-- history of the ratio has to survive that. +alter table public.promo_post + add column if not exists ownership text, + add column if not exists source_id uuid references public.promo_source(id) on delete set null, + -- True when the blend could not honour its target and drew from the other + -- side instead. Counted against fallback_policy.maxFallbackItemsPerDay. + add column if not exists via_fallback boolean not null default false; + +create index if not exists promo_post_blend_idx + on public.promo_post (list_id, created_at desc); + +-- ---------- promo_list: blend and fallback ---------- +alter table public.promo_list + -- Relative weights per ownership class. A 70/30 list converges on roughly + -- seven owned links for every three shared ones over a rolling window. + add column if not exists source_mix jsonb not null + default '{"owned": 70, "partner": 0, "shared": 30}'::jsonb, + + -- What to do when one side of the blend has nothing available. The default + -- keeps a list with no original content posting, while capping how far it + -- may drift into being a pure shared-content firehose. + add column if not exists fallback_policy jsonb not null + default '{"whenOwnedQueueEmpty": "use_shared", "whenSharedQueueEmpty": "use_owned", "maxFallbackItemsPerDay": 3}'::jsonb; + +-- ---------- RLS ---------- +-- promo_feed and promo_feed_item are shared across users and carry no +-- user_id, so they are readable only through a subscription the caller owns. +-- The worker uses the service role and bypasses all of this. + +alter table public.promo_feed enable row level security; +create policy "promo_feed readable via a subscribed list" + on public.promo_feed for select + using ( + exists ( + select 1 + from public.promo_source s + join public.promo_list l on l.id = s.list_id + where s.feed_id = promo_feed.id and l.user_id = auth.uid() + ) + ); + +alter table public.promo_feed_item enable row level security; +create policy "promo_feed_item readable via a subscribed list" + on public.promo_feed_item for select + using ( + exists ( + select 1 + from public.promo_source s + join public.promo_list l on l.id = s.list_id + where s.feed_id = promo_feed_item.feed_id and l.user_id = auth.uid() + ) + ); + +alter table public.promo_source enable row level security; +create policy "promo_source via owned list" + on public.promo_source for all + using ( + exists ( + select 1 from public.promo_list l + where l.id = list_id and l.user_id = auth.uid() + ) + ) + with check ( + exists ( + select 1 from public.promo_list l + where l.id = list_id and l.user_id = auth.uid() + ) + ); + +-- ---------- updated_at ---------- +create trigger promo_feed_updated_at + before update on public.promo_feed + for each row execute function public.promo_set_updated_at(); + +-- ---------- Grants ---------- +-- Feeds and their items are never written from the browser: only the worker +-- ingests, and only server actions (service role) create feed rows. +grant select on public.promo_feed to authenticated; +grant select on public.promo_feed_item to authenticated; +grant select, insert, update, delete on public.promo_source to authenticated; + +grant all on public.promo_feed to service_role; +grant all on public.promo_feed_item to service_role; +grant all on public.promo_source to service_role; diff --git a/tests/contract/mcp-promote.test.ts b/tests/contract/mcp-promote.test.ts index f008e671..1a3d0faa 100644 --- a/tests/contract/mcp-promote.test.ts +++ b/tests/contract/mcp-promote.test.ts @@ -23,6 +23,10 @@ describe("crawlproof MCP · promote module", () => { "generate_promo_post", "list_accounts", "post_to_socials", + "promote_add_feed_source", + "promote_add_keyword_source", + "promote_list_campaigns", + "promote_list_sources", "promote_url", ]); @@ -33,6 +37,15 @@ describe("crawlproof MCP · promote module", () => { expect(props?.url).toBeDefined(); expect(props?.account_ids).toBeDefined(); + // Content sources are reachable over MCP too, so an agent can build the + // same campaign a user would build in the dashboard. + const keywords = tools.find((t) => t.name === "promote_add_keyword_source"); + const keywordProps = (keywords?.inputSchema as { properties?: Record<string, unknown> }) + ?.properties; + expect(keywordProps?.campaign_id).toBeDefined(); + expect(keywordProps?.keywords).toBeDefined(); + expect(keywords?.description).toMatch(/one .*source per keyword/i); + await client.close(); await server.close(); }); diff --git a/tests/promote/blend.test.ts b/tests/promote/blend.test.ts new file mode 100644 index 00000000..d7847014 --- /dev/null +++ b/tests/promote/blend.test.ts @@ -0,0 +1,268 @@ +import { describe, it, expect } from "vitest"; +import { + chooseOwnership, + parseFallback, + parseMix, + rankByDeficit, + DEFAULT_MIX, + type BlendMix, + type FallbackPolicy, + type Ownership, +} from "@/lib/promote/blend"; + +const allAvailable = { owned: true, partner: true, shared: true }; + +const permissive: FallbackPolicy = { + whenOwnedQueueEmpty: "use_shared", + whenSharedQueueEmpty: "use_owned", + maxFallbackItemsPerDay: null, +}; + +describe("parseMix", () => { + it("reads a stored mix", () => { + expect(parseMix({ owned: 70, shared: 30 })).toEqual({ owned: 70, partner: 0, shared: 30 }); + }); + + it("falls back to the default when nothing is weighted", () => { + expect(parseMix({ owned: 0, shared: 0 })).toEqual(DEFAULT_MIX); + expect(parseMix(null)).toEqual(DEFAULT_MIX); + expect(parseMix("nonsense")).toEqual(DEFAULT_MIX); + }); + + it("ignores negative and non-numeric weights", () => { + expect(parseMix({ owned: -5, shared: "x", partner: 10 })).toEqual({ + owned: 0, + partner: 10, + shared: 0, + }); + }); +}); + +describe("parseFallback", () => { + it("keeps a valid policy", () => { + expect( + parseFallback({ + whenOwnedQueueEmpty: "pause", + whenSharedQueueEmpty: "use_any_available", + maxFallbackItemsPerDay: 5, + }), + ).toEqual({ + whenOwnedQueueEmpty: "pause", + whenSharedQueueEmpty: "use_any_available", + maxFallbackItemsPerDay: 5, + }); + }); + + it("preserves an explicit null cap as unlimited", () => { + expect(parseFallback({ maxFallbackItemsPerDay: null }).maxFallbackItemsPerDay).toBeNull(); + }); + + it("repairs unknown actions", () => { + expect(parseFallback({ whenOwnedQueueEmpty: "explode" }).whenOwnedQueueEmpty).toBe( + "use_shared", + ); + }); +}); + +describe("chooseOwnership", () => { + const mix: BlendMix = { owned: 70, partner: 0, shared: 30 }; + + it("opens a fresh list with the dominant class", () => { + const decision = chooseOwnership({ + mix, + posted: {}, + available: allAvailable, + fallback: permissive, + }); + expect(decision.ownership).toBe("owned"); + expect(decision.viaFallback).toBe(false); + expect(decision.reason).toBe("on_target"); + }); + + it("switches to the starved class once the leader is ahead of target", () => { + // 8 owned / 0 shared against a 70/30 target: shared is furthest behind. + const decision = chooseOwnership({ + mix, + posted: { owned: 8, shared: 0 }, + available: allAvailable, + fallback: permissive, + }); + expect(decision.ownership).toBe("shared"); + }); + + it("never draws from a class with no weight when the blend is satisfiable", () => { + const decision = chooseOwnership({ + mix, + posted: {}, + available: allAvailable, + fallback: permissive, + }); + expect(decision.ownership).not.toBe("partner"); + }); + + describe("when the owned queue is empty", () => { + const available = { owned: false, shared: true }; + + it("uses shared content by default, flagged as a fallback", () => { + const decision = chooseOwnership({ mix, posted: {}, available, fallback: permissive }); + expect(decision.ownership).toBe("shared"); + expect(decision.viaFallback).toBe(true); + expect(decision.reason).toBe("fallback"); + }); + + it("posts nothing when the policy says pause", () => { + const decision = chooseOwnership({ + mix, + posted: {}, + available, + fallback: { ...permissive, whenOwnedQueueEmpty: "pause" }, + }); + expect(decision.ownership).toBeNull(); + expect(decision.reason).toBe("fallback_disabled"); + }); + + it("stops once the daily fallback cap is reached", () => { + const fallback = { ...permissive, maxFallbackItemsPerDay: 3 }; + expect( + chooseOwnership({ mix, posted: {}, available, fallback, fallbackUsedToday: 2 }) + .ownership, + ).toBe("shared"); + const capped = chooseOwnership({ + mix, + posted: {}, + available, + fallback, + fallbackUsedToday: 3, + }); + expect(capped.ownership).toBeNull(); + expect(capped.reason).toBe("fallback_cap_reached"); + }); + + it("does not count against the cap when the target class is available", () => { + const decision = chooseOwnership({ + mix, + posted: {}, + available: allAvailable, + fallback: { ...permissive, maxFallbackItemsPerDay: 0 }, + fallbackUsedToday: 99, + }); + expect(decision.ownership).toBe("owned"); + expect(decision.viaFallback).toBe(false); + }); + }); + + it("falls back to owned when the shared queue is empty", () => { + const decision = chooseOwnership({ + mix, + posted: { owned: 9, shared: 0 }, + available: { owned: true, shared: false }, + fallback: permissive, + }); + expect(decision.ownership).toBe("owned"); + expect(decision.viaFallback).toBe(true); + }); + + it("reaches an unweighted class only under use_any_available", () => { + const onlyPartner = { owned: false, shared: false, partner: true }; + expect( + chooseOwnership({ mix, posted: {}, available: onlyPartner, fallback: permissive }) + .ownership, + ).toBeNull(); + expect( + chooseOwnership({ + mix, + posted: {}, + available: onlyPartner, + fallback: { ...permissive, whenOwnedQueueEmpty: "use_any_available" }, + }).ownership, + ).toBe("partner"); + }); + + it("posts nothing when no class has inventory", () => { + const decision = chooseOwnership({ + mix, + posted: {}, + available: {}, + fallback: permissive, + }); + expect(decision.ownership).toBeNull(); + expect(decision.reason).toBe("no_inventory"); + }); +}); + +describe("rankByDeficit", () => { + it("puts the class furthest below its target first", () => { + const mix: BlendMix = { owned: 70, partner: 0, shared: 30 }; + expect(rankByDeficit(mix, { owned: 10, shared: 0 })[0]).toBe("shared"); + expect(rankByDeficit(mix, { owned: 0, shared: 10 })[0]).toBe("owned"); + }); + + it("omits classes with no weight", () => { + expect(rankByDeficit({ owned: 100, partner: 0, shared: 0 }, {})).toEqual(["owned"]); + }); +}); + +describe("convergence — the acceptance criterion", () => { + // "A campaign can maintain a configured owned/shared publishing ratio." + // Simulate a campaign with unlimited inventory on both sides and check the + // realised ratio, which is what a user actually sees on their timeline. + function simulate(mix: BlendMix, ticks: number): Record<string, number> { + const posted: Partial<Record<Ownership, number>> = {}; + for (let i = 0; i < ticks; i++) { + const decision = chooseOwnership({ + mix, + posted, + available: allAvailable, + fallback: permissive, + }); + const key = decision.ownership!; + posted[key] = (posted[key] ?? 0) + 1; + } + return posted as Record<string, number>; + } + + it("converges on 70/30", () => { + const posted = simulate({ owned: 70, partner: 0, shared: 30 }, 100); + expect(posted.owned).toBe(70); + expect(posted.shared).toBe(30); + }); + + it("converges on 50/50", () => { + const posted = simulate({ owned: 50, partner: 0, shared: 50 }, 100); + expect(posted.owned).toBe(50); + expect(posted.shared).toBe(50); + }); + + it("converges on a three-way split", () => { + const posted = simulate({ owned: 50, partner: 20, shared: 30 }, 100); + expect(posted.owned).toBe(50); + expect(posted.partner).toBe(20); + expect(posted.shared).toBe(30); + }); + + it("never produces a long run of one class", () => { + // The failure mode weighted-random has: five shared posts in a row makes an + // account read as a content farm. + const mix: BlendMix = { owned: 70, partner: 0, shared: 30 }; + const posted: Partial<Record<Ownership, number>> = {}; + const sequence: Ownership[] = []; + for (let i = 0; i < 60; i++) { + const key = chooseOwnership({ + mix, + posted, + available: allAvailable, + fallback: permissive, + }).ownership!; + posted[key] = (posted[key] ?? 0) + 1; + sequence.push(key); + } + let longestRun = 1; + let run = 1; + for (let i = 1; i < sequence.length; i++) { + run = sequence[i] === sequence[i - 1] ? run + 1 : 1; + longestRun = Math.max(longestRun, run); + } + // 70/30 means owned legitimately posts twice in a row; four would be a bug. + expect(longestRun).toBeLessThanOrEqual(3); + }); +}); diff --git a/tests/promote/fake-supabase.ts b/tests/promote/fake-supabase.ts new file mode 100644 index 00000000..36fc6bd3 Binary files /dev/null and b/tests/promote/fake-supabase.ts differ diff --git a/tests/promote/feed-parse.test.ts b/tests/promote/feed-parse.test.ts new file mode 100644 index 00000000..c723d557 --- /dev/null +++ b/tests/promote/feed-parse.test.ts @@ -0,0 +1,138 @@ +import { describe, it, expect } from "vitest"; +import { readFileSync } from "node:fs"; +import { join } from "node:path"; +import { parseFeed, summarize } from "@/lib/promote/feedParse"; + +// A trimmed capture of the live https://rssamplifier.com/topics/bitcoin.rss — +// the exact shape a keyword source has to cope with, namespaces included. +const bitcoinFeed = readFileSync( + join(__dirname, "fixtures", "rssamplifier-bitcoin.rss"), + "utf8", +); + +describe("parseFeed — RSS Amplifier topic feed", () => { + const feed = parseFeed(bitcoinFeed, "https://rssamplifier.com/topics/bitcoin.rss"); + + it("reads the channel title", () => { + expect(feed.title).toBe("bitcoin — RSS Amplifier"); + }); + + it("returns every item in the document", () => { + expect(feed.items.length).toBe(4); + }); + + it("keeps the publisher's link, not the aggregator's", () => { + for (const item of feed.items) { + expect(item.url).toMatch(/^https?:\/\//); + expect(new URL(item.url).hostname).not.toBe("rssamplifier.com"); + } + }); + + it("captures title, guid and publish date", () => { + const first = feed.items[0]; + expect(first.title).toBe( + "Bitcoin catches a bid after weekly loss, on track for best day in over a month", + ); + expect(first.guid).toContain("investing.com"); + expect(first.publishedAt).toBe("2026-08-17T22:27:24.000Z"); + }); + + it("attributes the original publisher via dc:creator and <source>", () => { + const first = feed.items[0]; + expect(first.author).toBe("Investing.com"); + expect(first.sourceName).toBe("Cryptocurrency News"); + }); + + it("picks up a media:thumbnail image", () => { + expect(feed.items[0].imageUrl).toBe( + "https://content-media.investing.com/news/moved_LYNXMPEKA01G9_L.jpg", + ); + }); +}); + +describe("parseFeed — Atom", () => { + const atom = `<?xml version="1.0" encoding="utf-8"?> + <feed xmlns="http://www.w3.org/2005/Atom"> + <title>Example Blog + + Hello & welcome + + tag:example.com,2026:post-1 + 2026-08-01T10:00:00Z + Ada + <p>A short <b>intro</b>.</p> + + `; + + const feed = parseFeed(atom, "https://example.com/feed.xml"); + + it("resolves relative entry links against the feed URL", () => { + expect(feed.items[0].url).toBe("https://example.com/posts/hello"); + }); + + it("decodes entities in titles", () => { + expect(feed.items[0].title).toBe("Hello & welcome"); + }); + + it("strips markup out of the summary", () => { + expect(feed.items[0].summary).toBe("A short intro ."); + }); + + it("reads the atom author name and id", () => { + expect(feed.items[0].author).toBe("Ada"); + expect(feed.items[0].guid).toBe("tag:example.com,2026:post-1"); + }); +}); + +describe("parseFeed — hostile input", () => { + it("returns no items for a non-feed document rather than throwing", () => { + expect(parseFeed("not a feed", "https://a.com").items).toEqual( + [], + ); + }); + + it("returns no items for empty input", () => { + expect(parseFeed("", "https://a.com").items).toEqual([]); + }); + + it("skips entries whose link is unusable", () => { + const xml = ` + No link + Bad schemejavascript:alert(1) + Goodhttps://ok.example/post + `; + const items = parseFeed(xml, "https://a.com").items; + expect(items.map((i) => i.url)).toEqual(["https://ok.example/post"]); + }); + + it("honours the item cap", () => { + const items = Array.from( + { length: 30 }, + (_, i) => `https://e.com/${i}`, + ).join(""); + expect(parseFeed(`${items}`, "https://a.com", 10).items.length).toBe( + 10, + ); + }); +}); + +describe("summarize", () => { + it("returns null for empty prose", () => { + expect(summarize(null)).toBeNull(); + expect(summarize(" ")).toBeNull(); + }); + + it("drops script and style bodies", () => { + expect(summarize("Real text")).toBe( + "Real text", + ); + }); + + it("truncates long prose on a word boundary", () => { + const long = "word ".repeat(400); + const out = summarize(long)!; + expect(out.length).toBeLessThanOrEqual(601); + expect(out.endsWith("…")).toBe(true); + expect(out).not.toMatch(/wor…$/); + }); +}); diff --git a/tests/promote/fixtures/rssamplifier-bitcoin.rss b/tests/promote/fixtures/rssamplifier-bitcoin.rss new file mode 100644 index 00000000..8471b65e --- /dev/null +++ b/tests/promote/fixtures/rssamplifier-bitcoin.rss @@ -0,0 +1,53 @@ + + + + bitcoin — RSS Amplifier + https://rssamplifier.com/topics/bitcoin + Recent posts from the 345 feeds in the RSS Amplifier directory that cover bitcoin. + en + RSS Amplifier + Mon, 17 Aug 2026 22:27:24 GMT + + + Bitcoin catches a bid after weekly loss, on track for best day in over a month + https://www.investing.com/news/cryptocurrency-news/bitcoin-steadies-at-635k-iran-tensions-us-regulations-in-focus-4862213 + https://www.investing.com/news/cryptocurrency-news/bitcoin-steadies-at-635k-iran-tensions-us-regulations-in-focus-4862213 + Mon, 17 Aug 2026 22:27:24 GMT + Investing.com + Cryptocurrency News + + + + Bitcoin tests $65,600 resistance with weak trend: Live levels + https://www.investing.com/news/cryptocurrency-news/bitcoin-coiled-at-63591-with-volatility-at-6month-lows-live-levels-93CH-4862234 + https://www.investing.com/news/cryptocurrency-news/bitcoin-coiled-at-63591-with-volatility-at-6month-lows-live-levels-93CH-4862234 + Mon, 17 Aug 2026 19:19:17 GMT + Investing.com + Cryptocurrency News + + + + 🟠 O Bradesco vai vender Bitcoin + https://ascencriptonewsletter.substack.com/p/o-bradesco-vai-vender-bitcoin + https://ascencriptonewsletter.substack.com/p/o-bradesco-vai-vender-bitcoin + Mon, 17 Aug 2026 17:33:00 GMT + MAIS: Vazamento expõe 678 mil contribuintes na França | Setor quer supervisão do BC sobre stablecoins | Trezor alerta 13.689 clientes sobre phishing | Semana em Cripto + Ascen Cripto Newsletter + + + + Bitcoin hits $64K as gold gains while oil shakes off Trump Oman threat + https://cointelegraph.com/markets/bitcoin-hits-64k-as-gold-gains-while-oil-shakes-off-trump-oman-threat + https://cointelegraph.com/markets/bitcoin-hits-64k-as-gold-gains-while-oil-shakes-off-trump-oman-threat?utm_source=rss_feed&utm_medium=rss&utm_campaign=rss_partner_inbound + Mon, 17 Aug 2026 16:16:40 GMT + Bitcoin price strength saw BTC/USD pass $64,000 on 2% daily gains as gold pushed higher while US stocks wobbled on fresh US-Iran rhetoric. + Cointelegraph by William Suberg + Cointelegraph.com News + + + + diff --git a/tests/promote/ingest.test.ts b/tests/promote/ingest.test.ts new file mode 100644 index 00000000..c16d39b0 --- /dev/null +++ b/tests/promote/ingest.test.ts @@ -0,0 +1,298 @@ +import { describe, it, expect, beforeEach } from "vitest"; +import { readFileSync } from "node:fs"; +import { join } from "node:path"; +import { ingestDueFeeds, nextFetchAt, type FetchLike } from "@/lib/promote/ingest"; +import { makeFakeSupabase, resetIds, type FakeDb } from "./fake-supabase"; + +const bitcoinFeed = readFileSync( + join(__dirname, "fixtures", "rssamplifier-bitcoin.rss"), + "utf8", +); + +const NOW = new Date("2026-08-18T12:00:00.000Z"); + +// Unique constraints the real schema enforces, and that the fan-out relies on. +const CONSTRAINTS = [ + { table: "promo_feed_item", columns: ["feed_id", "url_hash"] }, + { table: "promo_link", columns: ["list_id", "url_hash"] }, +]; + +function response(body: string, init?: { status?: number; headers?: Record }) { + const headers = init?.headers ?? {}; + return { + ok: (init?.status ?? 200) < 400, + status: init?.status ?? 200, + headers: { get: (name: string) => headers[name.toLowerCase()] ?? null }, + text: async () => body, + }; +} + +function seed(overrides: Partial = {}): FakeDb { + return { + promo_feed: [ + { + id: "feed-bitcoin", + feed_url: "https://rssamplifier.com/topics/bitcoin.rss", + kind: "rssamplifier_topic", + etag: null, + last_modified: null, + fetch_interval_seconds: 900, + consecutive_failures: 0, + next_fetch_at: "2026-08-18T11:00:00.000Z", + }, + ], + promo_feed_item: [], + promo_source: [ + { + id: "source-a", + list_id: "list-a", + feed_id: "feed-bitcoin", + ownership: "shared", + enabled: true, + max_items_per_ingest: 10, + items_imported: 0, + }, + ], + promo_link: [], + ...overrides, + }; +} + +beforeEach(() => resetIds()); + +describe("ingestDueFeeds", () => { + it("fetches a due feed, stores its items and fans them out", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed, { headers: { etag: 'W/"abc"' } }); + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.feedsChecked).toBe(1); + expect(result.feedsFailed).toBe(0); + expect(result.itemsStored).toBe(4); + expect(result.linksCreated).toBe(4); + expect(db.promo_feed_item).toHaveLength(4); + expect(db.promo_link).toHaveLength(4); + }); + + it("stamps imported links with the subscription's ownership and provenance", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed); + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + const link = db.promo_link[0]; + expect(link.list_id).toBe("list-a"); + expect(link.source_id).toBe("source-a"); + expect(link.ownership).toBe("shared"); + expect(link.url_hash).toMatch(/^[a-f0-9]{64}$/); + expect(link.enabled).toBe(true); + // Attribution survives the trip, so shared content can be credited. + expect(link.source_name).toBeTruthy(); + }); + + it("stores the validators it was given, so the next poll is conditional", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => + response(bitcoinFeed, { + headers: { etag: 'W/"abc"', "last-modified": "Mon, 17 Aug 2026 22:27:24 GMT" }, + }); + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(db.promo_feed[0].etag).toBe('W/"abc"'); + expect(db.promo_feed[0].last_modified).toBe("Mon, 17 Aug 2026 22:27:24 GMT"); + }); + + it("sends the stored validators back on the next poll", async () => { + const withEtag = seed(); + withEtag.promo_feed[0].etag = 'W/"abc"'; + withEtag.promo_feed[0].last_modified = "Mon, 17 Aug 2026 22:27:24 GMT"; + const { client } = makeFakeSupabase(withEtag, CONSTRAINTS); + + const sent: Record[] = []; + const fetchImpl: FetchLike = async (_url, init) => { + sent.push(init?.headers ?? {}); + return response("", { status: 304 }); + }; + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(sent[0]["if-none-match"]).toBe('W/"abc"'); + expect(sent[0]["if-modified-since"]).toBe("Mon, 17 Aug 2026 22:27:24 GMT"); + }); + + it("treats 304 as success and does no work", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response("", { status: 304 }); + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.feedsUnchanged).toBe(1); + expect(result.feedsFailed).toBe(0); + expect(result.itemsStored).toBe(0); + expect(db.promo_link).toHaveLength(0); + expect(db.promo_feed[0].last_success_at).toBe(NOW.toISOString()); + }); + + it("does not import the same story twice across polls", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed); + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + // Feeds re-list the same entries every poll; make the feed due again. + db.promo_feed[0].next_fetch_at = "2026-08-18T11:00:00.000Z"; + const second = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(second.linksCreated).toBe(0); + expect(db.promo_link).toHaveLength(4); + expect(db.promo_feed_item).toHaveLength(4); + }); + + it("fetches once and fans out to every subscribing list", async () => { + const twoLists = seed(); + twoLists.promo_source.push({ + id: "source-b", + list_id: "list-b", + feed_id: "feed-bitcoin", + ownership: "shared", + enabled: true, + max_items_per_ingest: 10, + items_imported: 0, + }); + const { client, db } = makeFakeSupabase(twoLists, CONSTRAINTS); + + let fetches = 0; + const fetchImpl: FetchLike = async () => { + fetches++; + return response(bitcoinFeed); + }; + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + // The whole point of the shared registry. + expect(fetches).toBe(1); + expect(result.linksCreated).toBe(8); + expect(db.promo_link.filter((l) => l.list_id === "list-a")).toHaveLength(4); + expect(db.promo_link.filter((l) => l.list_id === "list-b")).toHaveLength(4); + }); + + it("skips disabled subscriptions", async () => { + const paused = seed(); + paused.promo_source[0].enabled = false; + const { client, db } = makeFakeSupabase(paused, CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed); + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.itemsStored).toBe(4); + expect(result.linksCreated).toBe(0); + expect(db.promo_link).toHaveLength(0); + }); + + it("honours a source's per-ingest cap", async () => { + const capped = seed(); + capped.promo_source[0].max_items_per_ingest = 2; + const { client, db } = makeFakeSupabase(capped, CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed); + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(db.promo_link).toHaveLength(2); + }); + + it("records the running import count on the source", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response(bitcoinFeed); + + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(db.promo_source[0].items_imported).toBe(4); + expect(db.promo_source[0].last_ingested_at).toBe(NOW.toISOString()); + }); + + it("leaves a feed that is not due yet alone", async () => { + const notDue = seed(); + notDue.promo_feed[0].next_fetch_at = "2026-08-18T13:00:00.000Z"; + const { client } = makeFakeSupabase(notDue, CONSTRAINTS); + + let fetches = 0; + const fetchImpl: FetchLike = async () => { + fetches++; + return response(bitcoinFeed); + }; + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + expect(result.feedsChecked).toBe(0); + expect(fetches).toBe(0); + }); + + it("claims a feed before fetching, so an overlapping pass cannot double-fetch", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => { + // Mid-fetch, the row must already be pushed past `now`. + expect(new Date(db.promo_feed[0].next_fetch_at).getTime()).toBeGreaterThan( + NOW.getTime(), + ); + return response(bitcoinFeed); + }; + await ingestDueFeeds(client, { fetchImpl, now: NOW }); + }); + + describe("failure handling", () => { + it("records an HTTP error and backs the feed off", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response("nope", { status: 500 }); + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.feedsFailed).toBe(1); + expect(db.promo_feed[0].consecutive_failures).toBe(1); + expect(db.promo_feed[0].last_error).toContain("500"); + }); + + it("treats a page that is not a feed as a failure", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => response("hi"); + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.feedsFailed).toBe(1); + expect(db.promo_feed[0].last_error).toContain("no usable entries"); + }); + + it("survives a thrown fetch without losing the pass", async () => { + const { client, db } = makeFakeSupabase(seed(), CONSTRAINTS); + const fetchImpl: FetchLike = async () => { + throw new Error("ECONNRESET"); + }; + + const result = await ingestDueFeeds(client, { fetchImpl, now: NOW }); + + expect(result.feedsFailed).toBe(1); + expect(db.promo_feed[0].last_error).toBe("ECONNRESET"); + }); + }); +}); + +describe("nextFetchAt", () => { + const feed = { fetch_interval_seconds: 900, consecutive_failures: 0 }; + + it("uses the plain interval while healthy", () => { + expect(nextFetchAt(feed, NOW, false)).toBe("2026-08-18T12:15:00.000Z"); + }); + + it("backs off geometrically while failing", () => { + expect(nextFetchAt({ ...feed, consecutive_failures: 1 }, NOW, true)).toBe( + "2026-08-18T12:15:00.000Z", + ); + expect(nextFetchAt({ ...feed, consecutive_failures: 3 }, NOW, true)).toBe( + "2026-08-18T13:00:00.000Z", + ); + }); + + it("caps the backoff at a day, so a recovered feed is noticed", () => { + const far = nextFetchAt({ ...feed, consecutive_failures: 99 }, NOW, true); + expect(new Date(far).getTime() - NOW.getTime()).toBeLessThanOrEqual(86_400_000); + }); +}); diff --git a/tests/promote/keywords.test.ts b/tests/promote/keywords.test.ts new file mode 100644 index 00000000..1d80a339 --- /dev/null +++ b/tests/promote/keywords.test.ts @@ -0,0 +1,100 @@ +import { describe, it, expect } from "vitest"; +import { + parseKeywords, + slugifyKeyword, + topicFeedUrl, + topicPageUrl, +} from "@/lib/promote/keywords"; + +describe("slugifyKeyword", () => { + it("lowercases a simple keyword", () => { + expect(slugifyKeyword("Bitcoin")).toBe("bitcoin"); + }); + + it("hyphenates multi-word keywords, matching how RSS Amplifier slugs topics", () => { + // Live check: /topics/artificial-intelligence.rss is a real feed, + // /topics/artificial%20intelligence.rss is not. + expect(slugifyKeyword("artificial intelligence")).toBe("artificial-intelligence"); + expect(slugifyKeyword(" AI Agent ")).toBe("ai-agent"); + }); + + it("strips punctuation and diacritics", () => { + expect(slugifyKeyword("C++ programming!")).toBe("c-programming"); + expect(slugifyKeyword("café culture")).toBe("cafe-culture"); + }); + + it("rejects input with nothing usable in it", () => { + expect(slugifyKeyword("")).toBeNull(); + expect(slugifyKeyword(" ")).toBeNull(); + expect(slugifyKeyword("!!!")).toBeNull(); + }); + + it("never emits a trailing hyphen when it truncates", () => { + const slug = slugifyKeyword("a".repeat(78) + " bb")!; + expect(slug.length).toBeLessThanOrEqual(80); + expect(slug.endsWith("-")).toBe(false); + }); +}); + +describe("parseKeywords", () => { + it("splits a comma-separated list into one source per keyword", () => { + expect(parseKeywords("bitcoin,blockchain,ethereum").map((k) => k.slug)).toEqual([ + "bitcoin", + "blockchain", + "ethereum", + ]); + }); + + it("splits on newlines too", () => { + expect(parseKeywords("bitcoin\nblockchain").map((k) => k.slug)).toEqual([ + "bitcoin", + "blockchain", + ]); + }); + + it("does NOT split on spaces — a phrase is one keyword", () => { + const parsed = parseKeywords("artificial intelligence"); + expect(parsed).toHaveLength(1); + expect(parsed[0].slug).toBe("artificial-intelligence"); + }); + + it("keeps the display label the user typed", () => { + expect(parseKeywords(" Artificial Intelligence ")[0]).toEqual({ + slug: "artificial-intelligence", + label: "Artificial Intelligence", + }); + }); + + it("deduplicates keywords that normalize to the same topic", () => { + expect(parseKeywords("Bitcoin, bitcoin , BITCOIN").map((k) => k.slug)).toEqual([ + "bitcoin", + ]); + }); + + it("drops empty entries from sloppy input", () => { + expect(parseKeywords("bitcoin,, ,\n\nethereum,").map((k) => k.slug)).toEqual([ + "bitcoin", + "ethereum", + ]); + }); + + it("returns nothing for empty input", () => { + expect(parseKeywords("")).toEqual([]); + }); +}); + +describe("topic URLs", () => { + it("builds the feed URL the spec calls for", () => { + expect(topicFeedUrl("bitcoin")).toBe("https://rssamplifier.com/topics/bitcoin.rss"); + }); + + it("builds the human-facing topic page URL", () => { + expect(topicPageUrl("bitcoin")).toBe("https://rssamplifier.com/topics/bitcoin"); + }); + + it("honours an override base without doubling the slash", () => { + expect(topicFeedUrl("bitcoin", "https://staging.example.com/")).toBe( + "https://staging.example.com/topics/bitcoin.rss", + ); + }); +}); diff --git a/tests/promote/normalize-url.test.ts b/tests/promote/normalize-url.test.ts new file mode 100644 index 00000000..144609a2 --- /dev/null +++ b/tests/promote/normalize-url.test.ts @@ -0,0 +1,147 @@ +import { describe, it, expect } from "vitest"; +import { + canonicalizeUrl, + normalizeUrlForIdentity, + normalizedTitleHash, + isTrackingParam, + urlHash, +} from "@/lib/promote/normalizeUrl"; + +describe("canonicalizeUrl — the form we publish", () => { + it("strips utm_* campaign tags", () => { + expect( + canonicalizeUrl("https://example.com/post?utm_source=x&utm_medium=social"), + ).toBe("https://example.com/post"); + }); + + it("strips network click ids", () => { + expect(canonicalizeUrl("https://example.com/p?fbclid=abc&gclid=def&igshid=g")).toBe( + "https://example.com/p", + ); + }); + + it("keeps parameters that select the resource", () => { + expect(canonicalizeUrl("https://example.com/search?q=bitcoin&page=2&utm_source=x")).toBe( + "https://example.com/search?q=bitcoin&page=2", + ); + }); + + it("keeps a bare ref, which routes on some sites", () => { + // CrawlProof's own short links carry ref_slug; dropping ref-ish params + // wholesale risks publishing a link that no longer resolves. + expect(canonicalizeUrl("https://example.com/p?ref=producthunt")).toBe( + "https://example.com/p?ref=producthunt", + ); + }); + + it("drops the fragment", () => { + expect(canonicalizeUrl("https://example.com/post#section-2")).toBe( + "https://example.com/post", + ); + }); + + it("preserves host case sensitivity of the path and the www prefix", () => { + expect(canonicalizeUrl("https://www.example.com/Post/Title")).toBe( + "https://www.example.com/Post/Title", + ); + }); + + it("rejects anything that is not a web link", () => { + expect(canonicalizeUrl("javascript:alert(1)")).toBeNull(); + expect(canonicalizeUrl("ftp://example.com/f")).toBeNull(); + expect(canonicalizeUrl("not a url")).toBeNull(); + expect(canonicalizeUrl("")).toBeNull(); + }); +}); + +describe("normalizeUrlForIdentity — the form we dedupe on", () => { + it("folds http and https to one identity", () => { + expect(normalizeUrlForIdentity("http://example.com/p")).toBe( + normalizeUrlForIdentity("https://example.com/p"), + ); + }); + + it("folds the www prefix away", () => { + expect(normalizeUrlForIdentity("https://www.example.com/p")).toBe( + normalizeUrlForIdentity("https://example.com/p"), + ); + }); + + it("folds a trailing slash away", () => { + expect(normalizeUrlForIdentity("https://example.com/p/")).toBe( + normalizeUrlForIdentity("https://example.com/p"), + ); + }); + + it("does not fold the root path away", () => { + expect(normalizeUrlForIdentity("https://example.com/")).toBe("https://example.com/"); + }); + + it("treats reordered query parameters as the same resource", () => { + expect(normalizeUrlForIdentity("https://e.com/s?b=2&a=1")).toBe( + normalizeUrlForIdentity("https://e.com/s?a=1&b=2"), + ); + }); + + it("keeps genuinely different resources apart", () => { + expect(normalizeUrlForIdentity("https://e.com/a")).not.toBe( + normalizeUrlForIdentity("https://e.com/b"), + ); + expect(normalizeUrlForIdentity("https://e.com/s?page=1")).not.toBe( + normalizeUrlForIdentity("https://e.com/s?page=2"), + ); + }); + + it("is case-insensitive on the host but not the path", () => { + expect(normalizeUrlForIdentity("https://EXAMPLE.com/Post")).toBe( + "https://example.com/Post", + ); + }); +}); + +describe("urlHash", () => { + it("agrees for URLs that differ only in tracking and shape", () => { + expect(urlHash("http://www.example.com/p/?utm_source=x")).toBe( + urlHash("https://example.com/p"), + ); + }); + + it("is a sha256 hex digest", () => { + expect(urlHash("https://example.com/p")).toMatch(/^[a-f0-9]{64}$/); + }); + + it("returns null rather than hashing garbage", () => { + expect(urlHash("not a url")).toBeNull(); + }); +}); + +describe("normalizedTitleHash", () => { + it("ignores dash style and case, so the same story matches", () => { + expect(normalizedTitleHash("Foo — Bar")).toBe(normalizedTitleHash("foo - bar")); + }); + + it("separates genuinely different headlines", () => { + expect(normalizedTitleHash("Bitcoin rises")).not.toBe( + normalizedTitleHash("Bitcoin falls"), + ); + }); + + it("returns null for an empty title", () => { + expect(normalizedTitleHash("")).toBeNull(); + expect(normalizedTitleHash(null)).toBeNull(); + }); +}); + +describe("isTrackingParam", () => { + it("matches whole analytics families by prefix", () => { + expect(isTrackingParam("utm_content")).toBe(true); + expect(isTrackingParam("mtm_campaign")).toBe(true); + expect(isTrackingParam("_hsenc")).toBe(true); + }); + + it("leaves content parameters alone", () => { + expect(isTrackingParam("q")).toBe(false); + expect(isTrackingParam("page")).toBe(false); + expect(isTrackingParam("id")).toBe(false); + }); +}); diff --git a/tests/promote/select-link.test.ts b/tests/promote/select-link.test.ts new file mode 100644 index 00000000..0ebd3698 --- /dev/null +++ b/tests/promote/select-link.test.ts @@ -0,0 +1,223 @@ +import { describe, it, expect, beforeEach } from "vitest"; +import { selectNextLink } from "@/lib/promote/selectLink"; +import { makeFakeSupabase, resetIds, type FakeDb } from "./fake-supabase"; + +const NOW = new Date("2026-08-18T12:00:00.000Z"); + +function link(over: Record = {}) { + return { + id: "link-1", + list_id: "list-a", + url: "https://example.com/a", + title: "A", + angle: null, + summary: null, + source_name: null, + source_id: null, + ownership: "owned", + enabled: true, + last_promoted_at: null, + times_promoted: 0, + ...over, + }; +} + +function db(over: Partial = {}): FakeDb { + return { promo_link: [], promo_post: [], ...over }; +} + +const list = { + id: "list-a", + source_mix: { owned: 70, partner: 0, shared: 30 }, + fallback_policy: { + whenOwnedQueueEmpty: "use_shared", + whenSharedQueueEmpty: "use_owned", + maxFallbackItemsPerDay: 3, + }, +}; + +beforeEach(() => resetIds()); + +describe("selectNextLink", () => { + it("returns nothing for an empty list", async () => { + const { client } = makeFakeSupabase(db()); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link).toBeNull(); + expect(selection.decision.reason).toBe("no_inventory"); + }); + + it("prefers the owned link on a fresh 70/30 list", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "owned-1", ownership: "owned" }), + link({ id: "shared-1", ownership: "shared" }), + ], + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("owned-1"); + expect(selection.decision.viaFallback).toBe(false); + }); + + it("switches to shared once owned is ahead of target", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "owned-1", ownership: "owned" }), + link({ id: "shared-1", ownership: "shared" }), + ], + promo_post: Array.from({ length: 8 }, (_, i) => ({ + id: `post-${i}`, + list_id: "list-a", + ownership: "owned", + status: "posted", + via_fallback: false, + created_at: `2026-08-18T0${i}:00:00.000Z`, + })), + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("shared-1"); + }); + + it("rotates within a class: least recently promoted first", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "old", ownership: "owned", last_promoted_at: "2026-08-18T09:00:00.000Z" }), + link({ id: "older", ownership: "owned", last_promoted_at: "2026-08-18T08:00:00.000Z" }), + link({ id: "never", ownership: "owned", last_promoted_at: null }), + ], + }), + ); + // A link that has never posted goes first. + expect((await selectNextLink(client, list, NOW)).link?.id).toBe("never"); + }); + + it("ignores disabled links", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "off", ownership: "owned", enabled: false }), + link({ id: "on", ownership: "shared" }), + ], + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("on"); + expect(selection.decision.viaFallback).toBe(true); + }); + + it("falls back to shared when the owned queue is empty", async () => { + const { client } = makeFakeSupabase( + db({ promo_link: [link({ id: "shared-1", ownership: "shared" })] }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("shared-1"); + expect(selection.decision.viaFallback).toBe(true); + expect(selection.decision.reason).toBe("fallback"); + }); + + it("stops falling back once the daily cap is spent", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [link({ id: "shared-1", ownership: "shared" })], + promo_post: Array.from({ length: 3 }, (_, i) => ({ + id: `fb-${i}`, + list_id: "list-a", + ownership: "shared", + status: "posted", + via_fallback: true, + created_at: "2026-08-18T10:00:00.000Z", + })), + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link).toBeNull(); + expect(selection.decision.reason).toBe("fallback_cap_reached"); + }); + + it("counts fallbacks in a rolling window, not forever", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [link({ id: "shared-1", ownership: "shared" })], + promo_post: Array.from({ length: 3 }, (_, i) => ({ + id: `fb-${i}`, + list_id: "list-a", + ownership: "shared", + status: "posted", + via_fallback: true, + // Two days ago: outside the window, so it must not still block. + created_at: "2026-08-16T10:00:00.000Z", + })), + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("shared-1"); + }); + + it("treats pre-source posts with no ownership as owned", async () => { + // Lists that predate content sources have promo_post rows with a null + // ownership. They were all hand-pasted links, so they must count as owned + // or the blend reads the history as 100% shared and never posts shared. + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "owned-1", ownership: "owned" }), + link({ id: "shared-1", ownership: "shared" }), + ], + promo_post: Array.from({ length: 9 }, (_, i) => ({ + id: `legacy-${i}`, + list_id: "list-a", + ownership: null, + status: "posted", + via_fallback: false, + created_at: `2026-08-18T0${i}:00:00.000Z`, + })), + }), + ); + const selection = await selectNextLink(client, list, NOW); + expect(selection.link?.id).toBe("shared-1"); + }); + + it("counts pending posts, so cookie-auth platforms do not skew the ratio", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "owned-1", ownership: "owned" }), + link({ id: "shared-1", ownership: "shared" }), + ], + promo_post: Array.from({ length: 8 }, (_, i) => ({ + id: `p-${i}`, + list_id: "list-a", + ownership: "owned", + status: "pending", + via_fallback: false, + created_at: `2026-08-18T0${i}:00:00.000Z`, + })), + }), + ); + expect((await selectNextLink(client, list, NOW)).link?.id).toBe("shared-1"); + }); + + it("ignores another list's history", async () => { + const { client } = makeFakeSupabase( + db({ + promo_link: [ + link({ id: "owned-1", ownership: "owned" }), + link({ id: "shared-1", ownership: "shared" }), + ], + promo_post: Array.from({ length: 20 }, (_, i) => ({ + id: `other-${i}`, + list_id: "list-other", + ownership: "owned", + status: "posted", + via_fallback: false, + created_at: "2026-08-18T10:00:00.000Z", + })), + }), + ); + expect((await selectNextLink(client, list, NOW)).link?.id).toBe("owned-1"); + }); +}); diff --git a/worker/index.ts b/worker/index.ts index 60331723..fa8bded9 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 { ingestDueFeeds } from "../lib/promote/ingest"; import { refreshCookieSessions } from "../lib/sp/sessionRefresh"; const supabaseUrl = process.env.NEXT_PUBLIC_SUPABASE_URL!; @@ -1385,6 +1386,23 @@ setInterval( 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 — +// polling is per feed, not per list. +async function promoteIngestSweep() { + const r = await ingestDueFeeds(supabase); + if (r.feedsChecked > 0) { + console.log( + `[worker] promote ingest feeds=${r.feedsChecked} unchanged=${r.feedsUnchanged} failed=${r.feedsFailed} items=${r.itemsStored} links=${r.linksCreated}`, + ); + } +} +setInterval( + () => promoteIngestSweep().catch((e) => console.error("[worker] promote ingest", e)), + 5 * 60_000, +); + // Port-drift scans: bridge queued rows to the "prober" BullMQ queue and // reconcile results (uptime-monitoring-prd.md §12). Short interval so the // Security tab's SSE stream gets snappy queued→running→done feedback. @@ -1448,5 +1466,6 @@ server.listen(port, bindHost, () => { sweep().catch(() => {}); socialFeedSweep().catch((e) => console.error("[worker] social feed sweep", e)); promoteSweep().catch((e) => console.error("[worker] promote sweep", e)); + promoteIngestSweep().catch((e) => console.error("[worker] promote ingest", e)); sessionRefreshSweep().catch((e) => console.error("[worker] session refresh sweep", e)); });