diff --git a/README.md b/README.md index 70f07e5..7e1e83c 100644 --- a/README.md +++ b/README.md @@ -79,6 +79,20 @@ const { records, cursor } = await contrail.query( await contrail.ingest(); ``` +### Persistent ingestion + +```ts +// Long-lived Jetstream connection with automatic batching and reconnection +const controller = new AbortController(); +await contrail.runPersistent({ + batchSize: 50, // flush every N events (default: 50) + flushIntervalMs: 5000, // or every N ms (default: 5000) + signal: controller.signal, +}); +``` + +Call `controller.abort()` for graceful shutdown — the current batch is flushed and the cursor is saved. + ### Discover users and backfill ```ts @@ -141,7 +155,59 @@ const db = new SqliteDatabase(new Database("data.db")); const contrail = new Contrail({ ...config, db }); ``` -## Running the example (Cloudflare Workers) +### PostgreSQL adapter (Node.js / server) + +```ts +import { createPostgresDatabase } from "contrail/postgres"; +import pg from "pg"; + +const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL }); +const db = createPostgresDatabase(pool); +const contrail = new Contrail({ ...config, db }); +``` + +PostgreSQL uses JSONB for record storage, tsvector generated columns for full-text search (instead of FTS5), and `BIGINT` for timestamp columns. + +## Testing + +### SQLite tests (default) + +```bash +pnpm install +pnpm test +``` + +All tests under `tests/` (except `postgres.test.ts`) run against an in-memory SQLite database with no external dependencies. + +### PostgreSQL tests + +PostgreSQL integration tests require a running PostgreSQL instance and a dedicated test database. + +```bash +# Create a test database (one-time setup) +createdb contrail_test + +# Run PostgreSQL tests +TEST_DATABASE_URL="postgresql://user:password@localhost:5432/contrail_test" pnpm test -- tests/postgres.test.ts +``` + +The test suite drops and recreates Contrail tables in `beforeEach` — it will **not** touch other databases or schemas. + +To run both SQLite and PostgreSQL tests together: + +```bash +TEST_DATABASE_URL="postgresql://user:password@localhost:5432/contrail_test" pnpm test +``` + +If `TEST_DATABASE_URL` is not set, PostgreSQL tests are automatically skipped. + +## Examples + +### PostgreSQL (Node.js) + +See [`examples/postgres/`](examples/postgres/) for a complete example with Docker Compose, persistent Jetstream ingestion, user discovery/backfill, and an HTTP API server. + +### Cloudflare Workers This repo includes a working example that indexes AT Protocol calendar events and RSVPs on Cloudflare Workers + D1. @@ -189,7 +255,7 @@ Ingestion runs automatically via cron (`*/1 * * * *`). Schema is auto-initialize | `references.*.field` | — | Field containing the target record's AT URI | | `queries` | `{}` | Custom query handlers (raw Response) | | `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 | +| `searchable` | disabled | Full-text search fields. SQLite uses FTS5 virtual tables; PostgreSQL uses tsvector generated columns with GIN indexes. Provide `string[]` to enable, omit to disable | ### Top-level options @@ -250,7 +316,7 @@ When using `createHandler`, all endpoints are available 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. Supports FTS5 syntax including prefix (`meetup*`), phrases (`"rust meetup"`), and boolean (`rust OR typescript`). Combinable with all other filters. +**Search** uses SQLite FTS5 or PostgreSQL tsvector 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 diff --git a/examples/postgres/README.md b/examples/postgres/README.md new file mode 100644 index 0000000..12e1308 --- /dev/null +++ b/examples/postgres/README.md @@ -0,0 +1,112 @@ +# Contrail — PostgreSQL Example + +A complete example of using Contrail to index AT Protocol calendar events and RSVPs with PostgreSQL. Includes persistent Jetstream ingestion (long-running listener) and user discovery/backfill. + +## Setup + +```bash +# Copy this folder to a new project +cp -r examples/postgres my-contrail-app +cd my-contrail-app + +# Install dependencies +npm install +``` + +> **Note:** The `contrail` dependency in `package.json` points at `github:flo-bit/contrail`. +> If you're using a fork with PostgreSQL support that hasn't been merged yet, update the +> dependency to point at your fork's branch: +> +> ```bash +> npm install github:your-username/contrail#your-branch +> ``` +> +> Or install from a local checkout: +> +> ```bash +> npm install /path/to/your/contrail +> ``` + +### Start PostgreSQL + +**Option A: Docker (recommended)** + +```bash +docker compose up -d +``` + +This starts PostgreSQL on port 5433 (to avoid conflicts with a local instance) with a `contrail` database ready to use. Override the port with `PG_PORT=5432 docker compose up -d`. + +**Option B: Native PostgreSQL** + +```bash +createdb contrail +``` + +Set `DATABASE_URL` to point at your local instance: + +```bash +export DATABASE_URL="postgresql://user:password@localhost:5432/contrail" +``` + +## Configure + +Edit `config.ts` to define your collections, queryable fields, relations, and references. See the [Contrail README](https://github.com/flo-bit/contrail) for all options. + +## Run + +### 1. Discover users and backfill records + +```bash +npm run sync +``` + +This finds users from ATProto relays and backfills their existing records from PDS. Safe to interrupt and restart — progress is saved per-DID in the database. + +### 2. Start persistent ingestion + +```bash +npm run ingest +``` + +This opens a long-lived Jetstream connection and continuously indexes new records as they appear on the network. Events are batched and flushed every 5 seconds (or every 50 events, whichever comes first). Handles reconnection automatically. + +Press `Ctrl+C` for graceful shutdown — the current batch is flushed and the cursor is saved so the next run picks up where it left off. + +### 3. Serve the XRPC API + +```bash +npm run serve +``` + +Your XRPC API is now available at `http://localhost:3000`: + +``` +# List events sorted by RSVP count +/xrpc/community.lexicon.calendar.event.listRecords?sort=rsvpsCount + +# Upcoming events with 10+ going RSVPs +/xrpc/community.lexicon.calendar.event.listRecords?startsAtMin=2026-03-16&rsvpsGoingCountMin=10 + +# Single event with hydrated RSVPs and profiles +/xrpc/community.lexicon.calendar.event.getRecord?uri=at://...&hydrateRsvps=10&profiles=true + +# Search events +/xrpc/community.lexicon.calendar.event.listRecords?search=meetup + +# RSVPs for a specific event +/xrpc/community.lexicon.calendar.rsvp.listRecords?subjectUri=at://... +``` + +## Running everything together + +In production you'd typically run sync once (or periodically), then keep `ingest` and `serve` running as separate processes: + +```bash +# Initial sync (run once, or periodically to discover new users) +npm run sync + +# In separate terminals (or use a process manager) +npm run ingest +npm run serve +``` diff --git a/examples/postgres/config.ts b/examples/postgres/config.ts new file mode 100644 index 0000000..31c3ca1 --- /dev/null +++ b/examples/postgres/config.ts @@ -0,0 +1,42 @@ +import type { ContrailConfig } from "contrail"; + +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", + field: "subject.uri", + }, + }, + }, + }, +}; diff --git a/examples/postgres/docker-compose.yml b/examples/postgres/docker-compose.yml new file mode 100644 index 0000000..7a881bf --- /dev/null +++ b/examples/postgres/docker-compose.yml @@ -0,0 +1,14 @@ +services: + postgres: + image: postgres:17 + environment: + POSTGRES_DB: contrail + POSTGRES_USER: contrail + POSTGRES_PASSWORD: contrail + ports: + - "${PG_PORT:-5433}:5432" + volumes: + - pgdata:/var/lib/postgresql/data + +volumes: + pgdata: diff --git a/examples/postgres/ingest.ts b/examples/postgres/ingest.ts new file mode 100644 index 0000000..b5dcdd6 --- /dev/null +++ b/examples/postgres/ingest.ts @@ -0,0 +1,42 @@ +/** + * Persistent Jetstream ingestion against PostgreSQL. + * + * Opens a long-lived Jetstream connection and continuously indexes new records. + * Events are batched and flushed periodically. Handles reconnection automatically. + * + * Usage: + * DATABASE_URL="postgresql://contrail:contrail@localhost:5432/contrail" npx tsx ingest.ts + */ +import pg from "pg"; +import { Contrail } from "contrail"; +import { createPostgresDatabase } from "contrail/postgres"; +import { config } from "./config"; + +const DATABASE_URL = + process.env.DATABASE_URL ?? + "postgresql://contrail:contrail@localhost:5433/contrail"; + +async function main() { + const pool = new pg.Pool({ connectionString: DATABASE_URL }); + const db = createPostgresDatabase(pool); + const contrail = new Contrail({ ...config, db }); + + const controller = new AbortController(); + const shutdown = () => { + console.log("\nShutting down gracefully..."); + controller.abort(); + }; + process.on("SIGTERM", shutdown); + process.on("SIGINT", shutdown); + + console.log("Starting persistent ingestion..."); + await contrail.runPersistent({ signal: controller.signal }); + + await pool.end(); + console.log("Done."); +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/examples/postgres/package.json b/examples/postgres/package.json new file mode 100644 index 0000000..70f6d02 --- /dev/null +++ b/examples/postgres/package.json @@ -0,0 +1,20 @@ +{ + "name": "contrail-postgres-example", + "version": "0.0.1", + "private": true, + "type": "module", + "scripts": { + "sync": "tsx sync.ts", + "ingest": "tsx ingest.ts", + "serve": "tsx serve.ts" + }, + "dependencies": { + "contrail": "github:flo-bit/contrail", + "pg": "^8.20.0" + }, + "devDependencies": { + "@types/pg": "^8.20.0", + "tsx": "^4.21.0", + "typescript": "^5.7.3" + } +} diff --git a/examples/postgres/serve.ts b/examples/postgres/serve.ts new file mode 100644 index 0000000..c3713c2 --- /dev/null +++ b/examples/postgres/serve.ts @@ -0,0 +1,61 @@ +/** + * Serve the Contrail XRPC API over HTTP using PostgreSQL. + * + * Usage: + * DATABASE_URL="postgresql://contrail:contrail@localhost:5432/contrail" npx tsx serve.ts + */ +import pg from "pg"; +import { createServer } from "node:http"; +import { Contrail } from "contrail"; +import { createHandler } from "contrail/server"; +import { createPostgresDatabase } from "contrail/postgres"; +import { config } from "./config"; + +const DATABASE_URL = + process.env.DATABASE_URL ?? + "postgresql://contrail:contrail@localhost:5433/contrail"; +const PORT = parseInt(process.env.PORT ?? "3000", 10); + +async function main() { + const pool = new pg.Pool({ connectionString: DATABASE_URL }); + const db = createPostgresDatabase(pool); + const contrail = new Contrail({ ...config, db }); + + await contrail.init(); + + const handle = createHandler(contrail); + + const server = createServer(async (req, res) => { + const url = new URL(req.url!, `http://localhost:${PORT}`); + const request = new Request(url.toString(), { + method: req.method, + headers: Object.entries(req.headers).reduce( + (h, [k, v]) => { + if (v) h[k] = Array.isArray(v) ? v.join(", ") : v; + return h; + }, + {} as Record, + ), + }); + + const response = await handle(request); + res.writeHead(response.status, Object.fromEntries(response.headers)); + res.end(await response.text()); + }); + + server.listen(PORT, () => { + console.log(`Contrail XRPC API listening on http://localhost:${PORT}`); + }); + + const shutdown = () => { + console.log("\nShutting down..."); + server.close(() => pool.end().then(() => process.exit(0))); + }; + process.on("SIGTERM", shutdown); + process.on("SIGINT", shutdown); +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/examples/postgres/sync.ts b/examples/postgres/sync.ts new file mode 100644 index 0000000..65d055d --- /dev/null +++ b/examples/postgres/sync.ts @@ -0,0 +1,71 @@ +/** + * Discover users and backfill their records into PostgreSQL. + * + * Safe to kill at any point — discovery and backfill cursors are saved + * per-DID in the database. Restarting resumes from where it left off. + * + * Usage: + * DATABASE_URL="postgresql://contrail:contrail@localhost:5432/contrail" npx tsx sync.ts + */ +import pg from "pg"; +import { Contrail } from "contrail"; +import { createPostgresDatabase } from "contrail/postgres"; +import { config } from "./config"; + +const DATABASE_URL = + process.env.DATABASE_URL ?? + "postgresql://contrail:contrail@localhost:5433/contrail"; + +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 pool = new pg.Pool({ connectionString: DATABASE_URL }); + const db = createPostgresDatabase(pool); + const contrail = new Contrail({ ...config, db }); + const syncStart = Date.now(); + + console.log(`=== Sync (PostgreSQL) ===\n`); + + 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: 10, + 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`); + + await pool.end(); +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/examples/postgres/tsconfig.json b/examples/postgres/tsconfig.json new file mode 100644 index 0000000..9b7d9df --- /dev/null +++ b/examples/postgres/tsconfig.json @@ -0,0 +1,13 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "ES2022", + "moduleResolution": "bundler", + "lib": ["ES2022"], + "strict": true, + "noEmit": true, + "skipLibCheck": true, + "isolatedModules": true + }, + "include": ["."] +}