From 2b1a9745f66b9ea9cce2a6b4c3aa805959a116f8 Mon Sep 17 00:00:00 2001 From: Samuel Newman Date: Fri, 7 Mar 2025 03:11:44 +0000 Subject: [PATCH] configurable jetstream instance --- packages/appview/src/ingestors/jetstream.ts | 10 +++++++--- packages/appview/src/lib/env.ts | 1 + packages/client/vite.config.ts | 11 ----------- 3 files changed, 8 insertions(+), 14 deletions(-) diff --git a/packages/appview/src/ingestors/jetstream.ts b/packages/appview/src/ingestors/jetstream.ts index 550dbaa..d9ac499 100644 --- a/packages/appview/src/ingestors/jetstream.ts +++ b/packages/appview/src/ingestors/jetstream.ts @@ -3,6 +3,7 @@ import pino from 'pino' import WebSocket from 'ws' import type { Database } from '#/db' +import { env } from '#/lib/env' export async function createJetstreamIngester(db: Database) { const logger = pino({ name: 'jetstream ingestion' }) @@ -19,6 +20,7 @@ export async function createJetstreamIngester(db: Database) { let lastCursorWrite = 0 return new Jetstream({ + instanceUrl: env.JETSTREAM_INSTANCE, logger, cursor: cursor?.seq || undefined, setCursor: async (seq) => { @@ -55,7 +57,6 @@ export async function createJetstreamIngester(db: Database) { ) if (!validatedRecord.success) return - // Store the status in our SQLite await db .insertInto('status') .values({ @@ -73,7 +74,6 @@ export async function createJetstreamIngester(db: Database) { ) .execute() } else if (evt.commit.operation === 'delete') { - // Remove the status from our SQLite await db.deleteFrom('status').where('uri', '=', uri).execute() } }, @@ -85,6 +85,7 @@ export async function createJetstreamIngester(db: Database) { } export class Jetstream { + private instanceUrl: string private logger: pino.Logger private handleEvent: (evt: JetstreamEvent) => Promise private onError: (err: unknown) => void @@ -95,6 +96,7 @@ export class Jetstream { private wantedCollections: string[] constructor({ + instanceUrl, logger, cursor, setCursor, @@ -102,6 +104,7 @@ export class Jetstream { onError, wantedCollections, }: { + instanceUrl: string logger: pino.Logger cursor?: number setCursor?: (seq: number) => Promise @@ -109,6 +112,7 @@ export class Jetstream { onError: (err: any) => void wantedCollections: string[] }) { + this.instanceUrl = instanceUrl this.logger = logger this.cursor = cursor this.setCursor = setCursor @@ -123,7 +127,7 @@ export class Jetstream { if (this.cursor !== undefined) { params.append('cursor', this.cursor.toString()) } - return `wss://jetstream.mozzius.dev/subscribe?${params.toString()}` + return `${this.instanceUrl}/subscribe?${params.toString()}` } start() { diff --git a/packages/appview/src/lib/env.ts b/packages/appview/src/lib/env.ts index 854bfd9..70aeac2 100644 --- a/packages/appview/src/lib/env.ts +++ b/packages/appview/src/lib/env.ts @@ -15,4 +15,5 @@ export const env = cleanEnv(process.env, { COOKIE_SECRET: str({ devDefault: '0'.repeat(32) }), SERVICE_DID: str({ default: undefined }), PUBLIC_URL: str({ devDefault: '' }), + JETSTREAM_INSTANCE: str({ default: 'wss://jetstream.mozzius.dev' }), }) diff --git a/packages/client/vite.config.ts b/packages/client/vite.config.ts index 7951d25..0d6fe9a 100644 --- a/packages/client/vite.config.ts +++ b/packages/client/vite.config.ts @@ -21,17 +21,6 @@ export default defineConfig({ '^/(xrpc|oauth|client-metadata\.json)/.*': { target: 'http://localhost:3001', changeOrigin: true, - configure: (proxy, _options) => { - proxy.on('error', (err, _req, _res) => { - console.log('PROXY ERROR', err); - }); - proxy.on('proxyReq', (proxyReq, req, _res) => { - console.log('PROXY REQUEST', req.method, req.url); - }); - proxy.on('proxyRes', (proxyRes, req, _res) => { - console.log('PROXY RESPONSE', req.method, req.url, proxyRes.statusCode); - }); - }, }, }, }, -- 2.51.2