diff --git a/lib/sp/browserPost.ts b/lib/sp/browserPost.ts index b4c31d56..dc2c0d65 100644 --- a/lib/sp/browserPost.ts +++ b/lib/sp/browserPost.ts @@ -30,6 +30,7 @@ import { type ImageStylePref, } from "@/lib/sp/imageGen"; import { reconcileOutreach } from "@/lib/sp/outreachReconcile"; +import { pickDefaultSubreddit } from "@/lib/sp/redditSubreddit"; export async function processBrowserPost(args: { postId: string; @@ -111,9 +112,12 @@ export async function processBrowserPost(args: { let result: { platformPostId: string; webUrl: string }; if (account.platform === "reddit") { - const subreddit = (claimed.subreddit as string | null) ?? ""; const title = (claimed.title as string | null) ?? text.slice(0, 300); - if (!subreddit) throw new Error("Subreddit is required for Reddit posts."); + // Fall back to a related, relatively open subreddit when none was stored + // (e.g. outreach/autopost rows) instead of failing the post. + const subreddit = + ((claimed.subreddit as string | null) ?? "").trim() || + pickDefaultSubreddit(title || text); result = await redditBrowserPost({ cookies, subreddit, title, text }); } else if (account.platform === "facebook_page") { result = await facebookBrowserPost({ diff --git a/lib/sp/browserSemaphore.ts b/lib/sp/browserSemaphore.ts new file mode 100644 index 00000000..a5be6c64 --- /dev/null +++ b/lib/sp/browserSemaphore.ts @@ -0,0 +1,58 @@ +// Global cap on concurrently running headless Chromium instances. +// +// Each browser-automated post (see platforms/browser.ts) spins up its own +// Chromium. Firing many at once — e.g. an outreach "post to all" that hits the +// worker's /sp/browser-post endpoint once per account — launched them all in +// parallel and exhausted the process thread limit, so every launch died with +// pthread_create: Resource temporarily unavailable (11) +// FATAL: Failed to start BrowserThread:IO +// This semaphore serializes launches down to a small ceiling; extra posts wait +// their turn instead of crashing. The slot is held for the whole lifetime of a +// browser (launch → close), so it caps *running* browsers, not just launches. +// +// Tune with SP_BROWSER_CONCURRENCY (default 2). It's a module singleton, so the +// cap is shared across every post handled by a given worker process. + +export class AsyncSemaphore { + private available: number; + private readonly waiters: Array<() => void> = []; + + constructor(max: number) { + this.available = Math.max(1, Math.floor(max)); + } + + async acquire(): Promise { + if (this.available > 0) { + this.available -= 1; + return; + } + await new Promise((resolve) => this.waiters.push(resolve)); + } + + release(): void { + const next = this.waiters.shift(); + if (next) { + // Hand the slot straight to the next waiter — never over-issue. + next(); + } else { + this.available += 1; + } + } + + // Acquire, run fn, and always release — even if fn throws. + async run(fn: () => Promise): Promise { + await this.acquire(); + try { + return await fn(); + } finally { + this.release(); + } + } +} + +function configuredMax(): number { + const raw = Number(process.env.SP_BROWSER_CONCURRENCY); + return Number.isFinite(raw) && raw >= 1 ? Math.floor(raw) : 2; +} + +export const browserSemaphore = new AsyncSemaphore(configuredMax()); diff --git a/lib/sp/platforms/browser.ts b/lib/sp/platforms/browser.ts index c2edc872..3920f9f6 100644 --- a/lib/sp/platforms/browser.ts +++ b/lib/sp/platforms/browser.ts @@ -15,6 +15,7 @@ // a problem. import { chromium, type Browser, type BrowserContext } from "playwright"; +import { browserSemaphore } from "@/lib/sp/browserSemaphore"; export type BrowserCookie = { name: string; @@ -39,15 +40,35 @@ async function launchContext(cookies: BrowserCookie[]): Promise<{ browser: Browser; ctx: BrowserContext; }> { - const browser = await chromium.launch({ headless: true }); - const ctx = await browser.newContext({ - userAgent: UA, - viewport: { width: 1280, height: 800 }, - locale: "en-US", - timezoneId: "America/New_York", - }); - await ctx.addCookies(cookies); - return { browser, ctx }; + // Hold a concurrency slot for the whole browser lifetime so we never run more + // than SP_BROWSER_CONCURRENCY headless Chromiums at once (see + // browserSemaphore). The slot is released when the browser disconnects, which + // fires from every platform function's `finally { browser.close() }`. + await browserSemaphore.acquire(); + let released = false; + const release = () => { + if (!released) { + released = true; + browserSemaphore.release(); + } + }; + try { + const browser = await chromium.launch({ headless: true }); + browser.once("disconnected", release); + const ctx = await browser.newContext({ + userAgent: UA, + viewport: { width: 1280, height: 800 }, + locale: "en-US", + timezoneId: "America/New_York", + }); + await ctx.addCookies(cookies); + return { browser, ctx }; + } catch (err) { + // Launch or context setup failed before we could attach the disconnect + // handler — free the slot so a failure can't leak capacity. + release(); + throw err; + } } export function parseCookies(raw: string): BrowserCookie[] { diff --git a/lib/sp/post.ts b/lib/sp/post.ts index 7b8865d3..18a9a6b7 100644 --- a/lib/sp/post.ts +++ b/lib/sp/post.ts @@ -9,6 +9,7 @@ import type { SupabaseClient } from "@supabase/supabase-js"; import { encryptSecret, decryptSecret } from "@/lib/sp/vault"; import { enqueueBrowserPost } from "@/lib/lx/workerClient"; +import { resolveSubreddit } from "@/lib/sp/redditSubreddit"; import { createBlueskyPost, createBlueskySession, @@ -108,6 +109,13 @@ export async function postViaAccount(args: { // Cookie-auth accounts post via Playwright in the worker. if (account.auth_mode === "cookie") { + // Reddit needs a subreddit, but outreach/autopost flows often don't supply + // one. Rather than fail, route to a related, relatively open subreddit + // (see resolveSubreddit) and persist the choice on the row. + const subreddit = + account.platform === "reddit" + ? resolveSubreddit(input.subreddit, input.title ?? text) + : input.subreddit ?? null; const { data: row, error: insErr } = await supabase .from("sp_post") .insert({ @@ -117,7 +125,7 @@ export async function postViaAccount(args: { source, rendered_text: text, rendered_media_url: input.mediaUrl ?? [], - subreddit: input.subreddit ?? null, + subreddit, title: input.title ?? null, scheduled_for: new Date().toISOString(), status: "queued_browser", diff --git a/lib/sp/redditSubreddit.ts b/lib/sp/redditSubreddit.ts new file mode 100644 index 00000000..0c5f706c --- /dev/null +++ b/lib/sp/redditSubreddit.ts @@ -0,0 +1,83 @@ +// Pick a fallback subreddit for a Reddit post that didn't specify one. +// +// Reddit posting is cookie/browser-based, and outreach (and some autoposts) +// don't carry a subreddit — which used to fail the post outright with +// "Subreddit is required for Reddit posts." Instead we route the post to a +// relatively open ("low moderation") subreddit that fits the content, so it +// actually goes out rather than erroring. +// +// This is a best-effort curated list — moderation strictness changes over time +// and can't be verified from here — so it's overridable with +// SP_REDDIT_DEFAULT_SUBS (comma-separated, most-preferred first). Keyword +// routing picks the most topically-relevant entry for the post text; when +// nothing matches (or an env override is set) it falls back to the first entry. + +type Curated = { sub: string; keywords: string[] }; + +// Ordered by fallback preference. r/SideProject is broadly open to sharing a +// tool/site; the others catch SEO / blogging / AI-search themed posts. +const CURATED: Curated[] = [ + { + sub: "SideProject", + keywords: ["launch", "built", "made", "project", "tool", "app", "startup", "saas"], + }, + { + sub: "juststart", + keywords: ["blog", "traffic", "website", "site", "content", "affiliate", "niche"], + }, + { + sub: "SEO", + keywords: ["seo", "search", "rank", "ranking", "google", "serp", "backlink", "keyword", "crawl"], + }, + { + sub: "artificial", + keywords: ["ai", "llm", "chatgpt", "gpt", "claude", "gemini", "perplexity", "aeo", "answer engine"], + }, +]; + +function normalizeSub(value: string | null | undefined): string { + return (value ?? "").trim().replace(/^\/?r\//i, "").trim(); +} + +function envSubs(): string[] { + return (process.env.SP_REDDIT_DEFAULT_SUBS ?? "") + .split(",") + .map((s) => normalizeSub(s)) + .filter(Boolean); +} + +// Choose a subreddit for the given post content. Prefers an env-configured +// list; otherwise keyword-routes over the curated list. Always returns a +// non-empty subreddit name (no leading "r/"). +export function pickDefaultSubreddit(content: string | null | undefined): string { + const overrides = envSubs(); + const text = (content ?? "").toLowerCase(); + + if (overrides.length > 0) { + // Env list has no keyword hints; use the first (most-preferred) entry. + return overrides[0]; + } + + let best = CURATED[0].sub; + let bestScore = 0; + for (const { sub, keywords } of CURATED) { + const score = keywords.reduce( + (acc, kw) => (text.includes(kw) ? acc + 1 : acc), + 0, + ); + if (score > bestScore) { + bestScore = score; + best = sub; + } + } + return best; +} + +// Return the given subreddit if one was supplied, otherwise a content-based +// fallback. Strips a leading "r/" either way. +export function resolveSubreddit( + supplied: string | null | undefined, + content: string | null | undefined, +): string { + return normalizeSub(supplied) || pickDefaultSubreddit(content); +} diff --git a/tests/sp/browser-semaphore.test.ts b/tests/sp/browser-semaphore.test.ts new file mode 100644 index 00000000..35131ca0 --- /dev/null +++ b/tests/sp/browser-semaphore.test.ts @@ -0,0 +1,69 @@ +import { describe, it, expect } from "vitest"; +import { AsyncSemaphore } from "@/lib/sp/browserSemaphore"; + +const tick = () => new Promise((r) => setTimeout(r, 0)); + +describe("AsyncSemaphore", () => { + it("never runs more than `max` tasks at once", async () => { + const sem = new AsyncSemaphore(2); + let active = 0; + let peak = 0; + const task = () => + sem.run(async () => { + active++; + peak = Math.max(peak, active); + await tick(); + active--; + }); + + await Promise.all(Array.from({ length: 8 }, task)); + + expect(peak).toBe(2); + expect(active).toBe(0); + }); + + it("serializes fully at max=1", async () => { + const sem = new AsyncSemaphore(1); + const order: number[] = []; + let active = 0; + let peak = 0; + await Promise.all( + [1, 2, 3].map((n) => + sem.run(async () => { + active++; + peak = Math.max(peak, active); + await tick(); + order.push(n); + active--; + }), + ), + ); + expect(peak).toBe(1); + expect(order).toEqual([1, 2, 3]); + }); + + it("releases the slot even when a task throws", async () => { + const sem = new AsyncSemaphore(1); + await expect( + sem.run(async () => { + throw new Error("boom"); + }), + ).rejects.toThrow("boom"); + + // If the slot leaked, this second acquire would hang forever. + let ran = false; + await sem.run(async () => { + ran = true; + }); + expect(ran).toBe(true); + }); + + it("clamps a max below 1 up to 1", async () => { + const sem = new AsyncSemaphore(0); + let ran = false; + await sem.run(async () => { + ran = true; + }); + expect(ran).toBe(true); + }); +}); diff --git a/tests/sp/reddit-subreddit.test.ts b/tests/sp/reddit-subreddit.test.ts new file mode 100644 index 00000000..120c74c9 --- /dev/null +++ b/tests/sp/reddit-subreddit.test.ts @@ -0,0 +1,51 @@ +import { afterEach, describe, it, expect } from "vitest"; +import { + pickDefaultSubreddit, + resolveSubreddit, +} from "@/lib/sp/redditSubreddit"; + +const ENV = "SP_REDDIT_DEFAULT_SUBS"; + +afterEach(() => { + delete process.env[ENV]; +}); + +describe("resolveSubreddit", () => { + it("keeps a supplied subreddit and strips a leading r/", () => { + expect(resolveSubreddit("r/webdev", "anything")).toBe("webdev"); + expect(resolveSubreddit("/r/SEO", "anything")).toBe("SEO"); + expect(resolveSubreddit("SideProject", "")).toBe("SideProject"); + }); + + it("falls back to a content-routed subreddit when none supplied", () => { + expect(resolveSubreddit("", "How I improved my Google search ranking")) + .toBe("SEO"); + expect(resolveSubreddit(null, "")).toBe("SideProject"); // default first entry + }); +}); + +describe("pickDefaultSubreddit", () => { + it("routes by topical keywords", () => { + expect(pickDefaultSubreddit("Cut my backlink audit crawl time")).toBe("SEO"); + expect(pickDefaultSubreddit("We just built and launched a new SaaS tool")).toBe( + "SideProject", + ); + expect(pickDefaultSubreddit("Using an LLM / ChatGPT for answer-engine AEO")).toBe( + "artificial", + ); + expect(pickDefaultSubreddit("Grew my blog traffic on a niche affiliate site")).toBe( + "juststart", + ); + }); + + it("defaults to the first curated entry when nothing matches", () => { + expect(pickDefaultSubreddit("completely unrelated text")).toBe("SideProject"); + expect(pickDefaultSubreddit("")).toBe("SideProject"); + }); + + it("honors the SP_REDDIT_DEFAULT_SUBS override (most-preferred first)", () => { + process.env[ENV] = "r/test, myovveride"; + expect(pickDefaultSubreddit("How I improved my Google ranking")).toBe("test"); + expect(resolveSubreddit("", "seo ranking backlink")).toBe("test"); + }); +});