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/);
  });
});