import { db, getPds, isDid, rkeyToUrl, fetchHandler } from "../utils.ts"; import { decodeFirst } from "@atcute/cbor"; import { DevAtcitiesRoute } from "../lexicons/index.ts"; import { is } from "@atcute/lexicons"; import { ComAtprotoSyncSubscribeRepos } from "@atcute/atproto"; import { Client } from "@atcute/client"; import { routes } from "../db/schema.ts"; import { and, eq } from "drizzle-orm"; const ErrorEvent = Symbol("error event"); export default function newRecords(db: db) { // https://pdsls.dev/at://did:plc:6msi3pj7krzih5qxqtryxlzw/com.atproto.lexicon.schema/com.atproto.sync.subscribeRepos const ws = new WebSocket( `${ Deno.env.get("ATPROTO_RELAY") || (() => { throw "ATPROTO_RELAY not set"; })() }/xrpc/com.atproto.sync.subscribeRepos` ); ws.addEventListener("close", (ev) => { throw ev; }); ws.addEventListener("error", (ev) => { throw ev; }); ws.addEventListener("open", () => console.log("Listening for new events.")); try { ws.addEventListener("message", async (ev) => { if (ev.data instanceof Blob) { try { const [header, remainder] = decodeFirst(await ev.data.bytes()); const [payload] = decodeFirst(remainder); // https://atproto.com/specs/event-stream#framing if (typeof header !== "object" || !header) return; if (!("op" in header) || typeof header.op !== "number") return; if (header.op === -1) { console.error(header); ws.close(); throw ErrorEvent; } // we only care about commits // identity events and similar are irrelevant if (!is(ComAtprotoSyncSubscribeRepos.commitSchema, payload)) { return; } if (!isDid(payload.repo)) { console.warn("Invalid did:", payload.repo); return; } const pds = await getPds(payload.repo); if (!pds) { return; } const client = new Client({ handler: fetchHandler({ service: pds, }), }); for (const op of payload.ops) { const [collection, rkey] = op.path.split("/"); if (!collection || !rkey) { console.warn("Invalid path", op.path); continue; } const path = rkeyToUrl(rkey); if (!path) continue; if (collection !== "dev.atcities.route") continue; switch (op.action) { case "create": case "update": { // hydrate data const { data, ok } = await client.get( "com.atproto.repo.getRecord", { params: { repo: payload.repo, collection: "dev.atcities.route", rkey, }, } ); if (!ok) { console.warn(data); continue; } if (!is(DevAtcitiesRoute.mainSchema, data.value)) { console.warn( "Invalid record:", `at://${payload.repo}/dev.atcities.route.${rkey}` ); continue; } if (data.value.page.$type !== "dev.atcities.route#blob") { continue; } // if this page exists in db, get the id so it can be replaced const id = ( await db .select({ id: routes.id, }) .from(routes) .where( and( eq(routes.did, payload.repo), eq(routes.url_route, path) ) ) ).at(0); await db.insert(routes).values({ did: payload.repo, url_route: path, id: id ? id.id : undefined, blob_cid: "ref" in data.value.page.blob ? data.value.page.blob.ref.$link : data.value.page.blob.cid, mime: data.value.page.blob.mimeType, }); continue; } case "delete": { const deleted = ( await db .select({ id: routes.id, cid: routes.blob_cid, }) .from(routes) .where( and( eq(routes.did, payload.repo), eq(routes.url_route, path) ) ) ).at(0); if (!deleted) continue; await db.delete(routes).where(eq(routes.id, deleted.id)); const cidInUse = ( await db .select({}) .from(routes) .where(eq(routes.blob_cid, deleted.cid)) ).length > 0; // purge cid if not in use if (!cidInUse) Deno.remove(`./blobs/${payload.repo}/${deleted.cid}`); continue; } } } } catch (e) { if (Error.isError(e)) { console.warn( "Invalid DAG-CBOR data:", e.name, e.message, e.cause, e.stack ); } else throw e; } } }); } catch (e) { // if and only if an error message was recived, reopen the connection if (e === ErrorEvent) newRecords(db); else throw e; } }