From f7041e4b02b6bd20264cceaa2f52ec83adeda33f Mon Sep 17 00:00:00 2001 From: Roscoe Rubin-Rottenberg Date: Thu, 4 Dec 2025 17:53:51 -0500 Subject: [PATCH] ingester --- README.md | 3 +- deno.json | 1 + deno.lock | 44 +++++++++++++++++++++++++ ingester/handlers/follow.ts | 43 ++++++++++++++++++++++++ ingester/handlers/like.ts | 46 ++++++++++++++++++++++++++ ingester/handlers/post.ts | 50 ++++++++++++++++++++++++++++ ingester/handlers/repost.ts | 46 ++++++++++++++++++++++++++ ingester/index.ts | 66 +++++++++++++++++++++++++++++++++++++ ingester/types.ts | 22 +++++++++++++ main.ts | 8 +++-- utils/env.ts | 2 +- 11 files changed, 325 insertions(+), 6 deletions(-) create mode 100644 ingester/handlers/follow.ts create mode 100644 ingester/handlers/like.ts create mode 100644 ingester/handlers/post.ts create mode 100644 ingester/handlers/repost.ts create mode 100644 ingester/index.ts create mode 100644 ingester/types.ts diff --git a/README.md b/README.md index 021bac1..5eaf023 100644 --- a/README.md +++ b/README.md @@ -46,8 +46,7 @@ instance: ``` SPRK_DB_NAME=feed-gen -SPRK_DB_HOST=localhost -SPRK_DB_PORT=27017 +SPRK_DB_URI=mongodb://localhost:27017 SPRK_DB_USER=username SPRK_DB_PASSWORD=password SPRK_FEEDGEN_DOMAIN=feeds.example.com diff --git a/deno.json b/deno.json index 915aba5..5a4f17f 100644 --- a/deno.json +++ b/deno.json @@ -7,6 +7,7 @@ "@atp/common": "jsr:@atp/common@^0.1.0-alpha.4", "@atp/identity": "jsr:@atp/identity@^0.1.0-alpha.1", "@atp/lexicon": "jsr:@atp/lexicon@^0.1.0-alpha.2", + "@atp/sync": "jsr:@atp/sync@^0.1.0-alpha.4", "@atp/syntax": "jsr:@atp/syntax@^0.1.0-alpha.2", "@atp/xrpc-server": "jsr:@atp/xrpc-server@^0.1.0-alpha.3", "@logtape/logtape": "jsr:@logtape/logtape@^1.2.2", diff --git a/deno.lock b/deno.lock index 25bd966..1c52130 100644 --- a/deno.lock +++ b/deno.lock @@ -9,8 +9,11 @@ "jsr:@atp/identity@~0.1.0-alpha.1": "0.1.0-alpha.1", "jsr:@atp/lexicon@~0.1.0-alpha.1": "0.1.0-alpha.2", "jsr:@atp/lexicon@~0.1.0-alpha.2": "0.1.0-alpha.2", + "jsr:@atp/repo@~0.1.0-alpha.2": "0.1.0-alpha.2", + "jsr:@atp/sync@~0.1.0-alpha.4": "0.1.0-alpha.4", "jsr:@atp/syntax@~0.1.0-alpha.1": "0.1.0-alpha.2", "jsr:@atp/syntax@~0.1.0-alpha.2": "0.1.0-alpha.2", + "jsr:@atp/xrpc-server@~0.1.0-alpha.2": "0.1.0-alpha.3", "jsr:@atp/xrpc-server@~0.1.0-alpha.3": "0.1.0-alpha.3", "jsr:@atp/xrpc@~0.1.0-alpha.2": "0.1.0-alpha.2", "jsr:@hono/hono@^4.10.7": "4.10.7", @@ -39,6 +42,7 @@ "npm:jose@^6.1.3": "6.1.3", "npm:mongoose@^8.20.1": "8.20.1", "npm:multiformats@^13.4.1": "13.4.1", + "npm:p-queue@^8.1.1": "8.1.1", "npm:rate-limiter-flexible@^2.4.2": "2.4.2", "npm:zod@^4.1.11": "4.1.13" }, @@ -88,6 +92,32 @@ "npm:zod" ] }, + "@atp/repo@0.1.0-alpha.2": { + "integrity": "6da50453bbd527a679237d15bc9569eb2195503189f9be9d3023060f3f89f44a", + "dependencies": [ + "jsr:@atp/bytes", + "jsr:@atp/common@~0.1.0-alpha.4", + "jsr:@atp/crypto@~0.1.0-alpha.2", + "jsr:@atp/lexicon@~0.1.0-alpha.2", + "jsr:@std/encoding", + "npm:@ipld/dag-cbor", + "npm:multiformats", + "npm:zod" + ] + }, + "@atp/sync@0.1.0-alpha.4": { + "integrity": "9b6aa6ccc9447270843272e40bfcb26520eddaf37f98202bcbab6c0bee0a602b", + "dependencies": [ + "jsr:@atp/common@~0.1.0-alpha.4", + "jsr:@atp/identity", + "jsr:@atp/lexicon@~0.1.0-alpha.2", + "jsr:@atp/repo", + "jsr:@atp/syntax@~0.1.0-alpha.1", + "jsr:@atp/xrpc-server@~0.1.0-alpha.2", + "npm:multiformats", + "npm:p-queue" + ] + }, "@atp/syntax@0.1.0-alpha.1": { "integrity": "9e2055cace77cf63a8c52a4a94c39492215e7135101db7bc2289ebad9bec1991" }, @@ -232,6 +262,9 @@ "tslib" ] }, + "eventemitter3@5.0.1": { + "integrity": "sha512-GWkBvjiSZK87ELrYOSESUYeVIc9mvLLf/nXalMOS5dYrgZq9o5OVkbZAVM06CVxYsCwH9BDZFPlQTlPA1j4ahA==" + }, "jose@6.1.3": { "integrity": "sha512-0TpaTfihd4QMNwrz/ob2Bp7X04yuxJkjRGi4aKmOqwhov54i6u79oCv7T+C7lo70MKH6BesI3vscD1yb/yzKXQ==" }, @@ -283,6 +316,16 @@ "multiformats@13.4.1": { "integrity": "sha512-VqO6OSvLrFVAYYjgsr8tyv62/rCQhPgsZUXLTqoFLSgdkgiUYKYeArbt1uWLlEpkjxQe+P0+sHlbPEte1Bi06Q==" }, + "p-queue@8.1.1": { + "integrity": "sha512-aNZ+VfjobsWryoiPnEApGGmf5WmNsCo9xu8dfaYamG5qaLP7ClhLN6NgsFe6SwJ2UbLEBK5dv9x8Mn5+RVhMWQ==", + "dependencies": [ + "eventemitter3", + "p-timeout" + ] + }, + "p-timeout@6.1.4": { + "integrity": "sha512-MyIV3ZA/PmyBN/ud8vV9XzwTrNtR4jFrObymZYnZqMmW0zA8Z17vnT0rBgFE/TlohB+YCHqXMgZzb3Csp49vqg==" + }, "punycode@2.3.1": { "integrity": "sha512-vYt7UD1U9Wg6138shLtLOvdAu+8DsC/ilFtEVHcH+wydcSpNE20AfSOduf6MkRFahL5FY7X1oU7nKVZFtfq8Fg==" }, @@ -329,6 +372,7 @@ "jsr:@atp/common@~0.1.0-alpha.4", "jsr:@atp/identity@~0.1.0-alpha.1", "jsr:@atp/lexicon@~0.1.0-alpha.2", + "jsr:@atp/sync@~0.1.0-alpha.4", "jsr:@atp/syntax@~0.1.0-alpha.2", "jsr:@atp/xrpc-server@~0.1.0-alpha.3", "jsr:@hono/hono@^4.10.7", diff --git a/ingester/handlers/follow.ts b/ingester/handlers/follow.ts new file mode 100644 index 0000000..3162773 --- /dev/null +++ b/ingester/handlers/follow.ts @@ -0,0 +1,43 @@ +import { CollectionHandler } from "../types.ts"; +import type { Record as FollowRecord } from "../../lex/types/so/sprk/graph/follow.ts"; + +/** Handler for so.sprk.graph.follow collection */ +const followHandler: CollectionHandler = { + collection: "so.sprk.graph.follow", + + handleInsert: async (ctx, evt) => { + const record = evt.record as FollowRecord; + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Follow.findOneAndUpdate( + { uri }, + { + uri, + cid: evt.cid, + authorDid: evt.did, + subject: record.subject, + createdAt: record.createdAt, + indexedAt: new Date().toISOString(), + }, + { upsert: true, new: true }, + ); + ctx.logger.debug("Indexed follow", { uri }); + } catch (err) { + ctx.logger.error("Failed to index follow", { uri, error: err }); + } + }, + + handleDelete: async (ctx, evt) => { + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Follow.deleteOne({ uri }); + ctx.logger.debug("Deleted follow", { uri }); + } catch (err) { + ctx.logger.error("Failed to delete follow", { uri, error: err }); + } + }, +}; + +export default followHandler; diff --git a/ingester/handlers/like.ts b/ingester/handlers/like.ts new file mode 100644 index 0000000..f7028b4 --- /dev/null +++ b/ingester/handlers/like.ts @@ -0,0 +1,46 @@ +import { CollectionHandler } from "../types.ts"; +import type { Record as LikeRecord } from "../../lex/types/so/sprk/feed/like.ts"; + +/** Handler for so.sprk.feed.like collection */ +const likeHandler: CollectionHandler = { + collection: "so.sprk.feed.like", + + handleInsert: async (ctx, evt) => { + const record = evt.record as LikeRecord; + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Like.findOneAndUpdate( + { uri }, + { + uri, + cid: evt.cid, + authorDid: evt.did, + subject: record.subject.uri, + subjectCid: record.subject.cid, + via: record.via?.uri ?? null, + viaCid: record.via?.cid ?? null, + createdAt: record.createdAt, + indexedAt: new Date().toISOString(), + }, + { upsert: true, new: true }, + ); + ctx.logger.debug("Indexed like", { uri }); + } catch (err) { + ctx.logger.error("Failed to index like", { uri, error: err }); + } + }, + + handleDelete: async (ctx, evt) => { + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Like.deleteOne({ uri }); + ctx.logger.debug("Deleted like", { uri }); + } catch (err) { + ctx.logger.error("Failed to delete like", { uri, error: err }); + } + }, +}; + +export default likeHandler; diff --git a/ingester/handlers/post.ts b/ingester/handlers/post.ts new file mode 100644 index 0000000..32cbe34 --- /dev/null +++ b/ingester/handlers/post.ts @@ -0,0 +1,50 @@ +import { CollectionHandler } from "../types.ts"; +import type { Record as PostRecord } from "../../lex/types/so/sprk/feed/post.ts"; + +/** Handler for so.sprk.feed.post collection */ +const postHandler: CollectionHandler = { + collection: "so.sprk.feed.post", + + handleInsert: async (ctx, evt) => { + const record = evt.record as PostRecord; + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Post.findOneAndUpdate( + { uri }, + { + uri, + cid: evt.cid, + authorDid: evt.did, + caption: record.caption, + media: record.media, + sound: record.sound, + langs: record.langs ?? [], + labels: (record.labels && "values" in record.labels) + ? record.labels.values + : [], + tags: record.tags ?? [], + createdAt: record.createdAt, + indexedAt: new Date().toISOString(), + }, + { upsert: true, new: true }, + ); + ctx.logger.debug("Indexed post", { uri }); + } catch (err) { + ctx.logger.error("Failed to index post", { uri, error: err }); + } + }, + + handleDelete: async (ctx, evt) => { + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Post.deleteOne({ uri }); + ctx.logger.debug("Deleted post", { uri }); + } catch (err) { + ctx.logger.error("Failed to delete post", { uri, error: err }); + } + }, +}; + +export default postHandler; diff --git a/ingester/handlers/repost.ts b/ingester/handlers/repost.ts new file mode 100644 index 0000000..ee37c35 --- /dev/null +++ b/ingester/handlers/repost.ts @@ -0,0 +1,46 @@ +import { CollectionHandler } from "../types.ts"; +import type { Record as RepostRecord } from "../../lex/types/so/sprk/feed/repost.ts"; + +/** Handler for so.sprk.feed.repost collection */ +const repostHandler: CollectionHandler = { + collection: "so.sprk.feed.repost", + + handleInsert: async (ctx, evt) => { + const record = evt.record as RepostRecord; + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Repost.findOneAndUpdate( + { uri }, + { + uri, + cid: evt.cid, + authorDid: evt.did, + subject: record.subject.uri, + subjectCid: record.subject.cid, + via: record.via?.uri ?? null, + viaCid: record.via?.cid ?? null, + createdAt: record.createdAt, + indexedAt: new Date().toISOString(), + }, + { upsert: true, new: true }, + ); + ctx.logger.debug("Indexed repost", { uri }); + } catch (err) { + ctx.logger.error("Failed to index repost", { uri, error: err }); + } + }, + + handleDelete: async (ctx, evt) => { + const uri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; + + try { + await ctx.db.models.Repost.deleteOne({ uri }); + ctx.logger.debug("Deleted repost", { uri }); + } catch (err) { + ctx.logger.error("Failed to delete repost", { uri, error: err }); + } + }, +}; + +export default repostHandler; diff --git a/ingester/index.ts b/ingester/index.ts new file mode 100644 index 0000000..ea91833 --- /dev/null +++ b/ingester/index.ts @@ -0,0 +1,66 @@ +import { Firehose } from "@atp/sync"; +import { IdResolver } from "@atp/identity"; +import { getLogger, Logger } from "@logtape/logtape"; +import { Database } from "../db/connection.ts"; +import { CollectionHandler, HandlerContext } from "./types.ts"; + +// collections +import postHandler from "./handlers/post.ts"; +// uncomment the following lines if you want to ingest likes, reposts, or follows +// import likeHandler from "./handlers/like.ts"; +// import repostHandler from "./handlers/repost.ts" +// import followHandler from "./handlers/follow.ts" + +export class Ingester { + idResolver: IdResolver; + logger: Logger; + firehose: Firehose; + db: Database; + handlers: Map; + + constructor( + db: Database, + handlers: CollectionHandler[] = [ + postHandler, + // uncomment the following lines to ingest likes, reposts, or follows + // likeHandler + // repostHandler + // followHandler + ], + idResolver?: IdResolver, + ) { + this.logger = getLogger(["feedgen", "ingester"]); + this.db = db; + this.idResolver = idResolver ?? new IdResolver(); + + // Build handler map for O(1) lookup + this.handlers = new Map(handlers.map((h) => [h.collection, h])); + + const ctx: HandlerContext = { + db: this.db, + logger: this.logger, + }; + + this.firehose = new Firehose({ + idResolver: this.idResolver, + handleEvent: async (evt) => { + if (!("collection" in evt)) return; + + const handler = this.handlers.get(evt.collection); + if (!handler) return; + + if (evt.event === "create" && handler.handleInsert) { + await handler.handleInsert(ctx, evt); + } else if (evt.event === "update" && handler.handleInsert) { + await handler.handleInsert(ctx, evt); + } else if (evt.event === "delete" && handler.handleDelete) { + await handler.handleDelete(ctx, evt); + } + }, + onError: (err) => { + this.logger.error("Firehose error", { error: err }); + }, + filterCollections: handlers.map((h) => h.collection), + }); + } +} diff --git a/ingester/types.ts b/ingester/types.ts new file mode 100644 index 0000000..011a06f --- /dev/null +++ b/ingester/types.ts @@ -0,0 +1,22 @@ +import { Event } from "@atp/sync"; +import { Logger } from "@logtape/logtape"; +import { Database } from "../db/connection.ts"; + +/** Context passed to collection handlers */ +export interface HandlerContext { + db: Database; + logger: Logger; +} + +/** Handler for a specific collection's events */ +export interface CollectionHandler { + collection: string; + handleInsert?: ( + ctx: HandlerContext, + evt: Event & { event: "create" | "update" }, + ) => Promise; + handleDelete?: ( + ctx: HandlerContext, + evt: Event & { event: "delete" }, + ) => Promise; +} diff --git a/main.ts b/main.ts index f6f2ef3..83faeac 100644 --- a/main.ts +++ b/main.ts @@ -11,6 +11,7 @@ import describeFeedGenerator from "./api/describeFeedGenerator.ts"; import getFeedSkeleton from "./api/getFeedSkeleton.ts"; import wellKnown from "./api/well-known.ts"; import health from "./api/health.ts"; +import { Ingester } from "./ingester/index.ts"; await configure({ sinks: { @@ -94,12 +95,15 @@ export async function startServer() { Deno.exit(1); } + const ingester = new Ingester(db); + ingester.firehose.start(); + const { SPRK_HOST, SPRK_PORT } = env; Deno.serve({ hostname: SPRK_HOST, port: SPRK_PORT, onListen: (info) => { - logger.info(`Server listening on ${info.hostname}:${info.port}`); + logger.info(`Server listening on http://${info.hostname}:${info.port}`); }, }, app.fetch); @@ -125,5 +129,3 @@ export async function startServer() { if (import.meta.main) { startServer(); } - -export default app; diff --git a/utils/env.ts b/utils/env.ts index bcee951..039cf61 100644 --- a/utils/env.ts +++ b/utils/env.ts @@ -8,7 +8,7 @@ export const env = { SPRK_FEEDGEN_DOMAIN: envStr("FEEDGEN_DOMAIN"), SPRK_DB_NAME: envStr("SPRK_DB_NAME"), - SPRK_DB_URI: envStr("SPRK_DB_URI") ?? "mongo://localhost:27017", + SPRK_DB_URI: envStr("SPRK_DB_URI") ?? "mongodb://localhost:27017", SPRK_DB_USER: envStr("SPRK_DB_USER"), SPRK_DB_PASS: envStr("SPRK_DB_PASS"), -- 2.51.2