import * as plugins from '../plugins.js'; import { createPinnedWebPushAgent, resolveWebPushEndpoint, WebPushEndpointPolicyError, type TWebPushDnsLookup, } from './helpers.webpush-endpoint.js'; export interface IWebPushTransportRequest { subscription: plugins.webpush.PushSubscription; payload: string; ttlSeconds: number; urgency: plugins.webpush.Urgency; topic?: string; vapid: { subject: string; publicKey: string; privateKey: string; }; } export interface IWebPushTransportResponse { statusCode: number; retryAfter?: string; } export interface IWebPushTransport { send(requestArg: IWebPushTransportRequest): Promise; } export type TWebPushTransportErrorCode = | 'PUSH_DNS_TIMEOUT' | 'PUSH_REQUEST_TIMEOUT' | 'PUSH_RESPONSE_TOO_LARGE' | 'PUSH_RESPONSE_INVALID' | 'PUSH_NETWORK_ERROR'; export class WebPushTransportError extends Error { public constructor( public readonly code: TWebPushTransportErrorCode, optionsArg: { cause?: unknown } = {}, ) { super(code, optionsArg); this.name = 'WebPushTransportError'; } } export interface IWebPushHttpTransportOptions { dnsLookup?: TWebPushDnsLookup; dnsTimeoutMs?: number; connectTimeoutMs?: number; requestTimeoutMs?: number; maxResponseBytes?: number; } const normalizeRetryAfter = (valueArg: string | string[] | undefined): string | undefined => { if (Array.isArray(valueArg)) return valueArg[0]?.slice(0, 128); return typeof valueArg === 'string' ? valueArg.slice(0, 128) : undefined; }; export class WebPushHttpTransport implements IWebPushTransport { private readonly dnsLookup?: TWebPushDnsLookup; private readonly dnsTimeoutMs: number; private readonly connectTimeoutMs: number; private readonly requestTimeoutMs: number; private readonly maxResponseBytes: number; public constructor(optionsArg: IWebPushHttpTransportOptions = {}) { this.dnsLookup = optionsArg.dnsLookup; this.dnsTimeoutMs = this.normalizePositiveInteger( optionsArg.dnsTimeoutMs, 5_000, 'dnsTimeoutMs', ); this.connectTimeoutMs = this.normalizePositiveInteger( optionsArg.connectTimeoutMs, 10_000, 'connectTimeoutMs', ); this.requestTimeoutMs = this.normalizePositiveInteger( optionsArg.requestTimeoutMs, 15_000, 'requestTimeoutMs', ); this.maxResponseBytes = this.normalizePositiveInteger( optionsArg.maxResponseBytes, 16_384, 'maxResponseBytes', ); } public async send(requestArg: IWebPushTransportRequest): Promise { const sendDeadline = Date.now() + this.requestTimeoutMs; let resolved: Awaited>; try { resolved = await new Promise>>( (resolve, reject) => { let settled = false; const settle = ( errorArg: unknown, valueArg?: Awaited>, ) => { if (settled) return; settled = true; clearTimeout(timer); if (errorArg) reject(errorArg); else resolve(valueArg!); }; const timer = setTimeout(() => { settle(new WebPushTransportError('PUSH_DNS_TIMEOUT')); }, Math.min(this.dnsTimeoutMs, this.requestTimeoutMs)); resolveWebPushEndpoint( requestArg.subscription.endpoint, this.dnsLookup, ).then( (valueArg) => settle(null, valueArg), (errorArg) => settle(errorArg), ); }, ); } catch (error: unknown) { if (error instanceof WebPushEndpointPolicyError) throw error; if (error instanceof WebPushTransportError) throw error; throw new WebPushTransportError('PUSH_NETWORK_ERROR', { cause: error }); } const remainingRequestTimeoutMs = sendDeadline - Date.now(); if (remainingRequestTimeoutMs <= 0) { throw new WebPushTransportError('PUSH_REQUEST_TIMEOUT'); } const agent = createPinnedWebPushAgent(resolved); try { const requestDetails = plugins.webpush.generateRequestDetails( requestArg.subscription, requestArg.payload, { TTL: requestArg.ttlSeconds, contentEncoding: 'aes128gcm', urgency: requestArg.urgency, ...(requestArg.topic ? { topic: requestArg.topic } : {}), vapidDetails: requestArg.vapid, }, ); if ( requestDetails.method !== 'POST' || requestDetails.endpoint !== requestArg.subscription.endpoint || !Buffer.isBuffer(requestDetails.body) ) { throw new WebPushTransportError('PUSH_RESPONSE_INVALID'); } return await new Promise((resolve, reject) => { let settled = false; let requestTimer: NodeJS.Timeout | undefined; const settle = ( errorArg: Error | null, responseArg?: IWebPushTransportResponse, ) => { if (settled) return; settled = true; if (requestTimer) clearTimeout(requestTimer); if (errorArg) reject(errorArg); else resolve(responseArg!); }; const pushRequest = plugins.https.request({ protocol: 'https:', hostname: resolved.hostname, servername: resolved.hostname, port: 443, path: `${resolved.url.pathname}${resolved.url.search}`, method: 'POST', headers: requestDetails.headers, agent, timeout: Math.min(this.connectTimeoutMs, remainingRequestTimeoutMs), maxHeaderSize: 16_384, setHost: true, }, (pushResponse) => { let responseBytes = 0; pushResponse.on('data', (chunkArg: Buffer | string) => { responseBytes += Buffer.byteLength(chunkArg); if (responseBytes > this.maxResponseBytes) { const error = new WebPushTransportError('PUSH_RESPONSE_TOO_LARGE'); settle(error); pushResponse.destroy(error); pushRequest.destroy(error); } }); pushResponse.once('end', () => { const statusCode = pushResponse.statusCode; if (!Number.isSafeInteger(statusCode) || statusCode! < 100 || statusCode! > 599) { settle(new WebPushTransportError('PUSH_RESPONSE_INVALID')); return; } settle(null, { statusCode: statusCode!, retryAfter: normalizeRetryAfter(pushResponse.headers['retry-after']), }); }); pushResponse.once('error', (errorArg) => { settle( errorArg instanceof WebPushTransportError ? errorArg : new WebPushTransportError('PUSH_NETWORK_ERROR', { cause: errorArg }), ); }); }); pushRequest.once('timeout', () => { const error = new WebPushTransportError('PUSH_REQUEST_TIMEOUT'); settle(error); pushRequest.destroy(error); }); pushRequest.once('error', (errorArg) => { settle( errorArg instanceof WebPushTransportError ? errorArg : new WebPushTransportError('PUSH_NETWORK_ERROR', { cause: errorArg }), ); }); requestTimer = setTimeout(() => { const error = new WebPushTransportError('PUSH_REQUEST_TIMEOUT'); settle(error); pushRequest.destroy(error); }, remainingRequestTimeoutMs); requestTimer.unref(); pushRequest.end(requestDetails.body); }); } finally { agent.destroy(); } } private normalizePositiveInteger( valueArg: number | undefined, defaultArg: number, labelArg: string, ): number { const value = valueArg ?? defaultArg; if (!Number.isSafeInteger(value) || value <= 0) { throw new Error(`Web Push transport ${labelArg} must be a positive integer`); } return value; } }