Skip to content
Merged
Show file tree
Hide file tree
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
10 changes: 10 additions & 0 deletions docs/LISTENER-CONFIGURATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,16 @@ Also used by the DB-backed retry scheduler for `baseDelayMs` / `multiplier` / `j
| `RETRY_SCHEDULER_BATCH_SIZE` | integer | `10` | No (defaulted) | Jobs per tick |
| `RETRY_MAX_DELAY_MS` | integer (ms) | `3600000` (1h) | No (defaulted) | Max backoff clamp for retry scheduler |

### Outbound webhook requests

| Name | Type | Default | Required | Purpose / effect |
|------|------|---------|----------|------------------|
| `WEBHOOK_TIMEOUT_MS` | integer (ms) | `10000` | No (defaulted) | Timeout for outbound webhook POSTs, applied independently of other network operations |

`WEBHOOK_TIMEOUT_MS` is validated at startup: non-numeric values abort startup with a `ConfigError`, and values below `1` or above `300000` (5 minutes) are rejected. When a webhook does not respond within the configured timeout the request is aborted and recorded as a distinct **timeout** failure (rather than a generic network or HTTP error), so it is visible separately in logs and retry handling.

**Recommended (guidance):** leave the `10000` ms default for typical endpoints; lower it when the receiver is expected to be fast and you want to fail over sooner, and raise it (up to `300000`) only for endpoints with a known long processing time.

### Scheduled notification scheduler

| Name | Type | Default | Required | Purpose / effect |
Expand Down
9 changes: 5 additions & 4 deletions listener/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -216,9 +216,11 @@ RETRY_SCHEDULER_PROCESSOR_ID=
# Number of retry jobs to process per scheduler tick.
RETRY_SCHEDULER_BATCH_SIZE=10

# Request timeout for outbound webhook delivery (ms). Timed-out requests are
# treated as failures and are subject to normal retry backoff.
# WEBHOOK_DELIVERY_TIMEOUT_MS=10000
# Timeout for outbound webhook requests (ms). Applies to the webhook delivery
# path independently of Discord/scheduler timeouts. Must be 1..300000.
# Timed-out requests are treated as failures and are subject to normal retry
# backoff.
WEBHOOK_TIMEOUT_MS=10000

# -----------------------------------------------------------------------------
# Retry Policy
Expand All @@ -232,7 +234,6 @@ RETRY_SCHEDULER_BATCH_SIZE=10
# configuration_error) are excluded from the default so a bad credential or a
# misconfigured URL is not retried pointlessly.
# RETRY_POLICY_RETRYABLE_FAILURE_TYPES=network_error,timeout,rate_limited,server_error,unknown

# -----------------------------------------------------------------------------
# Scheduled Notification Scheduler
# -----------------------------------------------------------------------------
Expand Down
27 changes: 27 additions & 0 deletions listener/src/config-schema.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -96,4 +96,31 @@ describe('Configuration Schema Validation (#694)', () => {
const errors = ConfigurationSchemaValidator.validate(invalidConfig, APP_CONFIG_SCHEMA);
expect(errors.some((e) => e.field === 'stellarRpcUrl' && e.message.includes('pattern'))).toBe(true);
});

it('rejects an out-of-range retryScheduler.webhookTimeoutMs', () => {
const tooSmall = {
...sampleValidConfig,
retryScheduler: { webhookTimeoutMs: 0 },
};
const tooLarge = {
...sampleValidConfig,
retryScheduler: { webhookTimeoutMs: 300001 },
};

const smallErrors = ConfigurationSchemaValidator.validate(tooSmall, APP_CONFIG_SCHEMA);
const largeErrors = ConfigurationSchemaValidator.validate(tooLarge, APP_CONFIG_SCHEMA);

expect(smallErrors.some((e) => e.field === 'retryScheduler.webhookTimeoutMs')).toBe(true);
expect(largeErrors.some((e) => e.field === 'retryScheduler.webhookTimeoutMs')).toBe(true);
});

it('accepts a valid retryScheduler.webhookTimeoutMs', () => {
const config = {
...sampleValidConfig,
retryScheduler: { webhookTimeoutMs: 10000 },
};

const errors = ConfigurationSchemaValidator.validate(config, APP_CONFIG_SCHEMA);
expect(errors.some((e) => e.field === 'retryScheduler.webhookTimeoutMs')).toBe(false);
});
});
9 changes: 9 additions & 0 deletions listener/src/config-schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
* - Validation errors identify the affected configuration field explicitly.
*/

import { MAX_WEBHOOK_TIMEOUT_MS } from './services/webhook-delivery-service';

export type FieldType = 'string' | 'number' | 'boolean' | 'array' | 'object';

export interface SchemaFieldRule {
Expand Down Expand Up @@ -223,6 +225,13 @@ export const APP_CONFIG_SCHEMA: ConfigSchema = {
batchSize: { type: 'number', min: 1 },
timingBufferMs: { type: 'number', min: 0 },
},
retryScheduler: {
enabled: { type: 'boolean' },
pollIntervalMs: { type: 'number', min: 1000 },
lockTimeoutMs: { type: 'number', min: 1000 },
batchSize: { type: 'number', min: 1 },
webhookTimeoutMs: { type: 'number', min: 1, max: MAX_WEBHOOK_TIMEOUT_MS },
},
rateLimit: {
enabled: { type: 'boolean' },
windowMs: { type: 'number', min: 1000 },
Expand Down
47 changes: 47 additions & 0 deletions listener/src/config.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -669,6 +669,53 @@ describe('Config validation', () => {
});
});
});

describe('WEBHOOK_TIMEOUT_MS', () => {
it('defaults to 10000 when unset', () => {
delete process.env.WEBHOOK_TIMEOUT_MS;

expect(loadConfig().retryScheduler?.webhookTimeoutMs).toBe(10000);
});

it('loads a configured webhook timeout', () => {
process.env.WEBHOOK_TIMEOUT_MS = '2500';

expect(loadConfig().retryScheduler?.webhookTimeoutMs).toBe(2500);
});

it('rejects a non-numeric webhook timeout at parse time', () => {
process.env.WEBHOOK_TIMEOUT_MS = 'soon';

expect(() => loadConfig()).toThrow(ConfigError);
expect(() => loadConfig()).toThrow(
'WEBHOOK_TIMEOUT_MS must be a valid integer, got "soon"'
);
});

it('rejects a zero webhook timeout', () => {
process.env.WEBHOOK_TIMEOUT_MS = '0';

const config = loadConfig();
expect(() => validateConfig(config)).toThrow(ConfigError);
expect(() => validateConfig(config)).toThrow('WEBHOOK_TIMEOUT_MS must be >= 1 ms');
});

it('rejects a negative webhook timeout', () => {
process.env.WEBHOOK_TIMEOUT_MS = '-50';

const config = loadConfig();
expect(() => validateConfig(config)).toThrow(ConfigError);
expect(() => validateConfig(config)).toThrow('WEBHOOK_TIMEOUT_MS must be >= 1 ms');
});

it('rejects an absurdly large webhook timeout', () => {
process.env.WEBHOOK_TIMEOUT_MS = '9999999';

const config = loadConfig();
expect(() => validateConfig(config)).toThrow(ConfigError);
expect(() => validateConfig(config)).toThrow('WEBHOOK_TIMEOUT_MS must be <= 300000 ms');
});
});
});
});

Expand Down
20 changes: 19 additions & 1 deletion listener/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,10 @@ import {
parseLogLevel,
} from './utils/logger';
import { DEFAULT_MAX_BODY_BYTES } from './middleware/body-limit';
import {
DEFAULT_WEBHOOK_TIMEOUT_MS,
MAX_WEBHOOK_TIMEOUT_MS,
} from './services/webhook-delivery-service';

export class ConfigError extends Error {
constructor(message: string) {
Expand Down Expand Up @@ -271,11 +275,15 @@ function loadRetrySchedulerConfig(policy: RetryPolicyOptions): RetrySchedulerOpt
multiplier: parseIntegerEnv('RETRY_MULTIPLIER', '2'),
maxDelayMs: parseIntegerEnv('RETRY_MAX_DELAY_MS', String(60 * 60 * 1000)),
jitter: trimEnv('RETRY_JITTER') !== 'false',
// Policy knobs are owned by the retry policy; fold them in so the scheduler
// Policy knobs are owned by the retry policy; fold them in so the scheduler
// and the in-memory queues agree on the attempt budget and on which failure
// types are worth retrying.
maxAttempts: policy.maxAttempts,
retryableFailureTypes: policy.retryableFailureTypes,
webhookTimeoutMs: parseIntegerEnv(
'WEBHOOK_TIMEOUT_MS',
String(DEFAULT_WEBHOOK_TIMEOUT_MS)
),
};
}

Expand Down Expand Up @@ -775,6 +783,16 @@ export function validateConfig(config: Config): void {
`RETRY_SCHEDULER_BATCH_SIZE must be >= 1 (received: ${config.retryScheduler.batchSize}).`,
);
}
if (config.retryScheduler.webhookTimeoutMs < 1) {
errors.push(
`WEBHOOK_TIMEOUT_MS must be >= 1 ms (received: ${config.retryScheduler.webhookTimeoutMs}).`,
);
} else if (config.retryScheduler.webhookTimeoutMs > MAX_WEBHOOK_TIMEOUT_MS) {
errors.push(
`WEBHOOK_TIMEOUT_MS must be <= ${MAX_WEBHOOK_TIMEOUT_MS} ms ` +
`(received: ${config.retryScheduler.webhookTimeoutMs}).`,
);
}
}

// ── Retry policy (#842) ───────────────────────────────────────────────────
Expand Down
48 changes: 47 additions & 1 deletion listener/src/services/retry-scheduler-webhook.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
* - Successful retries are marked COMPLETED and logged
*/

import { jest, describe, it, expect, beforeEach } from '@jest/globals';
import { jest, describe, it, expect, beforeEach, afterEach } from '@jest/globals';
import { RetryScheduler, RETRY_SCHEDULER_DEFAULTS } from './retry-scheduler';
import { WebhookDeliveryService } from './webhook-delivery-service';
import { NotificationStatus, NotificationType } from '../types/scheduled-notification';
Expand Down Expand Up @@ -457,4 +457,50 @@ describe('RetryScheduler — webhook retry queue', () => {
);
});
});

// ── Configurable webhook timeout ──────────────────────────────────────────
describe('webhook timeout configuration', () => {
const realFetch = global.fetch;

afterEach(() => {
(global as any).fetch = realFetch;
});

it('applies the configured webhookTimeoutMs to the outbound request', async () => {
const notification = makeWebhookNotification();
const repo = makeRepo({
fetchDueRetries: jest.fn().mockImplementation(() => Promise.resolve([notification])),
});

let aborted = false;
(global as any).fetch = jest.fn().mockImplementation((_url: any, init: any) =>
new Promise((_resolve, reject) => {
(init.signal as AbortSignal).addEventListener('abort', () => {
aborted = true;
const err = new Error('The operation was aborted.');
err.name = 'AbortError';
reject(err);
});
}),
);

const scheduler = new RetryScheduler(repo, {
...RETRY_SCHEDULER_DEFAULTS,
pollIntervalMs: 1000,
lockTimeoutMs: 1000,
batchSize: 1,
webhookTimeoutMs: 25,
});
await scheduler.runOnce();

expect(aborted).toBe(true);
expect(repo.markAsFailedOrRetry).toHaveBeenCalledWith(
10,
expect.objectContaining({ message: expect.stringContaining('timed out after 25ms') }),
expect.anything(),
expect.anything(),
expect.anything(),
);
});
});
});
15 changes: 9 additions & 6 deletions listener/src/services/retry-scheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { generateRequestId } from '../utils/request-id';
import { ScheduledNotificationRepository } from './scheduled-notification-repository';
import { ScheduledNotification, NotificationStatus } from '../types/scheduled-notification';
import { DiscordNotificationService } from './discord-notification';
import { WebhookDeliveryService } from './webhook-delivery-service';
import { WebhookDeliveryService, DEFAULT_WEBHOOK_TIMEOUT_MS } from './webhook-delivery-service';
import { getWorkerManager } from './worker-manager';
import { DeliveryReceiptRepository } from './delivery-receipt-repository';
import { DeliveryResult } from '../types/provider-capabilities';
Expand Down Expand Up @@ -36,7 +36,10 @@ export interface RetrySchedulerConfig {
maxDelayMs: number;
/** Add ±25 % random jitter to prevent thundering herd. Default: true. */
jitter: boolean;
/** Request timeout for outbound webhook delivery (ms). Default: 10 000. */
/**
* Timeout (ms) applied to outbound webhook requests (`WEBHOOK_TIMEOUT_MS`).
* Default: DEFAULT_WEBHOOK_TIMEOUT_MS.
*/
webhookTimeoutMs: number;
/**
* Hard ceiling on total delivery attempts, including the first one.
Expand All @@ -61,7 +64,7 @@ export const RETRY_SCHEDULER_DEFAULTS: RetrySchedulerConfig = {
multiplier: 2,
maxDelayMs: 60 * 60 * 1_000,
jitter: true,
webhookTimeoutMs: 10_000,
webhookTimeoutMs: DEFAULT_WEBHOOK_TIMEOUT_MS,
};

/**
Expand Down Expand Up @@ -128,9 +131,9 @@ export class RetryScheduler {
this.processorId = this.config.processorId ?? `retry-${uuidv4()}`;
this.repository = repository;
this.discordService = discordService ?? null;
this.webhookDeliveryService =
webhookDeliveryService ?? new WebhookDeliveryService({ timeoutMs: this.config.webhookTimeoutMs });
this.webhookDeliveryService = webhookDeliveryService ?? new WebhookDeliveryService();
this.webhookDeliveryService =
webhookDeliveryService ??
new WebhookDeliveryService({ timeoutMs: this.config.webhookTimeoutMs });
this.deliveryReceiptRepository = deliveryReceiptRepository;
}

Expand Down
55 changes: 55 additions & 0 deletions listener/src/services/webhook-delivery-service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -272,4 +272,59 @@ describe('WebhookDeliveryService', () => {
);
});
});

// ── Machine-readable failure classification ──────────────────────────────

describe('failure classification', () => {
it('labels a request timeout as "timeout"', async () => {
mockSendWebhook.mockRejectedValue(makeAbortError());

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBe('timeout');
});

it('labels a TimeoutError as "timeout"', async () => {
const err = new Error('request timed out');
err.name = 'TimeoutError';
mockSendWebhook.mockRejectedValue(err);

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBe('timeout');
});

it('labels a generic network error as "network"', async () => {
mockSendWebhook.mockRejectedValue(new Error('ECONNREFUSED'));

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBe('network');
expect(result.failureReason).not.toBe('timeout');
});

it('labels a 5xx response as "http_retryable"', async () => {
mockSendWebhook.mockResolvedValue(makeResponse(503, false));

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBe('http_retryable');
});

it('labels a 4xx response as "http_permanent"', async () => {
mockSendWebhook.mockResolvedValue(makeResponse(404, false));

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBe('http_permanent');
});

it('leaves failureReason undefined on success', async () => {
mockSendWebhook.mockResolvedValue(makeResponse(200));

const result = await service.deliver(TARGET_URL, PAYLOAD, REQUEST_ID);

expect(result.failureReason).toBeUndefined();
});
});
});
Loading