diff --git a/apps/jetstream/package.json b/apps/jetstream/package.json index 666b479..502977f 100644 --- a/apps/jetstream/package.json +++ b/apps/jetstream/package.json @@ -4,6 +4,7 @@ "type": "module", "scripts": { "dev": "bun --watch --env-file=../../.env index.ts", + "debug": "bun --inspect --env-file=../../.env index.ts", "prod": "bun index.ts", "prod:node": "node --experimental-strip-types index.ts", "prod:deno": "deno run --allow-env --allow-read --allow-ffi --allow-net --unstable-sloppy-imports --env-file=../../.env index.ts" diff --git a/apps/web/tsconfig.json b/apps/web/tsconfig.json index 8d86433..83677ab 100644 --- a/apps/web/tsconfig.json +++ b/apps/web/tsconfig.json @@ -6,6 +6,8 @@ "~astro/*": ["./src/components/astro/*"], "~/*": ["./*"], }, + "verbatimModuleSyntax": true + }, "extends": "astro/tsconfigs/strict", diff --git a/packages/atproto/context.ts b/packages/atproto/context.ts index 31bebae..4ebf603 100644 --- a/packages/atproto/context.ts +++ b/packages/atproto/context.ts @@ -6,7 +6,10 @@ import { config } from "./config"; import { appData } from "./db/schema"; export async function createAtContext() { - const dbClient = postgres(config.pgURL); + const dbClient = postgres(config.pgURL, { + idle_timeout: 20, + max_lifetime: 60 * 30, + }); const db = drizzle({ client: dbClient, }); diff --git a/packages/atproto/domain/extract-text-from-post.ts b/packages/atproto/domain/extract-text-from-post.ts index 40f5cbb..805b475 100644 --- a/packages/atproto/domain/extract-text-from-post.ts +++ b/packages/atproto/domain/extract-text-from-post.ts @@ -1,3 +1,4 @@ +import type { FeedPostWithUri } from "@andrioid/jetstream"; import { AppBskyEmbedExternal, AppBskyEmbedImages, @@ -7,7 +8,6 @@ import { import { eq } from "drizzle-orm"; import type { AtContext } from "../context"; import { postRecords } from "./post/post-record.table"; -import type { FeedPostWithUri } from "./queue-post"; export type ExtractedTextType = | "body" diff --git a/packages/atproto/domain/jetstream-did-list.ts b/packages/atproto/domain/jetstream-did-list.ts new file mode 100644 index 0000000..96bbdb2 --- /dev/null +++ b/packages/atproto/domain/jetstream-did-list.ts @@ -0,0 +1,17 @@ +import type { AtContext } from "../context"; +import { followTable } from "./user/user-follows.table"; + +export async function getDids( + ctx: AtContext +): Promise>> { + const wantedDids = await ctx.db + .selectDistinct({ did: followTable.follows }) + .from(followTable); + + const dids = wantedDids.map((d) => d.did); + if (dids.length === 0) { + console.warn("aborting jetstream connection, no dids requested"); + return []; + } + return dids; +} diff --git a/packages/atproto/domain/jetstream-subscription.ts b/packages/atproto/domain/jetstream-subscription.ts index f173d24..439b23a 100644 --- a/packages/atproto/domain/jetstream-subscription.ts +++ b/packages/atproto/domain/jetstream-subscription.ts @@ -2,27 +2,15 @@ import { Jetstream } from "@andrioid/jetstream"; import { subMinutes } from "date-fns"; import { desc } from "drizzle-orm"; import type { AtContext } from "../context"; +import { getDids } from "./jetstream-did-list"; import { postTable } from "./post/post.table"; import { queuePost } from "./queue-post"; import { queueRePost } from "./queue-repost"; -import { followTable } from "./user/user-follows.table"; export const LISTEN_NOTIFY_NEW_SUBSCRIBERS = "atproto.subscriber.update"; export async function listenForPosts(ctx: AtContext) { - async function getDids(): Promise>> { - const wantedDids = await ctx.db - .selectDistinct({ did: followTable.follows }) - .from(followTable); - - const dids = wantedDids.map((d) => d.did); - if (dids.length === 0) { - console.warn("aborting jetstream connection, no dids requested"); - return []; - } - return [...dids]; - } - const dids = await getDids(); + const dids = await getDids(ctx); const latestPost = await ctx.db .select() @@ -40,10 +28,21 @@ export async function listenForPosts(ctx: AtContext) { wantedCollections: ["app.bsky.feed.post", "app.bsky.feed.repost"], cursor: cursor?.toString(), }); + + let postCounter = 0; + const initialMem = process.memoryUsage().rss; + js.on({ event: "post", cb: async (msg) => { await queuePost(ctx, msg); + postCounter++; + const mem = process.memoryUsage().rss; + console.log( + `[jetstream] queued ${postCounter}`, + mem, + Number((mem / initialMem) * 100).toFixed(2) + ); }, }); js.on({ @@ -67,7 +66,7 @@ export async function listenForPosts(ctx: AtContext) { return; } - const newDids = await getDids(); + const newDids = await getDids(ctx); console.log("[jetstream] restarting jetstream, new subscribers"); js.setupSockets({ wantedDids: newDids, diff --git a/packages/atproto/domain/store-post-with-record.ts b/packages/atproto/domain/store-post-with-record.ts index 88e721e..ec4d868 100644 --- a/packages/atproto/domain/store-post-with-record.ts +++ b/packages/atproto/domain/store-post-with-record.ts @@ -66,5 +66,4 @@ export async function storePost(ctx: AtContext, post: FeedPostWithUri) { uri: post.uri, }) ); - console.log(`[queue] ${post.uri} stored`); } diff --git a/packages/jetstream/jetstream.ts b/packages/jetstream/jetstream.ts index 0b7e7bd..b000d8a 100644 --- a/packages/jetstream/jetstream.ts +++ b/packages/jetstream/jetstream.ts @@ -152,8 +152,6 @@ export class Jetstream { if (msg.kind !== "commit") return; // only kind supported atm if (msg.commit.operation !== "create") return; - // TODO: Broke these types being smart - fix tomorrow - switch (msg.commit.collection) { case "app.bsky.feed.post": //console.log("[jetstream post", msg);