From 1bd62dc53815581c61696caf0283084356307b95 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Andri=20=C3=93skarsson?= Date: Mon, 13 Jan 2025 16:54:01 +0100 Subject: [PATCH] worker update follows --- apps/worker/index.ts | 16 ++------------ apps/worker/package.json | 3 +-- .../atproto/domain/get-tech-following-feed.ts | 5 +++-- .../atproto/feeds/queries/following-all.ts | 3 ++- packages/atproto/worker/add-job.ts | 21 +++++++++++++++++++ {apps => packages/atproto}/worker/crontab | 0 .../atproto}/worker/graphile.config.ts | 4 +++- packages/atproto/worker/main.ts | 16 ++++++++++++++ .../tasks/classify-unclassified-posts.ts | 0 .../atproto}/worker/tasks/cleanup.ts | 0 .../tasks/fetch-missing-post-records.ts | 0 .../atproto}/worker/tasks/train-classifier.ts | 0 .../atproto/worker/tasks/update-followers.ts | 21 +++++++++++++++++++ tsconfig.json | 3 ++- 14 files changed, 71 insertions(+), 21 deletions(-) create mode 100644 packages/atproto/worker/add-job.ts rename {apps => packages/atproto}/worker/crontab (100%) rename {apps => packages/atproto}/worker/graphile.config.ts (72%) create mode 100644 packages/atproto/worker/main.ts rename {apps => packages/atproto}/worker/tasks/classify-unclassified-posts.ts (100%) rename {apps => packages/atproto}/worker/tasks/cleanup.ts (100%) rename {apps => packages/atproto}/worker/tasks/fetch-missing-post-records.ts (100%) rename {apps => packages/atproto}/worker/tasks/train-classifier.ts (100%) create mode 100644 packages/atproto/worker/tasks/update-followers.ts diff --git a/apps/worker/index.ts b/apps/worker/index.ts index 2ff92f6..5c523cb 100644 --- a/apps/worker/index.ts +++ b/apps/worker/index.ts @@ -1,18 +1,6 @@ -import { run } from "graphile-worker"; -import preset from "./graphile.config"; +import { workerProcess } from "@andrioid/atproto/worker/main"; -async function main() { - const runner = await run({ preset }); - - runner.events.on("job:complete", () => { - // Trigger Bun's GC manually after a job finishes - Bun.gc(false); - }); - - await runner.promise; -} - -main().catch((err) => { +workerProcess().catch((err) => { console.error(err); process.exit(1); }); diff --git a/apps/worker/package.json b/apps/worker/package.json index 13aa3be..a24fe29 100644 --- a/apps/worker/package.json +++ b/apps/worker/package.json @@ -8,8 +8,7 @@ "prod": "bun index.ts" }, "dependencies": { - "@andrioid/atproto": "workspace:*", - "graphile-worker": "^0.16.6" + "@andrioid/atproto": "workspace:*" }, "devDependencies": {}, "peerDependencies": { diff --git a/packages/atproto/domain/get-tech-following-feed.ts b/packages/atproto/domain/get-tech-following-feed.ts index 1cec333..f1eaf62 100644 --- a/packages/atproto/domain/get-tech-following-feed.ts +++ b/packages/atproto/domain/get-tech-following-feed.ts @@ -3,7 +3,7 @@ import { max } from "drizzle-orm"; import { and, desc, eq, gt, gte } from "drizzle-orm/expressions"; import type { FeedHandlerArgs, FeedHandlerOutput } from "../feeds"; import { fromCursor, toCursor } from "../helpers/cursor"; -import { getOrUpdateFollows } from "./get-or-update-follows"; +import { addJob } from "../worker/add-job"; import { repostTable } from "./post/post-reposts.table"; import { postScores } from "./post/post-scores.view"; import { postTable } from "./post/post.table"; @@ -14,7 +14,8 @@ export async function getTechFollowingFeed( ): Promise { const { ctx, actorDid, cursor, limit = 50 } = args; - await getOrUpdateFollows(ctx, actorDid); + await addJob("update-followers", { did: actorDid }); + //await getOrUpdateFollows(ctx, actorDid); const fls = ctx.db .select({ follows: followTable.follows }) diff --git a/packages/atproto/feeds/queries/following-all.ts b/packages/atproto/feeds/queries/following-all.ts index 7a148bf..239acaa 100644 --- a/packages/atproto/feeds/queries/following-all.ts +++ b/packages/atproto/feeds/queries/following-all.ts @@ -1,3 +1,4 @@ +/* import { desc } from "drizzle-orm"; import { union } from "drizzle-orm/pg-core"; import type { FeedHandlerArgs } from ".."; @@ -12,5 +13,5 @@ export function followingAllQuery(args: FeedHandlerArgs) { .orderBy(desc(postTable.lastMentioned)); return combinedQuery; } - +*/ // TODO: Read up on WITH clause and CTEs diff --git a/packages/atproto/worker/add-job.ts b/packages/atproto/worker/add-job.ts new file mode 100644 index 0000000..2c8f1d9 --- /dev/null +++ b/packages/atproto/worker/add-job.ts @@ -0,0 +1,21 @@ +import { config } from "@andrioid/atproto/config"; +import type { UpdateFollowersTaskType } from "@andrioid/atproto/worker/tasks/update-followers"; +import type { Job, TaskSpec } from "graphile-worker"; +import { makeWorkerUtils } from "graphile-worker"; + +export const workerUtils = await makeWorkerUtils({ + connectionString: config.pgURL, +}); + +type AddJobFnArgs = UpdateFollowersTaskType; + +type AddJobFn = (...args: AddJobFnArgs) => Promise; + +export const addJob: AddJobFn = (...args) => { + const taskSpec: TaskSpec = { + queueName: args[0], + maxAttempts: 2, + ...args[2], + }; + return workerUtils.addJob(args[0], args[1], args[2]); +}; diff --git a/apps/worker/crontab b/packages/atproto/worker/crontab similarity index 100% rename from apps/worker/crontab rename to packages/atproto/worker/crontab diff --git a/apps/worker/graphile.config.ts b/packages/atproto/worker/graphile.config.ts similarity index 72% rename from apps/worker/graphile.config.ts rename to packages/atproto/worker/graphile.config.ts index 9f2782a..8e03a80 100644 --- a/apps/worker/graphile.config.ts +++ b/packages/atproto/worker/graphile.config.ts @@ -1,5 +1,6 @@ import { config } from "@andrioid/atproto/config"; import { WorkerPreset } from "graphile-worker"; +import path from "node:path"; const preset: GraphileConfig.Preset = { extends: [WorkerPreset], @@ -9,7 +10,8 @@ const preset: GraphileConfig.Preset = { pollInterval: 2000, preparedStatements: true, schema: "graphile_worker", - crontabFile: "crontab", + crontabFile: path.join(import.meta.dirname, "./crontab"), + taskDirectory: path.join(import.meta.dirname, "./tasks"), concurrentJobs: 1, fileExtensions: [".ts"], }, diff --git a/packages/atproto/worker/main.ts b/packages/atproto/worker/main.ts new file mode 100644 index 0000000..0009aef --- /dev/null +++ b/packages/atproto/worker/main.ts @@ -0,0 +1,16 @@ +export * from "@andrioid/atproto/worker/add-job"; +import { run } from "graphile-worker"; +import preset from "./graphile.config"; + +export async function workerProcess() { + const runner = await run({ + preset, + }); + + runner.events.on("job:complete", () => { + // Trigger Bun's GC manually after a job finishes + Bun.gc(false); + }); + + await runner.promise; +} diff --git a/apps/worker/tasks/classify-unclassified-posts.ts b/packages/atproto/worker/tasks/classify-unclassified-posts.ts similarity index 100% rename from apps/worker/tasks/classify-unclassified-posts.ts rename to packages/atproto/worker/tasks/classify-unclassified-posts.ts diff --git a/apps/worker/tasks/cleanup.ts b/packages/atproto/worker/tasks/cleanup.ts similarity index 100% rename from apps/worker/tasks/cleanup.ts rename to packages/atproto/worker/tasks/cleanup.ts diff --git a/apps/worker/tasks/fetch-missing-post-records.ts b/packages/atproto/worker/tasks/fetch-missing-post-records.ts similarity index 100% rename from apps/worker/tasks/fetch-missing-post-records.ts rename to packages/atproto/worker/tasks/fetch-missing-post-records.ts diff --git a/apps/worker/tasks/train-classifier.ts b/packages/atproto/worker/tasks/train-classifier.ts similarity index 100% rename from apps/worker/tasks/train-classifier.ts rename to packages/atproto/worker/tasks/train-classifier.ts diff --git a/packages/atproto/worker/tasks/update-followers.ts b/packages/atproto/worker/tasks/update-followers.ts new file mode 100644 index 0000000..6fece18 --- /dev/null +++ b/packages/atproto/worker/tasks/update-followers.ts @@ -0,0 +1,21 @@ +import { createAtContext, getOrUpdateFollows } from "@andrioid/atproto"; +import type { TaskSpec } from "graphile-worker"; + +export type UpdateFollowersTaskType = [ + identifier: "update-followers", + payload: { + did: string; + }, + spec?: TaskSpec +]; + +export default async function updateFollowersTask( + payload: UpdateFollowersTaskType[1] +) { + if (!payload.did) { + console.warn("updateFollowersTask called without valid payload"); + return; + } + const ctx = await createAtContext(); + await getOrUpdateFollows(ctx, payload.did); +} diff --git a/tsconfig.json b/tsconfig.json index 238655f..462e7d1 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -23,5 +23,6 @@ "noUnusedLocals": false, "noUnusedParameters": false, "noPropertyAccessFromIndexSignature": false - } + }, + "exclude": ["node_modules/"] } -- 2.51.2