Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 36 additions & 5 deletions src/invoiceBatchProcessor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,14 @@ function isRateLimitError(error: unknown): boolean {
return error instanceof Error && /429|rate.?limit|too many requests/i.test(error.message);
}

function isValidPauseMs(value: number): boolean {
return Number.isInteger(value) && value >= 0;
}

function retryAfterMs(error: unknown, fallbackMs: number): number {
const withRetryAfter = error as { retryAfterMs?: number } | undefined;
return typeof withRetryAfter?.retryAfterMs === "number" ? withRetryAfter.retryAfterMs : fallbackMs;
const candidate = withRetryAfter?.retryAfterMs;
return typeof candidate === "number" && isValidPauseMs(candidate) ? candidate : fallbackMs;
}

export class InvoiceBatchProcessor {
Expand All @@ -76,14 +81,34 @@ export class InvoiceBatchProcessor {
const maxConcurrent = config.maxConcurrent ?? DEFAULT_MAX_CONCURRENT;
const rateLimitPauseMs = config.rateLimitPauseMs ?? DEFAULT_RATE_LIMIT_PAUSE_MS;

// Issue #1010 — reject invalid maxConcurrent before any dispatch.
if (!Number.isInteger(maxConcurrent) || maxConcurrent < 1) {
throw new RangeError("maxConcurrent must be a positive integer");
}
// Issue #1009 — reject invalid rateLimitPauseMs before any dispatch.
if (!isValidPauseMs(rateLimitPauseMs)) {
throw new RangeError("rateLimitPauseMs must be a non-negative integer");
}

let cursor = 0;
let pausedUntil = 0;
let slotSeq = 0;
const inFlight = new Map<number, Promise<{ slot: number; result: BatchInvoiceResult }>>();

// Issue #1009 — monotonic pause deadline that workers recheck after each
// sleep so a shorter hint from another worker cannot shorten the shared
// pause, and a worker that falls asleep before a longer hint arrives
// does not dispatch immediately on wake.
const waitForGlobalPause = async (): Promise<void> => {
while (true) {
const wait = pausedUntil - Date.now();
if (wait <= 0) return;
await sleep(wait);
}
};

const runOne = async (invoiceId: string): Promise<BatchInvoiceResult> => {
const wait = pausedUntil - Date.now();
if (wait > 0) await sleep(wait);
await waitForGlobalPause();

const amount = config.amounts[invoiceId];
if (amount === undefined) {
Expand All @@ -103,7 +128,10 @@ export class InvoiceBatchProcessor {
return result;
} catch (error) {
if (isRateLimitError(error)) {
pausedUntil = Date.now() + retryAfterMs(error, rateLimitPauseMs);
// Issue #1009 — take the max so the pause deadline only ever moves
// forward, never backward.
const nextPausedUntil = Date.now() + retryAfterMs(error, rateLimitPauseMs);
pausedUntil = Math.max(pausedUntil, nextPausedUntil);
}
const result: BatchInvoiceResult = {
invoiceId,
Expand All @@ -125,7 +153,10 @@ export class InvoiceBatchProcessor {
);
};

for (let i = 0; i < maxConcurrent; i++) launch();
// Issue #1010 — cap initial launches at invoice count so maxConcurrent
// values larger than the batch size don't stall or loop unnecessarily.
const initialLaunches = Math.min(maxConcurrent, invoiceIds.length);
for (let i = 0; i < initialLaunches; i++) launch();

while (inFlight.size > 0) {
const { slot, result } = await Promise.race(inFlight.values());
Expand Down