feat: Webhook-Warteschlange in Datenbank persistieren
Build and Push Multi-Platform Images / build-and-push (push) Successful in 31s

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 <noreply@anthropic.com>
This commit is contained in:
2026-06-29 17:54:16 +02:00
parent 66a2cccd20
commit 9718d6888a
8 changed files with 182 additions and 61 deletions
@@ -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';
@@ -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';
@@ -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;
}
@@ -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<void> {
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<void> {
await queryRunner.query(`DROP TABLE \`webhook_queue\``);
}
}
@@ -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<typeof createQueueRepoMock>;
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<void>((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);
@@ -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<number>();
private isProcessing = false;
constructor(
private readonly paperlessProcessor: PaperlessProcessorService,
@InjectRepository(Setting)
private readonly settingRepo: Repository<Setting>,
@InjectRepository(WebhookQueueItem)
private readonly queueRepo: Repository<WebhookQueueItem>,
) {}
/** Anzahl der aktuell wartenden Dokument-IDs. */
get size(): number {
return this.pending.size;
count(): Promise<number> {
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<void> {
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<void> {
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() };
}
/**
@@ -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 };
}
@@ -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],
})