Files
WeHub Mirror 6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

1239 lines
48 KiB
TypeScript

import dns from 'node:dns/promises'
import { Readable } from 'node:stream'
import zlib from 'node:zlib'
import http from 'http'
import https from 'https'
import type { LookupFunction } from 'net'
import { createLogger } from '@sim/logger'
import { preferIpv4, resolveHostAddresses } from '@sim/security/dns'
import { isLoopbackIp, isPrivateIp, isPrivateIpHost, unwrapIpv6Brackets } from '@sim/security/ssrf'
import { toError } from '@sim/utils/errors'
import { omit } from '@sim/utils/object'
import { HttpProxyAgent } from 'http-proxy-agent'
import { HttpsProxyAgent } from 'https-proxy-agent'
import * as ipaddr from 'ipaddr.js'
import {
Agent,
type Dispatcher,
type RequestInit as UndiciRequestInit,
request as undiciRequest,
} from 'undici'
import { isHosted, isPrivateDatabaseHostsAllowed } from '@/lib/core/config/env-flags'
import { type ValidationResult, validateExternalUrl } from '@/lib/core/security/input-validation'
import { nodeReadableToWebStream } from '@/lib/core/utils/node-stream'
import { PayloadSizeLimitError } from '@/lib/core/utils/stream-limits'
const logger = createLogger('InputValidation')
/**
* Result type for async URL validation with resolved IP
*/
export interface AsyncValidationResult extends ValidationResult {
resolvedIP?: string
originalHostname?: string
}
/**
* Validates a URL and resolves its DNS to prevent SSRF via DNS rebinding
*
* This function:
* 1. Performs basic URL validation (protocol, format)
* 2. Resolves the hostname to an IP address
* 3. Validates the resolved IP is not private/reserved
* 4. Returns the resolved IP for use in the actual request
*
* @param url - The URL to validate
* @param paramName - Name of the parameter for error messages
* @returns AsyncValidationResult with resolved IP for DNS pinning
*/
export async function validateUrlWithDNS(
url: string | null | undefined,
paramName = 'url',
options: { allowHttp?: boolean } = {}
): Promise<AsyncValidationResult> {
const basicValidation = validateExternalUrl(url, paramName, options)
if (!basicValidation.isValid) {
return basicValidation
}
const parsedUrl = new URL(url!)
const hostname = parsedUrl.hostname
const hostnameLower = hostname.toLowerCase()
const cleanHostname = unwrapIpv6Brackets(hostnameLower)
// Whole loopback range — see the matching note in input-validation.ts.
const isLocalhost = cleanHostname === 'localhost' || isLoopbackIp(cleanHostname)
try {
// Refused records are filtered rather than failing the whole host, matching
// createSsrfGuardedLookup below. Pinning to a surviving public address is
// just as safe as refusing outright, and rejecting the host would break a
// split-horizon resolver that answers with a private record alongside the
// public one — with no operator opt-out on this path.
const { addresses } = await resolveHostAddresses(cleanHostname)
const usable = addresses.filter(
(address) => !isPrivateIp(address) || (isLocalhost && !isHosted && isLoopbackIp(address))
)
if (usable.length === 0) {
logger.warn('URL resolves to blocked IP address', {
paramName,
hostname,
resolvedIP: addresses.find((address) => isPrivateIp(address)),
})
return {
isValid: false,
error: `${paramName} resolves to a blocked IP address`,
}
}
return {
isValid: true,
// Re-preferred over the surviving set so the pin is never an address the
// filter above just refused.
resolvedIP: preferIpv4(usable as [string, ...string[]]),
originalHostname: hostname,
}
} catch (error) {
logger.warn('DNS lookup failed for URL', {
paramName,
hostname,
error: toError(error).message,
})
return {
isValid: false,
error: `${paramName} hostname could not be resolved`,
}
}
}
/**
* Result of validating a user-supplied HTTP proxy URL.
*/
export interface ProxyValidationResult {
isValid: boolean
/** Proxy URL with hostname rewritten to the resolved IP (creds/port preserved) to pin the proxy connection. */
pinnedProxyUrl?: string
error?: string
}
/**
* Validates a user-supplied HTTP proxy URL and returns an IP-pinned form.
*
* When a request routes through a proxy, the TCP connection targets the proxy
* host (the proxy resolves the destination), so target-IP pinning no longer
* governs egress and the proxy URL becomes the SSRF surface. This function:
* 1. Enforces the `http:` scheme (raw TCP to the proxy, no TLS-to-proxy SNI to
* reconcile, so the host can be safely rewritten to an IP).
* 2. Resolves the proxy host's DNS and blocks private/reserved/loopback IPs via
* {@link validateUrlWithDNS}.
* 3. Pins the connection by rewriting the hostname to the resolved IP while
* preserving credentials/port, closing the DNS-rebinding (TOCTOU) window.
*
* @param proxyUrl - The proxy URL (e.g. `http://user:pass@host:port`)
*/
export async function validateAndPinProxyUrl(
proxyUrl: string | null | undefined
): Promise<ProxyValidationResult> {
if (!proxyUrl || typeof proxyUrl !== 'string') {
return { isValid: false, error: 'proxyUrl must be a string' }
}
let parsed: URL
try {
parsed = new URL(proxyUrl)
} catch {
return { isValid: false, error: 'proxyUrl must be a valid URL' }
}
if (parsed.protocol !== 'http:') {
return {
isValid: false,
error: 'proxyUrl must use http:// (https/socks proxies are not supported)',
}
}
const validation = await validateUrlWithDNS(proxyUrl, 'proxyUrl', { allowHttp: true })
if (!validation.isValid) {
return { isValid: false, error: validation.error }
}
const resolvedIP = validation.resolvedIP!
// validateUrlWithDNS permits loopback for self-hosted dev targets; a proxy governs
// egress, so loopback/private proxy hosts stay blocked unconditionally.
if (isPrivateIp(resolvedIP)) {
return { isValid: false, error: 'proxyUrl resolves to a blocked IP address' }
}
// Bracket IPv6 literals: assigning an unbracketed IPv6 address to URL.hostname
// is a no-op, which would leave the DNS hostname in place and reopen rebinding.
parsed.hostname = resolvedIP.includes(':') ? `[${resolvedIP}]` : resolvedIP
return { isValid: true, pinnedProxyUrl: parsed.toString() }
}
/**
* Validates a database hostname by resolving DNS and checking the resolved IP
* against private/reserved ranges to prevent SSRF via database connections.
*
* Unlike validateHostname (which enforces strict RFC hostname format), this
* function is permissive about hostname format to avoid breaking legitimate
* database hostnames (e.g. underscores in Docker/K8s service names). It only
* blocks localhost and private/reserved IPs.
*
* Self-hosted operators can set `ALLOW_PRIVATE_DATABASE_HOSTS` to reach databases
* on their private network (e.g. a Docker/Swarm service name that resolves to an
* internal IP). The opt-in only bypasses the private/reserved/loopback block; DNS
* is still resolved so the caller can pin the connection to the resolved IP. The
* bypass is never honored on the hosted platform (see {@link isPrivateDatabaseHostsAllowed}).
*
* @param host - The database hostname to validate
* @param paramName - Name of the parameter for error messages
* @returns AsyncValidationResult with resolved IP
*/
export async function validateDatabaseHost(
host: string | null | undefined,
paramName = 'host'
): Promise<AsyncValidationResult> {
if (!host) {
return { isValid: false, error: `${paramName} is required` }
}
const cleanHost = unwrapIpv6Brackets(host.toLowerCase())
if (cleanHost === 'localhost' && !isPrivateDatabaseHostsAllowed) {
return { isValid: false, error: `${paramName} cannot be localhost` }
}
if (isPrivateIpHost(cleanHost) && !isPrivateDatabaseHostsAllowed) {
return { isValid: false, error: `${paramName} cannot be a private IP address` }
}
try {
const { addresses, preferred } = await resolveHostAddresses(cleanHost)
const blockedAddress = isPrivateDatabaseHostsAllowed
? undefined
: addresses.find((candidate) => isPrivateIp(candidate))
if (blockedAddress !== undefined) {
logger.warn('Database host resolves to blocked IP address', {
paramName,
hostname: host,
resolvedIP: blockedAddress,
})
return {
isValid: false,
error: `${paramName} resolves to a blocked IP address`,
}
}
return {
isValid: true,
resolvedIP: preferred,
originalHostname: host,
}
} catch (error) {
logger.warn('DNS lookup failed for database host', {
paramName,
hostname: host,
error: toError(error).message,
})
return {
isValid: false,
error: `${paramName} hostname could not be resolved`,
}
}
}
/**
* Patterns run against the WHERE clause with string/identifier literals masked
* out (so an attacker cannot smuggle `OR 1` or `; DROP` inside a quoted value).
*
* The connector-literal rules below are intentionally `OR`-only: only an
* `OR <truthy>` term broadens a mutation to every row. `AND <number>` is a no-op
* for broadening and is also exactly what `BETWEEN low AND high` produces, so
* matching it would reject legitimate range filters (e.g. `id BETWEEN 1 AND 10`).
*/
const SQL_WHERE_MASKED_PATTERNS: readonly RegExp[] = [
/;\s*\w/, // stacked statement
/\bunion\s+(?:all\s+)?select\b/i,
/\binto\s+(?:out|dump)file\b/i,
/--/,
/\/\*/,
/\*\//,
/\b(?:sleep|pg_sleep|benchmark)\s*\(/i,
/\b(\w+)\s*=\s*\1\b/i, // same (unquoted) operand both sides: x=x, 1=1
/\b\d+(?:\.\d+)?\s*(?:=|==|<>|!=|<=|>=|<|>)\s*\d+(?:\.\d+)?\b/, // constant vs constant: 1=1, 1<2, 2>1
/\bor\s+(?:true|false)\b/i, // OR TRUE / OR FALSE
/\bor\s+\d+(?:\.\d+)?\b(?!\s*[=<>!+\-*/%])/i, // standalone truthy literal after OR: OR 1, OR 42
/^\s*(?:\d+(?:\.\d+)?|true|false)\s*$/i, // bare constant: "1" / "true" / "false"
]
/**
* Patterns run against the raw WHERE clause (need the literal contents intact),
* e.g. equality between two identical string literals.
*/
const SQL_WHERE_RAW_PATTERNS: readonly RegExp[] = [
/(['"])([^'"]*)\1\s*(?:=|==|<>|!=)\s*\1\2\1/, // 'a'='a' / "x"="x"
]
/**
* Replaces the contents of string literals ('...'), double-quoted and
* backtick-quoted identifiers with spaces (preserving length) so structural
* scans do not treat data inside quotes as SQL. Comments are intentionally left
* intact so comment-injection sequences are still detected.
*/
function maskSqlStringLiterals(sql: string): string {
let out = ''
let i = 0
while (i < sql.length) {
const ch = sql[i]
if (ch === "'" || ch === '"' || ch === '`') {
out += ' '
i++
while (i < sql.length && sql[i] !== ch) {
if (ch !== '`' && sql[i] === '\\') {
out += ' '
i += 2
continue
}
out += ' '
i++
}
if (i < sql.length) {
out += ' '
i++
}
continue
}
out += ch
i++
}
return out
}
/**
* Validates a free-form SQL `WHERE` condition for injection and always-true
* tautology patterns. Returns a {@link ValidationResult}; callers decide whether
* to throw or surface the error.
*
* IMPORTANT: this is **defense-in-depth, not a security boundary**. A free-form
* SQL condition cannot be exhaustively validated against every always-true
* expression (e.g. `OR 2 > 1`, `OR (1)`, `OR NOT 0`, `OR length(x) >= 0`). The
* real boundary is that the caller supplies their own database credentials and
* could run equivalent SQL directly (e.g. via a raw-SQL/execute operation). This
* guard stops the easy, obvious ways an injected condition broadens a mutation
* to every row; it is not a substitute for constraining untrusted input upstream.
*
* @param where - The WHERE clause condition (without the `WHERE` keyword)
* @param paramName - Label used in the error message
*/
export function validateSqlWhereClause(
where: string | null | undefined,
paramName = 'WHERE clause'
): ValidationResult {
if (typeof where !== 'string' || where.trim().length === 0) {
return { isValid: false, error: `${paramName} is required` }
}
const masked = maskSqlStringLiterals(where)
const matched =
SQL_WHERE_MASKED_PATTERNS.some((pattern) => pattern.test(masked)) ||
SQL_WHERE_RAW_PATTERNS.some((pattern) => pattern.test(where))
if (matched) {
return {
isValid: false,
error: `${paramName} contains a disallowed or always-true expression`,
}
}
return { isValid: true }
}
export interface SecureFetchOptions {
method?: string
headers?: Record<string, string>
body?: string | Buffer | Uint8Array
timeout?: number
maxRedirects?: number
/**
* Maximum bytes read from the response body. Defaults to
* {@link DEFAULT_MAX_RESPONSE_BYTES} — there is deliberately no "unlimited" mode, since
* many callers target a user-supplied host that can stream an endless body.
*/
maxResponseBytes?: number
signal?: AbortSignal
/** Drop the Authorization header when following a redirect, so it is not sent to the redirect target's origin. */
stripAuthOnRedirect?: boolean
/**
* Pre-validated, IP-pinned `http://` proxy URL (see {@link validateAndPinProxyUrl}).
* When set, the connection routes through this proxy and target-IP pinning is
* bypassed (the proxy resolves the target).
*/
proxyUrl?: string
}
export class SecureFetchHeaders {
private headers: Map<string, string>
private setCookies: string[]
constructor(headers: Record<string, string>, setCookies: string[] = []) {
this.headers = new Map(Object.entries(headers).map(([k, v]) => [k.toLowerCase(), v]))
this.setCookies = setCookies
}
get(name: string): string | null {
return this.headers.get(name.toLowerCase()) ?? null
}
/** Returns the raw `Set-Cookie` header values as an array. Each entry is one cookie. */
getSetCookie(): string[] {
return [...this.setCookies]
}
toRecord(): Record<string, string> {
const record: Record<string, string> = {}
for (const [key, value] of this.headers) {
record[key] = value
}
return record
}
[Symbol.iterator]() {
return this.headers.entries()
}
}
export interface SecureFetchResponse {
ok: boolean
status: number
statusText: string
headers: SecureFetchHeaders
body: ReadableStream<Uint8Array> | null
text: () => Promise<string>
json: () => Promise<unknown>
arrayBuffer: () => Promise<ArrayBuffer>
}
const DEFAULT_MAX_REDIRECTS = 5
/**
* Fail-safe ceiling applied by {@link secureFetchWithPinnedIP} when the caller does not
* pass `maxResponseBytes`. Many callers fetch a user-supplied host, so an omitted cap
* would let a malicious upstream stream an endless chunked body into memory until the
* process is OOM-killed. Set to the platform's largest legitimate payload (100MB, matching
* the upload limit); callers that need more must opt in explicitly, and callers handling
* small JSON should pass a much tighter cap.
*/
export const DEFAULT_MAX_RESPONSE_BYTES = 100 * 1024 * 1024
/** Response cap for JSON/control-plane proxies to user-supplied hosts. */
export const MAX_JSON_API_RESPONSE_BYTES = 10 * 1024 * 1024
function isRedirectStatus(status: number): boolean {
return status >= 300 && status < 400 && status !== 304
}
function isRetryableHttpStatus(status: number): boolean {
return status === 429 || (status >= 500 && status <= 599)
}
function resolveRedirectUrl(baseUrl: string, location: string): string {
try {
return new URL(location, baseUrl).toString()
} catch {
throw new Error(`Invalid redirect location: ${location}`)
}
}
/**
* Creates a DNS lookup function that always returns a pre-resolved IP address.
* Use this to prevent DNS rebinding (TOCTOU) attacks when connecting to
* user-controlled hostnames via non-HTTP protocols (SMTP, SSH, IMAP, etc.).
*/
export function createPinnedLookup(resolvedIP: string): LookupFunction {
const isIPv6 = resolvedIP.includes(':')
const family = isIPv6 ? 6 : 4
return (_hostname, options, callback) => {
if (options.all) {
callback(null, [{ address: resolvedIP, family }])
} else {
callback(null, resolvedIP, family)
}
}
}
/**
* DNS lookup that resolves normally and validates EVERY resolved address against
* the SSRF policy at socket-connect time (the LibreChat `getSSRFConnect` pattern).
* Private/reserved/loopback records are filtered out; if nothing publicly routable
* remains the connect fails. Because the check runs on each dial — including
* redirects and reconnects — there is no validated-then-trusted window for a DNS
* rebind to slip through, and unlike single-IP pinning the connector keeps the
* full public address set, so the OS/undici can fall back across addresses.
* IPv4 is ordered first (`verbatim: false`) — our egress is IPv4-only.
*/
export function createSsrfGuardedLookup(): LookupFunction {
return (hostname, options, callback) => {
dns
.lookup(hostname, { all: true, verbatim: false })
.then((addresses) => {
const usable = addresses.filter((entry) => !isPrivateIp(entry.address))
if (usable.length === 0) {
callback(
new Error(`Blocked by SSRF policy: ${hostname} has no publicly routable address`),
'',
4
)
return
}
if (options.all) callback(null, usable)
else callback(null, usable[0].address, usable[0].family)
})
.catch((error) => callback(toError(error), '', 4))
}
}
const MAX_GUARDED_REDIRECTS = 5
/**
* Rejects a redirect hop whose target is a private/reserved IP LITERAL. Node's
* `net.connect` bypasses the custom `lookup` for numeric hosts (`isIP(host)`
* short-circuits), so the connect-time guard never sees IP-literal dials —
* a 3xx to `http://169.254.169.254/` would otherwise connect directly. Hostname
* targets are covered by {@link createSsrfGuardedLookup} at connect time.
*/
function assertGuardedRedirectTarget(url: URL, allowedPinnedIp?: string): void {
if (url.protocol !== 'http:' && url.protocol !== 'https:') {
throw new Error(`Blocked by SSRF policy: redirect to unsupported protocol ${url.protocol}`)
}
const host = unwrapIpv6Brackets(url.hostname)
if (ipaddr.isValid(host) && isPrivateIp(host)) {
// The pinned-private carve-out permits exactly its own validated IP as a target (a
// self-hosted MCP on a private IP, or a same-host redirect that stays on it) — but nothing
// else private (a redirect to e.g. the cloud metadata IP is still blocked).
if (
allowedPinnedIp &&
ipaddr.isValid(allowedPinnedIp) &&
ipaddr.process(host).toString() === ipaddr.process(allowedPinnedIp).toString()
) {
return
}
throw new Error('Blocked by SSRF policy: redirect to a private or reserved address')
}
}
/**
* Manual, revalidating redirect follower used by the guarded fetch. Auto-follow
* is unsafe here on two counts the connect-time lookup cannot cover: IP-literal
* redirect targets bypass the lookup entirely (validated per hop instead), and
* undici retains CUSTOM request headers across cross-origin redirects (it strips
* only Authorization/Cookie) — so caller headers are dropped on any cross-origin
* hop. Exported for tests.
*/
export async function followRedirectsGuarded(
rawFetch: (url: string, init: UndiciRequestInit) => Promise<Response>,
input: string,
init: UndiciRequestInit,
options?: { allowRedirectToIp?: string }
): Promise<Response> {
let currentUrl = new URL(input)
// The initial URL gets the same IP-literal check as redirect hops, so the exported guard is
// self-contained even when a caller skips its own up-front validation. `allowRedirectToIp`
// (the pinned-private MCP carve-out's validated IP) permits that one private target — both the
// initial URL and any hop that stays on it — while everything else private stays blocked.
assertGuardedRedirectTarget(currentUrl, options?.allowRedirectToIp)
let method = (init.method ?? 'GET').toUpperCase()
let body = init.body
let headers = init.headers
for (let hop = 0; ; hop++) {
const response = await rawFetch(currentUrl.href, {
...init,
method,
body,
headers,
redirect: 'manual',
})
const status = response.status
const location = response.headers.get('location')
if (![301, 302, 303, 307, 308].includes(status) || !location) {
// `response.url` is already the final hop's URL (set per-request by the raw fetch); flag
// `redirected` too when at least one hop was followed, matching fetch semantics.
if (hop > 0)
Object.defineProperty(response, 'redirected', { value: true, configurable: true })
return response
}
// Cancel the redirect body up front so the throw paths below (hop cap, blocked
// target) can't leave a socket checked out on the long-lived Agent.
await response.body?.cancel().catch(() => {})
if (hop >= MAX_GUARDED_REDIRECTS) {
throw new Error(`Blocked by SSRF policy: more than ${MAX_GUARDED_REDIRECTS} redirects`)
}
const nextUrl = new URL(location, currentUrl)
assertGuardedRedirectTarget(nextUrl, options?.allowRedirectToIp)
// Per the fetch spec: 303 (and 301/302 on POST) switch to a bodyless GET, dropping
// the entity headers that described the removed body (a retained Content-Length /
// Content-Type on a bodyless GET is malformed and undici rejects it).
if (status === 303 || ((status === 301 || status === 302) && method === 'POST')) {
method = 'GET'
body = undefined
if (headers !== undefined) {
const sanitized = new Headers(headers as HeadersInit)
sanitized.delete('content-length')
sanitized.delete('content-type')
sanitized.delete('content-encoding')
sanitized.delete('transfer-encoding')
// double-cast-allowed: Headers is a valid undici HeadersInit at runtime but the DOM/undici types differ
headers = sanitized as unknown as UndiciRequestInit['headers']
}
}
if (nextUrl.origin !== currentUrl.origin) {
headers = undefined
// 307/308 preserve method+body; forwarding a body cross-origin can hand OAuth
// client secrets / tokens to an open-redirect target now that redirects really
// dial the new origin. No legitimate MCP/OAuth flow does this — refuse it.
if (body !== undefined && body !== null) {
throw new Error(
'Blocked by SSRF policy: cross-origin redirect would forward a request body'
)
}
}
currentUrl = nextUrl
}
}
/** Coerce a DOM/undici `HeadersInit` into the record shape undici `request` accepts. */
function toUndiciRequestHeaders(
headers: UndiciRequestInit['headers']
): Record<string, string> | undefined {
if (!headers) return undefined
const record: Record<string, string> = {}
if (Array.isArray(headers)) {
for (const [key, value] of headers as [string, string][]) {
if (value != null) record[key] = String(value)
}
return record
}
// Single cast (no `as unknown`): the optional `forEach` is satisfiable by both a plain
// record (absent) and a `Headers` instance (present), so it detects the iterable form.
const iterableHeaders = headers as {
forEach?: (cb: (value: string, key: string) => void) => void
}
if (typeof iterableHeaders.forEach === 'function') {
iterableHeaders.forEach((value, key) => {
record[key] = value
})
return record
}
for (const [key, value] of Object.entries(headers as Record<string, unknown>)) {
if (value != null) record[key] = Array.isArray(value) ? value.join(', ') : String(value)
}
return record
}
/** Coerce a DOM/undici body init into a value undici `request` accepts. */
function toUndiciRequestBody(
body: UndiciRequestInit['body']
): string | Buffer | Uint8Array | Readable | undefined {
if (body == null) return undefined
// fetch accepts URLSearchParams (form-encoded) and undici.request does not — the MCP SDK's
// OAuth token/refresh exchange sends one. Serialize it to its wire form.
if (body instanceof URLSearchParams) return body.toString()
if (body instanceof ArrayBuffer) return Buffer.from(body)
if (ArrayBuffer.isView(body) && !(body instanceof Uint8Array)) {
return Buffer.from(body.buffer, body.byteOffset, body.byteLength)
}
if (typeof (body as ReadableStream).getReader === 'function') {
// double-cast-allowed: DOM ReadableStream and the node:stream Web type differ but are structurally compatible at runtime
return Readable.fromWeb(body as unknown as Parameters<typeof Readable.fromWeb>[0])
}
// string, Uint8Array/Buffer, or Readable — passed through unchanged.
// double-cast-allowed: undici BodyInit is wider than what request() accepts; our guarded/pinned callers only send these
return body as unknown as string | Buffer | Uint8Array | Readable
}
/**
* Decompression transform for a `Content-Encoding`, or `null` to pass the body through.
* `undici.fetch` decodes the body automatically; `undici.request` does not, so this restores
* fetch parity for gzip/deflate/br responses (common behind CDNs). Unknown/absent encodings
* pass through untouched.
*/
function contentEncodingDecoder(
encoding: string
): zlib.Gunzip | zlib.Inflate | zlib.BrotliDecompress | null {
switch (encoding) {
case 'gzip':
case 'x-gzip':
return zlib.createGunzip()
case 'deflate':
return zlib.createInflate()
case 'br':
return zlib.createBrotliDecompress()
default:
return null
}
}
/**
* Streaming-safe replacement for `undiciFetch(url, { ...init, dispatcher })`.
*
* undici's `fetch` exposes the response body as a WHATWG `ReadableStream` whose
* bridge is broken under the Bun runtime (which the standalone server runs on):
* response headers arrive but `response.body` never yields data, hanging every
* incremental read — MCP SSE `tools/list`, provider streaming — to its timeout.
* undici's lower-level `request()` returns a Node `Readable` instead, which Bun
* implements natively and streams correctly; `Readable.toWeb` bridges it back to
* a spec `Response`. Buffered reads (`.json()`/`.text()`/`.arrayBuffer()`) behave
* identically on both runtimes, so this is a drop-in substitute.
*
* SSRF is unchanged: the same `dispatcher` (Agent carrying the guarded/pinned
* `connect.lookup`) governs every connection, and `maxResponseSize` still caps the
* body. Redirects are NOT followed here (`maxRedirections: 0`); the caller drives
* them via {@link followRedirectsGuarded}, exactly as it did over `fetch`'s
* `redirect: 'manual'`.
*/
async function undiciRequestAsResponse(
input: RequestInfo | URL,
init: RequestInit,
dispatcher: Dispatcher
): Promise<Response> {
let url: string
let effectiveInit = init as UndiciRequestInit
if (typeof Request !== 'undefined' && input instanceof Request) {
// A Request input carries its own method/headers/body/signal; lift them (explicit
// init fields win, per fetch semantics) so a guarded POST isn't downgraded to GET.
const bodyAllowed = input.method !== 'GET' && input.method !== 'HEAD'
effectiveInit = {
method: input.method,
headers: input.headers,
body: bodyAllowed ? await input.clone().arrayBuffer() : undefined,
signal: input.signal,
...(init as UndiciRequestInit),
// double-cast-allowed: DOM RequestInit and undici RequestInit differ in TS but match at runtime
} as unknown as UndiciRequestInit
url = input.url
} else {
url = typeof input === 'string' ? input : input instanceof URL ? input.href : input.url
}
const method = (effectiveInit.method ?? 'GET').toUpperCase()
const canHaveBody = method !== 'GET' && method !== 'HEAD'
const requestHeaders = toUndiciRequestHeaders(effectiveInit.headers) ?? {}
const requestBody = canHaveBody ? toUndiciRequestBody(effectiveInit.body) : undefined
// fetch auto-adds a form content-type for a URLSearchParams body; preserve that parity
// when the caller didn't set one (the MCP SDK does set it explicitly, but not every caller).
if (
canHaveBody &&
effectiveInit.body instanceof URLSearchParams &&
!Object.keys(requestHeaders).some((key) => key.toLowerCase() === 'content-type')
) {
requestHeaders['content-type'] = 'application/x-www-form-urlencoded;charset=UTF-8'
}
const { statusCode, headers, body } = await undiciRequest(url, {
method: method as Dispatcher.HttpMethod,
headers: requestHeaders,
body: requestBody,
signal: effectiveInit.signal ?? undefined,
dispatcher,
// No `maxRedirections`: request() does not auto-follow by default, so the caller's
// `followRedirectsGuarded` drives every hop with per-hop SSRF validation.
})
const responseHeaders = new Headers()
for (const [key, value] of Object.entries(headers)) {
if (Array.isArray(value)) for (const v of value) responseHeaders.append(key, v)
else if (value != null) responseHeaders.append(key, value)
}
// Null-body statuses (204/205/304) can't carry a body; drain undici's (empty) stream so its
// socket returns to the pool. Attach an error listener first so a socket reset mid-drain
// surfaces as a handled event, not an unhandled 'error' that crashes the process.
const isNullBody = statusCode === 204 || statusCode === 205 || statusCode === 304
if (isNullBody) {
body.on('error', () => {})
body.resume()
const response = new Response(null, { status: statusCode, headers: responseHeaders })
Object.defineProperty(response, 'url', { value: url, configurable: true })
return response
}
// Decode Content-Encoding like `fetch` does (`request()` returns raw bytes). `maxResponseSize`
// still caps the compressed wire bytes on `body`.
const contentEncoding = String(headers['content-encoding'] ?? '')
.toLowerCase()
.trim()
const decoder = contentEncodingDecoder(contentEncoding)
if (decoder) {
// The bridged body is now decoded; drop framing headers that would misdescribe it.
responseHeaders.delete('content-encoding')
responseHeaders.delete('content-length')
}
// Build the bridge over the stream the consumer reads (the decoder when decoding).
// `nodeReadableToWebStream` attaches its `error` listener synchronously, so wiring the pipe
// AFTER it means a synchronous zlib error (e.g. a server mislabeling a non-gzip body as gzip)
// is caught and rejects the reader instead of taking down the process.
const webBody = nodeReadableToWebStream(decoder ?? body)
if (decoder) {
body.once('error', (err) => decoder.destroy(err)) // forward maxResponseSize / socket reset
decoder.once('close', () => body.destroy()) // tear the source down so the socket can't leak
body.pipe(decoder)
}
try {
const response = new Response(webBody, { status: statusCode, headers: responseHeaders })
// undici.request never sets `url`; `fetch` did, and consumers rely on it (the MCP
// transport's response-cap wrapper copies it; the SDK resolves relative
// auth-metadata URLs against it). Preserve parity.
Object.defineProperty(response, 'url', { value: url, configurable: true })
return response
} catch (err) {
// `new Response` rejects an out-of-range status (a 1xx undici shouldn't surface, but
// defensively): destroy the source so its socket can't leak, then rethrow.
body.destroy()
throw err
}
}
/**
* Normalizes a `fetch(input, init)` call into a URL string + init. A `Request` input carries
* its own method/headers/body/signal; lift them into the init (explicit init fields win, per
* fetch semantics) so a manual redirect follower can't silently downgrade a POST Request to a
* bare GET or lose its headers.
*/
async function liftFetchArgs(
input: RequestInfo | URL,
init?: RequestInit
): Promise<{ target: string; effectiveInit: RequestInit }> {
const target = typeof input === 'string' ? input : input instanceof URL ? input.href : input.url
if (typeof Request !== 'undefined' && input instanceof Request) {
const bodyAllowed = input.method !== 'GET' && input.method !== 'HEAD'
return {
target,
effectiveInit: {
method: input.method,
headers: input.headers,
body: bodyAllowed ? await input.clone().arrayBuffer() : undefined,
signal: input.signal,
// Carry the Request's redirect mode so the pinned fetch honors `manual`/`error`
// instead of defaulting a `Request({ redirect: 'manual' })` to `follow`.
redirect: input.redirect,
...init,
},
}
}
return { target, effectiveInit: init ?? {} }
}
/**
* SSRF-guarded `fetch` + its `Agent` for outbound requests to user-controlled
* hosts: DNS resolves normally, and every socket connect validates the chosen
* addresses via {@link createSsrfGuardedLookup}; redirects are followed manually
* with per-hop validation (see {@link followRedirectsGuarded}) so IP-literal
* targets can't bypass the lookup and custom headers never cross origins. See
* {@link createPinnedFetchWithDispatcher} for the `maxResponseSize` semantics.
*/
export function createSsrfGuardedFetchWithDispatcher(options?: { maxResponseSize?: number }): {
fetch: typeof fetch
dispatcher: Agent
} {
const dispatcher = new Agent({
allowH2: false,
connect: { lookup: createSsrfGuardedLookup() },
...(options?.maxResponseSize !== undefined ? { maxResponseSize: options.maxResponseSize } : {}),
})
const rawFetch = (url: string, init: UndiciRequestInit): Promise<Response> =>
// double-cast-allowed: DOM RequestInit and undici RequestInit differ in TS but match at runtime
undiciRequestAsResponse(url, init as unknown as RequestInit, dispatcher)
const guarded = async (input: RequestInfo | URL, init?: RequestInit): Promise<Response> => {
const { target, effectiveInit } = await liftFetchArgs(input, init)
// double-cast-allowed: DOM RequestInit and undici RequestInit are structurally compatible at runtime but the TS types differ
return followRedirectsGuarded(rawFetch, target, effectiveInit as unknown as UndiciRequestInit)
}
return { fetch: guarded, dispatcher }
}
/**
* Builds a standard `fetch`-compatible function that pins every outbound
* connection to `resolvedIP`, preventing DNS-rebinding (TOCTOU) between URL
* validation and connection. The original hostname is preserved for TLS SNI and
* the `Host` header so it still matches the certificate. This is the single
* source of truth for pinned outbound fetches — both the LLM providers and the
* MCP transport consume it.
*
* Pass the returned function as the `fetch` option to the OpenAI/Anthropic SDKs
* (or call it directly) after validating the URL with {@link validateUrlWithDNS}
* and capturing `resolvedIP`. Because the pinned lookup always returns
* `resolvedIP` regardless of hostname, any redirect the server returns also
* connects to the validated IP — an attacker cannot rebind a redirect target to
* an internal address.
*
* The `Agent` is captured for the lifetime of the returned function, so repeated
* calls (e.g. a provider tool loop) reuse its keep-alive connections.
*
* `allowH2` opts the pinned Agent into HTTP/2 (ALPN-negotiated, h1.1 fallback).
* It defaults to `false` to leave existing consumers unchanged. Enabling it does
* not weaken pinning: the pinned `connect.lookup` forces every connection on the
* Agent to `resolvedIP` regardless of authority, so h2 connection coalescing can
* never reach an address other than the validated one.
*/
export function createPinnedFetch(
resolvedIP: string,
options?: { allowH2?: boolean }
): typeof fetch {
return createPinnedFetchWithDispatcher(resolvedIP, options).fetch
}
/**
* Same as {@link createPinnedFetch} but also returns the underlying `Agent` so a
* caller with a defined connection lifetime (e.g. a long-lived MCP transport) can
* tear the Agent down on close instead of waiting for its idle timeout. Closing
* the Agent is what releases any pooled keep-alive / HTTP/2 sockets it holds.
*
* `maxResponseSize` caps the (decoded) response body in bytes and makes undici reject
* with `UND_ERR_RES_EXCEEDED_MAX_SIZE` once exceeded — a DoS backstop for one-shot
* callers reading from a URL taken from untrusted metadata. Omit it (the default) to
* leave the response unbounded, which streaming consumers like the MCP transport need.
*/
export function createPinnedFetchWithDispatcher(
resolvedIP: string,
options?: { allowH2?: boolean; maxResponseSize?: number }
): { fetch: typeof fetch; dispatcher: Agent } {
const dispatcher = new Agent({
allowH2: options?.allowH2 ?? false,
connect: { lookup: createPinnedLookup(resolvedIP) },
...(options?.maxResponseSize !== undefined ? { maxResponseSize: options.maxResponseSize } : {}),
})
const rawFetch = (url: string, init: UndiciRequestInit): Promise<Response> =>
// double-cast-allowed: DOM RequestInit and undici RequestInit differ in TS but match at runtime
undiciRequestAsResponse(url, init as unknown as RequestInit, dispatcher)
// Requests go through `undici.request` (not `undici.fetch`) because fetch's streaming
// `response.body` never delivers under the Bun runtime the server runs on — the same bug
// {@link createSsrfGuardedFetchWithDispatcher} works around. Redirects are handled here (not
// by a caller's wrapper — the pinned fetch is passed straight to provider/A2A SDKs), honoring
// the request's `redirect` mode: `manual`/`error` must NOT transparently follow (e.g.
// `detectMcpAuthType` inspects the 3xx to classify auth). The default `follow` uses
// {@link followRedirectsGuarded}, which drops headers on cross-origin hops (so a redirect
// can't disclose a provider `api-key` to another origin) and stamps the final `response.url`.
// Every hop still dispatches through the pinned `Agent` (its `connect.lookup` forces
// `resolvedIP`), so a redirect can't escape to another address.
const pinned = async (input: RequestInfo | URL, init?: RequestInit): Promise<Response> => {
const { target, effectiveInit } = await liftFetchArgs(input, init)
const mode = effectiveInit.redirect ?? 'follow'
// double-cast-allowed: DOM RequestInit and undici RequestInit are structurally compatible at runtime but the TS types differ
const undiciInit = effectiveInit as unknown as UndiciRequestInit
if (mode === 'manual') {
return rawFetch(target, undiciInit)
}
if (mode === 'error') {
const response = await rawFetch(target, undiciInit)
const location = response.headers.get('location')
if (response.status >= 300 && response.status < 400 && location) {
await response.body?.cancel().catch(() => {})
throw new TypeError('Pinned fetch received an unexpected redirect (redirect: "error")')
}
return response
}
// Permit this pinned IP as a redirect/initial target even when it's private (the
// self-hosted MCP carve-out on a private/loopback IP, and same-host redirects that stay on
// it) — otherwise the guarded policy would block a self-hosted server reaching itself. Any
// OTHER private target (e.g. a redirect to the cloud metadata IP) is still blocked.
return followRedirectsGuarded(rawFetch, target, undiciInit, { allowRedirectToIp: resolvedIP })
}
return { fetch: pinned, dispatcher }
}
/**
* Performs a fetch with IP pinning to prevent DNS rebinding attacks.
* Uses the pre-resolved IP address while preserving the original hostname for TLS SNI.
* Follows redirects securely by validating each redirect target.
*
* The response body is always bounded — `options.maxResponseBytes` when supplied (and
* positive), otherwise {@link DEFAULT_MAX_RESPONSE_BYTES}. Exceeding the cap rejects with
* a {@link PayloadSizeLimitError} and destroys the socket.
*/
export async function secureFetchWithPinnedIP(
url: string,
resolvedIP: string,
options: SecureFetchOptions & { allowHttp?: boolean } = {},
redirectCount = 0
): Promise<SecureFetchResponse> {
const maxRedirects = options.maxRedirects ?? DEFAULT_MAX_REDIRECTS
const requestedMaxResponseBytes = options.maxResponseBytes
const maxResponseBytes =
typeof requestedMaxResponseBytes === 'number' && requestedMaxResponseBytes > 0
? requestedMaxResponseBytes
: DEFAULT_MAX_RESPONSE_BYTES
return new Promise((resolve, reject) => {
const parsed = new URL(url)
const isHttps = parsed.protocol === 'https:'
const defaultPort = isHttps ? 443 : 80
const port = parsed.port ? Number.parseInt(parsed.port, 10) : defaultPort
let agent: http.Agent
if (options.proxyUrl) {
// Proxy connection is already IP-pinned by validateAndPinProxyUrl; target-IP
// pinning is intentionally bypassed (the proxy resolves the target). https
// targets tunnel via CONNECT, http targets use absolute-URI forwarding.
agent = isHttps ? new HttpsProxyAgent(options.proxyUrl) : new HttpProxyAgent(options.proxyUrl)
} else {
const lookup = createPinnedLookup(resolvedIP)
const agentOptions: http.AgentOptions = { lookup }
agent = isHttps ? new https.Agent(agentOptions) : new http.Agent(agentOptions)
}
const { 'accept-encoding': _, ...sanitizedHeaders } = options.headers ?? {}
const requestOptions: http.RequestOptions = {
hostname: parsed.hostname,
port,
path: parsed.pathname + parsed.search,
method: options.method || 'GET',
headers: sanitizedHeaders,
agent,
timeout: options.timeout || 300000,
}
const protocol = isHttps ? https : http
const req = protocol.request(requestOptions, (res) => {
const statusCode = res.statusCode || 0
const location = res.headers.location
if (isRedirectStatus(statusCode) && location && redirectCount < maxRedirects) {
res.resume()
const redirectUrl = resolveRedirectUrl(url, location)
validateUrlWithDNS(redirectUrl, 'redirectUrl', { allowHttp: options.allowHttp })
.then((validation) => {
if (!validation.isValid) {
settledReject(new Error(`Redirect blocked: ${validation.error}`))
return
}
const redirectOptions = options.stripAuthOnRedirect
? {
...options,
headers: omit(options.headers ?? {}, ['Authorization', 'authorization']),
}
: options
return secureFetchWithPinnedIP(
redirectUrl,
validation.resolvedIP!,
redirectOptions,
redirectCount + 1
)
})
.then((response) => {
if (response) settledResolve(response)
})
.catch(settledReject)
return
}
if (isRedirectStatus(statusCode) && location && redirectCount >= maxRedirects) {
res.resume()
settledReject(new Error(`Too many redirects (max: ${maxRedirects})`))
return
}
const headersRecord: Record<string, string> = {}
let setCookieArray: string[] = []
for (const [key, value] of Object.entries(res.headers)) {
const lowerKey = key.toLowerCase()
if (lowerKey === 'set-cookie') {
if (Array.isArray(value)) {
setCookieArray = value
headersRecord[lowerKey] = value.join(', ')
} else if (typeof value === 'string') {
setCookieArray = [value]
headersRecord[lowerKey] = value
}
} else if (typeof value === 'string') {
headersRecord[lowerKey] = value
} else if (Array.isArray(value)) {
headersRecord[lowerKey] = value.join(', ')
}
}
// Responses that carry no body (HEAD, 204, 304) may still advertise the resource's full
// size in content-length. That is metadata, not a payload, so it must not trip the cap —
// otherwise a HEAD probe of a large file, or a conditional-GET 304, would fail spuriously.
const isBodylessResponse =
(requestOptions.method || 'GET').toUpperCase() === 'HEAD' ||
statusCode === 204 ||
statusCode === 304
const contentLength = headersRecord['content-length']
if (contentLength && !isBodylessResponse) {
const parsedLength = Number.parseInt(contentLength, 10)
if (Number.isFinite(parsedLength) && parsedLength > maxResponseBytes) {
cleanupAbort()
res.destroy()
req.destroy()
if (isRetryableHttpStatus(statusCode)) {
settledResolve({
ok: false,
status: statusCode,
statusText: res.statusMessage || '',
headers: new SecureFetchHeaders(headersRecord, setCookieArray),
body: null,
text: async () => '',
json: async () => ({}),
arrayBuffer: async () => new ArrayBuffer(0),
})
return
}
settledReject(
new PayloadSizeLimitError({
label: 'response body',
maxBytes: maxResponseBytes,
observedBytes: parsedLength,
})
)
return
}
}
let totalBytes = 0
const nodeRes = res
const body = new ReadableStream<Uint8Array>({
start(controller) {
nodeRes.on('data', (chunk: Buffer) => {
totalBytes += chunk.length
if (totalBytes > maxResponseBytes) {
cleanupAbort()
controller.error(
new PayloadSizeLimitError({
label: 'response body',
maxBytes: maxResponseBytes,
observedBytes: totalBytes,
})
)
nodeRes.destroy()
return
}
controller.enqueue(new Uint8Array(chunk))
})
nodeRes.on('end', () => {
cleanupAbort()
controller.close()
})
nodeRes.on('error', (err) => {
cleanupAbort()
controller.error(err)
})
},
cancel() {
cleanupAbort()
nodeRes.destroy()
},
})
let bodyBufferPromise: Promise<Buffer> | null = null
function readBodyAsBuffer(): Promise<Buffer> {
if (!bodyBufferPromise) {
bodyBufferPromise = (async () => {
const reader = body.getReader()
const buffers: Uint8Array[] = []
while (true) {
const { done, value } = await reader.read()
if (done) break
if (value) buffers.push(value)
}
return Buffer.concat(buffers.map((b) => Buffer.from(b)))
})()
}
return bodyBufferPromise
}
settledResolve({
ok: statusCode >= 200 && statusCode < 300,
status: statusCode,
statusText: res.statusMessage || '',
headers: new SecureFetchHeaders(headersRecord, setCookieArray),
body,
text: async () => (await readBodyAsBuffer()).toString('utf-8'),
json: async () => JSON.parse((await readBodyAsBuffer()).toString('utf-8')),
arrayBuffer: async () => {
const buf = await readBodyAsBuffer()
return buf.buffer.slice(buf.byteOffset, buf.byteOffset + buf.byteLength) as ArrayBuffer
},
})
})
let onAbort: (() => void) | null = null
const cleanupAbort = () => {
if (onAbort && options.signal) {
options.signal.removeEventListener('abort', onAbort)
onAbort = null
}
}
const settledResolve: typeof resolve = (value) => {
resolve(value)
}
const settledReject: typeof reject = (reason) => {
cleanupAbort()
reject(reason)
}
req.on('error', (error) => {
settledReject(error)
})
req.on('timeout', () => {
req.destroy()
settledReject(new Error(`Request timed out after ${requestOptions.timeout}ms`))
})
if (options.signal) {
if (options.signal.aborted) {
req.destroy()
settledReject(options.signal.reason ?? new Error('Aborted'))
return
}
onAbort = () => {
req.destroy()
settledReject(options.signal?.reason ?? new Error('Aborted'))
}
options.signal.addEventListener('abort', onAbort, { once: true })
}
if (options.body) {
req.write(options.body)
}
req.end()
})
}
/**
* Validates a URL and performs a secure fetch with DNS pinning in one call.
* Combines validateUrlWithDNS and secureFetchWithPinnedIP for convenience.
*
* @param url - The URL to fetch
* @param options - Fetch options (method, headers, body, etc.)
* @param paramName - Name of the parameter for error messages (default: 'url')
* @returns SecureFetchResponse
* @throws Error if URL validation fails
*/
export async function secureFetchWithValidation(
url: string,
options: SecureFetchOptions & { allowHttp?: boolean } = {},
paramName = 'url'
): Promise<SecureFetchResponse> {
const validation = await validateUrlWithDNS(url, paramName, {
allowHttp: options.allowHttp,
})
if (!validation.isValid) {
throw new Error(validation.error)
}
return secureFetchWithPinnedIP(url, validation.resolvedIP!, options)
}