diff --git a/lexicons/com.atproto.sync.subscribeRepos.json b/lexicons/com.atproto.sync.subscribeRepos.json deleted file mode 100644 index 69f8c6f..0000000 --- a/lexicons/com.atproto.sync.subscribeRepos.json +++ /dev/null @@ -1,356 +0,0 @@ -{ - "lexicon": 1, - "id": "com.atproto.sync.subscribeRepos", - "defs": { - "info": { - "type": "object", - "required": [ - "name" - ], - "properties": { - "name": { - "type": "string", - "knownValues": [ - "OutdatedCursor" - ] - }, - "message": { - "type": "string" - } - } - }, - "main": { - "type": "subscription", - "errors": [ - { - "name": "FutureCursor" - }, - { - "name": "ConsumerTooSlow", - "description": "If the consumer of the stream can not keep up with events, and a backlog gets too large, the server will drop the connection." - } - ], - "message": { - "schema": { - "refs": [ - "#commit", - "#sync", - "#identity", - "#account", - "#handle", - "#migrate", - "#tombstone", - "#info" - ], - "type": "union" - } - }, - "parameters": { - "type": "params", - "properties": { - "cursor": { - "type": "integer", - "description": "The last known event seq number to backfill from." - } - } - }, - "description": "Repository event stream, aka Firehose endpoint. Outputs repo commits with diff data, and identity update events, for all repositories on the current server. See the atproto specifications for details around stream sequencing, repo versioning, CAR diff format, and more. Public and does not require auth; implemented by PDS and Relay." - }, - "sync": { - "type": "object", - "required": [ - "seq", - "did", - "blocks", - "rev", - "time" - ], - "properties": { - "did": { - "type": "string", - "format": "did", - "description": "The account this repo event corresponds to. Must match that in the commit object." - }, - "rev": { - "type": "string", - "description": "The rev of the commit. This value must match that in the commit object." - }, - "seq": { - "type": "integer", - "description": "The stream sequence number of this message." - }, - "time": { - "type": "string", - "format": "datetime", - "description": "Timestamp of when this message was originally broadcast." - }, - "blocks": { - "type": "bytes", - "maxLength": 10000, - "description": "CAR file containing the commit, as a block. The CAR header must include the commit block CID as the first 'root'." - } - }, - "description": "Updates the repo to a new state, without necessarily including that state on the firehose. Used to recover from broken commit streams, data loss incidents, or in situations where upstream host does not know recent state of the repository." - }, - "commit": { - "type": "object", - "nullable": [ - "since" - ], - "required": [ - "seq", - "rebase", - "tooBig", - "repo", - "commit", - "rev", - "since", - "blocks", - "ops", - "blobs", - "time" - ], - "properties": { - "ops": { - "type": "array", - "items": { - "ref": "#repoOp", - "type": "ref", - "description": "List of repo mutation operations in this commit (eg, records created, updated, or deleted)." - }, - "maxLength": 200 - }, - "rev": { - "type": "string", - "format": "tid", - "description": "The rev of the emitted commit. Note that this information is also in the commit object included in blocks, unless this is a tooBig event." - }, - "seq": { - "type": "integer", - "description": "The stream sequence number of this message." - }, - "repo": { - "type": "string", - "format": "did", - "description": "The repo this event comes from. Note that all other message types name this field 'did'." - }, - "time": { - "type": "string", - "format": "datetime", - "description": "Timestamp of when this message was originally broadcast." - }, - "blobs": { - "type": "array", - "items": { - "type": "cid-link", - "description": "DEPRECATED -- will soon always be empty. List of new blobs (by CID) referenced by records in this commit." - } - }, - "since": { - "type": "string", - "format": "tid", - "description": "The rev of the last emitted commit from this repo (if any)." - }, - "blocks": { - "type": "bytes", - "maxLength": 2000000, - "description": "CAR file containing relevant blocks, as a diff since the previous repo state. The commit must be included as a block, and the commit block CID must be the first entry in the CAR header 'roots' list." - }, - "commit": { - "type": "cid-link", - "description": "Repo commit object CID." - }, - "rebase": { - "type": "boolean", - "description": "DEPRECATED -- unused" - }, - "tooBig": { - "type": "boolean", - "description": "DEPRECATED -- replaced by #sync event and data limits. Indicates that this commit contained too many ops, or data size was too large. Consumers will need to make a separate request to get missing data." - }, - "prevData": { - "type": "cid-link", - "description": "The root CID of the MST tree for the previous commit from this repo (indicated by the 'since' revision field in this message). Corresponds to the 'data' field in the repo commit object. NOTE: this field is effectively required for the 'inductive' version of firehose." - } - }, - "description": "Represents an update of repository state. Note that empty commits are allowed, which include no repo data changes, but an update to rev and signature." - }, - "handle": { - "type": "object", - "required": [ - "seq", - "did", - "handle", - "time" - ], - "properties": { - "did": { - "type": "string", - "format": "did" - }, - "seq": { - "type": "integer" - }, - "time": { - "type": "string", - "format": "datetime" - }, - "handle": { - "type": "string", - "format": "handle" - } - }, - "description": "DEPRECATED -- Use #identity event instead" - }, - "repoOp": { - "type": "object", - "nullable": [ - "cid" - ], - "required": [ - "action", - "path", - "cid" - ], - "properties": { - "cid": { - "type": "cid-link", - "description": "For creates and updates, the new record CID. For deletions, null." - }, - "path": { - "type": "string" - }, - "prev": { - "type": "cid-link", - "description": "For updates and deletes, the previous record CID (required for inductive firehose). For creations, field should not be defined." - }, - "action": { - "type": "string", - "knownValues": [ - "create", - "update", - "delete" - ] - } - }, - "description": "A repo operation, ie a mutation of a single record." - }, - "account": { - "type": "object", - "required": [ - "seq", - "did", - "time", - "active" - ], - "properties": { - "did": { - "type": "string", - "format": "did" - }, - "seq": { - "type": "integer" - }, - "time": { - "type": "string", - "format": "datetime" - }, - "active": { - "type": "boolean", - "description": "Indicates that the account has a repository which can be fetched from the host that emitted this event." - }, - "status": { - "type": "string", - "description": "If active=false, this optional field indicates a reason for why the account is not active.", - "knownValues": [ - "takendown", - "suspended", - "deleted", - "deactivated", - "desynchronized", - "throttled" - ] - } - }, - "description": "Represents a change to an account's status on a host (eg, PDS or Relay). The semantics of this event are that the status is at the host which emitted the event, not necessarily that at the currently active PDS. Eg, a Relay takedown would emit a takedown with active=false, even if the PDS is still active." - }, - "migrate": { - "type": "object", - "nullable": [ - "migrateTo" - ], - "required": [ - "seq", - "did", - "migrateTo", - "time" - ], - "properties": { - "did": { - "type": "string", - "format": "did" - }, - "seq": { - "type": "integer" - }, - "time": { - "type": "string", - "format": "datetime" - }, - "migrateTo": { - "type": "string" - } - }, - "description": "DEPRECATED -- Use #account event instead" - }, - "identity": { - "type": "object", - "required": [ - "seq", - "did", - "time" - ], - "properties": { - "did": { - "type": "string", - "format": "did" - }, - "seq": { - "type": "integer" - }, - "time": { - "type": "string", - "format": "datetime" - }, - "handle": { - "type": "string", - "format": "handle", - "description": "The current handle for the account, or 'handle.invalid' if validation fails. This field is optional, might have been validated or passed-through from an upstream source. Semantics and behaviors for PDS vs Relay may evolve in the future; see atproto specs for more details." - } - }, - "description": "Represents a change to an account's identity. Could be an updated handle, signing key, or pds hosting endpoint. Serves as a prod to all downstream services to refresh their identity cache." - }, - "tombstone": { - "type": "object", - "required": [ - "seq", - "did", - "time" - ], - "properties": { - "did": { - "type": "string", - "format": "did" - }, - "seq": { - "type": "integer" - }, - "time": { - "type": "string", - "format": "datetime" - } - }, - "description": "DEPRECATED -- Use #account event instead" - } - } -} \ No newline at end of file diff --git a/server/src/backfill/index.ts b/server/src/backfill/index.ts index 796c7ee..153d59e 100644 --- a/server/src/backfill/index.ts +++ b/server/src/backfill/index.ts @@ -2,7 +2,6 @@ import { drizzle } from "drizzle-orm/libsql"; import { routes } from "../db/schema.ts"; import * as schema from "../db/schema.ts"; import { db as db_type } from "../utils.ts"; -import { Client, simpleFetchHandler } from "@atcute/client"; import oldRecords from "./old-records.ts"; import newRecords from "./new-records.ts"; @@ -10,31 +9,17 @@ const db: db_type = drizzle( Deno.env.get("DB_FILE_NAME")! || (() => { throw "DB_FILE_NAME not set"; - })(), + })() ); -// we block access to the database till we've set up listeners etc -// set an initial value for ts sake; this will set asap -let res: (value: db_type) => void = (v) => console.log("failed to set", v); -export default new Promise((p_res) => { - res = p_res; -}); - -// clear old records +// clear old records & blobs // accounts could have deleted their sites or anything // so we should just nuke them +// blobs could be kept but would be a nightmare await db.delete(routes); +await Deno.remove("./blobs", { recursive: true }); -const relay = new Client({ - handler: simpleFetchHandler({ - service: Deno.env.get("ATPROTO_RELAY") || - (() => { - throw "ATPROTO_RELAY not set"; - })(), - }), -}); - -oldRecords(db, relay); +oldRecords(db); newRecords(db); -res(db); +export default db; diff --git a/server/src/backfill/new-records.ts b/server/src/backfill/new-records.ts index a436eec..2ba1c29 100644 --- a/server/src/backfill/new-records.ts +++ b/server/src/backfill/new-records.ts @@ -1,16 +1,14 @@ import { db, getPds, isDid, rkeyToUrl } from "../utils.ts"; import { decodeFirst } from "@atcute/cbor"; -import { - ComAtprotoSyncSubscribeRepos, - DevAtcitiesRoute, -} from "../lexicons/index.ts"; +import { DevAtcitiesRoute } from "../lexicons/index.ts"; import { is } from "@atcute/lexicons"; +import { ComAtprotoSyncSubscribeRepos } from "@atcute/atproto"; import { Client, simpleFetchHandler } from "@atcute/client"; import { routes } from "../db/schema.ts"; import { and, eq } from "drizzle-orm"; const ErrorEvent = Symbol("error event"); -export default async function newRecords(db: db) { +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( `${ @@ -27,6 +25,7 @@ export default async function newRecords(db: db) { ws.addEventListener("error", (ev) => { throw ev; }); + ws.addEventListener("open", () => console.log("Listening for new events.")); try { ws.addEventListener("message", async (ev) => { @@ -43,6 +42,8 @@ export default async function newRecords(db: db) { throw ErrorEvent; } + // we only care about commits + // identity events and similar are irrelevant if (!is(ComAtprotoSyncSubscribeRepos.commitSchema, payload)) { return; } @@ -50,48 +51,34 @@ export default async function newRecords(db: db) { console.warn("Invalid did:", payload.repo); return; } + const pds = await getPds(payload.repo); + if (!pds) { + return; + } - const pdsClient = (() => { - let client: Client | undefined; - return async () => { - if (client) { - return client; - } - if (!isDid(payload.repo)) { - console.warn("Invalid did:", payload.repo); - return; - } - - const pds = await getPds(payload.repo); - if (!pds) { - return; - } - client = new Client({ - handler: simpleFetchHandler({ - service: pds, - }), - }); - - return client; - }; - })(); + const client = new Client({ + handler: simpleFetchHandler({ + 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) { console.warn("rkey not valid!"); continue; } if (collection !== "dev.atcities.route") continue; + switch (op.action) { case "create": case "update": { - const client = await pdsClient(); - if (!client) { - console.warn("could not resolve pds for", payload.repo); - return; - } + // hydrate data const { data, ok } = await client.get( "com.atproto.repo.getRecord", { @@ -109,15 +96,18 @@ export default async function newRecords(db: db) { } if (!is(DevAtcitiesRoute.mainSchema, data.value)) { - console.warn("Invalid record"); + console.warn( + "Invalid record:", + `at://${payload.repo}/dev.atcities.route.${rkey}` + ); continue; } if (data.value.page.$type !== "dev.atcities.route#blob") { - console.warn("Unknown page type"); continue; } + // if this page exists in db, get the id so it can be replaced const id = ( await db .select({ @@ -198,5 +188,4 @@ export default async function newRecords(db: db) { if (e === ErrorEvent) newRecords(db); else throw e; } - ws.addEventListener("open", () => console.log("subscribed!!")); } diff --git a/server/src/backfill/old-records.ts b/server/src/backfill/old-records.ts index 43209e3..5377900 100644 --- a/server/src/backfill/old-records.ts +++ b/server/src/backfill/old-records.ts @@ -4,12 +4,13 @@ import { is } from "@atcute/lexicons"; import { DevAtcitiesRoute } from "../lexicons/index.ts"; import { routes } from "../db/schema.ts"; -async function index(did: `did:${"web" | "plc"}:${string}`, db: db) { +async function indexUser(did: `did:${"web" | "plc"}:${string}`, db: db) { const pds = await getPds(did); if (!pds) return console.error(did, "could not be resolved to a pds."); const pdsClient = new Client({ handler: simpleFetchHandler({ service: pds }), }); + let cursor: string | undefined; while (true) { const { data, ok } = await pdsClient.get("com.atproto.repo.listRecords", { @@ -27,6 +28,8 @@ async function index(did: `did:${"web" | "plc"}:${string}`, db: db) { for (const record of data.records) { if (is(DevAtcitiesRoute.mainSchema, record.value)) { + // ignore non blob page values + // this will be expanded in future if (record.value.page.$type === "dev.atcities.route#blob") { const url = rkeyToUrl(record.uri.split("/")[4]); if (!url) continue; @@ -49,7 +52,17 @@ async function index(did: `did:${"web" | "plc"}:${string}`, db: db) { } } -export default async function (db: db, relay: Client) { +export default async function (db: db) { + const relay = new Client({ + handler: simpleFetchHandler({ + service: + Deno.env.get("ATPROTO_RELAY") || + (() => { + throw "ATPROTO_RELAY not set"; + })(), + }), + }); + let repos: `did:${string}:${string}`[] = []; let cursor: string | undefined; while (true) { @@ -72,9 +85,12 @@ export default async function (db: db, relay: Client) { if (!cursor) break; } + const pending: Promise[] = []; for (const i in repos) { const did = repos[i]; console.log(`indexing ${Number(i) + 1}/${repos.length} ${did}`); - if (isDid(did)) index(did, db); + if (isDid(did)) pending.push(indexUser(did, db)); } + await Promise.all(pending); + console.log("Finished backfilling."); } diff --git a/server/src/lexicons/index.ts b/server/src/lexicons/index.ts index 00be4d2..31a2f85 100644 --- a/server/src/lexicons/index.ts +++ b/server/src/lexicons/index.ts @@ -1,2 +1 @@ -export * as ComAtprotoSyncSubscribeRepos from "./types/com/atproto/sync/subscribeRepos.ts"; export * as DevAtcitiesRoute from "./types/dev/atcities/route.ts"; diff --git a/server/src/lexicons/types/com/atproto/sync/subscribeRepos.ts b/server/src/lexicons/types/com/atproto/sync/subscribeRepos.ts deleted file mode 100644 index 94b164f..0000000 --- a/server/src/lexicons/types/com/atproto/sync/subscribeRepos.ts +++ /dev/null @@ -1,184 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; - -const _accountSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#account"), - ), - active: /*#__PURE__*/ v.boolean(), - did: /*#__PURE__*/ v.didString(), - seq: /*#__PURE__*/ v.integer(), - status: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.string< - | "deactivated" - | "deleted" - | "desynchronized" - | "suspended" - | "takendown" - | "throttled" - | (string & {}) - >(), - ), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _commitSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#commit"), - ), - blobs: /*#__PURE__*/ v.array(/*#__PURE__*/ v.cidLink()), - blocks: /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.bytes(), [ - /*#__PURE__*/ v.bytesSize(0, 2000000), - ]), - commit: /*#__PURE__*/ v.cidLink(), - get ops() { - return /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.array(repoOpSchema), [ - /*#__PURE__*/ v.arrayLength(0, 200), - ]); - }, - prevData: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidLink()), - rebase: /*#__PURE__*/ v.boolean(), - repo: /*#__PURE__*/ v.didString(), - rev: /*#__PURE__*/ v.tidString(), - seq: /*#__PURE__*/ v.integer(), - since: /*#__PURE__*/ v.nullable(/*#__PURE__*/ v.tidString()), - time: /*#__PURE__*/ v.datetimeString(), - tooBig: /*#__PURE__*/ v.boolean(), -}); -const _handleSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#handle"), - ), - did: /*#__PURE__*/ v.didString(), - handle: /*#__PURE__*/ v.handleString(), - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _identitySchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#identity"), - ), - did: /*#__PURE__*/ v.didString(), - handle: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.handleString()), - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _infoSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#info"), - ), - message: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), - name: /*#__PURE__*/ v.string<"OutdatedCursor" | (string & {})>(), -}); -const _mainSchema = /*#__PURE__*/ v.subscription( - "com.atproto.sync.subscribeRepos", - { - params: /*#__PURE__*/ v.object({ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - }), - get message() { - return /*#__PURE__*/ v.variant([ - accountSchema, - commitSchema, - handleSchema, - identitySchema, - infoSchema, - migrateSchema, - syncSchema, - tombstoneSchema, - ]); - }, - }, -); -const _migrateSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#migrate"), - ), - did: /*#__PURE__*/ v.didString(), - migrateTo: /*#__PURE__*/ v.nullable(/*#__PURE__*/ v.string()), - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _repoOpSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#repoOp"), - ), - action: /*#__PURE__*/ v.string< - "create" | "delete" | "update" | (string & {}) - >(), - cid: /*#__PURE__*/ v.nullable(/*#__PURE__*/ v.cidLink()), - path: /*#__PURE__*/ v.string(), - prev: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidLink()), -}); -const _syncSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#sync"), - ), - blocks: /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.bytes(), [ - /*#__PURE__*/ v.bytesSize(0, 10000), - ]), - did: /*#__PURE__*/ v.didString(), - rev: /*#__PURE__*/ v.string(), - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _tombstoneSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("com.atproto.sync.subscribeRepos#tombstone"), - ), - did: /*#__PURE__*/ v.didString(), - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); - -type account$schematype = typeof _accountSchema; -type commit$schematype = typeof _commitSchema; -type handle$schematype = typeof _handleSchema; -type identity$schematype = typeof _identitySchema; -type info$schematype = typeof _infoSchema; -type main$schematype = typeof _mainSchema; -type migrate$schematype = typeof _migrateSchema; -type repoOp$schematype = typeof _repoOpSchema; -type sync$schematype = typeof _syncSchema; -type tombstone$schematype = typeof _tombstoneSchema; - -export interface accountSchema extends account$schematype {} -export interface commitSchema extends commit$schematype {} -export interface handleSchema extends handle$schematype {} -export interface identitySchema extends identity$schematype {} -export interface infoSchema extends info$schematype {} -export interface mainSchema extends main$schematype {} -export interface migrateSchema extends migrate$schematype {} -export interface repoOpSchema extends repoOp$schematype {} -export interface syncSchema extends sync$schematype {} -export interface tombstoneSchema extends tombstone$schematype {} - -export const accountSchema = _accountSchema as accountSchema; -export const commitSchema = _commitSchema as commitSchema; -export const handleSchema = _handleSchema as handleSchema; -export const identitySchema = _identitySchema as identitySchema; -export const infoSchema = _infoSchema as infoSchema; -export const mainSchema = _mainSchema as mainSchema; -export const migrateSchema = _migrateSchema as migrateSchema; -export const repoOpSchema = _repoOpSchema as repoOpSchema; -export const syncSchema = _syncSchema as syncSchema; -export const tombstoneSchema = _tombstoneSchema as tombstoneSchema; - -export interface Account extends v.InferInput {} -export interface Commit extends v.InferInput {} -export interface Handle extends v.InferInput {} -export interface Identity extends v.InferInput {} -export interface Info extends v.InferInput {} -export interface Migrate extends v.InferInput {} -export interface RepoOp extends v.InferInput {} -export interface Sync extends v.InferInput {} -export interface Tombstone extends v.InferInput {} - -export interface $params extends v.InferInput {} -export type $message = v.InferInput; - -declare module "@atcute/lexicons/ambient" { - interface XRPCSubscriptions { - "com.atproto.sync.subscribeRepos": mainSchema; - } -}