feat: Mindest-Wartezeit (5000 ms) zwischen Einreihen und Verarbeitung
Build and Push Multi-Platform Images / build-and-push (push) Successful in 33s

processQueue verarbeitet den ältesten Warteschlangen-Eintrag erst, wenn er
mindestens MIN_QUEUE_AGE_MS (Default 5000 ms, ENV WEBHOOK_QUEUE_MIN_AGE_MS)
in der Tabelle lag. Ist der älteste Eintrag (FIFO) noch zu jung, bricht der
Tick ab und prüft beim nächsten Intervall erneut. Unit-Test ergänzt.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-22 11:03:10 +02:00
parent beaa1be4a5
commit 1c70473cef
2 changed files with 57 additions and 14 deletions
@@ -13,7 +13,7 @@ import { WebhookQueueService } from './webhook-queue.service';
* FIFO-Array simuliert: `INSERT IGNORE` (Dedup), FIFO-`find` und `delete`. * FIFO-Array simuliert: `INSERT IGNORE` (Dedup), FIFO-`find` und `delete`.
*/ */
function createQueueRepoMock() { function createQueueRepoMock() {
const store: number[] = []; const store: { documentId: number; createdAt: Date }[] = [];
return { return {
store, store,
createQueryBuilder: jest.fn(() => ({ createQueryBuilder: jest.fn(() => ({
@@ -22,8 +22,15 @@ function createQueueRepoMock() {
values: (v: { documentId: number }) => ({ values: (v: { documentId: number }) => ({
orIgnore: () => ({ orIgnore: () => ({
execute: () => { execute: () => {
const added = !store.includes(v.documentId); const added = !store.some((s) => s.documentId === v.documentId);
if (added) store.push(v.documentId); if (added) {
// Standardmäßig "alt genug" (60s), damit die Mindest-Wartezeit
// die Verhaltenstests nicht blockiert.
store.push({
documentId: v.documentId,
createdAt: new Date(Date.now() - 60_000),
});
}
return Promise.resolve({ return Promise.resolve({
raw: { affectedRows: added ? 1 : 0 }, raw: { affectedRows: added ? 1 : 0 },
}); });
@@ -33,15 +40,24 @@ function createQueueRepoMock() {
}), }),
}), }),
})), })),
find: jest.fn(() => find: jest.fn(() => {
Promise.resolve( const sorted = [...store].sort(
store.length (a, b) => a.createdAt.getTime() - b.createdAt.getTime(),
? [{ documentId: store[0], action: null, createdAt: new Date() }] );
return Promise.resolve(
sorted.length
? [
{
documentId: sorted[0].documentId,
action: null,
createdAt: sorted[0].createdAt,
},
]
: [], : [],
), );
), }),
delete: jest.fn((criteria: { documentId: number }) => { delete: jest.fn((criteria: { documentId: number }) => {
const idx = store.indexOf(criteria.documentId); const idx = store.findIndex((s) => s.documentId === criteria.documentId);
if (idx >= 0) store.splice(idx, 1); if (idx >= 0) store.splice(idx, 1);
return Promise.resolve({ affected: 1 }); return Promise.resolve({ affected: 1 });
}), }),
@@ -75,7 +91,7 @@ describe('WebhookQueueService', () => {
it('reiht jede ID nur einmal ein (Dedup)', async () => { it('reiht jede ID nur einmal ein (Dedup)', async () => {
await service.enqueue(5); await service.enqueue(5);
await service.enqueue(5); await service.enqueue(5);
expect(queueRepo.store).toEqual([5]); expect(queueRepo.store.map((s) => s.documentId)).toEqual([5]);
expect(await service.count()).toBe(1); expect(await service.count()).toBe(1);
}); });
@@ -88,13 +104,15 @@ describe('WebhookQueueService', () => {
expect(processor.processDocumentById).toHaveBeenCalledWith(1); expect(processor.processDocumentById).toHaveBeenCalledWith(1);
expect(processor.processDocumentById).toHaveBeenCalledWith(2); expect(processor.processDocumentById).toHaveBeenCalledWith(2);
expect(processor.processDocumentById).toHaveBeenCalledTimes(2); expect(processor.processDocumentById).toHaveBeenCalledTimes(2);
expect(queueRepo.store).toEqual([]); expect(queueRepo.store).toHaveLength(0);
}); });
it('entfernt die ID vor Verarbeitungsbeginn aus der Tabelle', async () => { it('entfernt die ID vor Verarbeitungsbeginn aus der Tabelle', async () => {
let containedWhileProcessing = true; let containedWhileProcessing = true;
processor.processDocumentById.mockImplementation((id: number) => { processor.processDocumentById.mockImplementation((id: number) => {
containedWhileProcessing = queueRepo.store.includes(id); containedWhileProcessing = queueRepo.store.some(
(s) => s.documentId === id,
);
return Promise.resolve({ processed: true }); return Promise.resolve({ processed: true });
}); });
@@ -120,7 +138,7 @@ describe('WebhookQueueService', () => {
await service.processQueue(); await service.processQueue();
expect(processor.processDocumentById).toHaveBeenCalledTimes(2); expect(processor.processDocumentById).toHaveBeenCalledTimes(2);
expect(queueRepo.store).toEqual([]); expect(queueRepo.store).toHaveLength(0);
}); });
it('startet keinen zweiten Durchlauf parallel (isProcessing-Guard)', async () => { it('startet keinen zweiten Durchlauf parallel (isProcessing-Guard)', async () => {
@@ -147,4 +165,19 @@ describe('WebhookQueueService', () => {
resolveFirst?.(); resolveFirst?.();
await firstRun; await firstRun;
}); });
it('verarbeitet einen Eintrag erst nach der Mindest-Wartezeit', async () => {
// Frisch eingereihter Eintrag (Alter ~0) darf noch nicht verarbeitet werden.
queueRepo.store.push({ documentId: 99, createdAt: new Date() });
await service.processQueue();
expect(processor.processDocumentById).not.toHaveBeenCalled();
expect(queueRepo.store).toHaveLength(1);
// Nach Überschreiten der Mindest-Wartezeit (5000 ms) wird verarbeitet.
queueRepo.store[0].createdAt = new Date(Date.now() - 6000);
await service.processQueue();
expect(processor.processDocumentById).toHaveBeenCalledWith(99);
expect(queueRepo.store).toHaveLength(0);
});
}); });
@@ -12,6 +12,11 @@ const LAST_WEBHOOK_CALL_TAG = 'last_webhook_call';
// Prüfintervall der Warteschlange (sehr kurz). Über ENV überschreibbar. // Prüfintervall der Warteschlange (sehr kurz). Über ENV überschreibbar.
const QUEUE_INTERVAL_MS = Number(process.env.WEBHOOK_QUEUE_INTERVAL_MS) || 1000; const QUEUE_INTERVAL_MS = Number(process.env.WEBHOOK_QUEUE_INTERVAL_MS) || 1000;
// Mindest-Verweildauer zwischen Einreihen und Verarbeitung. Ein Eintrag wird
// frühestens verarbeitet, wenn er so lange in der Warteschlange lag. Über ENV
// überschreibbar (Default 5000 ms).
const MIN_QUEUE_AGE_MS = Number(process.env.WEBHOOK_QUEUE_MIN_AGE_MS) || 5000;
interface WebhookStatusInfo { interface WebhookStatusInfo {
documentId: number | null; documentId: number | null;
action?: string; action?: string;
@@ -34,6 +39,8 @@ interface WebhookStatusInfo {
* erneutes Feuern während der Verarbeitung reiht sie wieder ein (ein weiterer * erneutes Feuern während der Verarbeitung reiht sie wieder ein (ein weiterer
* Lauf folgt danach). * Lauf folgt danach).
* - Es läuft immer nur eine Verarbeitung gleichzeitig (`isProcessing`-Guard). * - Es läuft immer nur eine Verarbeitung gleichzeitig (`isProcessing`-Guard).
* - Zwischen Einreihen und Verarbeitung liegen mindestens `MIN_QUEUE_AGE_MS`
* (Default 5000 ms); jüngere Einträge warten bis zum nächsten Tick.
* *
* Da die Warteschlange in der Datenbank liegt, überstehen ausstehende IDs einen * Da die Warteschlange in der Datenbank liegt, überstehen ausstehende IDs einen
* Neustart und werden nach dem Boot weiterverarbeitet. * Neustart und werden nach dem Boot weiterverarbeitet.
@@ -100,6 +107,9 @@ export class WebhookQueueService {
take: 1, take: 1,
}); });
if (!next) break; if (!next) break;
// Mindest-Wartezeit einhalten: ist der älteste Eintrag noch zu jung,
// sind es (FIFO) alle → diesen Tick beenden, beim nächsten erneut prüfen.
if (Date.now() - next.createdAt.getTime() < MIN_QUEUE_AGE_MS) break;
// ... und SOFORT (vor Verarbeitungsstart) aus der Tabelle entfernen. // ... und SOFORT (vor Verarbeitungsstart) aus der Tabelle entfernen.
await this.queueRepo.delete({ documentId: next.documentId }); await this.queueRepo.delete({ documentId: next.documentId });
await this.handleDocument(next.documentId); await this.handleDocument(next.documentId);