From 9718d6888afdebb6bd18a81f92576cc808349abd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Bj=C3=B6rn=20P=C3=B6ttker?= Date: Mon, 29 Jun 2026 17:54:16 +0200 Subject: [PATCH] feat: Webhook-Warteschlange in Datenbank persistieren MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Die Warteschlange liegt nicht mehr im Speicher, sondern in der neuen Tabelle webhook_queue (Entity WebhookQueueItem, documentId als Primärschlüssel). So überstehen ausstehende IDs einen Neustart und werden nach dem Boot weiterverarbeitet. - Neue Entity + Migration (CreateWebhookQueue); in data-source.ts/Barrel registriert. Dev: synchronize legt die Tabelle an; Prod: migrationsRun. - enqueue nutzt INSERT IGNORE (atomar, race-sicher) -> Dedup über den PK. - processQueue holt FIFO (createdAt ASC), entfernt die Zeile vor der Verarbeitung, arbeitet sequenziell (isProcessing-Guard). - getStatus/queueSize lesen die DB-Anzahl. Controller awaitet enqueue. - Unit-Test mit simuliertem FIFO-Repo: Dedup, Remove-on-start, Re-Enqueue, No-Parallel (deterministisch). Co-Authored-By: Claude Opus 4.8 --- paperless-backend/src/database/data-source.ts | 2 + .../src/database/entities/index.ts | 1 + .../entities/webhook-queue-item.entity.ts | 26 +++++ .../1782700000000-CreateWebhookQueue.ts | 25 +++++ .../src/webhook/webhook-queue.service.spec.ts | 104 +++++++++++++----- .../src/webhook/webhook-queue.service.ts | 76 +++++++------ .../src/webhook/webhook.controller.ts | 2 +- .../src/webhook/webhook.module.ts | 7 +- 8 files changed, 182 insertions(+), 61 deletions(-) create mode 100644 paperless-backend/src/database/entities/webhook-queue-item.entity.ts create mode 100644 paperless-backend/src/database/migrations/1782700000000-CreateWebhookQueue.ts diff --git a/paperless-backend/src/database/data-source.ts b/paperless-backend/src/database/data-source.ts index aed7b4d..952bf6b 100644 --- a/paperless-backend/src/database/data-source.ts +++ b/paperless-backend/src/database/data-source.ts @@ -26,6 +26,7 @@ import { CorrespondentEmailMapping, UserSettings, LabelPrintJob, + WebhookQueueItem, } from './entities'; // CLI-Kontext: .env laden (Laufzeit im Container liefert die Variablen via Docker, @@ -57,6 +58,7 @@ export const entities = [ CorrespondentEmailMapping, UserSettings, LabelPrintJob, + WebhookQueueItem, ]; const isProduction = process.env.NODE_ENV === 'production'; diff --git a/paperless-backend/src/database/entities/index.ts b/paperless-backend/src/database/entities/index.ts index 34082d7..75d7d18 100644 --- a/paperless-backend/src/database/entities/index.ts +++ b/paperless-backend/src/database/entities/index.ts @@ -21,3 +21,4 @@ export { InboxPostprocessingAction } from './inbox-postprocessing-action.entity' export { CorrespondentEmailMapping } from './correspondent-email-mapping.entity'; export { UserSettings } from './user-settings.entity'; export { LabelPrintJob } from './label-print-job.entity'; +export { WebhookQueueItem } from './webhook-queue-item.entity'; diff --git a/paperless-backend/src/database/entities/webhook-queue-item.entity.ts b/paperless-backend/src/database/entities/webhook-queue-item.entity.ts new file mode 100644 index 0000000..e3b6505 --- /dev/null +++ b/paperless-backend/src/database/entities/webhook-queue-item.entity.ts @@ -0,0 +1,26 @@ +import { + Entity, + PrimaryColumn, + Column, + CreateDateColumn, + Index, +} from 'typeorm'; + +/** + * Persistente Warteschlange der vom Paperless-Webhook gemeldeten Dokument-IDs. + * `documentId` ist Primärschlüssel und erzwingt damit, dass jede ID nur einmal + * in der Warteschlange steht (Dedup auf DB-Ebene). Die Tabelle übersteht + * Neustarts; ausstehende IDs werden nach dem Boot weiterverarbeitet. + */ +@Entity('webhook_queue') +export class WebhookQueueItem { + @PrimaryColumn({ type: 'int' }) + documentId!: number; + + @Column({ type: 'varchar', length: 100, nullable: true }) + action!: string | null; + + @Index() + @CreateDateColumn() + createdAt!: Date; +} diff --git a/paperless-backend/src/database/migrations/1782700000000-CreateWebhookQueue.ts b/paperless-backend/src/database/migrations/1782700000000-CreateWebhookQueue.ts new file mode 100644 index 0000000..aee2c23 --- /dev/null +++ b/paperless-backend/src/database/migrations/1782700000000-CreateWebhookQueue.ts @@ -0,0 +1,25 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +/** + * Legt die persistente Webhook-Warteschlange an (`webhook_queue`). + * `documentId` ist Primärschlüssel → jede Dokument-ID kommt nur einmal vor. + */ +export class CreateWebhookQueue1782700000000 implements MigrationInterface { + name = 'CreateWebhookQueue1782700000000'; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS \`webhook_queue\` ( + \`documentId\` int NOT NULL, + \`action\` varchar(100) NULL, + \`createdAt\` datetime(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6), + PRIMARY KEY (\`documentId\`), + INDEX \`IDX_webhook_queue_createdAt\` (\`createdAt\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP TABLE \`webhook_queue\``); + } +} diff --git a/paperless-backend/src/webhook/webhook-queue.service.spec.ts b/paperless-backend/src/webhook/webhook-queue.service.spec.ts index a92569c..469ad22 100644 --- a/paperless-backend/src/webhook/webhook-queue.service.spec.ts +++ b/paperless-backend/src/webhook/webhook-queue.service.spec.ts @@ -8,13 +8,51 @@ jest.mock('../paperless/paperless-processor.service', () => ({ import { WebhookQueueService } from './webhook-queue.service'; +/** + * Baut ein Mock-Repository, das die `webhook_queue`-Tabelle durch ein einfaches + * FIFO-Array simuliert: `INSERT IGNORE` (Dedup), FIFO-`find` und `delete`. + */ +function createQueueRepoMock() { + const store: number[] = []; + return { + store, + createQueryBuilder: jest.fn(() => ({ + insert: () => ({ + into: () => ({ + values: (v: { documentId: number }) => ({ + orIgnore: () => ({ + execute: () => { + const added = !store.includes(v.documentId); + if (added) store.push(v.documentId); + return Promise.resolve({ + raw: { affectedRows: added ? 1 : 0 }, + }); + }, + }), + }), + }), + }), + })), + find: jest.fn(() => + Promise.resolve( + store.length + ? [{ documentId: store[0], action: null, createdAt: new Date() }] + : [], + ), + ), + delete: jest.fn((criteria: { documentId: number }) => { + const idx = store.indexOf(criteria.documentId); + if (idx >= 0) store.splice(idx, 1); + return Promise.resolve({ affected: 1 }); + }), + count: jest.fn(() => Promise.resolve(store.length)), + }; +} + describe('WebhookQueueService', () => { let processor: { processDocumentById: jest.Mock }; - let settingRepo: { - findOneBy: jest.Mock; - create: jest.Mock; - save: jest.Mock; - }; + let settingRepo: { findOneBy: jest.Mock; create: jest.Mock; save: jest.Mock }; + let queueRepo: ReturnType; let service: WebhookQueueService; beforeEach(() => { @@ -26,71 +64,83 @@ describe('WebhookQueueService', () => { create: jest.fn((x: unknown) => x), save: jest.fn().mockResolvedValue(undefined), }; - service = new WebhookQueueService(processor as never, settingRepo as never); + queueRepo = createQueueRepoMock(); + service = new WebhookQueueService( + processor as never, + settingRepo as never, + queueRepo as never, + ); }); - it('reiht jede ID nur einmal ein (Dedup)', () => { - service.enqueue(5); - service.enqueue(5); - expect(service.size).toBe(1); + it('reiht jede ID nur einmal ein (Dedup)', async () => { + await service.enqueue(5); + await service.enqueue(5); + expect(queueRepo.store).toEqual([5]); + expect(await service.count()).toBe(1); }); it('verarbeitet alle eingereihten IDs und leert die Warteschlange', async () => { - service.enqueue(1); - service.enqueue(2); + await service.enqueue(1); + await service.enqueue(2); await service.processQueue(); expect(processor.processDocumentById).toHaveBeenCalledWith(1); expect(processor.processDocumentById).toHaveBeenCalledWith(2); expect(processor.processDocumentById).toHaveBeenCalledTimes(2); - expect(service.size).toBe(0); + expect(queueRepo.store).toEqual([]); }); - it('entfernt die ID vor Verarbeitungsbeginn aus der Liste', async () => { - let sizeWhileProcessing = -1; - processor.processDocumentById.mockImplementation(() => { - sizeWhileProcessing = service.size; + it('entfernt die ID vor Verarbeitungsbeginn aus der Tabelle', async () => { + let containedWhileProcessing = true; + processor.processDocumentById.mockImplementation((id: number) => { + containedWhileProcessing = queueRepo.store.includes(id); return Promise.resolve({ processed: true }); }); - service.enqueue(42); + await service.enqueue(42); await service.processQueue(); - // Beim Verarbeitungsstart war die einzige ID bereits entfernt. - expect(sizeWhileProcessing).toBe(0); + // Beim Verarbeitungsstart war die ID bereits aus der Tabelle entfernt. + expect(containedWhileProcessing).toBe(false); }); it('arbeitet eine während der Verarbeitung erneut eingereihte ID erneut ab', async () => { let firstRun = true; - processor.processDocumentById.mockImplementation((id: number) => { + processor.processDocumentById.mockImplementation(async (id: number) => { // Beim ersten Lauf feuert der Webhook erneut, während verarbeitet wird. if (firstRun && id === 7) { firstRun = false; - service.enqueue(7); + await service.enqueue(7); } - return Promise.resolve({ processed: true }); + return { processed: true }; }); - service.enqueue(7); + await service.enqueue(7); await service.processQueue(); expect(processor.processDocumentById).toHaveBeenCalledTimes(2); - expect(service.size).toBe(0); + expect(queueRepo.store).toEqual([]); }); it('startet keinen zweiten Durchlauf parallel (isProcessing-Guard)', async () => { let resolveFirst: (() => void) | undefined; + let signalStarted!: () => void; + const started = new Promise((res) => { + signalStarted = res; + }); processor.processDocumentById.mockImplementation( () => new Promise<{ processed: boolean }>((res) => { resolveFirst = () => res({ processed: true }); + signalStarted(); // erste Verarbeitung läuft und blockiert hier }), ); - service.enqueue(1); + await service.enqueue(1); const firstRun = service.processQueue(); // blockiert in processDocumentById - await service.processQueue(); // muss sofort zurückkehren + await started; // warten, bis der erste Lauf tatsächlich verarbeitet + await service.processQueue(); // muss sofort zurückkehren (isProcessing) expect(processor.processDocumentById).toHaveBeenCalledTimes(1); diff --git a/paperless-backend/src/webhook/webhook-queue.service.ts b/paperless-backend/src/webhook/webhook-queue.service.ts index 698a9cc..e7b5272 100644 --- a/paperless-backend/src/webhook/webhook-queue.service.ts +++ b/paperless-backend/src/webhook/webhook-queue.service.ts @@ -4,6 +4,7 @@ import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { PaperlessProcessorService } from '../paperless/paperless-processor.service'; import { Setting } from '../database/entities/setting.entity'; +import { WebhookQueueItem } from '../database/entities/webhook-queue-item.entity'; // Tag (Schlüssel) des Settings-Eintrags, der den letzten Webhook-Aufruf festhält. const LAST_WEBHOOK_CALL_TAG = 'last_webhook_call'; @@ -21,53 +22,64 @@ interface WebhookStatusInfo { /** * Entkoppelt den Webhook-Empfang von der Verarbeitung: Eingehende Dokument-IDs - * werden in eine deduplizierte Warteschlange gelegt und von einem separaten - * Intervall-Prozess **sequenziell** (ohne Überschneidung) abgearbeitet. So führt - * mehrfaches schnelles Speichern desselben Dokuments nicht zu parallelen Läufen. + * werden in eine **persistente**, deduplizierte Warteschlange (`webhook_queue`) + * gelegt und von einem separaten Intervall-Prozess **sequenziell** (ohne + * Überschneidung) abgearbeitet. So führt mehrfaches schnelles Speichern + * desselben Dokuments nicht zu parallelen Läufen. * * Verhalten: - * - Jede ID kommt nur **einmal** in der Liste vor (Dedup). - * - Beim Verarbeitungsstart wird die ID **sofort** aus der Liste entfernt; ein + * - Jede ID kommt nur **einmal** in der Warteschlange vor (Dedup via Primär- + * schlüssel `documentId` / `INSERT IGNORE`). + * - Beim Verarbeitungsstart wird die ID **sofort** aus der Tabelle entfernt; ein * erneutes Feuern während der Verarbeitung reiht sie wieder ein (ein weiterer * Lauf folgt danach). - * - Es läuft immer nur eine Verarbeitung gleichzeitig. + * - Es läuft immer nur eine Verarbeitung gleichzeitig (`isProcessing`-Guard). * - * Hinweis: Die Warteschlange liegt im Speicher; bei einem Neustart gehen noch - * nicht verarbeitete IDs verloren (Paperless müsste den Webhook erneut feuern). + * Da die Warteschlange in der Datenbank liegt, überstehen ausstehende IDs einen + * Neustart und werden nach dem Boot weiterverarbeitet. */ @Injectable() export class WebhookQueueService { private readonly logger = new Logger(WebhookQueueService.name); - private readonly pending = new Set(); private isProcessing = false; constructor( private readonly paperlessProcessor: PaperlessProcessorService, @InjectRepository(Setting) private readonly settingRepo: Repository, + @InjectRepository(WebhookQueueItem) + private readonly queueRepo: Repository, ) {} /** Anzahl der aktuell wartenden Dokument-IDs. */ - get size(): number { - return this.pending.size; + count(): Promise { + return this.queueRepo.count(); } /** * Reiht eine Dokument-ID zur Verarbeitung ein. Ist die ID bereits in der - * Warteschlange, wird sie nicht erneut hinzugefügt (Dedup). + * Warteschlange, wird sie nicht erneut hinzugefügt (Dedup über den Primär- + * schlüssel; `INSERT IGNORE` ist atomar und race-sicher). */ - enqueue(documentId: number, action?: string): void { - if (this.pending.has(documentId)) { + async enqueue(documentId: number, action?: string): Promise { + const result = await this.queueRepo + .createQueryBuilder() + .insert() + .into(WebhookQueueItem) + .values({ documentId, action: action ?? null }) + .orIgnore() + .execute(); + + const added = + ((result.raw as { affectedRows?: number })?.affectedRows ?? 0) > 0; + if (added) { + this.logger.log(`Dokument ${documentId} eingereiht.`); + } else { this.logger.log( `Dokument ${documentId} ist bereits in der Warteschlange – nicht erneut hinzugefügt.`, ); - } else { - this.pending.add(documentId); - this.logger.log( - `Dokument ${documentId} eingereiht (Warteschlange: ${this.pending.size}).`, - ); } - void this.recordStatus({ documentId, action, status: 'queued' }); + await this.recordStatus({ documentId, action, status: 'queued' }); } /** @@ -77,20 +89,20 @@ export class WebhookQueueService { */ @Interval(QUEUE_INTERVAL_MS) async processQueue(): Promise { - if (this.isProcessing || this.pending.size === 0) return; + if (this.isProcessing) return; this.isProcessing = true; try { // Solange Einträge vorhanden sind, einzeln und sequenziell abarbeiten. - while (this.pending.size > 0) { - // Nächste ID entnehmen und SOFORT (vor Verarbeitungsstart) entfernen. - let documentId: number | undefined; - for (const id of this.pending) { - documentId = id; - break; - } - if (documentId === undefined) break; - this.pending.delete(documentId); - await this.handleDocument(documentId); + for (;;) { + // Ältesten Eintrag (FIFO) holen ... + const [next] = await this.queueRepo.find({ + order: { createdAt: 'ASC' }, + take: 1, + }); + if (!next) break; + // ... und SOFORT (vor Verarbeitungsstart) aus der Tabelle entfernen. + await this.queueRepo.delete({ documentId: next.documentId }); + await this.handleDocument(next.documentId); } } finally { this.isProcessing = false; @@ -127,7 +139,7 @@ export class WebhookQueueService { lastCall = { raw: setting.Wert }; } } - return { lastCall, queueSize: this.pending.size }; + return { lastCall, queueSize: await this.queueRepo.count() }; } /** diff --git a/paperless-backend/src/webhook/webhook.controller.ts b/paperless-backend/src/webhook/webhook.controller.ts index 0de6e12..07eb00e 100644 --- a/paperless-backend/src/webhook/webhook.controller.ts +++ b/paperless-backend/src/webhook/webhook.controller.ts @@ -49,7 +49,7 @@ export class WebhookController { ); // Sofort einreihen und antworten; die eigentliche Verarbeitung übernimmt // der separate Queue-Prozess (sequenziell, dedupliziert). - this.webhookQueue.enqueue(documentId, payload.action); + await this.webhookQueue.enqueue(documentId, payload.action); return { status: 'queued', documentId }; } diff --git a/paperless-backend/src/webhook/webhook.module.ts b/paperless-backend/src/webhook/webhook.module.ts index 5ab970e..68af4c8 100644 --- a/paperless-backend/src/webhook/webhook.module.ts +++ b/paperless-backend/src/webhook/webhook.module.ts @@ -5,9 +5,14 @@ import { WebhookQueueService } from './webhook-queue.service'; import { PaperlessModule } from '../paperless/paperless.module'; import { AuthModule } from '../auth/auth.module'; import { Setting } from '../database/entities/setting.entity'; +import { WebhookQueueItem } from '../database/entities/webhook-queue-item.entity'; @Module({ - imports: [TypeOrmModule.forFeature([Setting]), PaperlessModule, AuthModule], + imports: [ + TypeOrmModule.forFeature([Setting, WebhookQueueItem]), + PaperlessModule, + AuthModule, + ], controllers: [WebhookController], providers: [WebhookQueueService], })