From 66a2cccd20558c6e10d4643bd5af9866249370ca Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Bj=C3=B6rn=20P=C3=B6ttker?= Date: Mon, 29 Jun 2026 17:09:23 +0200 Subject: [PATCH] =?UTF-8?q?feat:=20Webhook-Warteschlange=20f=C3=BCr=20sequ?= =?UTF-8?q?enzielle,=20deduplizierte=20Verarbeitung?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Der Webhook reiht eingehende Dokument-IDs nur noch in eine Warteschlange ein und antwortet sofort (status "queued"). Ein separater Intervall-Prozess (WebhookQueueService) arbeitet die IDs nacheinander ab – ohne Überschneidung, falls ein Dokument mehrfach kurz hintereinander gespeichert wird: - Jede ID kommt nur einmal in der Liste vor (Dedup). - Beim Verarbeitungsstart wird die ID sofort entfernt; ein erneutes Feuern während der Verarbeitung reiht sie wieder ein (ein weiterer Lauf folgt). - Es läuft immer nur eine Verarbeitung gleichzeitig (isProcessing-Guard). - Prüfintervall sehr kurz (WEBHOOK_QUEUE_INTERVAL_MS, Default 1000 ms). Status-Recording (last_webhook_call) + GET /api/webhook/status wandern in den Queue-Service und liefern zusätzlich die aktuelle Warteschlangen-Größe. Frontend-Tab zeigt Status "In Warteschlange" und die Queue-Größe an. Unit-Test deckt Dedup, Remove-on-start, Re-Enqueue und No-Parallel ab. Co-Authored-By: Claude Opus 4.8 --- .../src/webhook/webhook-queue.service.spec.ts | 100 +++++++++++ .../src/webhook/webhook-queue.service.ts | 168 ++++++++++++++++++ .../src/webhook/webhook.controller.ts | 111 ++---------- .../src/webhook/webhook.module.ts | 2 + paperless-frontend/src/api/webhook.ts | 3 +- paperless-frontend/src/pages/SettingsPage.tsx | 21 ++- 6 files changed, 299 insertions(+), 106 deletions(-) create mode 100644 paperless-backend/src/webhook/webhook-queue.service.spec.ts create mode 100644 paperless-backend/src/webhook/webhook-queue.service.ts diff --git a/paperless-backend/src/webhook/webhook-queue.service.spec.ts b/paperless-backend/src/webhook/webhook-queue.service.spec.ts new file mode 100644 index 0000000..a92569c --- /dev/null +++ b/paperless-backend/src/webhook/webhook-queue.service.spec.ts @@ -0,0 +1,100 @@ +// PaperlessProcessorService zieht über die Postprocessing-Kette das ESM-Paket +// "webdav" nach, das Jest nicht transformiert. Für diesen Unit-Test ersetzen wir +// das Modul durch eine Dummy-Klasse – der Service erhält seine Abhängigkeiten +// ohnehin als Mocks injiziert. +jest.mock('../paperless/paperless-processor.service', () => ({ + PaperlessProcessorService: class {}, +})); + +import { WebhookQueueService } from './webhook-queue.service'; + +describe('WebhookQueueService', () => { + let processor: { processDocumentById: jest.Mock }; + let settingRepo: { + findOneBy: jest.Mock; + create: jest.Mock; + save: jest.Mock; + }; + let service: WebhookQueueService; + + beforeEach(() => { + processor = { + processDocumentById: jest.fn().mockResolvedValue({ processed: true }), + }; + settingRepo = { + findOneBy: jest.fn().mockResolvedValue(null), + create: jest.fn((x: unknown) => x), + save: jest.fn().mockResolvedValue(undefined), + }; + service = new WebhookQueueService(processor as never, settingRepo as never); + }); + + it('reiht jede ID nur einmal ein (Dedup)', () => { + service.enqueue(5); + service.enqueue(5); + expect(service.size).toBe(1); + }); + + it('verarbeitet alle eingereihten IDs und leert die Warteschlange', async () => { + service.enqueue(1); + 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); + }); + + it('entfernt die ID vor Verarbeitungsbeginn aus der Liste', async () => { + let sizeWhileProcessing = -1; + processor.processDocumentById.mockImplementation(() => { + sizeWhileProcessing = service.size; + return Promise.resolve({ processed: true }); + }); + + service.enqueue(42); + await service.processQueue(); + + // Beim Verarbeitungsstart war die einzige ID bereits entfernt. + expect(sizeWhileProcessing).toBe(0); + }); + + it('arbeitet eine während der Verarbeitung erneut eingereihte ID erneut ab', async () => { + let firstRun = true; + processor.processDocumentById.mockImplementation((id: number) => { + // Beim ersten Lauf feuert der Webhook erneut, während verarbeitet wird. + if (firstRun && id === 7) { + firstRun = false; + service.enqueue(7); + } + return Promise.resolve({ processed: true }); + }); + + service.enqueue(7); + await service.processQueue(); + + expect(processor.processDocumentById).toHaveBeenCalledTimes(2); + expect(service.size).toBe(0); + }); + + it('startet keinen zweiten Durchlauf parallel (isProcessing-Guard)', async () => { + let resolveFirst: (() => void) | undefined; + processor.processDocumentById.mockImplementation( + () => + new Promise<{ processed: boolean }>((res) => { + resolveFirst = () => res({ processed: true }); + }), + ); + + service.enqueue(1); + const firstRun = service.processQueue(); // blockiert in processDocumentById + await service.processQueue(); // muss sofort zurückkehren + + expect(processor.processDocumentById).toHaveBeenCalledTimes(1); + + resolveFirst?.(); + await firstRun; + }); +}); diff --git a/paperless-backend/src/webhook/webhook-queue.service.ts b/paperless-backend/src/webhook/webhook-queue.service.ts new file mode 100644 index 0000000..698a9cc --- /dev/null +++ b/paperless-backend/src/webhook/webhook-queue.service.ts @@ -0,0 +1,168 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { Interval } from '@nestjs/schedule'; +import { InjectRepository } from '@nestjs/typeorm'; +import { Repository } from 'typeorm'; +import { PaperlessProcessorService } from '../paperless/paperless-processor.service'; +import { Setting } from '../database/entities/setting.entity'; + +// Tag (Schlüssel) des Settings-Eintrags, der den letzten Webhook-Aufruf festhält. +const LAST_WEBHOOK_CALL_TAG = 'last_webhook_call'; + +// Prüfintervall der Warteschlange (sehr kurz). Über ENV überschreibbar. +const QUEUE_INTERVAL_MS = Number(process.env.WEBHOOK_QUEUE_INTERVAL_MS) || 1000; + +interface WebhookStatusInfo { + documentId: number | null; + action?: string; + status: 'queued' | 'processed' | 'error' | 'bad-request'; + reason?: string; + message?: string; +} + +/** + * 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. + * + * Verhalten: + * - Jede ID kommt nur **einmal** in der Liste vor (Dedup). + * - Beim Verarbeitungsstart wird die ID **sofort** aus der Liste 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. + * + * Hinweis: Die Warteschlange liegt im Speicher; bei einem Neustart gehen noch + * nicht verarbeitete IDs verloren (Paperless müsste den Webhook erneut feuern). + */ +@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, + ) {} + + /** Anzahl der aktuell wartenden Dokument-IDs. */ + get size(): number { + return this.pending.size; + } + + /** + * Reiht eine Dokument-ID zur Verarbeitung ein. Ist die ID bereits in der + * Warteschlange, wird sie nicht erneut hinzugefügt (Dedup). + */ + enqueue(documentId: number, action?: string): void { + if (this.pending.has(documentId)) { + 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' }); + } + + /** + * Prüft in kurzen Abständen die Warteschlange und arbeitet vorhandene IDs + * nacheinander ab. Ein bereits laufender Durchlauf wird nicht doppelt + * gestartet (kein paralleles Verarbeiten). + */ + @Interval(QUEUE_INTERVAL_MS) + async processQueue(): Promise { + if (this.isProcessing || this.pending.size === 0) 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); + } + } finally { + this.isProcessing = false; + } + } + + private async handleDocument(documentId: number): Promise { + try { + await this.paperlessProcessor.processDocumentById(documentId); + this.logger.log(`Dokument ${documentId} aus Warteschlange verarbeitet.`); + await this.recordStatus({ documentId, status: 'processed' }); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + this.logger.error( + `Fehler bei der Verarbeitung von Dokument ${documentId}: ${message}`, + ); + await this.recordStatus({ documentId, status: 'error', message }); + } + } + + /** + * Liefert den letzten festgehaltenen Webhook-Status sowie die aktuelle + * Warteschlangen-Größe. + */ + async getStatus(): Promise<{ lastCall: unknown; queueSize: number }> { + const setting = await this.settingRepo.findOneBy({ + Tag: LAST_WEBHOOK_CALL_TAG, + }); + let lastCall: unknown = null; + if (setting?.Wert) { + try { + lastCall = JSON.parse(setting.Wert); + } catch { + lastCall = { raw: setting.Wert }; + } + } + return { lastCall, queueSize: this.pending.size }; + } + + /** + * Hält den letzten Webhook-Vorgang in der Settings-Tabelle fest + * (Tag "last_webhook_call"). Fehler beim Speichern werden nur geloggt. + */ + async recordStatus(info: WebhookStatusInfo): Promise { + try { + const value = JSON.stringify({ + at: new Date().toISOString(), + documentId: info.documentId, + action: info.action ?? null, + status: info.status, + ...(info.reason ? { reason: info.reason } : {}), + ...(info.message ? { message: info.message.slice(0, 100) } : {}), + }); + + let setting = await this.settingRepo.findOneBy({ + Tag: LAST_WEBHOOK_CALL_TAG, + }); + if (!setting) { + setting = this.settingRepo.create({ + Typ: 0, + Tag: LAST_WEBHOOK_CALL_TAG, + Wert: value, + }); + } else { + setting.Wert = value; + } + await this.settingRepo.save(setting); + } catch (err) { + this.logger.error( + 'Konnte Webhook-Status nicht in den Settings speichern', + err instanceof Error ? err.stack : String(err), + ); + } + } +} diff --git a/paperless-backend/src/webhook/webhook.controller.ts b/paperless-backend/src/webhook/webhook.controller.ts index 4630a2e..0de6e12 100644 --- a/paperless-backend/src/webhook/webhook.controller.ts +++ b/paperless-backend/src/webhook/webhook.controller.ts @@ -9,11 +9,8 @@ import { UseGuards, BadRequestException, } from '@nestjs/common'; -import { InjectRepository } from '@nestjs/typeorm'; -import { Repository } from 'typeorm'; import { ApiKeyGuard } from '../auth/api-key.guard'; -import { PaperlessProcessorService } from '../paperless/paperless-processor.service'; -import { Setting } from '../database/entities/setting.entity'; +import { WebhookQueueService } from './webhook-queue.service'; export interface PaperlessWebhookPayload { doc_url?: string; @@ -22,18 +19,11 @@ export interface PaperlessWebhookPayload { [key: string]: any; } -// Tag (Schlüssel) des Settings-Eintrags, der den letzten Webhook-Aufruf festhält. -const LAST_WEBHOOK_CALL_TAG = 'last_webhook_call'; - @Controller('api/webhook') export class WebhookController { private readonly logger = new Logger(WebhookController.name); - constructor( - private readonly paperlessProcessor: PaperlessProcessorService, - @InjectRepository(Setting) - private readonly settingRepo: Repository, - ) {} + constructor(private readonly webhookQueue: WebhookQueueService) {} @UseGuards(ApiKeyGuard) @Post('paperless') @@ -44,7 +34,7 @@ export class WebhookController { this.logger.warn( `Webhook ohne ermittelbare Dokument-ID: ${JSON.stringify(payload)}`, ); - await this.recordWebhookCall({ + await this.webhookQueue.recordStatus({ documentId: null, action: payload.action, status: 'bad-request', @@ -55,54 +45,22 @@ export class WebhookController { } this.logger.log( - `Webhook: action=${payload.action}, document=${documentId}`, + `Webhook: action=${payload.action}, document=${documentId} → eingereiht`, ); - - try { - const result = - await this.paperlessProcessor.processDocumentById(documentId); - const status = result.processed ? 'processed' : 'skipped'; - await this.recordWebhookCall({ - documentId, - action: payload.action, - status, - reason: result.reason, - }); - return { status, reason: result.reason }; - } catch (err) { - // Fehler tolerieren (wie der bisherige Cron) und 200 zurückgeben, - // damit Paperless keine Retry-Schleife startet. - const message = err instanceof Error ? err.message : String(err); - this.logger.error( - `Fehler bei Webhook-Verarbeitung von Dokument ${documentId}: ${message}`, - ); - await this.recordWebhookCall({ - documentId, - action: payload.action, - status: 'error', - message, - }); - return { status: 'error', message }; - } + // Sofort einreihen und antworten; die eigentliche Verarbeitung übernimmt + // der separate Queue-Prozess (sequenziell, dedupliziert). + this.webhookQueue.enqueue(documentId, payload.action); + return { status: 'queued', documentId }; } /** - * Liefert Informationen zum letzten Webhook-Aufruf (Zeitpunkt, Dokument, - * Ergebnis). Über die globalen Guards per JWT oder API-Key zugänglich. + * Liefert Informationen zum letzten Webhook-Vorgang (Zeitpunkt, Dokument, + * Ergebnis) und die aktuelle Warteschlangen-Größe. Über die globalen Guards + * per JWT oder API-Key zugänglich. */ @Get('status') - async getWebhookStatus(): Promise<{ lastCall: unknown }> { - const setting = await this.settingRepo.findOneBy({ - Tag: LAST_WEBHOOK_CALL_TAG, - }); - if (!setting?.Wert) { - return { lastCall: null }; - } - try { - return { lastCall: JSON.parse(setting.Wert) }; - } catch { - return { lastCall: { raw: setting.Wert } }; - } + async getWebhookStatus(): Promise<{ lastCall: unknown; queueSize: number }> { + return this.webhookQueue.getStatus(); } /** @@ -124,47 +82,4 @@ export class WebhookController { } return null; } - - /** - * Hält den letzten Webhook-Aufruf in der Settings-Tabelle fest - * (Tag "last_webhook_call"). Fehler beim Speichern werden nur geloggt und - * beeinflussen die Webhook-Antwort nicht. - */ - private async recordWebhookCall(info: { - documentId: number | null; - action?: string; - status: string; - reason?: string; - message?: string; - }): Promise { - try { - const value = JSON.stringify({ - at: new Date().toISOString(), - documentId: info.documentId, - action: info.action ?? null, - status: info.status, - ...(info.reason ? { reason: info.reason } : {}), - ...(info.message ? { message: info.message.slice(0, 100) } : {}), - }); - - let setting = await this.settingRepo.findOneBy({ - Tag: LAST_WEBHOOK_CALL_TAG, - }); - if (!setting) { - setting = this.settingRepo.create({ - Typ: 0, - Tag: LAST_WEBHOOK_CALL_TAG, - Wert: value, - }); - } else { - setting.Wert = value; - } - await this.settingRepo.save(setting); - } catch (err) { - this.logger.error( - 'Konnte letzten Webhook-Aufruf nicht in den Settings speichern', - err instanceof Error ? err.stack : String(err), - ); - } - } } diff --git a/paperless-backend/src/webhook/webhook.module.ts b/paperless-backend/src/webhook/webhook.module.ts index 5a280d7..5ab970e 100644 --- a/paperless-backend/src/webhook/webhook.module.ts +++ b/paperless-backend/src/webhook/webhook.module.ts @@ -1,6 +1,7 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { WebhookController } from './webhook.controller'; +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'; @@ -8,5 +9,6 @@ import { Setting } from '../database/entities/setting.entity'; @Module({ imports: [TypeOrmModule.forFeature([Setting]), PaperlessModule, AuthModule], controllers: [WebhookController], + providers: [WebhookQueueService], }) export class WebhookModule {} diff --git a/paperless-frontend/src/api/webhook.ts b/paperless-frontend/src/api/webhook.ts index d3d971f..db45565 100644 --- a/paperless-frontend/src/api/webhook.ts +++ b/paperless-frontend/src/api/webhook.ts @@ -4,13 +4,14 @@ export interface WebhookLastCall { at: string; // ISO-Zeitstempel documentId: number | null; action: string | null; - status: 'processed' | 'skipped' | 'error' | 'bad-request' | string; + status: 'queued' | 'processed' | 'error' | 'bad-request' | string; reason?: string; message?: string; } export interface WebhookStatus { lastCall: WebhookLastCall | null; + queueSize?: number; } export const webhookApi = { diff --git a/paperless-frontend/src/pages/SettingsPage.tsx b/paperless-frontend/src/pages/SettingsPage.tsx index d26ab52..4ac127c 100644 --- a/paperless-frontend/src/pages/SettingsPage.tsx +++ b/paperless-frontend/src/pages/SettingsPage.tsx @@ -2774,8 +2774,8 @@ function WebhookStatusTab() { const renderStatusBadge = (s: string) => { switch (s) { + case 'queued': return ; case 'processed': return ; - case 'skipped': return ; case 'error': return ; case 'bad-request': return ; default: return ; @@ -2783,6 +2783,7 @@ function WebhookStatusTab() { }; const lastCall = status?.lastCall ?? null; + const queueSize = status?.queueSize ?? 0; return ( <> @@ -2791,13 +2792,19 @@ function WebhookStatusTab() { Paperless-NGX ruft nach dem Bearbeiten eines Dokuments den Webhook{' '} /api/webhook/paperless auf - (per API-Key authentifiziert). Hier siehst du den zuletzt - verarbeiteten Aufruf. Eine vollständige Historie steht in den - Backend-Logs. + (per API-Key authentifiziert). Die ID wird in eine Warteschlange + gelegt und von einem separaten Prozess nacheinander verarbeitet. + Hier siehst du den zuletzt festgehaltenen Vorgang. Eine vollständige + Historie steht in den Backend-Logs. - + + + 0 ? 'processing' : 'default'}> + Warteschlange: {queueSize} + +