diff --git a/AGENTS.md b/AGENTS.md deleted file mode 100644 index 47d5ef9..0000000 --- a/AGENTS.md +++ /dev/null @@ -1,71 +0,0 @@ -# CLAUDE.md - -This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository. - -## What is Contrail? - -Contrail is an AT Protocol (Bluesky) data indexer that runs on Cloudflare Workers + D1. You define collections in `src/config.ts`, and it automatically handles Jetstream ingestion, PDS backfill, user discovery via relays, and serves typed XRPC query endpoints. - -## Commands - -```bash -pnpm install # Install dependencies -pnpm dev # Start wrangler dev server (port 8787) -pnpm dev:auto # Start dev with auto-ingestion every 60s -pnpm generate # Regenerate lexicons + types from config -pnpm generate:pull # Pull lexicons from network, detect fields, generate types -pnpm typecheck # Run tsc --noEmit -pnpm ingest # Manually trigger ingestion cycle (requires dev server running) -pnpm sync # Discover users + backfill records from PDS (requires dev server) -pnpm deploy # Deploy to Cloudflare Workers -``` - -## Architecture - -### Entry points - -- **`src/adapters/worker.ts`** — Cloudflare Worker entry point. Handles HTTP fetch (Hono app) and scheduled cron (ingestion cycle). -- **`src/adapters/sqlite.ts`** — Local SQLite adapter wrapping better-sqlite3 to match the D1 `Database` interface. Excluded from tsconfig for Worker builds. - -### Core modules (`src/core/`) - -- **`types.ts`** — Central types (`Database`, `ContrailConfig`, `RecordRow`, `IngestEvent`) and config resolution logic. The `Database`/`Statement` interface abstracts over D1 and SQLite. -- **`jetstream.ts`** — Connects to Bluesky Jetstream WebSocket endpoints, collects events, runs full ingest cycles (load cursor → ingest → apply → save cursor → refresh identities). -- **`backfill.ts`** — Two functions: `backfillUser` (pages through a user's PDS via `com.atproto.repo.listRecords`) and `discoverDIDs` (discovers users via `com.atproto.sync.listReposByCollection` from relay endpoints). -- **`identity.ts`** — DID/handle resolution and caching in the `identities` table. -- **`client.ts`** — Creates authenticated AT Protocol XRPC clients for a given DID. -- **`queryable.generated.ts`** — Auto-generated file mapping collections to their queryable fields. **Do not edit manually** — run `pnpm generate`. - -### Database (`src/core/db/`) - -- **`schema.ts`** — Schema initialization: base tables (records, counts, backfills, discovery, cursor, identities) + dynamic indexes derived from queryable fields and relations. -- **`records.ts`** — Core DB operations: cursor management, applying ingest events (upsert/delete + relation count updates), querying records with filters/pagination. - -### Router (`src/core/router/`) - -- **`index.ts`** — Hono app factory with CORS, health checks, admin routes, and per-collection routes. -- **`collection.ts`** — Registers dynamic XRPC endpoints per collection: `getRecords`, `getRecord`, plus custom queries. -- **`admin.ts`** — Admin endpoints (`sync`, `getCursor`, `getOverview`) protected by `ADMIN_SECRET`. -- **`hydrate.ts`** — Hydration: embedding related records into responses. -- **`profiles.ts`** — Profile resolution: attaching handle + profile data to responses. - -### Code generation (`scripts/generate-lexicons.ts`) - -This script reads `src/config.ts`, auto-detects queryable fields from lexicon JSON files, generates XRPC lexicon definitions under `lexicons-generated/`, updates `lex.config.js`, and writes `src/core/queryable.generated.ts`. The `lex-cli generate` step then produces TypeScript types under `src/lexicon-types/`. - -### Key patterns - -- **Config-driven**: `src/config.ts` is the single source of truth. Adding a collection there drives schema creation, endpoint registration, Jetstream subscriptions, and code generation. -- **Database abstraction**: `Database`/`Statement` interfaces in `types.ts` let the same core logic run against both Cloudflare D1 and local SQLite. -- **Dependent collections**: Collections with `discover: false` (e.g., profiles) only index records for DIDs already known from discoverable collections. -- **Relations**: Materialized counts in the `counts` table, updated on every ingest event. `groupBy` splits counts by a field value. - -## Testing workflow - -Always follow TDD when making changes: - -1. **Write tests first** — add or update tests that cover the new behavior. Run them and confirm they **fail**. -2. **Update code** — implement the change. -3. **Run tests again** — confirm all tests **pass** (including existing ones). - -Tests live in `tests/` and use Vitest. Run with `npx vitest run`. diff --git a/README.md b/README.md index 1ea2be9..70f07e5 100644 --- a/README.md +++ b/README.md @@ -3,27 +3,49 @@ > [!WARNING] > Work in progress! Pre-alpha, expect breaking changes. -Define collections — get automatic Jetstream ingestion, PDS backfill, user discovery, and typed XRPC endpoints. Runs on Cloudflare Workers + D1. +A library for indexing AT Protocol records. Define collections — get automatic Jetstream ingestion, PDS backfill, user discovery, and typed XRPC endpoints. Works with Cloudflare Workers + D1, SvelteKit, Node.js, or any JavaScript runtime. -## Quickstart +## Install -## Config +```bash +npm install github:flo-bit/contrail +``` -Edit `src/config.ts` — this is the only file you need to touch: +## Usage ```ts -export const config: ContrailConfig = { - namespace: "com.example", // your reverse-domain namespace +import { Contrail } from "contrail"; + +const contrail = new Contrail({ + namespace: "com.example", + db, // any Database-compatible instance (D1, SQLite, etc.) collections: { "community.lexicon.calendar.event": { + queryable: { + mode: {}, // string → equality filter (?mode=online) + name: {}, // string → equality filter (?name=...) + startsAt: { type: "range" }, // range → min/max filters (?startsAtMin=...&startsAtMax=...) + endsAt: { type: "range" }, + }, + searchable: ["name", "description"], relations: { rsvps: { collection: "community.lexicon.calendar.rsvp", - groupBy: "status", // materialized counts by status + groupBy: "status", + count: true, + groups: { + interested: "community.lexicon.calendar.rsvp#interested", + going: "community.lexicon.calendar.rsvp#going", + notgoing: "community.lexicon.calendar.rsvp#notgoing", + }, }, }, }, "community.lexicon.calendar.rsvp": { + queryable: { + status: {}, + "subject.uri": {}, + }, references: { event: { collection: "community.lexicon.calendar.event", @@ -32,14 +54,107 @@ export const config: ContrailConfig = { }, }, }, +}); + +await contrail.init(); +``` + +### Query records + +```ts +const { records, cursor } = await contrail.query( + "community.lexicon.calendar.event", + { + filters: { mode: "in-person" }, + sort: { countType: "community.lexicon.calendar.rsvp", direction: "desc" }, + limit: 20, + } +); +``` + +### Ingest from Jetstream + +```ts +// Run one ingestion cycle (catches up to present, then stops) +await contrail.ingest(); +``` + +### Discover users and backfill + +```ts +// Find users from relays +await contrail.discover(); + +// Backfill their records from PDS +await contrail.backfill({ concurrency: 100 }); + +// Or both in one call +await contrail.sync({ concurrency: 100 }); +``` + +### Notify of immediate updates + +```ts +// After writing to a user's PDS, tell Contrail to fetch it now +await contrail.notify("at://did:plc:abc/community.lexicon.calendar.rsvp/123"); + +// Batch up to 25 URIs +await contrail.notify([uri1, uri2, uri3]); +``` + +### HTTP handler (XRPC endpoints) + +Mount the full XRPC API in any framework: + +```ts +import { createHandler } from "contrail/server"; + +const handle = createHandler(contrail); +// handle: (Request, db?) => Promise +``` + +**SvelteKit:** + +```ts +// src/routes/xrpc/[...path]/+server.ts +export const GET = ({ request }) => handle(request); +export const POST = ({ request }) => handle(request); +``` + +**Cloudflare Worker:** + +```ts +export default { + async fetch(request, env) { + return handle(request, env.DB); + }, }; ``` -### Dev +### SQLite adapter (Node.js / local dev) + +```ts +import { SqliteDatabase } from "contrail/sqlite"; +import Database from "better-sqlite3"; + +const db = new SqliteDatabase(new Database("data.db")); +const contrail = new Contrail({ ...config, db }); +``` + +## Running the example (Cloudflare Workers) + +This repo includes a working example that indexes AT Protocol calendar events and RSVPs on Cloudflare Workers + D1. + +### Setup ```bash pnpm install pnpm generate:pull # pull lexicons from network, auto-detect fields, generate types +``` + +### Dev + +```bash pnpm sync # discover users and backfill records from PDS pnpm dev:auto # start wrangler dev with auto-ingestion ``` @@ -48,35 +163,27 @@ pnpm dev:auto # start wrangler dev with auto-ingestion ```bash npx wrangler d1 create contrail -# Add database_id to wrangler.toml +# Add database_id to wrangler.jsonc pnpm deploy -# to sync in production, run it locally but set your d1 to remote, then run -pnpm sync +pnpm sync # discover + backfill against prod D1 ``` Ingestion runs automatically via cron (`*/1 * * * *`). Schema is auto-initialized. - -### What's auto-detected from lexicons - -When you run `pnpm generate`, queryable fields are derived from each collection's lexicon: - -- **String fields** → equality filter (`?status=going`) -- **Datetime/integer fields** → range filters (`?startsAtMin=2026-03-16&startsAtMax=2026-04-01`) -- **StrongRef fields** → `.uri` equality filter (`?subjectUri=at://...`) - -You can override any auto-detected field by specifying `queryable` manually in config. +## Config ### Collection options | Option | Default | Description | |--------|---------|-------------| -| `queryable` | auto-detected | Override auto-detected queryable fields | +| `queryable` | `{}` | Fields exposed as query filters. `{}` = string equality, `{ type: "range" }` = min/max | | `discover` | `true` | Find users via relays. `false` = only track known DIDs | | `relations` | `{}` | Many-to-one relationships with materialized counts | | `relations.*.field` | `"subject.uri"` | Field in the related record to match against | | `relations.*.match` | `"uri"` | Match against parent's `"uri"` or `"did"` | | `relations.*.groupBy` | — | Split counts by this field's value | +| `relations.*.groups` | — | Group value mappings (e.g. `{ going: "collection#going" }`) | +| `relations.*.count` | `true` | Enable materialized count columns on the parent | | `references` | `{}` | Forward references to other collections for hydration | | `references.*.collection` | — | Target collection NSID | | `references.*.field` | — | Field containing the target record's AT URI | @@ -84,13 +191,25 @@ You can override any auto-detected field by specifying `queryable` manually in c | `pipelineQueries` | `{}` | Custom query handlers that go through the standard filter/sort/hydration pipeline | | `searchable` | disabled | FTS5 search fields. Provide `string[]` to enable, omit to disable | +### Top-level options + +| Option | Default | Description | +|--------|---------|-------------| +| `namespace` | — | Your reverse-domain namespace (e.g. `"com.example"`) | +| `collections` | — | Collection configurations | +| `profiles` | `["app.bsky.actor.profile"]` | Profile collection NSIDs | +| `relays` | Bluesky relays | Relay URLs for user discovery | +| `jetstreams` | Bluesky Jetstream | Jetstream URLs for real-time ingestion | +| `feeds` | — | Personalized feed configurations | +| `logger` | `console` | Logger instance (`{ log, warn, error }`) | + ### Profiles `profiles` is a top-level config array of collection NSIDs that contain profile records (rkey `self`). Defaults to `["app.bsky.actor.profile"]`. These are auto-added to `collections` with `{ discover: false }`. Use `?profiles=true` on any endpoint to include a `profiles` map in the response, keyed by DID, with handle and profile record data. -## API +## XRPC API -All endpoints at `/xrpc/{nsid}.{method}`: +When using `createHandler`, all endpoints are available at `/xrpc/{nsid}.{method}`: | Endpoint | Description | |----------|-------------| @@ -98,10 +217,8 @@ All endpoints at `/xrpc/{nsid}.{method}`: | `{collection}.getRecord` | Get single record by URI | | `{namespace}.getProfile` | Get a user's profile by DID or handle | | `{namespace}.notifyOfUpdate` | Notify of a record change for immediate indexing | -| `{namespace}.admin.sync` | Discover + backfill (requires `ADMIN_SECRET`) | -| `{namespace}.admin.getCursor` | Current cursor position | -| `{namespace}.admin.getOverview` | All collections summary | -| `{namespace}.admin.reset` | Delete all data (requires `ADMIN_SECRET`) | +| `{namespace}.getCursor` | Current cursor position | +| `{namespace}.getOverview` | All collections summary | ### Query parameters @@ -111,7 +228,7 @@ All endpoints at `/xrpc/{nsid}.{method}`: |-------|---------|-------------| | `actor` | `?actor=did:plc:...` or `?actor=alice.bsky.social` | Filter by DID or handle (triggers on-demand backfill) | | `profiles` | `?profiles=true` | Include profile + identity info keyed by DID | -| `search` | `?search=meetup` | Full-text search across searchable fields (FTS5, ranked). Requires `searchable` to be configured. | +| `search` | `?search=meetup` | Full-text search across searchable fields (FTS5, ranked) | | `{field}` | `?status=going` | Equality filter on queryable string field | | `{field}Min` | `?startsAtMin=2026-03-16` | Range minimum (datetime/integer fields) | | `{field}Max` | `?endsAtMax=2026-04-01` | Range maximum (datetime/integer fields) | @@ -133,7 +250,7 @@ All endpoints at `/xrpc/{nsid}.{method}`: ?sort=rsvpsGoingCount&order=asc # by going count ascending ``` -**Search** uses SQLite FTS5 for ranked full-text search. To enable, set `searchable: ["field1", "field2"]` on a collection. Results are ranked by relevance (BM25) with `time_us` as tiebreaker. Supports FTS5 syntax including prefix (`meetup*`), phrases (`"rust meetup"`), and boolean (`rust OR typescript`). Combinable with all other filters. +**Search** uses SQLite FTS5 for ranked full-text search. To enable, set `searchable: ["field1", "field2"]` on a collection. Supports FTS5 syntax including prefix (`meetup*`), phrases (`"rust meetup"`), and boolean (`rust OR typescript`). Combinable with all other filters. ``` ?search=meetup # basic search @@ -161,7 +278,7 @@ All endpoints at `/xrpc/{nsid}.{method}`: # Single event with counts, RSVPs, and profiles /xrpc/community.lexicon.calendar.event.getRecord?uri=at://did:plc:.../community.lexicon.calendar.event/...&hydrateRsvps=10&profiles=true -# Search for events by name/description (requires searchable config) +# Search for events by name/description /xrpc/community.lexicon.calendar.event.listRecords?search=meetup&profiles=true # RSVPs for a specific event, with the referenced event embedded @@ -170,13 +287,13 @@ All endpoints at `/xrpc/{nsid}.{method}`: ## Notify of Updates -By default, Contrail ingests from Jetstream every minute. If your app writes to a user's PDS and needs the change reflected immediately, call `notifyOfUpdate` right after the write: +By default, Contrail ingests from Jetstream every minute (in the Worker example). If your app writes to a user's PDS and needs the change reflected immediately, use `contrail.notify()` or call the XRPC endpoint: ```ts -// User creates an RSVP via their PDS -const { uri } = await agent.createRecord({ ... }); +// Programmatic +await contrail.notify(uri); -// Tell Contrail to fetch and index it now +// Or via HTTP await fetch("https://your-contrail.workers.dev/xrpc/com.example.notifyOfUpdate", { method: "POST", headers: { "Content-Type": "application/json" }, @@ -188,22 +305,12 @@ Contrail fetches the record from the user's PDS and figures out what to do: | PDS returns | Already indexed? | Action | |---|---|---| -| Record (new CID) | No | **Create** — indexes it, updates relation counts | -| Record (new CID) | Yes | **Update** — upserts the record | +| Record (new CID) | No | **Create** — indexes it, recounts relations | +| Record (new CID) | Yes | **Update** — upserts the record, recounts relations | | Record (same CID) | Yes | **Skip** — nothing changed | -| 404 | Yes | **Delete** — removes it, decrements counts | +| 404 | Yes | **Delete** — removes it, recounts relations | | 404 | No | **No-op** | -You can also batch up to 25 URIs in one request: - -```ts -await fetch(".../xrpc/com.example.notifyOfUpdate", { - method: "POST", - headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ uris: [uri1, uri2, uri3] }), -}); -``` - When Jetstream later delivers the same event, the duplicate is detected by CID and skipped. ## Typesafe Client Usage @@ -259,6 +366,6 @@ const response = await rpc.get("community.lexicon.calendar.event.getRecords", { }); if (response.ok) { - console.log(response.data.records); // typed + console.log(response.data.records); // typed } ``` diff --git a/src/config.ts b/examples/cloudflare-workers/config.ts similarity index 57% rename from src/config.ts rename to examples/cloudflare-workers/config.ts index 3a0f5fb..9cb9265 100644 --- a/src/config.ts +++ b/examples/cloudflare-workers/config.ts @@ -1,19 +1,36 @@ -import type { ContrailConfig } from "./core/types"; +import type { ContrailConfig } from "../../src/index"; export const config: ContrailConfig = { namespace: "rsvp.atmo", collections: { "community.lexicon.calendar.event": { + queryable: { + mode: {}, + name: {}, + status: {}, + startsAt: { type: "range" }, + endsAt: { type: "range" }, + createdAt: { type: "range" }, + }, searchable: ["name", "description"], relations: { rsvps: { collection: "community.lexicon.calendar.rsvp", groupBy: "status", count: true, + groups: { + interested: "community.lexicon.calendar.rsvp#interested", + going: "community.lexicon.calendar.rsvp#going", + notgoing: "community.lexicon.calendar.rsvp#notgoing", + }, }, }, }, "community.lexicon.calendar.rsvp": { + queryable: { + status: {}, + "subject.uri": {}, + }, references: { event: { collection: "community.lexicon.calendar.event", diff --git a/examples/cloudflare-workers/generate.ts b/examples/cloudflare-workers/generate.ts new file mode 100644 index 0000000..e100806 --- /dev/null +++ b/examples/cloudflare-workers/generate.ts @@ -0,0 +1,19 @@ +/** + * Generates lexicon files and lex.config.js from config. + * + * Usage: npx tsx examples/cloudflare-workers/generate.ts + */ + +import { join, dirname } from "path"; +import { fileURLToPath } from "url"; +import { config } from "./config"; +import { generateLexicons } from "../../src/generate"; + +const ROOT_DIR = join(dirname(fileURLToPath(import.meta.url)), "../.."); + +generateLexicons({ + config, + rootDir: ROOT_DIR, + outputDir: join(ROOT_DIR, "lexicons-generated"), + writeRuntimeFiles: true, +}); diff --git a/examples/cloudflare-workers/sync.ts b/examples/cloudflare-workers/sync.ts new file mode 100644 index 0000000..323af32 --- /dev/null +++ b/examples/cloudflare-workers/sync.ts @@ -0,0 +1,70 @@ +/** + * Example: CLI sync script using Contrail as a library. + * + * Usage: + * npx tsx examples/cloudflare-workers/sync.ts # local D1 + * npx tsx examples/cloudflare-workers/sync.ts --remote # prod D1 + */ +import { Contrail } from "../../src/index"; +import { config } from "./config"; + +// For Cloudflare D1, use wrangler's getPlatformProxy: +import { getPlatformProxy } from "wrangler"; + +function elapsed(start: number): string { + const ms = Date.now() - start; + if (ms < 1000) return `${ms}ms`; + if (ms < 60_000) return `${(ms / 1000).toFixed(1)}s`; + const mins = Math.floor(ms / 60_000); + const secs = ((ms % 60_000) / 1000).toFixed(0); + return `${mins}m ${secs}s`; +} + +async function main() { + const remote = process.argv.includes("--remote"); + const syncStart = Date.now(); + + console.log(`=== Sync (${remote ? "remote/prod" : "local"} D1) ===\n`); + + const { env, dispose } = await getPlatformProxy<{ DB: D1Database }>({ + environment: remote ? "production" : undefined, + }); + + const contrail = new Contrail({ ...config, db: env.DB }); + + try { + await contrail.init(); + + console.log("--- Discovery ---"); + const discoveryStart = Date.now(); + const discovered = await contrail.discover(); + console.log(` Done: ${discovered.length} users in ${elapsed(discoveryStart)}\n`); + + console.log("--- Backfill ---"); + const backfillStart = Date.now(); + const total = await contrail.backfill({ + concurrency: 100, + onProgress: ({ records, usersComplete, usersTotal, usersFailed }) => { + const secs = (Date.now() - backfillStart) / 1000; + const rate = secs > 0 ? Math.round(records / secs) : 0; + const failStr = usersFailed > 0 ? ` | ${usersFailed} failed` : ""; + process.stdout.write( + `\r ${records} records | ${usersComplete}/${usersTotal} users | ${rate}/s | ${elapsed(backfillStart)}${failStr} ` + ); + }, + }); + process.stdout.write("\n"); + console.log(` Done: ${total} records in ${elapsed(backfillStart)}\n`); + + console.log(`=== Finished in ${elapsed(syncStart)} ===`); + console.log(` Discovered: ${discovered.length} users`); + console.log(` Backfilled: ${total} records`); + } finally { + await dispose(); + } +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/examples/cloudflare-workers/worker.ts b/examples/cloudflare-workers/worker.ts new file mode 100644 index 0000000..729d07f --- /dev/null +++ b/examples/cloudflare-workers/worker.ts @@ -0,0 +1,33 @@ +/** + * Example: Cloudflare Worker using Contrail as a library. + */ +import { Contrail } from "../../src/index"; +import { createHandler } from "../../src/server"; +import { config } from "./config"; + +const contrail = new Contrail(config); +const handle = createHandler(contrail); + +let initialized = false; + +export default { + async fetch(request: Request, env: { DB: D1Database }) { + if (!initialized) { + await contrail.init(env.DB); + initialized = true; + } + return handle(request, env.DB); + }, + + async scheduled( + _event: ScheduledEvent, + env: { DB: D1Database }, + ctx: ExecutionContext + ) { + if (!initialized) { + await contrail.init(env.DB); + initialized = true; + } + ctx.waitUntil(contrail.ingest({}, env.DB)); + }, +}; diff --git a/package.json b/package.json index 6f985f9..f77da72 100644 --- a/package.json +++ b/package.json @@ -1,17 +1,23 @@ { "name": "contrail", - "version": "1.0.0", + "version": "0.0.2", "private": true, + "type": "module", + "exports": { + ".": "./src/index.ts", + "./server": "./src/server.ts", + "./sqlite": "./src/adapters/sqlite.ts" + }, "scripts": { "dev": "wrangler dev --test-scheduled", "dev:auto": "tsx scripts/dev.ts", "deploy": "wrangler deploy", "clean": "tsx scripts/clean.ts", - "generate": "tsx scripts/generate-lexicons.ts && lex-cli generate", - "generate:pull": "tsx scripts/generate-lexicons.ts && lex-cli pull && tsx scripts/generate-lexicons.ts && lex-cli pull && lex-cli generate", + "generate": "tsx examples/cloudflare-workers/generate.ts", + "generate:pull": "tsx examples/cloudflare-workers/generate.ts && lex-cli pull && tsx examples/cloudflare-workers/generate.ts && lex-cli pull", "typecheck": "tsc --noEmit", "ingest": "curl -s http://localhost:8787/__scheduled?cron=*/1+*+*+*+*", - "sync": "tsx scripts/sync.ts", + "sync": "tsx examples/cloudflare-workers/sync.ts", "test": "vitest run", "test:watch": "vitest" }, diff --git a/scripts/clean.ts b/scripts/clean.ts index f3ffa5f..144384d 100644 --- a/scripts/clean.ts +++ b/scripts/clean.ts @@ -5,9 +5,10 @@ */ import { rmSync } from "fs"; -import { join } from "path"; +import { join, dirname } from "path"; +import { fileURLToPath } from "url"; -const ROOT = join(__dirname, ".."); +const ROOT = join(dirname(fileURLToPath(import.meta.url)), ".."); const targets = [ ".wrangler", @@ -15,7 +16,6 @@ const targets = [ "lexicons-pulled", "lexicons-generated", "src/lexicon-types", - "src/core/queryable.generated.ts", ]; for (const target of targets) { diff --git a/scripts/generate-lexicons.ts b/scripts/generate-lexicons.ts deleted file mode 100644 index 02891ec..0000000 --- a/scripts/generate-lexicons.ts +++ /dev/null @@ -1,18 +0,0 @@ -/** - * Generates lexicon files, lex.config.js, and queryable.generated.ts from config. - * - * Usage: npx tsx scripts/generate-lexicons.ts - */ - -import { join } from "path"; -import { config } from "../src/config"; -import { generateLexicons } from "../src/generate"; - -const ROOT_DIR = join(__dirname, ".."); - -generateLexicons({ - config, - rootDir: ROOT_DIR, - outputDir: join(ROOT_DIR, "lexicons-generated"), - writeRuntimeFiles: true, -}); diff --git a/scripts/sync.ts b/scripts/sync.ts deleted file mode 100644 index 9806ad7..0000000 --- a/scripts/sync.ts +++ /dev/null @@ -1,124 +0,0 @@ -/** - * Sync: discover users from relays and backfill their records from PDS. - * Runs directly against D1 via wrangler bindings — no dev server needed. - * - * Usage: - * npx tsx scripts/sync.ts # local D1 - * npx tsx scripts/sync.ts --remote # prod D1 - */ - -import { getPlatformProxy } from "wrangler"; -import { config as rawConfig } from "../src/config"; -import { resolveConfig, validateConfig, getCollectionNames } from "../src/core/types"; -import { initSchema } from "../src/core/db"; -import { discoverDIDs, backfillAll } from "../src/core/backfill"; - -const config = resolveConfig(rawConfig); -validateConfig(config); - -function elapsed(start: number): string { - const ms = Date.now() - start; - if (ms < 1000) return `${ms}ms`; - if (ms < 60_000) return `${(ms / 1000).toFixed(1)}s`; - const mins = Math.floor(ms / 60_000); - const secs = ((ms % 60_000) / 1000).toFixed(0); - return `${mins}m ${secs}s`; -} - -async function main() { - const remote = process.argv.includes("--remote"); - const syncStart = Date.now(); - - console.log(`=== Sync (${remote ? "remote/prod" : "local"} D1) ===\n`); - - const { env, dispose } = await getPlatformProxy<{ DB: D1Database }>({ - environment: remote ? "production" : undefined, - }); - const db = env.DB; - - try { - await initSchema(db, config); - - // Phase 1: Discover all DIDs from relays - console.log("--- Discovery ---"); - const discoveryStart = Date.now(); - const allDiscovered = new Set(); - while (true) { - const dids = await discoverDIDs(db, config, Infinity); - if (dids.length === 0) break; - for (const did of dids) allDiscovered.add(did); - console.log(` Found ${allDiscovered.size} unique users so far`); - } - console.log(` Done: ${allDiscovered.size} users in ${elapsed(discoveryStart)}\n`); - - // Ensure dependent collections have backfill entries for all known DIDs - const dependentCollections = getCollectionNames(config).filter( - (col) => config.collections[col]?.discover === false - ); - for (const depCol of dependentCollections) { - await db - .prepare( - `INSERT OR IGNORE INTO backfills (did, collection, completed) - SELECT i.did, ?, 0 FROM identities i - LEFT JOIN backfills b ON b.did = i.did AND b.collection = ? - WHERE b.did IS NULL` - ) - .bind(depCol, depCol) - .run(); - } - - // Phase 2: Backfill all pending records - console.log("--- Backfill ---"); - const backfillStart = Date.now(); - - const pending = await db - .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 0") - .first<{ count: number }>(); - const pendingCount = pending?.count ?? 0; - - const uniqueUsers = await db - .prepare("SELECT COUNT(DISTINCT did) as count FROM backfills WHERE completed = 0") - .first<{ count: number }>(); - - const userCount = uniqueUsers?.count ?? 0; - console.log(` ${pendingCount} pending collection backfills for ${userCount} users`); - - const total = await backfillAll(db, config, { - concurrency: 100, - onProgress: ({ records, usersComplete, usersTotal, usersFailed }) => { - const secs = (Date.now() - backfillStart) / 1000; - const rate = secs > 0 ? Math.round(records / secs) : 0; - const failStr = usersFailed > 0 ? ` | ${usersFailed} failed` : ""; - process.stdout.write( - `\r ${records} records | ${usersComplete}/${usersTotal} users | ${rate}/s | ${elapsed(backfillStart)}${failStr} ` - ); - }, - }); - process.stdout.write("\n"); - - console.log(` Done: ${total} records in ${elapsed(backfillStart)}\n`); - - // Summary - const finalRemaining = await db - .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 0") - .first<{ count: number }>(); - const failed = await db - .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 1 AND retries > 0") - .first<{ count: number }>(); - - console.log(`=== Finished in ${elapsed(syncStart)} ===`); - console.log(` Discovered: ${allDiscovered.size} users`); - console.log(` Backfilled: ${total} records`); - if ((finalRemaining?.count ?? 0) > 0) - console.log(` Remaining: ${finalRemaining!.count} backfills`); - if ((failed?.count ?? 0) > 0) - console.log(` Failed: ${failed!.count} (exceeded retries)`); - } finally { - await dispose(); - } -} - -main().catch((err) => { - console.error(err); - process.exit(1); -}); diff --git a/src/adapters/worker.ts b/src/adapters/worker.ts deleted file mode 100644 index bcba483..0000000 --- a/src/adapters/worker.ts +++ /dev/null @@ -1,38 +0,0 @@ -import { createApp } from "../core/router"; -import { initSchema } from "../core/db"; -import { runIngestCycle } from "../core/jetstream"; -import { config as rawConfig } from "../config"; -import { validateConfig, resolveConfig } from "../core/types"; - -const config = resolveConfig(rawConfig); -validateConfig(config); - -let schemaReady = false; - -async function ensureSchema(db: D1Database): Promise { - if (!schemaReady) { - await initSchema(db, config); - schemaReady = true; - } -} - -interface Env { - DB: D1Database; - ADMIN_SECRET?: string; -} - -export default { - async fetch(request: Request, env: Env): Promise { - await ensureSchema(env.DB); - return createApp(env.DB, config, env.ADMIN_SECRET).fetch(request); - }, - - async scheduled( - _event: ScheduledEvent, - env: Env, - ctx: ExecutionContext - ): Promise { - await ensureSchema(env.DB); - ctx.waitUntil(runIngestCycle(env.DB, config)); - }, -}; diff --git a/src/contrail.ts b/src/contrail.ts new file mode 100644 index 0000000..fd5a0b9 --- /dev/null +++ b/src/contrail.ts @@ -0,0 +1,91 @@ +import type { ContrailConfig, Database, ResolvedContrailConfig } from "./core/types"; +import { resolveConfig, validateConfig } from "./core/types"; +import { initSchema } from "./core/db/schema"; +import { queryRecords } from "./core/db/records"; +import type { QueryOptions, SortOption } from "./core/db/records"; +import { runIngestCycle } from "./core/jetstream"; +import { discoverDIDs, backfillAll } from "./core/backfill"; +import type { BackfillAllOptions, BackfillProgress } from "./core/backfill"; +import { processNotifyUris } from "./core/router/notify"; +import type { NotifyResult } from "./core/router/notify"; + +export interface ContrailOptions extends ContrailConfig { + db?: Database; +} + +export class Contrail { + readonly config: ResolvedContrailConfig; + private _db?: Database; + + constructor(options: ContrailOptions) { + const { db, ...configInput } = options; + this.config = resolveConfig(configInput); + validateConfig(this.config); + this._db = db; + } + + private getDb(db?: Database): Database { + const d = db ?? this._db; + if (!d) throw new Error("No database provided. Pass db to constructor or to this method."); + return d; + } + + /** Initialize the database schema. Must be called before other operations. */ + async init(db?: Database): Promise { + await initSchema(this.getDb(db), this.config); + } + + /** Query records from a collection. */ + async query( + collection: string, + options?: Omit, + db?: Database + ) { + return queryRecords(this.getDb(db), this.config, { collection, ...options }); + } + + /** Run one Jetstream ingestion cycle (catches up to present, then stops). */ + async ingest(options?: { timeoutMs?: number }, db?: Database): Promise { + await runIngestCycle(this.getDb(db), this.config, options?.timeoutMs); + } + + /** Discover users from relays. Returns discovered DIDs. */ + async discover(db?: Database): Promise { + const d = this.getDb(db); + const allDiscovered = new Set(); + while (true) { + const dids = await discoverDIDs(d, this.config, Infinity); + if (dids.length === 0) break; + for (const did of dids) allDiscovered.add(did); + } + return [...allDiscovered]; + } + + /** Backfill pending users' records from their PDS. */ + async backfill( + options?: BackfillAllOptions, + db?: Database + ): Promise { + return backfillAll(this.getDb(db), this.config, options); + } + + /** Discover + backfill in one call. */ + async sync( + options?: BackfillAllOptions, + db?: Database + ): Promise<{ discovered: number; backfilled: number }> { + const d = this.getDb(db); + const discovered = await this.discover(d); + const backfilled = await this.backfill(options, d); + return { discovered: discovered.length, backfilled }; + } + + /** Immediately fetch and index specific records from their PDS. */ + async notify( + uris: string | string[], + db?: Database + ): Promise { + const uriList = Array.isArray(uris) ? uris : [uris]; + return processNotifyUris(this.getDb(db), this.config, uriList); + } +} diff --git a/src/core/backfill.ts b/src/core/backfill.ts index 3d61827..17cfb73 100644 --- a/src/core/backfill.ts +++ b/src/core/backfill.ts @@ -412,9 +412,7 @@ async function fetchPage( `fetchPage(${relay}, ${collection})` ); } catch (err) { - console.error( - `Discovery failed for ${collection} from ${relay} after retries: ${err}` - ); + // Discovery page fetch failed after retries — skip this relay return null; } } diff --git a/src/core/db/records.ts b/src/core/db/records.ts index 1a03ad5..e339abe 100644 --- a/src/core/db/records.ts +++ b/src/core/db/records.ts @@ -1,5 +1,6 @@ import type { ContrailConfig, + ResolvedContrailConfig, RelationConfig, Database, Statement, @@ -8,7 +9,6 @@ import type { RecordSource, } from "../types"; import { getNestedValue, getRelationField, countColumnName, getFeedFollowCollections, recordsTableName } from "../types"; -import { resolvedRelationsMap } from "../queryable.generated"; import { getSearchableFields, ftsTableName, buildFtsContent } from "../search"; // --- Counts --- @@ -90,9 +90,9 @@ function buildCountStatements( // Grouped counts if (rel.groupBy) { - const mapping = (resolvedRelationsMap as Record)[parentCollection]?.[relationName]; + const mapping = (config as ResolvedContrailConfig)._resolved?.relations[parentCollection]?.[relationName]; if (mapping?.groups) { - for (const [, fullToken] of Object.entries(mapping.groups as Record)) { + for (const [, fullToken] of Object.entries(mapping.groups)) { const groupCol = countColumnName(fullToken); setClauses.push( `${groupCol} = (SELECT COUNT(*) FROM ${childTable} WHERE json_extract(record, '$.${field}') = ? AND json_extract(record, '$.${rel.groupBy}') = ?)` @@ -376,7 +376,7 @@ function getCountColumns(config: ContrailConfig, collection: string): { type: st const colConfig = config.collections[collection]; if (!colConfig?.relations) return []; const columns: { type: string; column: string }[] = []; - const relMap = resolvedRelationsMap[collection] ?? {}; + const relMap = (config as ResolvedContrailConfig)._resolved?.relations[collection] ?? {}; for (const [relName, rel] of Object.entries(colConfig.relations)) { if (rel.count === false) continue; diff --git a/src/core/db/schema.ts b/src/core/db/schema.ts index 63b3044..93699d7 100644 --- a/src/core/db/schema.ts +++ b/src/core/db/schema.ts @@ -1,8 +1,11 @@ -import type { ContrailConfig, Database } from "../types"; -import { getRelationField, countColumnName, recordsTableName } from "../types"; -import { resolvedQueryable, resolvedRelationsMap } from "../queryable.generated"; +import type { ContrailConfig, Database, ResolvedContrailConfig, ResolvedMaps } from "../types"; +import { getRelationField, countColumnName, recordsTableName, resolveConfig } from "../types"; import { getSearchableFields, ftsTableName } from "../search"; +function getResolved(config: ContrailConfig): ResolvedMaps { + return (config as ResolvedContrailConfig)._resolved ?? resolveConfig(config)._resolved; +} + const BASE_SCHEMA = ` CREATE TABLE IF NOT EXISTS backfills ( did TEXT NOT NULL, @@ -59,10 +62,11 @@ function buildCollectionTables(config: ContrailConfig): string[] { } function buildDynamicIndexes(config: ContrailConfig): string[] { + const resolved = getResolved(config); const indexes: string[] = []; for (const [collection, colConfig] of Object.entries(config.collections)) { const table = recordsTableName(collection); - const queryable = resolvedQueryable[collection] ?? colConfig.queryable ?? {}; + const queryable = resolved.queryable[collection] ?? colConfig.queryable ?? {}; for (const field of Object.keys(queryable)) { const idxName = `idx_${sanitizeName(collection)}_${sanitizeName(field)}`; indexes.push( @@ -84,12 +88,13 @@ function buildDynamicIndexes(config: ContrailConfig): string[] { } function buildCountColumns(config: ContrailConfig): string[] { + const resolved = getResolved(config); const stmts: string[] = []; const addedColumns = new Map>(); // table → columns for (const [collection, colConfig] of Object.entries(config.collections)) { const table = recordsTableName(collection); - const relMap = resolvedRelationsMap[collection] ?? {}; + const relMap = resolved.relations[collection] ?? {}; if (!addedColumns.has(table)) addedColumns.set(table, new Set()); const tableColumns = addedColumns.get(table)!; diff --git a/src/core/identity.ts b/src/core/identity.ts index 417fdb8..05040fd 100644 --- a/src/core/identity.ts +++ b/src/core/identity.ts @@ -1,5 +1,5 @@ import type { Did } from "@atcute/lexicons"; -import type { Database } from "./types"; +import type { Database, Logger } from "./types"; import { isDid, isHandle } from "@atcute/lexicons/syntax"; import { resolvePDS } from "./client"; @@ -82,8 +82,8 @@ export async function resolveIdentities( try { const identity = await fetchAndSave(db, did); map.set(did, identity); - } catch (err) { - console.warn(`Failed to resolve identity for ${did}: ${err}`); + } catch { + // Silently skip unresolvable identities } } @@ -152,8 +152,8 @@ export async function refreshStaleIdentities( for (const did of toRefresh) { try { await fetchAndSave(db, did); - } catch (err) { - console.warn(`Failed to refresh identity for ${did}: ${err}`); + } catch { + // Silently skip unresolvable identities } } } diff --git a/src/core/jetstream.ts b/src/core/jetstream.ts index 997ebc3..7d6261a 100644 --- a/src/core/jetstream.ts +++ b/src/core/jetstream.ts @@ -1,5 +1,5 @@ import { JetstreamSubscription } from "@atcute/jetstream"; -import type { ContrailConfig, IngestEvent, Database } from "./types"; +import type { ContrailConfig, IngestEvent, Database, Logger } from "./types"; import { getCollectionNames, getDependentCollections, DEFAULT_FEED_MAX_ITEMS } from "./types"; import { initSchema, getLastCursor, saveCursor, applyEvents, pruneFeedItems } from "./db"; import { refreshStaleIdentities } from "./identity"; @@ -12,12 +12,17 @@ let schemaInitialized = false; let lastFeedPruneMs = 0; const FEED_PRUNE_INTERVAL_MS = 60 * 60 * 1000; // 1 hour +function getLogger(config: ContrailConfig): Logger { + return config.logger ?? console; +} + export async function ingestEvents( config: ContrailConfig, cursor: number | null, safetyTimeoutMs: number = 25_000, knownDids?: Set ): Promise<{ events: IngestEvent[]; lastCursor: number | null }> { + const log = getLogger(config); const startTimeUs = Date.now() * 1000; const deadline = Date.now() + safetyTimeoutMs; const collected: IngestEvent[] = []; @@ -31,15 +36,15 @@ export async function ingestEvents( wantedCollections: collections, ...(cursor !== null ? { cursor } : {}), onConnectionOpen() { - console.log("Connected to Jetstream"); + log.log("Connected to Jetstream"); }, onConnectionClose(event) { - console.log( + log.log( `Disconnected from Jetstream: ${event.code} ${event.reason}` ); }, onConnectionError(event) { - console.error("Jetstream error:", event.error); + log.error("Jetstream error:", event.error); }, }); @@ -75,12 +80,12 @@ export async function ingestEvents( } if (event.time_us >= startTimeUs) { - console.log("Caught up to present, stopping ingestion"); + log.log("Caught up to present, stopping ingestion"); break; } if (Date.now() >= deadline) { - console.log("Safety timeout reached, stopping ingestion"); + log.log("Safety timeout reached, stopping ingestion"); break; } } @@ -95,6 +100,8 @@ export async function runIngestCycle( config: ContrailConfig, timeoutMs: number = 25_000 ): Promise { + const log = getLogger(config); + if (!schemaInitialized) { await initSchema(db, config); schemaInitialized = true; @@ -103,7 +110,7 @@ export async function runIngestCycle( const cursor = await getLastCursor(db); const collections = getCollectionNames(config); - console.log( + log.log( `Starting ingestion. Cursor: ${cursor ?? "none"}, Collections: ${collections.join(", ")}` ); @@ -114,14 +121,14 @@ export async function runIngestCycle( if (dependentCollections.length > 0) { if (cachedKnownDids) { knownDids = cachedKnownDids; - console.log(`Using cached known DIDs (${knownDids.size} users)`); + log.log(`Using cached known DIDs (${knownDids.size} users)`); } else { const result = await db .prepare("SELECT did FROM identities") .all<{ did: string }>(); knownDids = new Set((result.results ?? []).map((r) => r.did)); cachedKnownDids = knownDids; - console.log(`Loaded ${knownDids.size} known DIDs from database`); + log.log(`Loaded ${knownDids.size} known DIDs from database`); } } @@ -132,7 +139,7 @@ export async function runIngestCycle( knownDids ); - console.log(`Received ${events.length} events from Jetstream`); + log.log(`Received ${events.length} events from Jetstream`); for (let i = 0; i < events.length; i += BATCH_SIZE) { const batch = events.slice(i, i + BATCH_SIZE); @@ -145,13 +152,13 @@ export async function runIngestCycle( try { await refreshStaleIdentities(db, uniqueDids); } catch (err) { - console.warn(`Identity refresh failed: ${err}`); + log.warn(`Identity refresh failed: ${err}`); } } if (lastCursor !== null) { await saveCursor(db, lastCursor); - console.log(`Saved cursor: ${lastCursor}`); + log.log(`Saved cursor: ${lastCursor}`); } // Prune feed items hourly @@ -160,9 +167,9 @@ export async function runIngestCycle( ...Object.values(config.feeds).map((f) => f.maxItems ?? DEFAULT_FEED_MAX_ITEMS) ); const pruned = await pruneFeedItems(db, maxItems); - if (pruned > 0) console.log(`Pruned ${pruned} old feed items`); + if (pruned > 0) log.log(`Pruned ${pruned} old feed items`); lastFeedPruneMs = Date.now(); } - console.log(`Ingestion complete. Stored ${events.length} events.`); + log.log(`Ingestion complete. Stored ${events.length} events.`); } diff --git a/src/core/router/admin.ts b/src/core/router/admin.ts index 4dd4a28..c6b2e64 100644 --- a/src/core/router/admin.ts +++ b/src/core/router/admin.ts @@ -1,4 +1,4 @@ -import type { Hono, Context, Next } from "hono"; +import type { Hono } from "hono"; import type { ContrailConfig, Database } from "../types"; import { getCollectionNames, recordsTableName } from "../types"; import { getLastCursor } from "../db"; @@ -6,25 +6,11 @@ import { getLastCursor } from "../db"; export function registerAdminRoutes( app: Hono, db: Database, - config: ContrailConfig, - adminSecret?: string + config: ContrailConfig ): void { - const requireAdmin = async (c: Context, next: Next) => { - if (adminSecret) { - const auth = c.req.header("Authorization"); - if (auth !== `Bearer ${adminSecret}`) - return c.json({ error: "Unauthorized" }, 401); - } else { - const url = new URL(c.req.url); - if (url.hostname !== "localhost" && url.hostname !== "127.0.0.1") - return c.json({ error: "ADMIN_SECRET not configured" }, 403); - } - await next(); - }; - const ns = config.namespace; - app.get(`/xrpc/${ns}.admin.getCursor`, async (c) => { + app.get(`/xrpc/${ns}.getCursor`, async (c) => { const cursor = await getLastCursor(db); if (cursor === null) return c.json({ cursor: null }); @@ -36,7 +22,7 @@ export function registerAdminRoutes( }); }); - app.get(`/xrpc/${ns}.admin.getOverview`, async (c) => { + app.get(`/xrpc/${ns}.getOverview`, async (c) => { const collections: { collection: string; records: number; unique_users: number }[] = []; for (const collection of getCollectionNames(config)) { @@ -54,11 +40,4 @@ export function registerAdminRoutes( collections, }); }); - - app.get(`/xrpc/${ns}.admin.reset`, requireAdmin, async (c) => { - const collectionTables = getCollectionNames(config).map(recordsTableName); - const tables = [...collectionTables, "backfills", "discovery", "cursor", "identities"]; - await db.batch(tables.map((t) => db.prepare(`DELETE FROM ${t}`))); - return c.json({ ok: true }); - }); } diff --git a/src/core/router/collection.ts b/src/core/router/collection.ts index 8bcbcf3..6766099 100644 --- a/src/core/router/collection.ts +++ b/src/core/router/collection.ts @@ -1,7 +1,6 @@ import type { Hono } from "hono"; -import type { ContrailConfig, Database, RecordRow, QueryableField, RecordSource } from "../types"; +import type { ContrailConfig, ResolvedContrailConfig, Database, RecordRow, QueryableField, RecordSource } from "../types"; import { getCollectionNames, countColumnName, recordsTableName } from "../types"; -import { resolvedQueryable, resolvedRelationsMap } from "../queryable.generated"; import { queryRecords } from "../db"; import type { SortOption } from "../db/records"; import { backfillUser } from "../backfill"; @@ -24,7 +23,7 @@ export async function runPipeline( const relations = colConfig.relations ?? {}; const references = colConfig.references ?? {}; const queryableFields: Record = - resolvedQueryable[collection] ?? colConfig.queryable ?? {}; + (config as ResolvedContrailConfig)._resolved?.queryable[collection] ?? colConfig.queryable ?? {}; const limit = parseIntParam(params.get("limit"), 50); const cursor = params.get("cursor") || undefined; @@ -58,7 +57,7 @@ export async function runPipeline( } const countFilters: Record = {}; - const relMap = resolvedRelationsMap[collection] ?? {}; + const relMap = (config as ResolvedContrailConfig)._resolved?.relations[collection] ?? {}; for (const [relName, rel] of Object.entries(relations)) { const totalMin = parseIntParam(params.get(`${relName}CountMin`)); if (totalMin != null) countFilters[rel.collection] = totalMin; @@ -138,7 +137,7 @@ export async function runPipeline( const formattedRecords: FormattedRecord[] = rows.map((row) => { const formatted = formatRecord(row); - flattenCounts(formatted, row.counts, collection, relations); + flattenCounts(formatted, row.counts, relations); const h = hydrates[row.uri]; if (h) { for (const [relName, groups] of Object.entries(h)) { @@ -203,8 +202,8 @@ export function registerCollectionRoutes( if (!row) return c.json({ error: "Record not found" }, 404); const formatted = formatRecord({ ...row, collection }); - const counts = extractCounts(row, collection, relations); - if (counts) flattenCounts(formatted, counts, collection, relations); + const counts = extractCounts(row, relations); + if (counts) flattenCounts(formatted, counts, relations); const params = new URL(c.req.url).searchParams; const wantProfilesSingle = params.get("profiles") === "true"; @@ -277,21 +276,18 @@ export function registerCollectionRoutes( function extractCounts( row: any, - collection: string, relations: Record ): Record | undefined { - const relMap = resolvedRelationsMap[collection] ?? {}; const counts: Record = {}; - for (const [relName, rel] of Object.entries(relations)) { + for (const [, rel] of Object.entries(relations)) { if (rel.count === false) continue; const totalCol = countColumnName(rel.collection); const val = row[totalCol]; if (val != null && val !== 0) counts[rel.collection] = val; - const mapping = relMap[relName]; - if (mapping) { - for (const [, fullToken] of Object.entries(mapping.groups)) { + if (rel.groups) { + for (const [, fullToken] of Object.entries(rel.groups as Record)) { const groupCol = countColumnName(fullToken); const gval = row[groupCol]; if (gval != null && gval !== 0) counts[fullToken] = gval; @@ -305,24 +301,19 @@ function extractCounts( function flattenCounts( formatted: FormattedRecord, counts: Record | undefined, - collection: string, relations: Record ): void { if (!counts) return; - const relMap = resolvedRelationsMap[collection] ?? {}; const capitalize = (s: string) => s.charAt(0).toUpperCase() + s.slice(1); const collectionToRelName: Record = {}; const tokenToField: Record = {}; - for (const [relName, mapping] of Object.entries(relMap)) { - collectionToRelName[mapping.collection] = relName; - for (const [shortName, fullToken] of Object.entries(mapping.groups)) { - tokenToField[fullToken] = `${relName}${capitalize(shortName)}Count`; - } - } for (const [relName, rel] of Object.entries(relations)) { - if (!collectionToRelName[rel.collection]) { - collectionToRelName[rel.collection] = relName; + collectionToRelName[rel.collection] = relName; + if (rel.groups) { + for (const [shortName, fullToken] of Object.entries(rel.groups as Record)) { + tokenToField[fullToken] = `${relName}${capitalize(shortName)}Count`; + } } } diff --git a/src/core/router/index.ts b/src/core/router/index.ts index 09988a1..59f96e4 100644 --- a/src/core/router/index.ts +++ b/src/core/router/index.ts @@ -11,8 +11,7 @@ import { backfillUser } from "../backfill"; export function createApp( db: Database, - config: ContrailConfig, - adminSecret?: string + config: ContrailConfig ): Hono { const app = new Hono(); app.use("*", cors()); @@ -43,7 +42,7 @@ export function createApp( return c.json(profile); }); - registerAdminRoutes(app, db, config, adminSecret); + registerAdminRoutes(app, db, config); registerCollectionRoutes(app, db, config); registerFeedRoutes(app, db, config); registerNotifyRoute(app, db, config); diff --git a/src/core/router/notify.ts b/src/core/router/notify.ts index e8f681f..92164f9 100644 --- a/src/core/router/notify.ts +++ b/src/core/router/notify.ts @@ -35,6 +35,108 @@ async function fetchRecordFromPDS( return { value: data.value, cid: data.cid }; } +export interface NotifyResult { + indexed: number; + deleted: number; + errors?: string[]; +} + +/** + * Process notify URIs: fetch from PDS, detect changes, apply events. + * Shared by both the Hono route and the Contrail.notify() method. + */ +export async function processNotifyUris( + db: Database, + config: ContrailConfig, + uris: string[] +): Promise { + const events: IngestEvent[] = []; + const errors: string[] = []; + + for (const uri of uris) { + const parsed = parseAtUri(uri); + if (!parsed) { + errors.push(`invalid AT URI: ${uri}`); + continue; + } + + // Only accept collections we're tracking + if (!config.collections[parsed.collection]) { + errors.push(`collection not tracked: ${parsed.collection}`); + continue; + } + + const pds = await getPDS(parsed.did as Did, db); + if (!pds) { + errors.push(`could not resolve PDS for ${parsed.did}`); + continue; + } + + const result = await fetchRecordFromPDS( + pds, + parsed.did, + parsed.collection, + parsed.rkey + ); + + const now = Date.now() * 1000; // microseconds + + // Check if this record already exists locally + const table = recordsTableName(parsed.collection); + const existing = await db + .prepare(`SELECT cid FROM ${table} WHERE uri = ?`) + .bind(uri) + .first<{ cid: string | null }>(); + + if (result) { + if (existing?.cid === result.cid) { + // Same CID — nothing changed + continue; + } + + events.push({ + uri, + did: parsed.did, + collection: parsed.collection, + rkey: parsed.rkey, + operation: existing ? "update" : "create", + cid: result.cid, + record: JSON.stringify(result.value), + time_us: now, + indexed_at: now, + }); + } else if (existing) { + // Record gone from PDS but exists locally — delete it. + const existingRecord = await db + .prepare(`SELECT record FROM ${table} WHERE uri = ?`) + .bind(uri) + .first<{ record: string | null }>(); + + events.push({ + uri, + did: parsed.did, + collection: parsed.collection, + rkey: parsed.rkey, + operation: "delete", + cid: null, + record: existingRecord?.record ?? null, + time_us: now, + indexed_at: now, + }); + } + } + + if (events.length > 0) { + await applyEvents(db, events, config); + } + + return { + indexed: events.filter((e) => e.operation === "create" || e.operation === "update").length, + deleted: events.filter((e) => e.operation === "delete").length, + errors: errors.length > 0 ? errors : undefined, + }; +} + export function registerNotifyRoute( app: Hono, db: Database, @@ -58,93 +160,7 @@ export function registerNotifyRoute( return c.json({ error: "max 25 URIs per request" }, 400); } - const events: IngestEvent[] = []; - const errors: string[] = []; - - for (const uri of uris) { - const parsed = parseAtUri(uri); - if (!parsed) { - errors.push(`invalid AT URI: ${uri}`); - continue; - } - - // Only accept collections we're tracking - if (!config.collections[parsed.collection]) { - errors.push(`collection not tracked: ${parsed.collection}`); - continue; - } - - const pds = await getPDS(parsed.did as Did, db); - if (!pds) { - errors.push(`could not resolve PDS for ${parsed.did}`); - continue; - } - - const result = await fetchRecordFromPDS( - pds, - parsed.did, - parsed.collection, - parsed.rkey - ); - - const now = Date.now() * 1000; // microseconds - - // Check if this record already exists locally - const table = recordsTableName(parsed.collection); - const existing = await db - .prepare(`SELECT cid FROM ${table} WHERE uri = ?`) - .bind(uri) - .first<{ cid: string | null }>(); - - if (result) { - if (existing?.cid === result.cid) { - // Same CID — nothing changed, skip to avoid double-counting - continue; - } - - events.push({ - uri, - did: parsed.did, - collection: parsed.collection, - rkey: parsed.rkey, - // Both "update" and "create" trigger a full recount of related records. - operation: existing ? "update" : "create", - cid: result.cid, - record: JSON.stringify(result.value), - time_us: now, - indexed_at: now, - }); - } else if (existing) { - // Record gone from PDS but exists locally — delete it. - // We need the old record data so buildCountStatements can decrement counts. - const existingRecord = await db - .prepare(`SELECT record FROM ${table} WHERE uri = ?`) - .bind(uri) - .first<{ record: string | null }>(); - - events.push({ - uri, - did: parsed.did, - collection: parsed.collection, - rkey: parsed.rkey, - operation: "delete", - cid: null, - record: existingRecord?.record ?? null, - time_us: now, - indexed_at: now, - }); - } - // If not on PDS and not local, nothing to do - } - - if (events.length > 0) { - await applyEvents(db, events, config); - } - - return c.json({ - indexed: events.filter((e) => e.operation === "create" || e.operation === "update").length, - deleted: events.filter((e) => e.operation === "delete").length, - errors: errors.length > 0 ? errors : undefined, - }); + const result = await processNotifyUris(db, config, uris); + return c.json(result); }); } diff --git a/src/core/types.ts b/src/core/types.ts index b1fb76c..1512e13 100644 --- a/src/core/types.ts +++ b/src/core/types.ts @@ -24,6 +24,8 @@ export interface RelationConfig { groupBy?: string; /** Enable materialized count columns on the parent. Defaults to true. */ count?: boolean; + /** Pre-resolved group mappings: shortName → full token (e.g. { going: "community.lexicon.calendar.rsvp#going" }). Auto-computed from groupBy if omitted. */ + groups?: Record; } /** A forward reference: this collection's records point at another collection. */ @@ -85,6 +87,12 @@ export const DEFAULT_RELAYS = [ "https://relay1.us-east.bsky.network" ]; +export interface Logger { + log(...args: any[]): void; + warn(...args: any[]): void; + error(...args: any[]): void; +} + export interface ContrailConfig { namespace: string; collections: Record; @@ -92,12 +100,29 @@ export interface ContrailConfig { relays?: string[]; jetstreams?: string[]; feeds?: Record; + logger?: Logger; +} + +export interface ResolvedRelation { + collection: string; + groupBy: string; + groups: Record; // shortName → full token value +} + +export interface ResolvedMaps { + queryable: Record>; + relations: Record>; +} + +/** Config after resolveConfig() — has computed queryable/relation maps attached. */ +export interface ResolvedContrailConfig extends ContrailConfig { + _resolved: ResolvedMaps; } /** - * Resolve config: apply defaults and auto-add profile collections. + * Resolve config: apply defaults, auto-add profile collections, compute queryable maps. */ -export function resolveConfig(config: ContrailConfig): ContrailConfig { +export function resolveConfig(config: ContrailConfig): ResolvedContrailConfig { const profiles = config.profiles ?? DEFAULT_PROFILES; const collections = { ...config.collections }; for (const col of profiles) { @@ -115,13 +140,47 @@ export function resolveConfig(config: ContrailConfig): ContrailConfig { } } - return { + const base = { ...config, collections, profiles, jetstreams: config.jetstreams ?? DEFAULT_JETSTREAMS, relays: config.relays ?? DEFAULT_RELAYS, + logger: config.logger ?? console, }; + + return { + ...base, + _resolved: _resolveQueryableMaps(base), + }; +} + +function _resolveQueryableMaps(config: ContrailConfig): ResolvedMaps { + const queryable: Record> = {}; + const relations: Record> = {}; + + for (const [collection, colConfig] of Object.entries(config.collections)) { + if (colConfig.queryable) { + queryable[collection] = colConfig.queryable; + } + + if (colConfig.relations) { + for (const [relName, rel] of Object.entries(colConfig.relations)) { + if (!rel.groupBy) continue; + const groups: Record = rel.groups ? { ...rel.groups } : {}; + if (Object.keys(groups).length > 0) { + if (!relations[collection]) relations[collection] = {}; + relations[collection][relName] = { + collection: rel.collection, + groupBy: rel.groupBy, + groups, + }; + } + } + } + } + + return { queryable, relations }; } export function getFeedFollowCollections(config: ContrailConfig): string[] { diff --git a/src/generate.ts b/src/generate.ts index 21a43cf..c71838e 100644 --- a/src/generate.ts +++ b/src/generate.ts @@ -611,10 +611,6 @@ export function generateLexicons(options: GenerateOptions): Record> = ${JSON.stringify(resolvedQueryableMap, null, 2)};\n\nexport interface ResolvedRelation {\n collection: string;\n groupBy: string;\n groups: Record; // shortName → full token value\n}\n\nexport const resolvedRelationsMap: Record> = ${JSON.stringify(resolvedRelationsMap, null, 2)};\n`; - writeFileSync(join(rootDir, "src", "core", "queryable.generated.ts"), queryableContent); - log("Generated src/core/queryable.generated.ts"); } log("\nDone!"); diff --git a/src/index.ts b/src/index.ts new file mode 100644 index 0000000..4834e8f --- /dev/null +++ b/src/index.ts @@ -0,0 +1,28 @@ +export { Contrail } from "./contrail"; +export type { ContrailOptions } from "./contrail"; + +export type { + ContrailConfig, + CollectionConfig, + RelationConfig, + ReferenceConfig, + QueryableField, + FeedConfig, + Database, + Statement, + Logger, + IngestEvent, + RecordRow, + RecordSource, + ResolvedContrailConfig, + ResolvedMaps, + ResolvedRelation, + CustomQueryHandler, + PipelineQueryHandler, +} from "./core/types"; + +export { resolveConfig, validateConfig } from "./core/types"; + +export type { QueryOptions, SortOption } from "./core/db/records"; +export type { BackfillProgress, BackfillAllOptions } from "./core/backfill"; +export type { NotifyResult } from "./core/router/notify"; diff --git a/src/server.ts b/src/server.ts new file mode 100644 index 0000000..2bd69c1 --- /dev/null +++ b/src/server.ts @@ -0,0 +1,32 @@ +import type { Contrail } from "./contrail"; +import type { Database } from "./core/types"; +import { createApp } from "./core/router"; + +/** + * Create an HTTP handler from a Contrail instance. + * Returns a standard (Request, db?) => Promise function. + * + * Usage: + * const handle = createHandler(contrail); + * // SvelteKit: export const GET = ({ request }) => handle(request); + * // Workers: return handle(request, env.DB); + */ +export function createHandler( + contrail: Contrail +): (request: Request, db?: Database) => Promise | Response { + // Cache the Hono app when db is bound at construction + let cachedApp: ReturnType | null = null; + + return (request: Request, db?: Database) => { + const d = db ?? (contrail as any)._db; + if (!d) throw new Error("No database provided. Pass db to Contrail constructor or to handler."); + + // If db is the same bound instance, reuse the Hono app + if (!db && !cachedApp) { + cachedApp = createApp(d, contrail.config); + } + + const app = db ? createApp(d, contrail.config) : cachedApp!; + return app.fetch(request); + }; +} diff --git a/tests/helpers.ts b/tests/helpers.ts index e53e936..38a0a11 100644 --- a/tests/helpers.ts +++ b/tests/helpers.ts @@ -1,12 +1,13 @@ import { createSqliteDatabase } from "../src/adapters/sqlite"; -import type { Database, ContrailConfig } from "../src/core/types"; +import type { Database, ResolvedContrailConfig } from "../src/core/types"; +import { resolveConfig } from "../src/core/types"; import { initSchema } from "../src/core/db/schema"; export function createTestDb(): Database { return createSqliteDatabase(":memory:"); } -export const TEST_CONFIG: ContrailConfig = { +export const TEST_CONFIG: ResolvedContrailConfig = resolveConfig({ namespace: "com.example", collections: { "community.lexicon.calendar.event": { @@ -19,6 +20,11 @@ export const TEST_CONFIG: ContrailConfig = { rsvps: { collection: "community.lexicon.calendar.rsvp", groupBy: "status", + groups: { + interested: "community.lexicon.calendar.rsvp#interested", + going: "community.lexicon.calendar.rsvp#going", + notgoing: "community.lexicon.calendar.rsvp#notgoing", + }, }, }, }, @@ -31,7 +37,7 @@ export const TEST_CONFIG: ContrailConfig = { }, }, }, -}; +}); export async function createTestDbWithSchema(): Promise { const db = createTestDb(); diff --git a/tests/search.test.ts b/tests/search.test.ts index 375419e..71ac87b 100644 --- a/tests/search.test.ts +++ b/tests/search.test.ts @@ -1,10 +1,11 @@ import { describe, it, expect, beforeEach } from "vitest"; -import type { Database, ContrailConfig } from "../src/core/types"; +import type { Database } from "../src/core/types"; +import { resolveConfig } from "../src/core/types"; import { createTestDb, makeEvent } from "./helpers"; import { initSchema } from "../src/core/db/schema"; import { applyEvents, queryRecords } from "../src/core/db/records"; -const SEARCH_CONFIG: ContrailConfig = { +const SEARCH_CONFIG = resolveConfig({ namespace: "com.example", collections: { "community.lexicon.calendar.event": { @@ -31,7 +32,7 @@ const SEARCH_CONFIG: ContrailConfig = { searchable: false, // disabled }, }, -}; +}); let db: Database; diff --git a/tsconfig.json b/tsconfig.json index 565f91b..330eb14 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -12,6 +12,6 @@ "resolveJsonModule": true, "isolatedModules": true }, - "include": ["src"], + "include": ["src", "examples"], "exclude": ["src/adapters/sqlite.ts", "src/generate.ts"] } diff --git a/wrangler.jsonc b/wrangler.jsonc index f62e354..73bc147 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -1,7 +1,7 @@ { "$schema": "node_modules/wrangler/config-schema.json", "name": "contrail", - "main": "src/adapters/worker.ts", + "main": "examples/cloudflare-workers/worker.ts", "compatibility_date": "2025-12-25", "observability": { "enabled": true