Skip to content

Commit 4ebc348

Browse files
authored
Merge pull request #297 from lorenzo-romano/feat/issue-165-webhook-lag
feat(webhooks): add dead-letter status and lag metrics to the delivery queue (#165)
2 parents a406074 + 08a2aba commit 4ebc348

4 files changed

Lines changed: 135 additions & 7 deletions

File tree

apps/web/src/app/api/webhooks/deliver/route.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { NextResponse } from 'next/server';
22
import { withClient, ensureSchema } from '@/lib/db';
3-
import { deliverDue } from '@/lib/webhooks';
3+
import { deliverDue, pendingDue } from '@/lib/webhooks';
44

55
export const dynamic = 'force-dynamic';
66
export const maxDuration = 30;
@@ -23,7 +23,11 @@ export async function GET(request: Request) {
2323
try {
2424
const result = await withClient(async (client) => {
2525
await ensureSchema(client);
26-
return deliverDue(client);
26+
const outcome = await deliverDue(client);
27+
// Remaining lag after this run: the signal a scheduler uses to scale
28+
// consumer frequency to backlog (#165).
29+
const lag = await pendingDue(client);
30+
return { ...outcome, lag };
2731
});
2832
return NextResponse.json({ success: true, ...result });
2933
} catch (error: unknown) {

apps/web/src/app/api/webhooks/route.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,9 @@ export async function GET() {
1414
configured: false,
1515
pending: 0,
1616
failed: 0,
17+
deadLetter: 0,
1718
delivered: 0,
19+
lag: 0,
1820
recentFailed: [],
1921
});
2022
}

apps/web/src/lib/webhooks.test.ts

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,8 @@ import {
1010
deliverDue,
1111
enqueueWebhookDelivery,
1212
payloadFromRow,
13+
pendingDue,
14+
webhookSummary,
1315
} from './webhooks';
1416

1517
describe('shouldRetry', () => {
@@ -176,3 +178,78 @@ describe('deliverDue — a sleeping host cannot stall the caller past the budget
176178
expect(elapsed).toBeLessThan(5_000);
177179
});
178180
});
181+
182+
describe('pendingDue — the lag signal a consumer fleet scales on (#165)', () => {
183+
it('counts only pending rows whose retry time has passed (or never set)', async () => {
184+
const queries: string[] = [];
185+
const query = vi.fn(async (sql: string) => {
186+
queries.push(sql);
187+
return { rows: [{ count: '7' }] };
188+
});
189+
190+
const lag = await pendingDue({ query } as never, { now: new Date('2026-08-01T00:00:00Z') });
191+
192+
expect(lag).toBe(7);
193+
const sql = queries[0];
194+
expect(sql).toContain("status = 'pending'");
195+
// Null next_retry_at (never attempted) is always due.
196+
expect(sql).toContain('next_retry_at IS NULL OR next_retry_at <= $1::timestamptz');
197+
});
198+
});
199+
200+
describe('webhookSummary', () => {
201+
it('reports lag and the dead-letter count alongside the existing tallies', async () => {
202+
const query = vi.fn(async (sql: string) => {
203+
if (/^SELECT status, count/.test(sql)) {
204+
return {
205+
rows: [
206+
{ status: 'pending', n: '2' },
207+
{ status: 'delivered', n: '5' },
208+
{ status: 'dead_letter', n: '3' },
209+
],
210+
};
211+
}
212+
if (/^SELECT count\(\*\)::text AS count/.test(sql)) {
213+
return { rows: [{ count: '2' }] };
214+
}
215+
return { rows: [] };
216+
});
217+
218+
const summary = await webhookSummary({ query } as never);
219+
220+
expect(summary.pending).toBe(2);
221+
expect(summary.delivered).toBe(5);
222+
expect(summary.deadLetter).toBe(3);
223+
expect(summary.lag).toBe(2);
224+
expect(summary.recentFailed).toEqual([]);
225+
});
226+
227+
it('lists dead-lettered deliveries in recentFailed for operator inspection', async () => {
228+
const query = vi.fn(async (sql: string) => {
229+
if (/^SELECT status, count/.test(sql)) return { rows: [{ status: 'dead_letter', n: '1' }] };
230+
if (/^SELECT count\(\*\)::text AS count/.test(sql)) return { rows: [{ count: '0' }] };
231+
if (/^SELECT id, payment_tx_hash/.test(sql)) {
232+
return {
233+
rows: [
234+
{
235+
id: '9',
236+
payment_tx_hash: 'a'.repeat(64),
237+
status: 'dead_letter',
238+
attempts: 8,
239+
last_status_code: 503,
240+
last_error: 'HTTP 503',
241+
updated_at: new Date('2026-08-01T00:00:00.000Z'),
242+
},
243+
],
244+
};
245+
}
246+
return { rows: [] };
247+
});
248+
249+
const summary = await webhookSummary({ query } as never);
250+
251+
expect(summary.deadLetter).toBe(1);
252+
expect(summary.recentFailed[0].status).toBe('dead_letter');
253+
expect(summary.recentFailed[0].attempts).toBe(8);
254+
});
255+
});

apps/web/src/lib/webhooks.ts

Lines changed: 50 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,12 @@ export const DELIVERY_WINDOW_MS = 24 * 60 * 60 * 1000;
1515
export const ATTEMPT_TIMEOUT_MS = 2_000;
1616
export const MAX_BACKOFF_MS = 60 * 60 * 1000;
1717

18-
export type DeliveryStatus = 'pending' | 'delivering' | 'delivered' | 'failed';
18+
export type DeliveryStatus =
19+
| 'pending'
20+
| 'delivering'
21+
| 'delivered'
22+
| 'failed'
23+
| 'dead_letter';
1924

2025
export interface PaymentPayload {
2126
tx_hash: string;
@@ -274,7 +279,9 @@ export async function deliverDue(
274279
transportError,
275280
});
276281
if (terminal.status === 'delivered') delivered++;
277-
else if (terminal.status === 'failed') failed++;
282+
// A dead-lettered row is terminal, not a retry — count it as failed here
283+
// so the run's tallies reflect deliveries that gave up (#165).
284+
else if (terminal.status === 'failed' || terminal.status === 'dead_letter') failed++;
278285
else retried++;
279286
}
280287

@@ -308,7 +315,11 @@ async function recordAttempt(
308315
let status: DeliveryStatus;
309316
if (ok) status = 'delivered';
310317
else if (next) status = 'pending';
311-
else status = 'failed';
318+
// A delivery that exhausts its attempt budget (or its 24h delivery window)
319+
// is dead-lettered: it is kept for operator inspection and never retried
320+
// again (#165). This is the queue's explicit dead-letter state, distinct
321+
// from a transient 'failed' row.
322+
else status = 'dead_letter';
312323

313324
await client.query(
314325
`INSERT INTO webhook_attempts (delivery_id, attempt_number, status_code, error)
@@ -332,10 +343,37 @@ async function recordAttempt(
332343
return { id: input.id, status, statusCode: input.statusCode, error: input.error };
333344
}
334345

346+
/**
347+
* The queue's lag: deliveries that are due for (re)delivery right now.
348+
*
349+
* This is the signal a consumer fleet auto-scales on (#165) — when lag stays
350+
* high, schedule more frequent or overlapping `/api/webhooks/deliver` runs;
351+
* when it is zero, the queue is drained. Rows with a null `next_retry_at`
352+
* (never attempted) are always due; the rest are due once their retry time
353+
* has passed.
354+
*/
355+
export async function pendingDue(
356+
client: Client,
357+
opts: { now?: Date } = {},
358+
): Promise<number> {
359+
const now = (opts.now ?? new Date()).toISOString();
360+
const res = await client.query<{ count: string }>(
361+
`SELECT count(*)::text AS count
362+
FROM webhook_deliveries
363+
WHERE status = 'pending'
364+
AND (next_retry_at IS NULL OR next_retry_at <= $1::timestamptz)`,
365+
[now],
366+
);
367+
return Number(res.rows[0]?.count ?? 0);
368+
}
369+
335370
export async function webhookSummary(client: Client): Promise<{
336371
pending: number;
337372
failed: number;
373+
deadLetter: number;
338374
delivered: number;
375+
/** Deliveries due right now — the lag a consumer fleet scales on (#165). */
376+
lag: number;
339377
recentFailed: Array<{
340378
id: number;
341379
paymentTxHash: string;
@@ -349,7 +387,12 @@ export async function webhookSummary(client: Client): Promise<{
349387
const counts = await client.query<{ status: string; n: string }>(
350388
`SELECT status, count(*)::text AS n FROM webhook_deliveries GROUP BY status`,
351389
);
352-
const byStatus: Record<string, number> = { pending: 0, failed: 0, delivered: 0 };
390+
const byStatus: Record<string, number> = {
391+
pending: 0,
392+
failed: 0,
393+
delivered: 0,
394+
dead_letter: 0,
395+
};
353396
for (const row of counts.rows) byStatus[row.status] = Number(row.n);
354397

355398
const recent = await client.query<{
@@ -363,15 +406,17 @@ export async function webhookSummary(client: Client): Promise<{
363406
}>(
364407
`SELECT id, payment_tx_hash, status, attempts, last_status_code, last_error, updated_at
365408
FROM webhook_deliveries
366-
WHERE status = 'failed'
409+
WHERE status IN ('failed', 'dead_letter')
367410
ORDER BY updated_at DESC
368411
LIMIT 20`,
369412
);
370413

371414
return {
372415
pending: (byStatus.pending ?? 0) + (byStatus.delivering ?? 0),
373416
failed: byStatus.failed ?? 0,
417+
deadLetter: byStatus.dead_letter ?? 0,
374418
delivered: byStatus.delivered ?? 0,
419+
lag: await pendingDue(client),
375420
recentFailed: recent.rows.map((row) => ({
376421
id: Number(row.id),
377422
paymentTxHash: row.payment_tx_hash,

0 commit comments

Comments
 (0)