Something went wrong. Try again.
an app to share curated trails sidetrail.app
atproto nextjs react rsc
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295import { 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<void> { 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<void> { 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<void> { 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<void> { 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<void> { 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<void> { await db.delete(completions).where(eq(completions.uri, uri));}
// ============================================================================// Account operations// ============================================================================
export async function deleteAllContentByDid(db: IngesterDb, did: string): Promise<void> { 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<void> { await db .insert(accounts) .values({ did, active: 1, seq: 0, updatedAt: new Date() }) .onConflictDoNothing();}
export async function handleAccountEvent( db: IngesterDb, evt: AccountEvent["account"],): Promise<void> { 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<void> { 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<string, unknown>;
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; } }}