fix: Warteschlange blieb stehen – Altersvergleich in die DB verlagert
Build and Push Multi-Platform Images / build-and-push (push) Successful in 32s

Der Mindest-Wartezeit-Check verglich Date.now() (Node-Uhr) mit dem aus der DB
gelesenen createdAt. Ohne gesetzte timezone (mysql2-Default 'local') wird
createdAt bei abweichender DB-Zeitzone als zukünftig interpretiert → Alter
negativ → jeder Eintrag galt als "zu jung" und wurde nie verarbeitet.

Der Vergleich läuft jetzt in der DB (createdAt <= NOW(6) - INTERVAL 5s via
TypeORM Raw), unabhängig von Zeitzone/Uhr des Node-Prozesses. Bereits
wartende Einträge werden dadurch nach dem Deploy sofort abgearbeitet.
Test-Mock bildet den DB-Alters-Filter nach.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-22 12:07:51 +02:00
parent 7b2a79be2a
commit 095cc4bb02
2 changed files with 22 additions and 14 deletions
@@ -41,16 +41,17 @@ function createQueueRepoMock() {
}), }),
})), })),
find: jest.fn(() => { find: jest.fn(() => {
const sorted = [...store].sort( // Simuliert den DB-seitigen Alters-Filter (createdAt <= NOW() - 5000 ms).
(a, b) => a.createdAt.getTime() - b.createdAt.getTime(), const eligible = [...store]
); .filter((s) => Date.now() - s.createdAt.getTime() >= 5000)
.sort((a, b) => a.createdAt.getTime() - b.createdAt.getTime());
return Promise.resolve( return Promise.resolve(
sorted.length eligible.length
? [ ? [
{ {
documentId: sorted[0].documentId, documentId: eligible[0].documentId,
action: null, action: null,
createdAt: sorted[0].createdAt, createdAt: eligible[0].createdAt,
}, },
] ]
: [], : [],
@@ -1,7 +1,7 @@
import { Injectable, Logger } from '@nestjs/common'; import { Injectable, Logger } from '@nestjs/common';
import { Interval } from '@nestjs/schedule'; import { Interval } from '@nestjs/schedule';
import { InjectRepository } from '@nestjs/typeorm'; import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm'; import { Raw, Repository } from 'typeorm';
import { PaperlessProcessorService } from '../paperless/paperless-processor.service'; import { PaperlessProcessorService } from '../paperless/paperless-processor.service';
import { Setting } from '../database/entities/setting.entity'; import { Setting } from '../database/entities/setting.entity';
import { WebhookQueueItem } from '../database/entities/webhook-queue-item.entity'; import { WebhookQueueItem } from '../database/entities/webhook-queue-item.entity';
@@ -99,18 +99,25 @@ export class WebhookQueueService {
if (this.isProcessing) return; if (this.isProcessing) return;
this.isProcessing = true; this.isProcessing = true;
try { try {
// Solange Einträge vorhanden sind, einzeln und sequenziell abarbeiten. // Solange (alte genug) Einträge vorhanden sind, sequenziell abarbeiten.
for (;;) { for (;;) {
// Ältesten Eintrag (FIFO) holen ... // Ältesten Eintrag (FIFO) holen, der die Mindest-Wartezeit erfüllt.
// Der Altersvergleich läuft bewusst in der DB (NOW() vs. createdAt) und
// ist damit unabhängig von Zeitzone/Uhr des Node-Prozesses ein Node/DB-
// Zeitversatz würde sonst den Vergleich verfälschen (Einträge nie "alt
// genug"). Nur ein numerisches Konstantenliteral wird interpoliert.
const [next] = await this.queueRepo.find({ const [next] = await this.queueRepo.find({
where: {
createdAt: Raw(
(alias) =>
`${alias} <= (NOW(6) - INTERVAL ${MIN_QUEUE_AGE_MS * 1000} MICROSECOND)`,
),
},
order: { createdAt: 'ASC' }, order: { createdAt: 'ASC' },
take: 1, take: 1,
}); });
if (!next) break; if (!next) break; // Warteschlange leer oder noch nichts alt genug
// Mindest-Wartezeit einhalten: ist der älteste Eintrag noch zu jung, // SOFORT (vor Verarbeitungsstart) aus der Tabelle entfernen.
// 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.
await this.queueRepo.delete({ documentId: next.documentId }); await this.queueRepo.delete({ documentId: next.documentId });
await this.handleDocument(next.documentId); await this.handleDocument(next.documentId);
} }