diff --git a/packages/appview/src/db.ts b/packages/appview/src/db.ts index fc177ec..bffd2c9 100644 --- a/packages/appview/src/db.ts +++ b/packages/appview/src/db.ts @@ -13,6 +13,7 @@ export type DatabaseSchema = { status: Status auth_session: AuthSession auth_state: AuthState + cursor: Cursor } export type Status = { @@ -33,6 +34,11 @@ export type AuthState = { state: AuthStateJson } +export type Cursor = { + id: number + seq: number +} + type AuthStateJson = string type AuthSessionJson = string @@ -47,6 +53,19 @@ const migrationProvider: MigrationProvider = { }, } +migrations['002'] = { + async up(db: Kysely) { + await db.schema + .createTable('cursor') + .addColumn('id', 'integer', (col) => col.primaryKey()) + .addColumn('seq', 'integer', (col) => col.notNull()) + .execute() + }, + async down(db: Kysely) { + await db.schema.dropTable('cursor').execute() + }, +} + migrations['001'] = { async up(db: Kysely) { await db.schema diff --git a/packages/appview/src/index.ts b/packages/appview/src/index.ts index 93b5304..55d25fb 100644 --- a/packages/appview/src/index.ts +++ b/packages/appview/src/index.ts @@ -47,7 +47,7 @@ export class Server { // Create the atproto utilities const oauthClient = await createClient(db) const baseIdResolver = createIdResolver() - const ingester = createIngester(db, baseIdResolver) + const ingester = await createIngester(db, baseIdResolver) const resolver = createBidirectionalResolver(baseIdResolver) const ctx = { db, diff --git a/packages/appview/src/ingester.ts b/packages/appview/src/ingester.ts index 43aeae8..3a42261 100644 --- a/packages/appview/src/ingester.ts +++ b/packages/appview/src/ingester.ts @@ -1,14 +1,43 @@ import { IdResolver } from '@atproto/identity' -import { Firehose, type Event } from '@atproto/sync' +import { Firehose, MemoryRunner, type Event } from '@atproto/sync' import { XyzStatusphereStatus } from '@statusphere/lexicon' import pino from 'pino' import type { Database } from '#/db' -export function createIngester(db: Database, idResolver: IdResolver) { +export async function createIngester(db: Database, idResolver: IdResolver) { const logger = pino({ name: 'firehose ingestion' }) + + const cursor = await db + .selectFrom('cursor') + .where('id', '=', 1) + .select('seq') + .executeTakeFirst() + + logger.info(`start cursor: ${cursor?.seq}`) + + // For throttling cursor writes + let lastCursorWrite = 0 + + const runner = new MemoryRunner({ + startCursor: cursor?.seq || undefined, + setCursor: async (seq) => { + const now = Date.now() + + if (now - lastCursorWrite >= 10000) { + lastCursorWrite = now + await db + .updateTable('cursor') + .set({ seq }) + .where('id', '=', 1) + .execute() + } + }, + }) + return new Firehose({ idResolver, + runner, handleEvent: async (evt: Event) => { // Watch for write events if (evt.event === 'create' || evt.event === 'update') {