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
1239 lines
48 KiB
TypeScript
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)
|
|
}
|