import { eq } from "drizzle-orm"; const log = process.env.NODE_ENV === "test" ? () => {} : console.log; import type { NodePgDatabase } from "drizzle-orm/node-postgres"; import { trails, walks, completions, accounts } from "@sidetrail/db"; import type { JetstreamEvent, AccountEvent } from "./jetstream.js"; import { validateRecord, type IndexedCollection } from "./lexicons.js"; export const COLLECTIONS = [ "app.sidetrail.trail", "app.sidetrail.walk", "app.sidetrail.completion", ]; export type IngesterDb = NodePgDatabase; // ============================================================================ // Trail operations // ============================================================================ export async function upsertTrail( db: IngesterDb, uri: string, cid: string, authorDid: string, rkey: string, record: unknown, createdAt: string, ): Promise { await db .insert(trails) .values({ uri, cid, authorDid, rkey, record: record as typeof trails.$inferInsert.record, createdAt: new Date(createdAt), indexedAt: new Date(), }) .onConflictDoUpdate({ target: trails.uri, set: { cid, record: record as typeof trails.$inferInsert.record, indexedAt: new Date(), }, }); } export async function deleteTrail(db: IngesterDb, uri: string): Promise { await db.delete(trails).where(eq(trails.uri, uri)); } // ============================================================================ // Walk operations // ============================================================================ export async function upsertWalk( db: IngesterDb, uri: string, cid: string, authorDid: string, rkey: string, trailUri: string, record: unknown, createdAt: string, ): Promise { await db .insert(walks) .values({ uri, cid, authorDid, rkey, trailUri, record: record as typeof walks.$inferInsert.record, createdAt: new Date(createdAt), indexedAt: new Date(), }) .onConflictDoUpdate({ target: walks.uri, set: { cid, record: record as typeof walks.$inferInsert.record, indexedAt: new Date(), }, }); } export async function deleteWalk(db: IngesterDb, uri: string): Promise { await db.delete(walks).where(eq(walks.uri, uri)); } // ============================================================================ // Completion operations // ============================================================================ export async function upsertCompletion( db: IngesterDb, uri: string, cid: string, authorDid: string, rkey: string, trailUri: string, record: unknown, createdAt: string, ): Promise { await db .insert(completions) .values({ uri, cid, authorDid, rkey, trailUri, record: record as typeof completions.$inferInsert.record, createdAt: new Date(createdAt), indexedAt: new Date(), }) .onConflictDoUpdate({ target: completions.uri, set: { cid, record: record as typeof completions.$inferInsert.record, indexedAt: new Date(), }, }); } export async function deleteCompletion(db: IngesterDb, uri: string): Promise { await db.delete(completions).where(eq(completions.uri, uri)); } // ============================================================================ // Account operations // ============================================================================ export async function deleteAllContentByDid(db: IngesterDb, did: string): Promise { await Promise.all([ db.delete(trails).where(eq(trails.authorDid, did)), db.delete(walks).where(eq(walks.authorDid, did)), db.delete(completions).where(eq(completions.authorDid, did)), ]); } export async function ensureAccount(db: IngesterDb, did: string): Promise { await db .insert(accounts) .values({ did, active: 1, seq: 0, updatedAt: new Date() }) .onConflictDoNothing(); } export async function handleAccountEvent( db: IngesterDb, evt: AccountEvent["account"], ): Promise { const { did, active, seq, status } = evt; const [existing] = await db .select({ seq: accounts.seq }) .from(accounts) .where(eq(accounts.did, did)) .limit(1); if (!existing) { return; } if (existing.seq >= seq) { return; } await db .update(accounts) .set({ active: active ? 1 : 0, status: active ? null : status, seq, updatedAt: new Date(), }) .where(eq(accounts.did, did)); if (!active && (status === "takendown" || status === "deleted")) { log(`Account ${did} ${status} - deleting all content`); await deleteAllContentByDid(db, did); } else if (!active) { log(`Account ${did} ${status ?? "inactive"} - marking inactive`); } else if (active) { log(`Account ${did} reactivated`); } } // ============================================================================ // Event handler // ============================================================================ export async function handleEvent(db: IngesterDb, evt: JetstreamEvent): Promise { if (evt.kind === "account") { await handleAccountEvent(db, evt.account); return; } if (evt.kind === "identity") return; if (evt.kind !== "commit") return; const { commit } = evt; const { collection, rkey } = commit; if (!COLLECTIONS.includes(collection)) return; // ATProto spec: commits from inactive accounts should be ignored const [accountStatus] = await db .select({ active: accounts.active }) .from(accounts) .where(eq(accounts.did, evt.did)) .limit(1); if (accountStatus && !accountStatus.active) { return; } const uri = `at://${evt.did}/${collection}/${rkey}`; if (commit.operation === "delete") { switch (collection) { case "app.sidetrail.trail": await deleteTrail(db, uri); break; case "app.sidetrail.walk": await deleteWalk(db, uri); break; case "app.sidetrail.completion": await deleteCompletion(db, uri); break; } return; } const record = commit.record as Record; const validation = validateRecord(collection as IndexedCollection, record); if (!validation.success) { log(`Rejecting invalid ${collection} ${uri}: ${validation.reason}`); return; } await ensureAccount(db, evt.did); switch (collection) { case "app.sidetrail.trail": await upsertTrail( db, uri, commit.cid, evt.did, rkey, record, (record.createdAt as string) || new Date().toISOString(), ); break; case "app.sidetrail.walk": { const trailRef = record.trail as { uri: string } | undefined; const trailUri = trailRef?.uri || ""; await upsertWalk( db, uri, commit.cid, evt.did, rkey, trailUri, record, (record.createdAt as string) || new Date().toISOString(), ); break; } case "app.sidetrail.completion": { const trailRef = record.trail as { uri: string } | undefined; const trailUri = trailRef?.uri || ""; await upsertCompletion( db, uri, commit.cid, evt.did, rkey, trailUri, record, (record.createdAt as string) || new Date().toISOString(), ); break; } } }