import type { HttpLibrary, RequestContext, } from "@kubernetes/client-node/dist/gen/http/http.js"; import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js"; import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js"; import http from "node:http"; import https from "node:https"; type TransportOptions = { maxConcurrent: number; minIntervalMs: number; maxRetries: number; baseRetryMs: number; maxRetryMs: number; random: () => number; }; const DefaultTransportOptions: TransportOptions = { maxConcurrent: 4, minIntervalMs: 100, maxRetries: 3, baseRetryMs: 250, maxRetryMs: 5000, random: Math.random, }; function abortError(): Error { const error = new Error("Request aborted"); error.name = "AbortError"; return error; } export function delay(ms: number, signal?: AbortSignal): Promise { if (signal?.aborted) return Promise.reject(abortError()); return new Promise((resolve, reject) => { const timer = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); const onAbort = () => { clearTimeout(timer); reject(abortError()); }; signal?.addEventListener("abort", onAbort, { once: true }); }); } export function parseRetryAfter( value: string | undefined, now = Date.now(), ): number | undefined { if (!value) return; const seconds = Number(value); if (Number.isFinite(seconds) && seconds >= 0) return seconds * 1000; const date = Date.parse(value); if (Number.isNaN(date)) return; return Math.max(0, date - now); } class RequestLimiter { private active = 0; private nextRequestAt = 0; private readonly queue: Array<{ signal?: AbortSignal; onAbort?: () => void; resolve: (release: () => void) => void; reject: (error: Error) => void; }> = []; constructor( private readonly maxConcurrent: number, private readonly minIntervalMs: number, ) {} acquire(signal?: AbortSignal): Promise<() => void> { if (signal?.aborted) return Promise.reject(abortError()); return new Promise((resolve, reject) => { const entry: (typeof this.queue)[number] = { signal, resolve, reject }; entry.onAbort = () => { const index = this.queue.indexOf(entry); if (index === -1) return; this.queue.splice(index, 1); reject(abortError()); }; signal?.addEventListener("abort", entry.onAbort, { once: true }); this.queue.push(entry); this.drain(); }); } private drain(): void { while (this.active < this.maxConcurrent && this.queue.length > 0) { const entry = this.queue.shift()!; entry.signal?.removeEventListener("abort", entry.onAbort!); if (entry.signal?.aborted) { entry.reject(abortError()); continue; } this.active += 1; const now = Date.now(); const waitMs = Math.max(0, this.nextRequestAt - now); this.nextRequestAt = Math.max(now, this.nextRequestAt) + this.minIntervalMs; void delay(waitMs, entry.signal) .then(() => { let released = false; entry.resolve(() => { if (released) return; released = true; this.active -= 1; this.drain(); }); }) .catch((error) => { this.active -= 1; entry.reject( error instanceof Error ? error : new Error(String(error)), ); this.drain(); }); } } } function sendOnce(request: RequestContext): Promise { return new Promise((resolve, reject) => { const url = new URL(request.getUrl()); const transport = url.protocol === "http:" ? http : https; const signal = request.getSignal(); const req = transport.request( url, { method: request.getHttpMethod().toString(), headers: request.getHeaders(), agent: request.getAgent() as never, }, (response) => { const chunks: Buffer[] = []; response.on("data", (chunk) => { chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); }); response.on("end", () => { signal?.removeEventListener("abort", onAbort); const buffer = Buffer.concat(chunks); const headers: Record = {}; for (const [key, value] of Object.entries(response.headers)) { if (Array.isArray(value)) headers[key] = value.join(", "); else if (value !== undefined) headers[key] = String(value); } resolve( new ResponseContext(response.statusCode ?? 0, headers, { text: async () => buffer.toString("utf8"), binary: async () => buffer, }), ); }); }, ); const onAbort = () => req.destroy(abortError()); if (signal?.aborted) onAbort(); else signal?.addEventListener("abort", onAbort, { once: true }); req.on("error", (error) => { signal?.removeEventListener("abort", onAbort); reject(error); }); const body = request.getBody(); if (body !== undefined && body !== null) { req.write(body as string | Uint8Array); } req.end(); }); } export function createKubernetesHttpLibrary( overrides: Partial = {}, ): HttpLibrary { const options = { ...DefaultTransportOptions, ...overrides }; const limiter = new RequestLimiter( options.maxConcurrent, options.minIntervalMs, ); return { send(request) { const result = (async () => { for (let attempt = 0; ; attempt += 1) { const release = await limiter.acquire(request.getSignal()); let response: ResponseContext; try { response = await sendOnce(request); } finally { release(); } if ( response.httpStatusCode !== 429 || attempt >= options.maxRetries ) { return response; } const retryAfter = parseRetryAfter(response.headers["retry-after"]); const exponential = Math.min( options.maxRetryMs, options.baseRetryMs * 2 ** attempt, ); const backoff = retryAfter ?? exponential * (0.8 + options.random() * 0.4); await delay(backoff, request.getSignal()); } })(); return from(result); }, }; }