diff --git a/.changeset/fluffy-camels-sniff.md b/.changeset/fluffy-camels-sniff.md new file mode 100644 index 0000000..f03c5b5 --- /dev/null +++ b/.changeset/fluffy-camels-sniff.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": patch +--- + +update cli diff --git a/apps/sveltekit-cloudflare-workers/package.json b/apps/sveltekit-cloudflare-workers/package.json index b9924c5..150c246 100644 --- a/apps/sveltekit-cloudflare-workers/package.json +++ b/apps/sveltekit-cloudflare-workers/package.json @@ -5,7 +5,7 @@ "type": "module", "scripts": { "dev": "vite dev", - "build": "contrail-lex generate && contrail-lex types && vite build && tsx scripts/append-scheduled.ts", + "build": "contrail-lex generate && contrail-lex types && vite build && contrail append-scheduled", "generate": "contrail-lex generate", "generate:pull": "contrail-lex all", "backfill": "contrail backfill", diff --git a/apps/sveltekit-cloudflare-workers/scripts/append-scheduled.ts b/apps/sveltekit-cloudflare-workers/scripts/append-scheduled.ts deleted file mode 100644 index c9bd1cd..0000000 --- a/apps/sveltekit-cloudflare-workers/scripts/append-scheduled.ts +++ /dev/null @@ -1,29 +0,0 @@ -/** - * Post-build script: appends a `scheduled` handler to the SvelteKit worker output. - * - * SvelteKit's adapter-cloudflare doesn't support the `scheduled` export natively - * (see https://github.com/sveltejs/kit/issues/4841). This script patches the - * generated _worker.js to add one that self-calls the /api/cron endpoint. - */ -import { readFileSync, writeFileSync } from 'fs'; -import { join, dirname } from 'path'; -import { fileURLToPath } from 'url'; - -const root = join(dirname(fileURLToPath(import.meta.url)), '..'); -const workerPath = join(root, '.svelte-kit', 'cloudflare', '_worker.js'); - -let code = readFileSync(workerPath, 'utf-8'); - -code += ` -// --- Appended by scripts/append-scheduled.ts --- -worker_default.scheduled = async function (event, env, ctx) { - const req = new Request('http://localhost/api/cron', { - method: 'POST', - headers: { 'X-Cron-Secret': env.CRON_SECRET || '' } - }); - ctx.waitUntil(this.fetch(req, env, ctx)); -}; -`; - -writeFileSync(workerPath, code); -console.log('Appended scheduled handler to _worker.js'); diff --git a/docs/08-labels.md b/docs/08-labels.md index 5dad220..39a873e 100644 --- a/docs/08-labels.md +++ b/docs/08-labels.md @@ -88,7 +88,7 @@ Account-level labels (subject = bare DID) hydrate onto profiles. They appear ins |---|---|---| | Cron-driven | `contrail.ingestLabels()` | Cloudflare Workers — one drain per cron tick | | Persistent | `contrail.runPersistentLabels()` | Node / long-lived servers — one socket per labeler, auto-reconnect | -| One-shot backfill | `pnpm contrail labels-backfill [--remote]` | Local script, drains until each labeler reports caught up | +| One-shot backfill | `pnpm contrail backfill --only labels [--remote]` | Local script, drains until each labeler reports caught up. (`pnpm contrail backfill` runs both records and labels.) | When `config.labels` is set, the bundled `createWorker` already calls `ingestLabels()` from `scheduled()` alongside `ingest()` — no boilerplate. diff --git a/docs/frameworks/sveltekit-cloudflare.md b/docs/frameworks/sveltekit-cloudflare.md index 691ca62..d03a488 100644 --- a/docs/frameworks/sveltekit-cloudflare.md +++ b/docs/frameworks/sveltekit-cloudflare.md @@ -23,8 +23,6 @@ src/ xrpc/[...path]/+server.ts # mounts all contrail XRPC endpoints api/cron/+server.ts # hit by the cron trigger (see below) wrangler.jsonc -scripts/ - append-scheduled.ts # workaround — see "Cron" below ``` ## 1. Declare the config @@ -119,7 +117,7 @@ export const load: PageServerLoad = async ({ platform, locals }) => { ## 5. Cron ingest — the workaround -SvelteKit's `@sveltejs/adapter-cloudflare` doesn't expose a `scheduled()` export on the generated worker ([issue #4841](https://github.com/sveltejs/kit/issues/4841)). The easiest fix: an HTTP endpoint that does the ingest, plus a post-build script that appends a `scheduled` handler calling that endpoint. +SvelteKit's `@sveltejs/adapter-cloudflare` doesn't expose a `scheduled()` export on the generated worker ([issue #4841](https://github.com/sveltejs/kit/issues/4841)). The fix is an HTTP endpoint that does the ingest, plus a post-build patch on `_worker.js` that appends a `scheduled` handler calling it. The patch is what `contrail append-scheduled` does. **Endpoint:** @@ -139,41 +137,17 @@ export const POST: RequestHandler = async ({ request, platform }) => { }; ``` -**Post-build patch:** - -```ts -// scripts/append-scheduled.ts -import { readFileSync, writeFileSync } from "fs"; -import { join, dirname } from "path"; -import { fileURLToPath } from "url"; - -const root = join(dirname(fileURLToPath(import.meta.url)), ".."); -const workerPath = join(root, ".svelte-kit", "cloudflare", "_worker.js"); - -writeFileSync( - workerPath, - readFileSync(workerPath, "utf-8") + - ` -worker_default.scheduled = async function (event, env, ctx) { - const req = new Request("http://localhost/api/cron", { - method: "POST", - headers: { "X-Cron-Secret": env.CRON_SECRET ?? "" }, - }); - ctx.waitUntil(this.fetch(req, env, ctx)); -}; -` -); -``` - -**Wire it into `build`:** +**Wire `contrail append-scheduled` into your `build` script:** ```jsonc // package.json "scripts": { - "build": "vite build && tsx scripts/append-scheduled.ts" + "build": "vite build && contrail append-scheduled" } ``` +`contrail append-scheduled` patches `.svelte-kit/cloudflare/_worker.js` to append a `scheduled()` export that POSTs to `/api/cron` with `env.CRON_SECRET`. Override with `--worker `, `--cron-path `, or `--secret-env ` if your project diverges. + `CRON_SECRET` is any random string — generate one, set it as a secret with `wrangler secret put CRON_SECRET`. The cron handler self-auths with it so nobody external can trigger your ingest. ## 6. Wrangler config @@ -237,5 +211,5 @@ From now on: - **Top-level await in `$lib/contrail/index.ts`** will fail to bundle — use the lazy `ensureInit` pattern above. - **`ensureInit` is per-isolate, not global.** Cloudflare cold-starts spin new isolates; each one pays one init call on its first request. `contrail.init()` is idempotent so this is safe, just not instant. -- **SvelteKit's `adapter-cloudflare` regenerates `_worker.js` on every build**, so the `append-scheduled.ts` patch has to run *after* `vite build`. Don't try to put it in `prebuild`. +- **SvelteKit's `adapter-cloudflare` regenerates `_worker.js` on every build**, so `contrail append-scheduled` has to run *after* `vite build`. Don't try to put it in `prebuild`. - **`D1Database` type in platform env** needs `@cloudflare/workers-types` in `devDependencies` and `types` in your tsconfig. diff --git a/packages/contrail/package.json b/packages/contrail/package.json index 5571988..2adafed 100644 --- a/packages/contrail/package.json +++ b/packages/contrail/package.json @@ -73,6 +73,7 @@ "@atcute/jetstream": "^1.0.2", "@atcute/lexicons": "^1.2.9", "@atcute/xrpc-server": "^0.1.12", + "cac": "^7.0.0", "hono": "^4.12.8", "jiti": "^2.4.0" }, diff --git a/packages/contrail/src/cli.ts b/packages/contrail/src/cli.ts index 505b5f4..8b1978f 100644 --- a/packages/contrail/src/cli.ts +++ b/packages/contrail/src/cli.ts @@ -1,358 +1,28 @@ #!/usr/bin/env node /** - * contrail — CLI for one-off operations against a contrail deployment. + * contrail — CLI entrypoint. * - * Usage: - * contrail backfill [--remote] [--binding DB] [--config path] - * contrail refresh [--remote] [--ignore-window 60s] [--by-collection] - * contrail dev [--binding DB] [--cron '*\/1 * * * *'] - * - * Config auto-detects at contrail.config.ts, src/contrail.config.ts, - * src/lib/contrail.config.ts, or app/contrail.config.ts (first match wins). + * Subcommand implementations live in `./cli/commands/`; this file just wires + * them into a single cac instance. See `contrail --help` for usage. */ -import { spawn } from "node:child_process"; -import { createInterface } from "node:readline/promises"; -import { stdin as input, stdout as output } from "node:process"; -import { findConfigFile, loadConfig, CONFIG_CANDIDATES_MESSAGE } from "./cli-config.js"; -import { backfillAll, refresh } from "./workers/backfill.js"; -import { Contrail } from "./contrail.js"; -import type { CollectionStats, RefreshResult } from "./core/refresh.js"; -import type { Database } from "./core/types.js"; - -type Subcommand = "backfill" | "refresh" | "labels-backfill" | "dev" | "help"; - -const USAGE = `contrail [options] - -Subcommands: - backfill One-time bulk load from each known DID's PDS (resumable) - refresh Fresh sweep: reconcile PDS vs DB, report missing + stale - labels-backfill One-shot drain per configured labeler — runs catch-up cycles - until each labeler has no more pending events (resumable) - dev Local wrangler dev + auto-trigger cron + backfill/refresh prompts - help Print this message - -Options (backfill): - --config Path to Contrail config file (TS or JS). - --root Project root for auto-detection. Default: CWD. - --remote Use production D1 bindings. - --binding D1 binding name in wrangler.jsonc. Default: "DB". - --concurrency Passed to contrail.backfillAll(). Default: 100. +import { cac } from "cac"; +import { registerBackfill } from "./cli/commands/backfill.js"; +import { registerRefresh } from "./cli/commands/refresh.js"; +import { registerDev } from "./cli/commands/dev.js"; +import { registerAppendScheduled } from "./cli/commands/append-scheduled.js"; -Options (refresh): - --config Path to Contrail config file (TS or JS). - --root Project root for auto-detection. Default: CWD. - --remote Use production D1 bindings. - --binding D1 binding name in wrangler.jsonc. Default: "DB". - --concurrency Passed to contrail.refresh(). Default: 50. - --ignore-window Seconds of grace for stale-update detection. Default: 60. - --by-collection Print per-collection stats, not just totals. +const cli = cac("contrail"); -Options (dev): - --config Path to Contrail config file (TS or JS). - --root Project root for auto-detection. Default: CWD. - --binding D1 binding name in wrangler.jsonc. Default: "DB". - --cron Cron expression to fire against /__scheduled. Default: "*/1 * * * *". - --stale-after Prompt to run refresh if the ingest cursor is older than this. Default: 60. - --yes Accept all prompts without asking (CI-friendly). -`; +registerBackfill(cli); +registerRefresh(cli); +registerDev(cli); +registerAppendScheduled(cli); -interface Args { - cmd: Subcommand; - config?: string; - root: string; - remote: boolean; - binding: string; - concurrency?: number; - ignoreWindowMs?: number; - byCollection: boolean; - cron?: string; - staleAfterMin?: number; - yes: boolean; -} - -function parseArgs(argv: string[]): Args { - const args = argv.slice(2); - const cmd = (args.shift() ?? "help") as Subcommand; - let config: string | undefined; - let root = process.cwd(); - let remote = false; - let binding = "DB"; - let concurrency: number | undefined; - let ignoreWindowMs: number | undefined; - let byCollection = false; - let cron: string | undefined; - let staleAfterMin: number | undefined; - let yes = false; - for (let i = 0; i < args.length; i++) { - const a = args[i]; - if (a === "--config") config = args[++i]; - else if (a === "--root") root = args[++i]; - else if (a === "--remote") remote = true; - else if (a === "--binding") binding = args[++i]; - else if (a === "--concurrency") concurrency = parseInt(args[++i], 10); - else if (a === "--ignore-window") ignoreWindowMs = parseInt(args[++i], 10) * 1000; - else if (a === "--by-collection") byCollection = true; - else if (a === "--cron") cron = args[++i]; - else if (a === "--stale-after") staleAfterMin = parseInt(args[++i], 10); - else if (a === "--yes" || a === "-y") yes = true; - else if (a === "-h" || a === "--help") - return { - cmd: "help", - config, - root, - remote, - binding, - concurrency, - ignoreWindowMs, - byCollection, - cron, - staleAfterMin, - yes, - }; - } - return { cmd, config, root, remote, binding, concurrency, ignoreWindowMs, byCollection, cron, staleAfterMin, yes }; -} - -function resolveConfigPath(opts: Args): string | null { - const p = findConfigFile(opts.root, opts.config); - if (!p) { - console.error( - "Could not find a Contrail config. Pass --config or place one at\n" + - ` ${CONFIG_CANDIDATES_MESSAGE}` - ); - } - return p; -} +cli.help(); -async function promptYesNo(question: string, defaultYes: boolean, autoYes: boolean): Promise { - if (autoYes) return true; - // Non-TTY → default-decline; the user isn't there to answer. - if (!input.isTTY) return false; - const rl = createInterface({ input, output }); - try { - const hint = defaultYes ? "[Y/n]" : "[y/N]"; - const ans = (await rl.question(`${question} ${hint} `)).trim().toLowerCase(); - if (ans === "") return defaultYes; - return ans === "y" || ans === "yes"; - } finally { - rl.close(); - } +try { + cli.parse(); +} catch (err) { + console.error(err); + process.exit(1); } - -function formatStats(s: CollectionStats): string { - return `${s.missing} missing, ${s.staleUpdates} stale updates, ${s.inSync} in sync`; -} - -function printRefreshReport(result: RefreshResult, byCollection: boolean): void { - console.log(""); - if (byCollection) { - console.log("by collection:"); - const entries = Object.entries(result.byCollection).sort(([a], [b]) => - a.localeCompare(b) - ); - for (const [nsid, stats] of entries) { - if (stats.missing === 0 && stats.staleUpdates === 0 && stats.inSync === 0) continue; - console.log(` ${nsid}`); - console.log(` ${formatStats(stats)}`); - } - console.log(""); - } - console.log("total:"); - console.log(` ${formatStats(result.total)}`); - console.log(` ${result.usersScanned} users scanned` + - (result.usersFailed ? `, ${result.usersFailed} failed` : "") + - ` in ${(result.elapsedMs / 1000).toFixed(1)}s`); - console.log(` (ignore window: ${(result.ignoreWindowMs / 1000).toFixed(0)}s)`); -} - -async function cmdDev(opts: Args): Promise { - const configPath = resolveConfigPath(opts); - if (!configPath) return 1; - const config = await loadConfig(configPath); - - // Pre-flight: connect to local D1 via wrangler's platform proxy, inspect - // state, prompt. Then dispose before starting wrangler dev so the two - // miniflare processes don't fight over the sqlite file. - const { getPlatformProxy } = await import("wrangler"); - const { env, dispose } = await getPlatformProxy(); - const db = (env as Record)[opts.binding] as Database | undefined; - - if (db) { - const contrail = new Contrail(config); - await contrail.init(db); - - const hasBackfilled = await db - .prepare("SELECT 1 FROM backfills WHERE completed = 1 LIMIT 1") - .first(); - - if (!hasBackfilled) { - console.log("no backfilled users in the local DB yet."); - if (await promptYesNo("run backfill now? (takes a few minutes)", true, opts.yes)) { - await contrail.backfillAll({ concurrency: opts.concurrency ?? 100 }, db); - } - } else { - const staleAfterMs = (opts.staleAfterMin ?? 60) * 60_000; - const row = await db - .prepare("SELECT time_us FROM cursor WHERE id = 1") - .first<{ time_us: number }>(); - if (row?.time_us) { - const ageMs = Date.now() - Math.floor(row.time_us / 1000); - if (ageMs > staleAfterMs) { - const hrs = (ageMs / 3_600_000).toFixed(1); - console.log(`ingest cursor is ${hrs}h old — you may have missed events.`); - if (await promptYesNo("run refresh first?", true, opts.yes)) { - await contrail.refresh({ ignoreWindowMs: opts.ignoreWindowMs }, db); - } - } - } - } - } - - await dispose(); - - // Start wrangler dev + fire /__scheduled on a loop so the cron actually - // runs in local dev (wrangler's cron scheduler only works in deployed - // production; --test-scheduled enables the manual-trigger endpoint). - const cron = opts.cron ?? "*/1 * * * *"; - const cronUrl = `http://localhost:8787/__scheduled?cron=${encodeURIComponent(cron)}`; - - const wrangler = spawn("npx", ["wrangler", "dev", "--test-scheduled"], { - stdio: "inherit", - shell: process.platform === "win32", - cwd: opts.root, - }); - - // Give wrangler a few seconds to bind the port, then start the cron loop. - const kickoff = setTimeout(() => { - fetch(cronUrl).catch(() => {}); // fire-and-forget; first run - console.log(`\nauto-ingest: firing ${cronUrl} every 60s\n`); - }, 3_000); - - const interval = setInterval(() => { - fetch(cronUrl).catch(() => {}); // wrangler may be restarting; swallow - }, 60_000); - - const cleanup = () => { - clearTimeout(kickoff); - clearInterval(interval); - try { - wrangler.kill("SIGINT"); - } catch { - /* ignore */ - } - }; - process.on("SIGINT", cleanup); - process.on("SIGTERM", cleanup); - - return new Promise((resolve) => { - wrangler.on("exit", (code) => { - clearTimeout(kickoff); - clearInterval(interval); - resolve(code ?? 0); - }); - }); -} - -async function main(): Promise { - const opts = parseArgs(process.argv); - - if (opts.cmd === "help") { - process.stdout.write(USAGE); - return 0; - } - - if (opts.cmd === "backfill") { - const configPath = resolveConfigPath(opts); - if (!configPath) return 1; - const config = await loadConfig(configPath); - await backfillAll({ - config, - remote: opts.remote, - binding: opts.binding, - concurrency: opts.concurrency ?? 100, - }); - return 0; - } - - if (opts.cmd === "refresh") { - const configPath = resolveConfigPath(opts); - if (!configPath) return 1; - const config = await loadConfig(configPath); - const result = await refresh({ - config, - remote: opts.remote, - binding: opts.binding, - concurrency: opts.concurrency ?? 50, - ignoreWindowMs: opts.ignoreWindowMs, - }); - printRefreshReport(result, opts.byCollection); - return 0; - } - - if (opts.cmd === "labels-backfill") { - const configPath = resolveConfigPath(opts); - if (!configPath) return 1; - const config = await loadConfig(configPath); - if (!config.labels || config.labels.sources.length === 0) { - console.error("No labels configured (config.labels.sources is empty)."); - return 1; - } - const { getPlatformProxy } = await import("wrangler"); - const { env, dispose } = await getPlatformProxy({ - environment: opts.remote ? "production" : undefined, - }); - try { - const db = (env as Record)[opts.binding] as Database | undefined; - if (!db) { - console.error(`No binding named "${opts.binding}" in wrangler env.`); - return 1; - } - const contrail = new Contrail(config); - await contrail.init(db); - - // Run cycles until each cycle drains nothing new — measured by the - // labeler_cursors not advancing across two consecutive cycles. - const before = new Map(); - let stable = 0; - while (stable < 2) { - const rows = ( - await db - .prepare("SELECT did, cursor FROM labeler_cursors") - .all<{ did: string; cursor: number }>() - ).results ?? []; - for (const r of rows) before.set(r.did, r.cursor); - await contrail.ingestLabels({ timeoutMs: 60_000 }, db); - const after = ( - await db - .prepare("SELECT did, cursor FROM labeler_cursors") - .all<{ did: string; cursor: number }>() - ).results ?? []; - let advanced = false; - for (const r of after) { - if ((before.get(r.did) ?? -1) !== r.cursor) { - advanced = true; - break; - } - } - stable = advanced ? 0 : stable + 1; - } - console.log("labels-backfill: caught up"); - return 0; - } finally { - await dispose(); - } - } - - if (opts.cmd === "dev") return cmdDev(opts); - - console.error(USAGE); - return 1; -} - -main().then( - (code) => process.exit(code), - (err) => { - console.error(err); - process.exit(1); - } -); diff --git a/packages/contrail/src/cli/commands/append-scheduled.ts b/packages/contrail/src/cli/commands/append-scheduled.ts new file mode 100644 index 0000000..9a6fce9 --- /dev/null +++ b/packages/contrail/src/cli/commands/append-scheduled.ts @@ -0,0 +1,83 @@ +import { readFileSync, writeFileSync } from "node:fs"; +import { resolve as resolvePath, relative as relativePath } from "node:path"; +import type { CAC } from "cac"; + +interface AppendScheduledOpts { + root: string; + worker: string; + cronPath: string; + secretEnv: string; +} + +const MARKER = "// contrail: scheduled handler"; + +/** + * Patch a SvelteKit/adapter-cloudflare worker bundle to expose a `scheduled()` + * handler. The adapter doesn't surface scheduled exports natively + * (sveltejs/kit#4841), so we append one that POSTs to an in-app cron endpoint + * (typically /api/cron) using a shared secret. + */ +export function registerAppendScheduled(cli: CAC): void { + cli + .command( + "append-scheduled", + "Patch a SvelteKit/adapter-cloudflare _worker.js with a scheduled() handler that hits an HTTP cron endpoint. Run after `vite build` in your build script." + ) + .option("--root ", "Project root for resolving --worker", { + default: process.cwd(), + }) + .option("--worker ", "Path to the generated worker bundle", { + default: ".svelte-kit/cloudflare/_worker.js", + }) + .option( + "--cron-path ", + "In-app path the scheduled handler hits", + { default: "/api/cron" } + ) + .option("--secret-env ", "Env var holding the cron secret", { + default: "CRON_SECRET", + }) + .action((options: AppendScheduledOpts) => { + const workerPath = resolvePath(options.root, options.worker); + + let code: string; + try { + code = readFileSync(workerPath, "utf-8"); + } catch (err) { + const e = err as NodeJS.ErrnoException; + if (e.code === "ENOENT") { + console.error( + `contrail append-scheduled: worker bundle not found at ${workerPath}\n` + + "Run your framework build (e.g. `vite build`) first, or pass --worker ." + ); + process.exit(1); + } + throw err; + } + + if (code.includes(MARKER)) { + console.log( + `contrail append-scheduled: ${relativePath(options.root, workerPath)} already patched, skipping.` + ); + return; + } + + const patched = + code + + ` +${MARKER} — patched by \`contrail append-scheduled\` +worker_default.scheduled = async function (event, env, ctx) { +\tconst req = new Request("http://localhost${options.cronPath}", { +\t\tmethod: "POST", +\t\theaders: { "X-Cron-Secret": env.${options.secretEnv} || "" } +\t}); +\tctx.waitUntil(this.fetch(req, env, ctx)); +}; +`; + + writeFileSync(workerPath, patched); + console.log( + `contrail append-scheduled: appended scheduled() → ${options.cronPath} to ${relativePath(options.root, workerPath)}` + ); + }); +} diff --git a/packages/contrail/src/cli/commands/backfill.ts b/packages/contrail/src/cli/commands/backfill.ts new file mode 100644 index 0000000..8b8747a --- /dev/null +++ b/packages/contrail/src/cli/commands/backfill.ts @@ -0,0 +1,82 @@ +import type { CAC } from "cac"; +import { backfillAll, labelsBackfillAll } from "../../workers/backfill.js"; +import { resolveAndLoadConfig } from "../shared.js"; + +interface BackfillOpts { + config?: string; + root?: string; + remote?: boolean; + binding: string; + concurrency: number; + only?: string; +} + +const VALID_ONLY = ["records", "labels"] as const; + +export function registerBackfill(cli: CAC): void { + cli + .command( + "backfill", + "One-time bulk load: records from each known DID's PDS, then labels from each configured labeler. Resumable." + ) + .option("--config ", "Path to Contrail config file (TS or JS)") + .option("--root ", "Project root for auto-detection (default: CWD)") + .option("--remote", "Use production D1 bindings") + .option("--binding ", "D1 binding name in wrangler.jsonc", { + default: "DB", + }) + .option( + "--concurrency ", + "Concurrency for record backfill (labels are per-labeler serial)", + { default: 100 } + ) + .option( + "--only ", + "Run only one half: 'records' or 'labels'. Default: both." + ) + .action(async (options: BackfillOpts) => { + const only = options.only; + if (only !== undefined && !VALID_ONLY.includes(only as never)) { + console.error( + `--only must be one of: ${VALID_ONLY.join(", ")} (got "${only}")` + ); + process.exit(1); + } + + const config = await resolveAndLoadConfig(options); + const wrangler = { + config, + remote: !!options.remote, + binding: options.binding, + }; + + const runRecords = only !== "labels"; + const runLabels = only !== "records"; + + if (runRecords) { + await backfillAll({ + ...wrangler, + concurrency: Number(options.concurrency), + }); + } + + if (runLabels) { + const labelsConfigured = + !!config.labels && config.labels.sources.length > 0; + if (!labelsConfigured) { + if (only === "labels") { + console.error( + "No labels configured (config.labels.sources is empty)." + ); + process.exit(1); + } + // --only not set and labels aren't configured: skip silently. + } else { + const result = await labelsBackfillAll(wrangler); + console.log( + `labels: caught up after ${result.cycles} cycle${result.cycles === 1 ? "" : "s"}` + ); + } + } + }); +} diff --git a/packages/contrail/src/cli/commands/dev.ts b/packages/contrail/src/cli/commands/dev.ts new file mode 100644 index 0000000..6b5a9ee --- /dev/null +++ b/packages/contrail/src/cli/commands/dev.ts @@ -0,0 +1,153 @@ +import { spawn } from "node:child_process"; +import type { CAC } from "cac"; +import { Contrail } from "../../contrail.js"; +import type { Database } from "../../core/types.js"; +import { promptYesNo, resolveAndLoadConfig } from "../shared.js"; + +interface DevOpts { + config?: string; + root: string; + binding: string; + cron: string; + concurrency: number; + ignoreWindow?: number; + staleAfter: number; + yes?: boolean; +} + +export function registerDev(cli: CAC): void { + cli + .command( + "dev", + "Local wrangler dev + auto-trigger cron + backfill/refresh prompts" + ) + .option("--config ", "Path to Contrail config file (TS or JS)") + .option("--root ", "Project root for auto-detection (default: CWD)", { + default: process.cwd(), + }) + .option("--binding ", "D1 binding name in wrangler.jsonc", { + default: "DB", + }) + .option( + "--cron ", + "Cron expression to fire against /__scheduled (default: every minute)", + { default: "*/1 * * * *" } + ) + .option( + "--concurrency ", + "Concurrency passed to backfill if prompted (default: 100)", + { default: 100 } + ) + .option( + "--ignore-window ", + "Refresh ignore-window in seconds, if prompted (default: server default)" + ) + .option( + "--stale-after ", + "Prompt to refresh if the ingest cursor is older than this (default: 60)", + { default: 60 } + ) + .option("--yes, -y", "Accept all prompts without asking (CI-friendly)") + .action(async (options: DevOpts) => { + const config = await resolveAndLoadConfig(options); + + // Pre-flight: connect to local D1 via wrangler's platform proxy, inspect + // state, prompt. Then dispose before starting wrangler dev so the two + // miniflare processes don't fight over the sqlite file. + const { getPlatformProxy } = await import("wrangler"); + const { env, dispose } = await getPlatformProxy(); + const db = (env as Record)[options.binding] as + | Database + | undefined; + + if (db) { + const contrail = new Contrail(config); + await contrail.init(db); + + const hasBackfilled = await db + .prepare("SELECT 1 FROM backfills WHERE completed = 1 LIMIT 1") + .first(); + + if (!hasBackfilled) { + console.log("no backfilled users in the local DB yet."); + if ( + await promptYesNo( + "run backfill now? (takes a few minutes)", + true, + !!options.yes + ) + ) { + await contrail.backfillAll( + { concurrency: Number(options.concurrency) }, + db + ); + } + } else { + const staleAfterMs = Number(options.staleAfter) * 60_000; + const row = await db + .prepare("SELECT time_us FROM cursor WHERE id = 1") + .first<{ time_us: number }>(); + if (row?.time_us) { + const ageMs = Date.now() - Math.floor(row.time_us / 1000); + if (ageMs > staleAfterMs) { + const hrs = (ageMs / 3_600_000).toFixed(1); + console.log( + `ingest cursor is ${hrs}h old — you may have missed events.` + ); + if (await promptYesNo("run refresh first?", true, !!options.yes)) { + const ignoreWindowMs = + options.ignoreWindow !== undefined + ? Number(options.ignoreWindow) * 1000 + : undefined; + await contrail.refresh({ ignoreWindowMs }, db); + } + } + } + } + } + + await dispose(); + + // Start wrangler dev + fire /__scheduled on a loop so the cron actually + // runs in local dev (wrangler's cron scheduler only works in deployed + // production; --test-scheduled enables the manual-trigger endpoint). + const cronUrl = `http://localhost:8787/__scheduled?cron=${encodeURIComponent(options.cron)}`; + + const wrangler = spawn("npx", ["wrangler", "dev", "--test-scheduled"], { + stdio: "inherit", + shell: process.platform === "win32", + cwd: options.root, + }); + + // Give wrangler a few seconds to bind the port, then start the cron loop. + const kickoff = setTimeout(() => { + fetch(cronUrl).catch(() => {}); // fire-and-forget; first run + console.log(`\nauto-ingest: firing ${cronUrl} every 60s\n`); + }, 3_000); + + const interval = setInterval(() => { + fetch(cronUrl).catch(() => {}); // wrangler may be restarting; swallow + }, 60_000); + + const cleanup = () => { + clearTimeout(kickoff); + clearInterval(interval); + try { + wrangler.kill("SIGINT"); + } catch { + /* ignore */ + } + }; + process.on("SIGINT", cleanup); + process.on("SIGTERM", cleanup); + + const code = await new Promise((resolve) => { + wrangler.on("exit", (c) => { + clearTimeout(kickoff); + clearInterval(interval); + resolve(c ?? 0); + }); + }); + process.exit(code); + }); +} diff --git a/packages/contrail/src/cli/commands/refresh.ts b/packages/contrail/src/cli/commands/refresh.ts new file mode 100644 index 0000000..3c5764e --- /dev/null +++ b/packages/contrail/src/cli/commands/refresh.ts @@ -0,0 +1,50 @@ +import type { CAC } from "cac"; +import { refresh } from "../../workers/backfill.js"; +import { printRefreshReport, resolveAndLoadConfig } from "../shared.js"; + +interface RefreshOpts { + config?: string; + root?: string; + remote?: boolean; + binding: string; + concurrency: number; + ignoreWindow?: number; + byCollection?: boolean; +} + +export function registerRefresh(cli: CAC): void { + cli + .command( + "refresh", + "Fresh sweep: reconcile PDS vs DB, report missing + stale" + ) + .option("--config ", "Path to Contrail config file (TS or JS)") + .option("--root ", "Project root for auto-detection (default: CWD)") + .option("--remote", "Use production D1 bindings") + .option("--binding ", "D1 binding name in wrangler.jsonc", { + default: "DB", + }) + .option("--concurrency ", "Passed to contrail.refresh()", { + default: 50, + }) + .option( + "--ignore-window ", + "Seconds of grace for stale-update detection (default: 60)" + ) + .option("--by-collection", "Print per-collection stats, not just totals") + .action(async (options: RefreshOpts) => { + const config = await resolveAndLoadConfig(options); + const ignoreWindowMs = + options.ignoreWindow !== undefined + ? Number(options.ignoreWindow) * 1000 + : undefined; + const result = await refresh({ + config, + remote: !!options.remote, + binding: options.binding, + concurrency: Number(options.concurrency), + ignoreWindowMs, + }); + printRefreshReport(result, !!options.byCollection); + }); +} diff --git a/packages/contrail/src/cli/shared.ts b/packages/contrail/src/cli/shared.ts new file mode 100644 index 0000000..ecb5d2e --- /dev/null +++ b/packages/contrail/src/cli/shared.ts @@ -0,0 +1,90 @@ +/** + * Helpers shared across `contrail` subcommands. + */ +import { createInterface } from "node:readline/promises"; +import { stdin as input, stdout as output } from "node:process"; +import { + findConfigFile, + loadConfig, + CONFIG_CANDIDATES_MESSAGE, +} from "../cli-config.js"; +import type { ContrailConfig } from "../core/types.js"; +import type { CollectionStats, RefreshResult } from "../core/refresh.js"; + +export interface ConfigOpts { + config?: string; + root?: string; +} + +/** + * Resolve and load a ContrailConfig from CLI options. Exits with code 1 if no + * config file is found, since every command except `append-scheduled` needs one. + */ +export async function resolveAndLoadConfig( + opts: ConfigOpts +): Promise { + const root = opts.root ?? process.cwd(); + const path = findConfigFile(root, opts.config); + if (!path) { + console.error( + "Could not find a Contrail config. Pass --config or place one at\n" + + ` ${CONFIG_CANDIDATES_MESSAGE}` + ); + process.exit(1); + } + return loadConfig(path); +} + +/** + * Yes/no prompt that respects --yes (auto-accept) and falls back to the + * provided default in non-TTY environments — the user isn't there to answer. + */ +export async function promptYesNo( + question: string, + defaultYes: boolean, + autoYes: boolean +): Promise { + if (autoYes) return true; + if (!input.isTTY) return false; + const rl = createInterface({ input, output }); + try { + const hint = defaultYes ? "[Y/n]" : "[y/N]"; + const ans = (await rl.question(`${question} ${hint} `)).trim().toLowerCase(); + if (ans === "") return defaultYes; + return ans === "y" || ans === "yes"; + } finally { + rl.close(); + } +} + +export function formatStats(s: CollectionStats): string { + return `${s.missing} missing, ${s.staleUpdates} stale updates, ${s.inSync} in sync`; +} + +export function printRefreshReport( + result: RefreshResult, + byCollection: boolean +): void { + console.log(""); + if (byCollection) { + console.log("by collection:"); + const entries = Object.entries(result.byCollection).sort(([a], [b]) => + a.localeCompare(b) + ); + for (const [nsid, stats] of entries) { + if (stats.missing === 0 && stats.staleUpdates === 0 && stats.inSync === 0) + continue; + console.log(` ${nsid}`); + console.log(` ${formatStats(stats)}`); + } + console.log(""); + } + console.log("total:"); + console.log(` ${formatStats(result.total)}`); + console.log( + ` ${result.usersScanned} users scanned` + + (result.usersFailed ? `, ${result.usersFailed} failed` : "") + + ` in ${(result.elapsedMs / 1000).toFixed(1)}s` + ); + console.log(` (ignore window: ${(result.ignoreWindowMs / 1000).toFixed(0)}s)`); +} diff --git a/packages/contrail/src/workers/backfill.ts b/packages/contrail/src/workers/backfill.ts index 5f0ae6c..49edc70 100644 --- a/packages/contrail/src/workers/backfill.ts +++ b/packages/contrail/src/workers/backfill.ts @@ -1,12 +1,14 @@ /** * Wrangler-backed helpers for Cloudflare Workers deployments: * - * - `backfillAll` — one-time bulk load from scratch (uses the `backfills` - * state table to resume across runs). + * - `backfillAll` — one-time bulk record load from scratch (uses the + * `backfills` state table to resume across runs). + * - `labelsBackfillAll` — drain pending events per configured labeler in + * repeated cycles until each labeler's cursor stops advancing. * - `refresh` — reconcile every known DID's PDS against our DB, report * what's missing or stale. Use after outages or long idle periods. * - * Both dynamically import `wrangler` (optional peer dep) and wire it to + * All dynamically import `wrangler` (optional peer dep) and wire it to * the user's D1 binding, then dispose the proxy on exit. */ import { Contrail } from "../contrail.js"; @@ -95,3 +97,65 @@ export async function refresh( ) ); } + +export interface LabelsBackfillAllViaWranglerOptions extends WranglerCommon { + /** Per-cycle subscribe timeout passed to `contrail.ingestLabels()`. Default: 60s. */ + cycleTimeoutMs?: number; + /** Called after each cycle with whether any cursor advanced. */ + onCycle?: (info: { cycle: number; advanced: boolean }) => void; +} + +export interface LabelsBackfillAllResult { + /** Total ingest cycles run before any labeler stopped advancing twice in a row. */ + cycles: number; + /** Whether labels were configured at all — false means we no-op'd. */ + ran: boolean; +} + +/** + * Drain pending events from each configured labeler. Runs `ingestLabels` + * in a loop, checking the `labeler_cursors` table after each cycle, and + * stops once two consecutive cycles fail to advance any cursor. + */ +export async function labelsBackfillAll( + opts: LabelsBackfillAllViaWranglerOptions +): Promise { + return withWrangler(opts, async (contrail, db) => { + if (!opts.config.labels || opts.config.labels.sources.length === 0) { + return { cycles: 0, ran: false }; + } + const timeoutMs = opts.cycleTimeoutMs ?? 60_000; + let cycles = 0; + let stable = 0; + while (stable < 2) { + const before = new Map(); + const beforeRows = + ( + await db + .prepare("SELECT did, cursor FROM labeler_cursors") + .all<{ did: string; cursor: number }>() + ).results ?? []; + for (const r of beforeRows) before.set(r.did, r.cursor); + + await contrail.ingestLabels({ timeoutMs }, db); + cycles++; + + const afterRows = + ( + await db + .prepare("SELECT did, cursor FROM labeler_cursors") + .all<{ did: string; cursor: number }>() + ).results ?? []; + let advanced = false; + for (const r of afterRows) { + if ((before.get(r.did) ?? -1) !== r.cursor) { + advanced = true; + break; + } + } + stable = advanced ? 0 : stable + 1; + opts.onCycle?.({ cycle: cycles, advanced }); + } + return { cycles, ran: true }; + }); +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 5c07e90..8e648bf 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -373,6 +373,9 @@ importers: '@atcute/xrpc-server': specifier: ^0.1.12 version: 0.1.12 + cac: + specifier: ^7.0.0 + version: 7.0.0 hono: specifier: ^4.12.8 version: 4.12.15 @@ -2134,6 +2137,10 @@ packages: resolution: {integrity: sha512-b6Ilus+c3RrdDk+JhLKUAQfzzgLEPy6wcXqS7f/xe1EETvsDP6GORG7SFuOs6cID5YkqchW/LXZbX5bc8j7ZcQ==} engines: {node: '>=8'} + cac@7.0.0: + resolution: {integrity: sha512-tixWYgm5ZoOD+3g6UTea91eow5z6AAHaho3g0V9CNSNb45gM8SmflpAc+GRd1InC4AqN/07Unrgp56Y94N9hJQ==} + engines: {node: '>=20.19.0'} + chai@6.2.2: resolution: {integrity: sha512-NUPRluOfOiTKBKvWPtSD4PhFvWCqOi0BGStNWs57X9js7XGTprSmFoz5F0tWhR4WPjNeR9jXqdC7/UpSJTnlRg==} engines: {node: '>=18'} @@ -5221,6 +5228,8 @@ snapshots: cac@6.7.14: {} + cac@7.0.0: {} + chai@6.2.2: {} chardet@2.1.1: {}