Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/flo-bit/contrail. atproto backend in a bottle flo-bit.dev/contrail
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520import { mkdir, readFile, readdir, rm, writeFile } from "node:fs/promises";import { basename, dirname, isAbsolute, relative, resolve } from "node:path";import { fileURLToPath } from "node:url";import { performance } from "node:perf_hooks";import { createRequire } from "node:module";import { Contrail, type ContrailConfig, type Database,} from "@atmo-dev/contrail";import { createSqliteDatabase } from "@atmo-dev/contrail/sqlite";import { getPlatformProxy } from "wrangler";
const APP_DIR = resolve(dirname(fileURLToPath(import.meta.url)), "..");const CONFIG_DIR = resolve(APP_DIR, "configs");const CACHE_DIR = resolve(APP_DIR, ".cache");const RESULTS_DIR = resolve(APP_DIR, "results");const WRANGLER_CONFIG = resolve(APP_DIR, "wrangler.jsonc");const require = createRequire(import.meta.url);
type BenchmarkBackend = "d1" | "sqlite";
interface Options { config: string; backend: BenchmarkBackend; concurrency: number; pdsConcurrency: number; didsPerPds: number; maxAttempts: number; keepCache: boolean; excludeDids: string[]; validationLexicons?: string;}
function positiveInteger(raw: string | undefined, name: string, fallback: number): number { if (raw === undefined) return fallback; const value = Number(raw); if (!Number.isInteger(value) || value <= 0) { throw new Error(`${name} must be a positive integer (got ${raw})`); } return value;}
function parseArgs(argv: string[]): Options { let config: string | undefined; let backend: BenchmarkBackend = "d1"; let concurrencyRaw: string | undefined; let pdsConcurrencyRaw: string | undefined; let didsPerPdsRaw: string | undefined; let maxAttemptsRaw: string | undefined; let keepCache = false; let validationLexicons: string | undefined; const excludeDids: string[] = [];
for (let index = 0; index < argv.length; index++) { const arg = argv[index]; if (arg === "--config") config = argv[++index]; else if (arg.startsWith("--config=")) config = arg.slice("--config=".length); else if (arg === "--backend" || arg.startsWith("--backend=")) { const value = arg === "--backend" ? argv[++index] : arg.slice("--backend=".length); if (value !== "d1" && value !== "sqlite") { throw new Error(`--backend must be d1 or sqlite (got ${value})`); } backend = value; } else if (arg === "--concurrency") concurrencyRaw = argv[++index]; else if (arg.startsWith("--concurrency=")) { concurrencyRaw = arg.slice("--concurrency=".length); } else if (arg === "--pds-concurrency") pdsConcurrencyRaw = argv[++index]; else if (arg.startsWith("--pds-concurrency=")) { pdsConcurrencyRaw = arg.slice("--pds-concurrency=".length); } else if (arg === "--dids-per-pds") didsPerPdsRaw = argv[++index]; else if (arg.startsWith("--dids-per-pds=")) { didsPerPdsRaw = arg.slice("--dids-per-pds=".length); } else if (arg === "--max-attempts") maxAttemptsRaw = argv[++index]; else if (arg.startsWith("--max-attempts=")) { maxAttemptsRaw = arg.slice("--max-attempts=".length); } else if (arg === "--exclude-did") excludeDids.push(argv[++index]); else if (arg.startsWith("--exclude-did=")) { excludeDids.push(arg.slice("--exclude-did=".length)); } else if (arg === "--validation-lexicons") { validationLexicons = argv[++index]; } else if (arg.startsWith("--validation-lexicons=")) { validationLexicons = arg.slice("--validation-lexicons=".length); } else if (arg === "--keep-cache") keepCache = true; else if (arg === "--help" || arg === "-h") { console.log(`Usage: pnpm bench --config <file> [options]
Options: --config <file> JSON config path or filename under configs/ (required) --backend <name> d1 (default) or native sqlite --concurrency <n> Concurrent identity resolutions (default: 100) --pds-concurrency <n> Concurrent PDS hosts (default: 20) --dids-per-pds <n> Concurrent accounts per PDS (default: 3) --max-attempts <n> Immediate attempts per failed account (default: 1) --exclude-did <did> Skip an actor after discovery (repeatable) --validation-lexicons <dir> Enable strict Lexicon + CID validation with JSON docs --keep-cache Keep the disposable local D1 after the run`); process.exit(0); } else { throw new Error(`Unknown argument: ${arg}`); } }
if (!config) throw new Error("--config is required"); return { config, backend, concurrency: positiveInteger(concurrencyRaw, "--concurrency", 100), pdsConcurrency: positiveInteger(pdsConcurrencyRaw, "--pds-concurrency", 20), didsPerPds: positiveInteger(didsPerPdsRaw, "--dids-per-pds", 3), maxAttempts: positiveInteger(maxAttemptsRaw, "--max-attempts", 1), keepCache, excludeDids, validationLexicons, };}
async function resolveConfigPath(input: string): Promise<string> { const candidates = isAbsolute(input) ? [input] : [resolve(process.cwd(), input), resolve(CONFIG_DIR, input)]; for (const candidate of candidates) { try { await readFile(candidate); return candidate; } catch { // Try the next location. } } throw new Error(`Config not found: ${input}`);}
async function loadLexiconDirectory(input: string): Promise<{ path: string; documents: object[];}> { const path = isAbsolute(input) ? input : resolve(process.cwd(), input); const entries = await readdir(path, { recursive: true, withFileTypes: true }); const files = entries .filter((entry) => entry.isFile() && entry.name.endsWith(".json")) .map((entry) => resolve(entry.parentPath, entry.name)) .sort(); if (files.length === 0) throw new Error(`No Lexicon JSON files found: ${path}`); const documents = await Promise.all( files.map(async (file) => JSON.parse(await readFile(file, "utf8")) as object), ); return { path, documents };}
async function installedPackageVersion( name: string, packageRequire: NodeJS.Require = require,): Promise<string | null> { try { const packagePath = packageRequire.resolve(`${name}/package.json`); const metadata = JSON.parse(await readFile(packagePath, "utf8")) as { version?: unknown; }; return typeof metadata.version === "string" ? metadata.version : null; } catch { return null; }}
function elapsed(start: number): number { return Math.round((performance.now() - start) * 100) / 100;}
function safeName(path: string): string { return basename(path) .replace(/\.config\.json$/i, "") .replace(/\.json$/i, "") .replace(/[^a-z0-9_-]+/gi, "-");}
interface FetchMetric { requests: number; errors: number; total_ms: number; max_ms: number; statuses: Record<string, number>;}
function instrumentFetch(): { metrics: Map<string, FetchMetric>; restore(): void; maxActive(): number;} { const original = globalThis.fetch; const metrics = new Map<string, FetchMetric>(); let active = 0; let peakActive = 0;
globalThis.fetch = async (input, init) => { let key = "invalid-url"; try { const raw = input instanceof Request ? input.url : String(input); const url = new URL(raw); const operation = url.pathname.split("/").pop() || url.pathname; key = `${url.host}/${operation}`; } catch { // Keep the fallback key. }
const metric = metrics.get(key) ?? { requests: 0, errors: 0, total_ms: 0, max_ms: 0, statuses: {}, }; metrics.set(key, metric); metric.requests++; active++; peakActive = Math.max(peakActive, active); const start = performance.now(); try { const response = await original(input, init); const status = String(response.status); metric.statuses[status] = (metric.statuses[status] ?? 0) + 1; return response; } catch (error) { metric.errors++; throw error; } finally { const duration = performance.now() - start; metric.total_ms += duration; metric.max_ms = Math.max(metric.max_ms, duration); active--; } };
return { metrics, restore() { globalThis.fetch = original; }, maxActive() { return peakActive; }, };}
async function main(): Promise<void> { const options = parseArgs(process.argv.slice(2)); const configPath = await resolveConfigPath(options.config); const config = JSON.parse(await readFile(configPath, "utf8")) as ContrailConfig; const validation = options.validationLexicons ? await loadLexiconDirectory(options.validationLexicons) : null; if (validation) { config.validation = { strict: true, verifyCid: true, }; for (const collection of Object.values(config.collections)) { collection.validate = true; } } const wranglerPackagePath = options.backend === "d1" ? require.resolve("wrangler/package.json") : null; const runtime = { node: process.version, wrangler: options.backend === "d1" ? await installedPackageVersion("wrangler") : null, workerd: wranglerPackagePath ? await installedPackageVersion("workerd", createRequire(wranglerPackagePath)) : null, sqlite: options.backend === "sqlite" ? (process.versions.sqlite ?? null) : null, }; const name = `${safeName(configPath)}${ options.backend === "sqlite" ? "-sqlite" : "" }${validation ? "-validated" : ""}`; const cachePath = resolve( CACHE_DIR, `${name}-r${options.concurrency}-h${options.pdsConcurrency}-d${options.didsPerPds}`, );
if (options.backend === "d1") { await rm(cachePath, { recursive: true, force: true }); await mkdir(cachePath, { recursive: true }); } await mkdir(RESULTS_DIR, { recursive: true });
console.log(`config: ${configPath}`); console.log( `backend: ${ options.backend === "d1" ? "fresh local D1" : "fresh in-memory native SQLite" }`, ); console.log( options.backend === "d1" ? `runtime: Wrangler ${runtime.wrangler ?? "unknown"}, workerd ${runtime.workerd ?? "unknown"}` : `runtime: Node ${runtime.node}, SQLite ${runtime.sqlite ?? "unknown"}`, ); console.log(`resolution: ${options.concurrency}`); console.log(`PDS hosts: ${options.pdsConcurrency}`); console.log(`DIDs / PDS: ${options.didsPerPds}`); console.log(`max attempts: ${options.maxAttempts}`); console.log(`excluded: ${options.excludeDids.length} actors`); console.log(`validation: ${validation ? `${validation.documents.length} Lexicons + CID` : "disabled"}`); if (options.backend === "d1") console.log(`cache: reset ${cachePath}`);
const startedAt = new Date(); const totalStart = performance.now(); let proxy: Awaited<ReturnType<typeof getPlatformProxy>> | undefined; let bindingMs = 0; let initMs = 0; let discoveryMs = 0; let backfillMs = 0; let discovered = 0; let acceptedRecords = 0; let backfillMetrics: any; let overview: any; let diagnostics: any; let fetchInstrumentation: ReturnType<typeof instrumentFetch> | undefined;
try { let phaseStart = performance.now(); let db: Database; if (options.backend === "sqlite") { db = createSqliteDatabase(":memory:"); } else { proxy = await getPlatformProxy({ configPath: WRANGLER_CONFIG, persist: { path: cachePath }, remoteBindings: false, envFiles: [], }); db = (proxy.env as Record<string, unknown>).DB as Database; } bindingMs = elapsed(phaseStart); fetchInstrumentation = instrumentFetch(); const contrail = new Contrail({ ...config, db, ...(validation ? { lexicons: validation.documents } : {}), });
phaseStart = performance.now(); await contrail.init(); initMs = elapsed(phaseStart);
phaseStart = performance.now(); const dids = await contrail.discover(); discovered = dids.length; discoveryMs = elapsed(phaseStart); console.log(`discovered: ${discovered} accounts in ${(discoveryMs / 1000).toFixed(2)}s`);
for (const did of new Set(options.excludeDids)) { await db.prepare("DELETE FROM backfills WHERE did = ?").bind(did).run(); }
phaseStart = performance.now(); const backfillStart = phaseStart; let lastProgressAt = 0; acceptedRecords = await contrail.backfill({ concurrency: options.concurrency, pdsConcurrency: options.pdsConcurrency, didsPerPds: options.didsPerPds, maxAttempts: options.maxAttempts, onMetrics(metrics) { backfillMetrics = metrics; }, onProgress(progress) { const now = Date.now(); if (now - lastProgressAt < 2_000) return; lastProgressAt = now; const seconds = ((performance.now() - backfillStart) / 1000).toFixed(1); console.log( `progress: +${seconds}s ${progress.usersComplete}/${progress.usersTotal} accounts, ` + `${progress.records} accepted, ${progress.usersFailed} failed`, ); }, }); backfillMs = elapsed(phaseStart);
const response = await contrail .app() .fetch(new Request("http://benchmark/status")); if (!response.ok) throw new Error(`Status request failed: ${response.status}`); overview = await response.json(); diagnostics = await contrail.diagnostics(); } finally { fetchInstrumentation?.restore(); await proxy?.dispose(); if (options.backend === "d1" && !options.keepCache) { await rm(cachePath, { recursive: true, force: true }); } }
const totalMs = elapsed(totalStart); const completedAt = new Date(); const indexedRecords = Number(overview.total_records ?? 0); const relativeConfig = relative(APP_DIR, configPath); const network = Object.fromEntries( [...(fetchInstrumentation?.metrics ?? new Map())].map(([key, metric]) => [ key, { ...metric, total_ms: Math.round(metric.total_ms * 100) / 100, max_ms: Math.round(metric.max_ms * 100) / 100, }, ]), ); const result = { format: "contrail.backfill-benchmark", version: 1, config: relativeConfig.startsWith("..") ? configPath : relativeConfig, backend: options.backend === "d1" ? "wrangler-local-d1" : "native-sqlite", runtime, options: { concurrency: options.concurrency, pdsConcurrency: options.pdsConcurrency, didsPerPds: options.didsPerPds, maxAttempts: options.maxAttempts, excludedDids: options.excludeDids, validation: validation ? { lexicons: relative(APP_DIR, validation.path), documents: validation.documents.length, strict: true, verifyCid: true, } : null, }, started_at: startedAt.toISOString(), completed_at: completedAt.toISOString(), timings_ms: { binding: bindingMs, init: initMs, discovery: discoveryMs, backfill: backfillMs, total: totalMs, }, throughput: { accepted_records_per_second: backfillMs > 0 ? Math.round((acceptedRecords / (backfillMs / 1000)) * 100) / 100 : 0, indexed_records_per_second: totalMs > 0 ? Math.round((indexedRecords / (totalMs / 1000)) * 100) / 100 : 0, }, discovered_accounts: discovered, accepted_records: acceptedRecords, indexed_records: indexedRecords, peak_rss_kib: process.resourceUsage().maxRSS, phases: backfillMetrics, network: { max_concurrent: fetchInstrumentation?.maxActive() ?? 0, requests: network, }, backfill: overview.backfill, collections: overview.collections, ingest_diagnostics: diagnostics, };
const timestamp = completedAt.toISOString().replace(/[:.]/g, "-"); const resultPath = resolve( RESULTS_DIR, `${name}-r${options.concurrency}-h${options.pdsConcurrency}-d${options.didsPerPds}-${timestamp}.json`, ); await writeFile(resultPath, `${JSON.stringify(result, null, 2)}\n`);
console.log(""); console.log(`backfill: ${(backfillMs / 1000).toFixed(2)}s`); console.log(`total: ${(totalMs / 1000).toFixed(2)}s`); console.log(`indexed: ${indexedRecords} records`); console.log( `throughput: ${result.throughput.accepted_records_per_second.toFixed(2)} accepted records/s`, ); console.log( `accounts: ${overview.backfill.accounts.complete} complete, ` + `${overview.backfill.accounts.pending} pending, ` + `${overview.backfill.accounts.retrying} retrying, ` + `${overview.backfill.accounts.failed} failed`, ); console.log(`network max: ${result.network.max_concurrent} concurrent requests`); const rejected = (diagnostics ?? []).reduce( (total: number, diagnostic: { total?: number }) => total + Number(diagnostic.total ?? 0), 0, ); console.log(`rejections: ${rejected} aggregate admission decisions`); if (backfillMetrics) { console.log( `phases: resolution ${(backfillMetrics.resolution_ms / 1000).toFixed(2)}s, ` + `derived ${(backfillMetrics.derived_rebuild_ms / 1000).toFixed(2)}s`, ); for (const [collection, metric] of Object.entries( backfillMetrics.collections, ) as Array<[string, any]>) { console.log( `collection: ${collection} — ${metric.fetched_records} fetched, ` + `${metric.accepted_records} accepted, ${metric.requests} requests, ` + `${(metric.fetch_ms / 1000).toFixed(2)}s fetch, ` + `${(metric.projection_and_checkpoint_ms / 1000).toFixed(2)}s project/checkpoint`, ); } } for (const [key, metric] of Object.entries(network)) { console.log( `network: ${key} — ${metric.requests} requests, ` + `${(metric.total_ms / 1000).toFixed(2)}s cumulative, ` + `${(metric.max_ms / 1000).toFixed(2)}s max`, ); } console.log(`result: ${resultPath}`);
if (overview.backfill.state !== "complete") { console.error(`Benchmark ended with backfill state: ${overview.backfill.state}`); process.exitCode = 1; }}
main().catch((error) => { console.error(error instanceof Error ? error.stack : error); process.exitCode = 1;});