diff --git a/.env b/.env new file mode 100644 index 0000000..71a7dae --- /dev/null +++ b/.env @@ -0,0 +1,5 @@ +DATABASE_URL=file:./label-watcher.db +MIGRATIONS_FOLDER=drizzle +NOTIFY_SMTP_URL=smtps://resend:.... +NOTIFY_SENDER_EMAIL=yougotmail@pdsmoover.com +LOG_LEVEL=debug diff --git a/.env.example b/.env.example index 91ec23f..406cd3a 100644 --- a/.env.example +++ b/.env.example @@ -1,3 +1,5 @@ DATABASE_URL=file:./label-watcher.db MIGRATIONS_FOLDER=drizzle NOTIFY_SMTP_URL=smtps://resend:.... +NOTIFY_SENDER_EMAIL=yougotmail@pdsmoover.com +LOG_LEVEL=info diff --git a/drizzle/0001_workable_leech.sql b/drizzle/0001_workable_leech.sql new file mode 100644 index 0000000..eaf6785 --- /dev/null +++ b/drizzle/0001_workable_leech.sql @@ -0,0 +1 @@ +ALTER TABLE `labels_applied` ADD `labeler` text NOT NULL; \ No newline at end of file diff --git a/drizzle/meta/0001_snapshot.json b/drizzle/meta/0001_snapshot.json new file mode 100644 index 0000000..83818f1 --- /dev/null +++ b/drizzle/meta/0001_snapshot.json @@ -0,0 +1,170 @@ +{ + "version": "6", + "dialect": "sqlite", + "id": "368e7243-8e16-4e27-a7af-a906d976d75e", + "prevId": "5a06d315-4901-4f9e-aab2-1c8dd03bb750", + "tables": { + "labeler_cursors": { + "name": "labeler_cursors", + "columns": { + "labeler_id": { + "name": "labeler_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "cursor": { + "name": "cursor", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "labeler_cursors_labeler_id_unique": { + "name": "labeler_cursors_labeler_id_unique", + "columns": [ + "labeler_id" + ], + "isUnique": true + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "labels_applied": { + "name": "labels_applied", + "columns": { + "id": { + "name": "id", + "type": "integer", + "primaryKey": true, + "notNull": true, + "autoincrement": true + }, + "did": { + "name": "did", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "label": { + "name": "label", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "labeler": { + "name": "labeler", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "action": { + "name": "action", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "negated": { + "name": "negated", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": false + }, + "date_applied": { + "name": "date_applied", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": { + "labels_applied_did_watched_repos_did_fk": { + "name": "labels_applied_did_watched_repos_did_fk", + "tableFrom": "labels_applied", + "tableTo": "watched_repos", + "columnsFrom": [ + "did" + ], + "columnsTo": [ + "did" + ], + "onDelete": "no action", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "watched_repos": { + "name": "watched_repos", + "columns": { + "did": { + "name": "did", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "pds_host": { + "name": "pds_host", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "active": { + "name": "active", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "date_first_seen": { + "name": "date_first_seen", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "watched_repos_did_unique": { + "name": "watched_repos_did_unique", + "columns": [ + "did" + ], + "isUnique": true + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + } + }, + "views": {}, + "enums": {}, + "_meta": { + "schemas": {}, + "tables": {}, + "columns": {} + }, + "internal": { + "indexes": {} + } +} \ No newline at end of file diff --git a/drizzle/meta/_journal.json b/drizzle/meta/_journal.json index 17320ce..c203c01 100644 --- a/drizzle/meta/_journal.json +++ b/drizzle/meta/_journal.json @@ -8,6 +8,13 @@ "when": 1771615394802, "tag": "0000_crazy_wallflower", "breakpoints": true + }, + { + "idx": 1, + "version": "6", + "when": 1771625783853, + "tag": "0001_workable_leech", + "breakpoints": true } ] } \ No newline at end of file diff --git a/package.json b/package.json index 8d9e113..9fa985d 100644 --- a/package.json +++ b/package.json @@ -5,7 +5,7 @@ "description": "", "main": "index.js", "scripts": { - "start": "tsc && node dist/index.js | pino-pretty", + "start": "tsc && node --env-file=.env dist/index.js | pino-pretty", "db:generate": "drizzle-kit generate", "db:migrate": "drizzle-kit migrate", "db:studio": "drizzle-kit studio" diff --git a/src/db/schema.ts b/src/db/schema.ts index aae7cc8..148a10e 100644 --- a/src/db/schema.ts +++ b/src/db/schema.ts @@ -14,6 +14,7 @@ export const labelsApplied = sqliteTable("labels_applied", { .notNull() .references(() => watchedRepos.did), label: text("label").notNull(), + labeler: text("labeler").notNull(), action: text("action").notNull(), negated: integer("negated", { mode: "boolean" }).default(false).notNull(), dateApplied: integer("date_applied", { mode: "timestamp" }).notNull(), diff --git a/src/handlers/handleNewLabel.ts b/src/handlers/handleNewLabel.ts index ebb4cd4..c9a8c42 100644 --- a/src/handlers/handleNewLabel.ts +++ b/src/handlers/handleNewLabel.ts @@ -3,29 +3,58 @@ import type { LabelerConfig } from "../types/settings.js"; import { logger } from "../logger.js"; import type { LibSQLDatabase } from "drizzle-orm/libsql"; import * as schema from "../db/schema.js"; +import { count, eq } from "drizzle-orm"; export const handleNewLabel = async ( config: LabelerConfig, label: Label, db: LibSQLDatabase, ) => { - // TODO: MAKE SURE TO CHECK NEG - logger.info({ host: config.host }, "From"); - let labledDate = new Date(label.cts); - if (config.labels[label.val]) { - logger.info( - { action: config.labels[label.val]?.action }, - "Listed label found. Performing the action", + try { + // TODO: MAKE SURE TO CHECK NEG + let labledDate = new Date(label.cts); + logger.debug( + { + labeler: config.host, + val: label.val, + uri: label.uri, + neg: label.neg, + date: labledDate, + }, + "Label", ); + + let labelConfig = config.labels[label.val]; + if (labelConfig) { + const isRepoWatched = await db + .select() + .from(schema.watchedRepos) + .where(eq(schema.watchedRepos.did, label.uri)) + .limit(1); + + if (isRepoWatched.length > 0) { + logger.info( + { action: config.labels[label.val]?.action }, + `Listed label: ${label.val} found. Performing the action against: ${label.uri}`, + ); + + await db.insert(schema.labelsApplied).values({ + did: label.uri, + label: label.val, + labeler: config.host, + action: labelConfig.action, + negated: label.neg ?? false, + dateApplied: labledDate, + }); + + return; + } + logger.warn( + { action: config.labels[label.val]?.action }, + "Listed label found but repo is not watched. Skipping", + ); + } + } catch (error) { + logger.error({ error }, "Error handling new label"); } - logger.info( - { - src: label.src, - val: label.val, - uri: label.uri, - neg: label.neg, - date: labledDate, - }, - "Label", - ); }; diff --git a/src/handlers/lablerSubscriber.ts b/src/handlers/lablerSubscriber.ts index aa8b729..29500c8 100644 --- a/src/handlers/lablerSubscriber.ts +++ b/src/handlers/lablerSubscriber.ts @@ -48,9 +48,12 @@ export const labelerSubscriber = ( } case "com.atproto.label.subscribeLabels#labels": { for (const label of message.labels) { - queue.add(async () => { - await handleNewLabel(config, label, db); - }); + // We only care about labels for identities, not content for now + if (label.uri.startsWith("did:")) { + queue.add(async () => { + await handleNewLabel(config, label, db); + }); + } } break; } diff --git a/src/index.ts b/src/index.ts index 3b9ffd8..d05f1e0 100644 --- a/src/index.ts +++ b/src/index.ts @@ -45,13 +45,23 @@ logger.info("Identity queue backfill and completion complete."); const lastCursors = await db.select().from(labelerCursor); // Sets up the subscribers to the labelers -const labelSubscribers = Object.entries(settings.labeler).map(([_, config]) => { - let lastCursorRow = lastCursors.find( - (cursor) => cursor.labelerId === config.host, - ); - let lastCursor = lastCursorRow?.cursor ?? undefined; - return labelerSubscriber(config, lastCursor, db, labelQueue); -}); +const labelSubscribers = Object.entries(settings.labeler) + .map(([_, config]) => { + if (config.labels == undefined) { + logger.info( + { host: config.host }, + "No labels to watch not starting subscriber for this one", + ); + return null; + } + + let lastCursorRow = lastCursors.find( + (cursor) => cursor.labelerId === config.host, + ); + let lastCursor = lastCursorRow?.cursor ?? undefined; + return labelerSubscriber(config, lastCursor, db, labelQueue); + }) + .filter((x) => x !== null); const pdsSubscribers = Object.entries(settings.pds) .map(([_, config]) => { diff --git a/src/logger.ts b/src/logger.ts index ed26420..29b52d7 100644 --- a/src/logger.ts +++ b/src/logger.ts @@ -1,3 +1,5 @@ import pino from "pino"; -export const logger = pino(); +export const logger = pino({ + level: process.env.LOG_LEVEL ?? "info", +});