From baa5f6447a3bd3761030d0fd885745c71f6a2005 Mon Sep 17 00:00:00 2001 From: "Sven (AAMOS AI)" Date: Sun, 16 Aug 2026 21:43:19 +0700 Subject: [PATCH] feat(worker): refund AI scan quota on definitive job failure + sanitize DLQ jobId --- apps/worker/src/index.ts | 7 +++- apps/worker/src/processors/quota-refund.ts | 39 ++++++++++++++++++++++ 2 files changed, 45 insertions(+), 1 deletion(-) create mode 100644 apps/worker/src/processors/quota-refund.ts diff --git a/apps/worker/src/index.ts b/apps/worker/src/index.ts index 3b12c67..19f6d99 100644 --- a/apps/worker/src/index.ts +++ b/apps/worker/src/index.ts @@ -4,6 +4,7 @@ import IORedis from "ioredis"; import { createServer } from "node:http"; import { createContext } from "./context.js"; import { processScanJob } from "./processors/scans.js"; +import { refundAiScanForFailedJob } from "./processors/quota-refund.js"; import { processModerateRecipe } from "./processors/moderation.js"; import { processTranslateRecipe } from "./processors/translation.js"; import { processGenerateWeekPlan } from "./processors/weekplan.js"; @@ -186,6 +187,10 @@ worker.on("failed", async (job, err) => { console.error(`✗ ${jobType} (${job?.id}): ${err.message} (attempt ${attemptsMade}/${attempts})`); if (job && attemptsMade >= attempts) { + const _sjid = (job.data as { scanJobId?: unknown })?.scanJobId; + if (_sjid) { + try { await refundAiScanForFailedJob(ctx, String(_sjid)); } catch (e) { console.error("Kvot-refund misslyckades:", e); } + } try { await deadLetterQueue.add( job.name, @@ -197,7 +202,7 @@ worker.on("failed", async (job, err) => { _failedAt: new Date().toISOString(), }, { - jobId: `${instanceId}:${job.id}`, + jobId: `${instanceId}__${job.id}`.replaceAll(":", "_"), removeOnComplete: { age: 30 * 24 * 60 * 60 }, // 30 dagar removeOnFail: { age: 90 * 24 * 60 * 60 }, // 90 dagar }, diff --git a/apps/worker/src/processors/quota-refund.ts b/apps/worker/src/processors/quota-refund.ts new file mode 100644 index 0000000..208a387 --- /dev/null +++ b/apps/worker/src/processors/quota-refund.ts @@ -0,0 +1,39 @@ +import { and, eq, sql } from "drizzle-orm"; +import { schema } from "@app/database"; +import type { WorkerContext } from "../context.js"; + +/** + * Återbetalar AI-skanningskvot när ett skann-jobb dör definitivt. + * Speglar consumeAiScan (api): kvot debiteras på household-ägaren (billingUser). + * Barcode-skanningar drar ingen kvot och hoppas över. Golvat på 0. + */ +export async function refundAiScanForFailedJob(ctx: WorkerContext, scanJobId: string): Promise { + const [job] = await ctx.db + .select({ userId: schema.scanJobs.userId, scanType: schema.scanJobs.scanType }) + .from(schema.scanJobs) + .where(eq(schema.scanJobs.id, scanJobId)) + .limit(1); + if (!job || job.scanType === "barcode") return; + + let billingUserId = job.userId; + const [member] = await ctx.db + .select({ householdId: schema.householdMembers.householdId }) + .from(schema.householdMembers) + .where(eq(schema.householdMembers.userId, job.userId)) + .orderBy(schema.householdMembers.joinedAt) + .limit(1); + if (member?.householdId) { + const [owner] = await ctx.db + .select({ userId: schema.householdMembers.userId }) + .from(schema.householdMembers) + .where(and(eq(schema.householdMembers.householdId, member.householdId), eq(schema.householdMembers.role, "owner"))) + .limit(1); + if (owner?.userId) billingUserId = owner.userId; + } + + const month = new Date().toISOString().slice(0, 7); + await ctx.db + .update(schema.aiUsageCounters) + .set({ aiScans: sql`GREATEST(${schema.aiUsageCounters.aiScans} - 1, 0)` }) + .where(and(eq(schema.aiUsageCounters.userId, billingUserId), eq(schema.aiUsageCounters.month, month))); +}