From 1317e8d1dc42a12ece181a0a4282d0448e43952c Mon Sep 17 00:00:00 2001 From: Samuel Newman Date: Fri, 7 Mar 2025 03:08:14 +0000 Subject: [PATCH] add jetstream --- lexicons/xyz/statusphere/getStatuses.json | 1 - packages/appview/src/context.ts | 7 +- packages/appview/src/db.ts | 15 ++ packages/appview/src/index.ts | 7 +- .../{ingester.ts => ingestors/firehose.ts} | 5 +- packages/appview/src/ingestors/index.ts | 2 + packages/appview/src/ingestors/jetstream.ts | 213 ++++++++++++++++++ packages/appview/src/lexicons/lexicons.ts | 3 - .../types/xyz/statusphere/getStatuses.ts | 1 - packages/client/src/components/StatusForm.tsx | 2 +- packages/client/src/components/StatusList.tsx | 13 +- packages/client/src/pages/HomePage.tsx | 17 +- packages/client/vite.config.ts | 11 + packages/lexicon/src/lexicons.ts | 3 - .../src/types/xyz/statusphere/getStatuses.ts | 1 - 15 files changed, 272 insertions(+), 29 deletions(-) rename packages/appview/src/{ingester.ts => ingestors/firehose.ts} (96%) create mode 100644 packages/appview/src/ingestors/index.ts create mode 100644 packages/appview/src/ingestors/jetstream.ts diff --git a/lexicons/xyz/statusphere/getStatuses.json b/lexicons/xyz/statusphere/getStatuses.json index d5cea78..156515a 100644 --- a/lexicons/xyz/statusphere/getStatuses.json +++ b/lexicons/xyz/statusphere/getStatuses.json @@ -22,7 +22,6 @@ "type": "object", "required": ["statuses"], "properties": { - "cursor": { "type": "string" }, "statuses": { "type": "array", "items": { diff --git a/packages/appview/src/context.ts b/packages/appview/src/context.ts index 8094bb2..6d5c01a 100644 --- a/packages/appview/src/context.ts +++ b/packages/appview/src/context.ts @@ -2,13 +2,14 @@ import { OAuthClient } from '@atproto/oauth-client-node' import { Firehose } from '@atproto/sync' import pino from 'pino' -import { Database } from './db' -import { BidirectionalResolver } from './id-resolver' +import { Database } from '#/db' +import { BidirectionalResolver } from '#/id-resolver' +import { Jetstream } from '#/ingestors' // Application state passed to the router and elsewhere export type AppContext = { db: Database - ingester: Firehose + ingester: Firehose | Jetstream logger: pino.Logger oauthClient: OAuthClient resolver: BidirectionalResolver diff --git a/packages/appview/src/db.ts b/packages/appview/src/db.ts index bffd2c9..069bef6 100644 --- a/packages/appview/src/db.ts +++ b/packages/appview/src/db.ts @@ -53,6 +53,11 @@ const migrationProvider: MigrationProvider = { }, } +migrations['003'] = { + async up(db: Kysely) {}, + async down(_db: Kysely) {}, +} + migrations['002'] = { async up(db: Kysely) { await db.schema @@ -60,6 +65,16 @@ migrations['002'] = { .addColumn('id', 'integer', (col) => col.primaryKey()) .addColumn('seq', 'integer', (col) => col.notNull()) .execute() + + // Insert initial cursor values: + // id=1 is for firehose, id=2 is for jetstream + await db + .insertInto('cursor' as never) + .values([ + { id: 1, seq: 0 }, + { id: 2, seq: 0 }, + ]) + .execute() }, async down(db: Kysely) { await db.schema.dropTable('cursor').execute() diff --git a/packages/appview/src/index.ts b/packages/appview/src/index.ts index fc574e2..699b28b 100644 --- a/packages/appview/src/index.ts +++ b/packages/appview/src/index.ts @@ -14,7 +14,7 @@ import { AppContext } from '#/context' import { createDb, migrateToLatest } from '#/db' import * as error from '#/error' import { createBidirectionalResolver, createIdResolver } from '#/id-resolver' -import { createIngester } from '#/ingester' +import { createFirehoseIngester, createJetstreamIngester } from '#/ingestors' import { createServer } from '#/lexicons' import { env } from '#/lib/env' @@ -36,7 +36,8 @@ export class Server { // Create the atproto utilities const oauthClient = await createClient(db) const baseIdResolver = createIdResolver() - const ingester = await createIngester(db, baseIdResolver) + const ingester = await createJetstreamIngester(db) + // Alternative: const ingester = await createFirehoseIngester(db, baseIdResolver) const resolver = createBidirectionalResolver(baseIdResolver) const ctx = { db, @@ -103,7 +104,7 @@ export class Server { }) } } else { - server.xrpc.router.set('trust proxy', true) + app.set('trust proxy', true) } // Use the port from env (should be 3001 for the API server) diff --git a/packages/appview/src/ingester.ts b/packages/appview/src/ingestors/firehose.ts similarity index 96% rename from packages/appview/src/ingester.ts rename to packages/appview/src/ingestors/firehose.ts index 3a42261..6e002a1 100644 --- a/packages/appview/src/ingester.ts +++ b/packages/appview/src/ingestors/firehose.ts @@ -5,7 +5,10 @@ import pino from 'pino' import type { Database } from '#/db' -export async function createIngester(db: Database, idResolver: IdResolver) { +export async function createFirehoseIngester( + db: Database, + idResolver: IdResolver, +) { const logger = pino({ name: 'firehose ingestion' }) const cursor = await db diff --git a/packages/appview/src/ingestors/index.ts b/packages/appview/src/ingestors/index.ts new file mode 100644 index 0000000..21132e7 --- /dev/null +++ b/packages/appview/src/ingestors/index.ts @@ -0,0 +1,2 @@ +export * from './jetstream' +export * from './firehose' diff --git a/packages/appview/src/ingestors/jetstream.ts b/packages/appview/src/ingestors/jetstream.ts new file mode 100644 index 0000000..550dbaa --- /dev/null +++ b/packages/appview/src/ingestors/jetstream.ts @@ -0,0 +1,213 @@ +import { XyzStatusphereStatus } from '@statusphere/lexicon' +import pino from 'pino' +import WebSocket from 'ws' + +import type { Database } from '#/db' + +export async function createJetstreamIngester(db: Database) { + const logger = pino({ name: 'jetstream ingestion' }) + + const cursor = await db + .selectFrom('cursor') + .where('id', '=', 2) + .select('seq') + .executeTakeFirst() + + logger.info(`start cursor: ${cursor?.seq}`) + + // For throttling cursor writes + let lastCursorWrite = 0 + + return new Jetstream({ + logger, + cursor: cursor?.seq || undefined, + setCursor: async (seq) => { + const now = Date.now() + + if (now - lastCursorWrite >= 30000) { + lastCursorWrite = now + logger.info(`writing cursor: ${seq}`) + await db + .updateTable('cursor') + .set({ seq }) + .where('id', '=', 2) + .execute() + } + }, + handleEvent: async (evt) => { + // ignore account and identity events + if ( + evt.kind !== 'commit' || + evt.commit.collection !== 'xyz.statusphere.status' + ) + return + + const now = new Date() + const uri = `at://${evt.did}/${evt.commit.collection}/${evt.commit.rkey}` + + if ( + (evt.commit.operation === 'create' || + evt.commit.operation === 'update') && + XyzStatusphereStatus.isRecord(evt.commit.record) + ) { + const validatedRecord = XyzStatusphereStatus.validateRecord( + evt.commit.record, + ) + if (!validatedRecord.success) return + + // Store the status in our SQLite + await db + .insertInto('status') + .values({ + uri, + authorDid: evt.did, + status: validatedRecord.value.status, + createdAt: validatedRecord.value.createdAt, + indexedAt: now.toISOString(), + }) + .onConflict((oc) => + oc.column('uri').doUpdateSet({ + status: validatedRecord.value.status, + indexedAt: now.toISOString(), + }), + ) + .execute() + } else if (evt.commit.operation === 'delete') { + // Remove the status from our SQLite + await db.deleteFrom('status').where('uri', '=', uri).execute() + } + }, + onError: (err) => { + logger.error({ err }, 'error during jetstream ingestion') + }, + wantedCollections: ['xyz.statusphere.status'], + }) +} + +export class Jetstream { + private logger: pino.Logger + private handleEvent: (evt: JetstreamEvent) => Promise + private onError: (err: unknown) => void + private setCursor?: (seq: number) => Promise + private cursor?: number + private ws?: WebSocket + private isStarted = false + private wantedCollections: string[] + + constructor({ + logger, + cursor, + setCursor, + handleEvent, + onError, + wantedCollections, + }: { + logger: pino.Logger + cursor?: number + setCursor?: (seq: number) => Promise + handleEvent: (evt: any) => Promise + onError: (err: any) => void + wantedCollections: string[] + }) { + this.logger = logger + this.cursor = cursor + this.setCursor = setCursor + this.handleEvent = handleEvent + this.onError = onError + this.wantedCollections = wantedCollections + } + + constructUrlWithQuery = (): string => { + const params = new URLSearchParams() + params.append('wantedCollections', this.wantedCollections.join(',')) + if (this.cursor !== undefined) { + params.append('cursor', this.cursor.toString()) + } + return `wss://jetstream.mozzius.dev/subscribe?${params.toString()}` + } + + start() { + if (this.isStarted) return + this.isStarted = true + this.ws = new WebSocket(this.constructUrlWithQuery()) + + this.ws.on('open', () => { + this.logger.info('Jetstream connection opened.') + }) + + this.ws.on('message', async (data) => { + try { + const event: JetstreamEvent = JSON.parse(data.toString()) + + // Update cursor if provided + if (event.time_us !== undefined && this.setCursor) { + await this.setCursor(event.time_us) + } + + await this.handleEvent(event) + } catch (err) { + this.onError(err) + } + }) + + this.ws.on('error', (err) => { + this.onError(err) + }) + + this.ws.on('close', (code, reason) => { + this.logger.error(`Jetstream closed. Code: ${code}, Reason: ${reason}`) + this.isStarted = false + }) + } + + destroy() { + if (this.ws) { + this.ws.close() + this.isStarted = false + } + } +} + +type JetstreamEvent = { + did: string + time_us: number +} & (CommitEvent | AccountEvent | IdentityEvent) + +type CommitEvent = { + kind: 'commit' + commit: + | { + operation: 'create' | 'update' + record: T + rev: string + collection: string + rkey: string + cid: string + } + | { + operation: 'delete' + rev: string + collection: string + rkey: string + } +} + +type IdentityEvent = { + kind: 'identity' + identity: { + did: string + handle: string + seq: number + time: string + } +} + +type AccountEvent = { + kind: 'account' + account: { + active: boolean + did: string + seq: number + time: string + } +} diff --git a/packages/appview/src/lexicons/lexicons.ts b/packages/appview/src/lexicons/lexicons.ts index d3b9645..c6c06cc 100644 --- a/packages/appview/src/lexicons/lexicons.ts +++ b/packages/appview/src/lexicons/lexicons.ts @@ -79,9 +79,6 @@ export const schemaDict = { type: 'object', required: ['statuses'], properties: { - cursor: { - type: 'string', - }, statuses: { type: 'array', items: { diff --git a/packages/appview/src/lexicons/types/xyz/statusphere/getStatuses.ts b/packages/appview/src/lexicons/types/xyz/statusphere/getStatuses.ts index 0c203f6..bd456fe 100644 --- a/packages/appview/src/lexicons/types/xyz/statusphere/getStatuses.ts +++ b/packages/appview/src/lexicons/types/xyz/statusphere/getStatuses.ts @@ -21,7 +21,6 @@ export interface QueryParams { export type InputSchema = undefined export interface OutputSchema { - cursor?: string statuses: XyzStatusphereDefs.StatusView[] } diff --git a/packages/client/src/components/StatusForm.tsx b/packages/client/src/components/StatusForm.tsx index 99408e7..9aaa8a2 100644 --- a/packages/client/src/components/StatusForm.tsx +++ b/packages/client/src/components/StatusForm.tsx @@ -5,7 +5,7 @@ import { XyzStatusphereDefs } from '@statusphere/lexicon' import useAuth from '#/hooks/useAuth' import api from '#/services/api' -const STATUS_OPTIONS = [ +export const STATUS_OPTIONS = [ '👍', '👎', '💙', diff --git a/packages/client/src/components/StatusList.tsx b/packages/client/src/components/StatusList.tsx index ea4050e..e9520f9 100644 --- a/packages/client/src/components/StatusList.tsx +++ b/packages/client/src/components/StatusList.tsx @@ -2,6 +2,7 @@ import { useEffect } from 'react' import { useQuery } from '@tanstack/react-query' import api from '#/services/api' +import { STATUS_OPTIONS } from './StatusForm' const StatusList = () => { // Use React Query to fetch and cache statuses @@ -23,11 +24,19 @@ const StatusList = () => { // Destructure data const statuses = data?.statuses || [] + + // Get a random emoji from the STATUS_OPTIONS array + const randomEmoji = STATUS_OPTIONS[Math.floor(Math.random() * STATUS_OPTIONS.length)] if (isPending && !data) { return ( -
- Loading statuses... +
+
+ {randomEmoji} +
+
+ Loading statuses... +
) } diff --git a/packages/client/src/pages/HomePage.tsx b/packages/client/src/pages/HomePage.tsx index 52be5a6..b3c4de7 100644 --- a/packages/client/src/pages/HomePage.tsx +++ b/packages/client/src/pages/HomePage.tsx @@ -1,22 +1,19 @@ import Header from '#/components/Header' -import StatusForm from '#/components/StatusForm' +import StatusForm, { STATUS_OPTIONS } from '#/components/StatusForm' import StatusList from '#/components/StatusList' import { useAuth } from '#/hooks/useAuth' const HomePage = () => { const { user, loading, error } = useAuth() + // Get a random emoji from the STATUS_OPTIONS array + const randomEmoji = + STATUS_OPTIONS[Math.floor(Math.random() * STATUS_OPTIONS.length)] + if (loading) { return ( -
-
-

- Loading Statusphere... -

-

- Setting up your experience -

-
+
+
{randomEmoji}
) } diff --git a/packages/client/vite.config.ts b/packages/client/vite.config.ts index 0d6fe9a..7951d25 100644 --- a/packages/client/vite.config.ts +++ b/packages/client/vite.config.ts @@ -21,6 +21,17 @@ 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); + }); + }, }, }, }, diff --git a/packages/lexicon/src/lexicons.ts b/packages/lexicon/src/lexicons.ts index d3b9645..c6c06cc 100644 --- a/packages/lexicon/src/lexicons.ts +++ b/packages/lexicon/src/lexicons.ts @@ -79,9 +79,6 @@ export const schemaDict = { type: 'object', required: ['statuses'], properties: { - cursor: { - type: 'string', - }, statuses: { type: 'array', items: { diff --git a/packages/lexicon/src/types/xyz/statusphere/getStatuses.ts b/packages/lexicon/src/types/xyz/statusphere/getStatuses.ts index ef5affc..3e24c9a 100644 --- a/packages/lexicon/src/types/xyz/statusphere/getStatuses.ts +++ b/packages/lexicon/src/types/xyz/statusphere/getStatuses.ts @@ -20,7 +20,6 @@ export interface QueryParams { export type InputSchema = undefined export interface OutputSchema { - cursor?: string statuses: XyzStatusphereDefs.StatusView[] } -- 2.51.2