Spaces:
Runtime error
Runtime error
File size: 8,295 Bytes
46252cd | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 | import { Job } from 'bullmq';
import { Repository } from 'typeorm';
import { ConfigService } from '@nestjs/config';
import { WebhookProcessor } from './webhook.processor';
import { Webhook } from '../../webhook/entities/webhook.entity';
import { WebhookDeliveryFailure } from '../../webhook/entities/webhook-delivery-failure.entity';
import { HookManager } from '../../../core/hooks';
import { WebhookJobData } from '../../webhook/webhook.service';
import { fetch as undiciFetch } from 'undici';
// Delivery goes through undici's fetch (via the SSRF-pinning helper), so mock that, not global fetch.
jest.mock('undici', () => {
const actual = jest.requireActual<typeof import('undici')>('undici');
return { __esModule: true, ...actual, fetch: jest.fn() };
});
/**
* Regression coverage for the production (QUEUE_ENABLED) webhook delivery path, which was
* previously untested. Covers the success path, the off-by-one final-attempt gate, the
* retry-count header, and the redirect refusal when SSRF protection is on.
*/
describe('WebhookProcessor', () => {
let processor: WebhookProcessor;
let repo: { update: jest.Mock };
let failureRepo: { insert: jest.Mock };
let hookManager: { execute: jest.Mock };
let configService: { get: jest.Mock };
let mockFetch: jest.Mock;
const origProtect = process.env.WEBHOOK_SSRF_PROTECT;
const makeJob = (overrides: Partial<WebhookJobData> = {}, attemptsMade = 0): Job<WebhookJobData> =>
({
id: 'job-1',
attemptsMade,
data: {
webhookId: 'wh-1',
url: 'https://8.8.8.8/hook', // IP literal → SSRF guard needs no DNS lookup
event: 'message.received',
payload: {
event: 'message.received',
timestamp: '',
sessionId: 'sess-1',
idempotencyKey: 'k',
deliveryId: 'd',
data: {},
},
headers: { 'Content-Type': 'application/json' },
attempt: 1,
maxRetries: 3,
...overrides,
},
}) as unknown as Job<WebhookJobData>;
beforeEach(() => {
repo = { update: jest.fn().mockResolvedValue({ affected: 1 }) };
failureRepo = { insert: jest.fn().mockResolvedValue({}) };
hookManager = { execute: jest.fn().mockResolvedValue({ continue: true, data: {} }) };
configService = { get: jest.fn((key: string, def?: unknown) => (key === 'webhook.timeout' ? 25000 : def)) };
processor = new WebhookProcessor(
repo as unknown as Repository<Webhook>,
failureRepo as unknown as Repository<WebhookDeliveryFailure>,
hookManager as unknown as HookManager,
configService as unknown as ConfigService,
);
// The merged delivery path uses withSafeFetch (undici), so mock undici's fetch, not global.fetch.
mockFetch = undiciFetch as jest.Mock;
process.env.WEBHOOK_SSRF_PROTECT = 'false'; // delivery-logic tests; redirect test flips it on
});
afterEach(() => {
mockFetch.mockReset();
if (origProtect === undefined) delete process.env.WEBHOOK_SSRF_PROTECT;
else process.env.WEBHOOK_SSRF_PROTECT = origProtect;
});
it('uses the configured WEBHOOK_TIMEOUT for the request abort signal (not a hardcoded 10s)', async () => {
const timeoutSpy = jest.spyOn(AbortSignal, 'timeout');
mockFetch.mockResolvedValue({ ok: true, status: 200 });
await processor.process(makeJob());
expect(configService.get).toHaveBeenCalledWith('webhook.timeout', 10000);
expect(timeoutSpy).toHaveBeenCalledWith(25000);
timeoutSpy.mockRestore();
});
it('on success updates lastTriggeredAt and fires webhook:delivered', async () => {
mockFetch.mockResolvedValue({ ok: true, status: 200 });
const result = await processor.process(makeJob());
expect(result.success).toBe(true);
expect(repo.update).toHaveBeenCalledTimes(1);
const updateArgs = repo.update.mock.calls[0] as unknown as [string, { lastTriggeredAt: Date }];
expect(updateArgs[0]).toBe('wh-1');
expect(updateArgs[1].lastTriggeredAt).toBeInstanceOf(Date);
expect(hookManager.execute).toHaveBeenCalledWith('webhook:delivered', expect.anything(), expect.anything());
});
it('sets X-OpenWA-Retry-Count to the attempt number', async () => {
mockFetch.mockResolvedValue({ ok: true, status: 200 });
await processor.process(makeJob({}, 2));
const call = mockFetch.mock.calls[0] as unknown as [string, { headers: Record<string, string> }];
expect(call[1].headers['X-OpenWA-Retry-Count']).toBe('2');
});
it('throws on a non-ok response WITHOUT firing webhook:error before the final attempt', async () => {
mockFetch.mockResolvedValue({ ok: false, status: 500, statusText: 'Server Error' });
await expect(processor.process(makeJob({ maxRetries: 3 }, 0))).rejects.toThrow();
expect(hookManager.execute).not.toHaveBeenCalledWith('webhook:error', expect.anything(), expect.anything());
});
it('fires webhook:error only on the final attempt', async () => {
mockFetch.mockResolvedValue({ ok: false, status: 500, statusText: 'Server Error' });
// attemptsMade=2, maxRetries=3 -> attemptsMade+1 >= maxRetries -> final
await expect(processor.process(makeJob({ maxRetries: 3 }, 2))).rejects.toThrow();
expect(hookManager.execute).toHaveBeenCalledWith('webhook:error', expect.anything(), expect.anything());
});
it('persists a durable delivery-failure record on the final attempt (with parsed HTTP status)', async () => {
mockFetch.mockResolvedValue({ ok: false, status: 503, statusText: 'Service Unavailable' });
await expect(
processor.process(makeJob({ maxRetries: 3, webhookId: 'wh-x', url: 'https://8.8.8.8/h' }, 2)),
).rejects.toThrow();
expect(failureRepo.insert).toHaveBeenCalledTimes(1);
expect(failureRepo.insert).toHaveBeenCalledWith(
expect.objectContaining({
webhookId: 'wh-x',
url: 'https://8.8.8.8/h',
sessionId: 'sess-1',
attempts: 3,
lastStatusCode: 503,
lastError: 'HTTP 503: Service Unavailable',
}),
);
});
it('does NOT persist a delivery-failure record before the final attempt', async () => {
mockFetch.mockResolvedValue({ ok: false, status: 500, statusText: 'Server Error' });
await expect(processor.process(makeJob({ maxRetries: 3 }, 0))).rejects.toThrow();
expect(failureRepo.insert).not.toHaveBeenCalled();
});
it('refuses to follow a redirect when SSRF protection is on', async () => {
process.env.WEBHOOK_SSRF_PROTECT = 'true';
mockFetch.mockResolvedValue({ ok: false, status: 0, type: 'opaqueredirect' });
await expect(processor.process(makeJob({ maxRetries: 1 }, 0))).rejects.toThrow();
expect(mockFetch).toHaveBeenCalledWith('https://8.8.8.8/hook', expect.objectContaining({ redirect: 'manual' }));
expect(repo.update).not.toHaveBeenCalled(); // never treated as delivered
});
// A literal link-local IP triggers the SSRF guard synchronously before any fetch/DNS, so this is
// fully offline. The webhook:error hook payload and the durable DLQ row must both carry the generic
// message — the resolved internal IP is a recon oracle. The server-side logger.error keeps full detail.
it('redacts the resolved internal IP from the webhook:error payload and DLQ row on an SSRF block', async () => {
process.env.WEBHOOK_SSRF_PROTECT = 'true';
// final attempt (attemptsMade=0, maxRetries=1 → 1 >= 1) so the hook + DLQ fire
await expect(processor.process(makeJob({ url: 'https://169.254.169.254/h', maxRetries: 1 }, 0))).rejects.toThrow();
expect(mockFetch).not.toHaveBeenCalled(); // blocked before any network
const hookCalls = hookManager.execute.mock.calls as unknown as Array<[string, { error: string }, unknown]>;
const errorHookCall = hookCalls.find(c => c[0] === 'webhook:error');
expect(errorHookCall).toBeDefined();
expect(errorHookCall![1].error).toBe('Destination address is not allowed');
expect(errorHookCall![1].error).not.toMatch(/169\.254\.169\.254/);
expect(failureRepo.insert).toHaveBeenCalledTimes(1);
const inserted = (failureRepo.insert.mock.calls[0] as unknown[])[0] as { lastError: string };
expect(inserted.lastError).toBe('Destination address is not allowed');
expect(inserted.lastError).not.toMatch(/169\.254\.169\.254/);
});
});
|