#![warn(unused_must_use)]
#[path = "AsyncHTTP.rs"]
pub mod async_http;
#[path = "CertificateInfo.rs"]
pub mod certificate_info;
#[path = "TlsInfo.rs"]
pub mod tls_info;
#[path = "Decompressor.rs"]
pub mod decompressor;
#[path = "H2Client.rs"]
pub mod h2_client;
pub use bun_http_types::h2 as h2_frame_parser;
#[path = "H3Client.rs"]
pub mod h3_client;
#[path = "HeaderBuilder.rs"]
pub mod header_builder;
#[path = "HeaderValueIterator.rs"]
pub mod header_value_iterator;
#[path = "Headers.rs"]
pub mod headers;
#[path = "HTTPCertError.rs"]
pub mod http_cert_error;
#[path = "HTTPContext.rs"]
pub mod http_context;
#[path = "HTTPRequestBody.rs"]
pub mod http_request_body;
#[path = "HTTPThread.rs"]
pub mod http_thread;
#[path = "InitError.rs"]
pub mod init_error;
#[path = "InternalState.rs"]
pub mod internal_state;
#[path = "lshpack.rs"]
pub mod lshpack;
#[path = "ProxyTunnel.rs"]
pub mod proxy_tunnel;
#[path = "SendFile.rs"]
pub mod send_file;
#[path = "Signals.rs"]
pub mod signals;
#[path = "ThreadSafeStreamBuffer.rs"]
pub mod thread_safe_stream_buffer;
#[path = "websocket.rs"]
pub mod websocket;
#[path = "websocket_http_client.rs"]
pub mod websocket_http_client;
#[path = "zlib.rs"]
pub mod zlib;
pub use async_http::AsyncHTTP;
pub use certificate_info::CertificateInfo;
pub use tls_info::BunTlsInfo;
pub use decompressor::Decompressor;
pub use header_builder::HeaderBuilder;
pub use headers::{Headers, HeadersExt};
pub use http_cert_error::HTTPCertError;
pub use http_context::{HTTPContext, HTTPSocket};
pub use http_request_body::HTTPRequestBody;
pub use http_thread::HttpThread as HTTPThread;
pub use http_thread::{defer_shutdown_reclaim, shutdown_for_exit};
pub use internal_state::InternalState;
pub use proxy_tunnel::ProxyTunnel;
pub use send_file::SendFile;
pub use signals::Signals;
pub use thread_safe_stream_buffer::ThreadSafeStreamBuffer;
#[path = "ssl_config.rs"]
pub mod ssl_config;
pub use ssl_config::SSLConfig;
pub use bun_uws::ssl_wrapper;
pub use bun_uws::ssl_wrapper::SSLWrapper;
pub use h2_client as h2;
pub use h2_client as H2;
pub use h3_client as h3;
pub use h3_client as H3;
pub use http_context as new_http_context;
pub type NewHTTPContext<const SSL: bool> = http_context::HTTPContext<SSL>;
pub type NewHttpContext<const SSL: bool> = http_context::HTTPContext<SSL>;
pub type HttpsContext = http_context::HTTPContext<true>;
pub type HttpContext = http_context::HTTPContext<false>;
pub type HttpClient<'a> = HTTPClient<'a>;
pub type HttpThread = HTTPThread;
pub type AsyncHttp<'a> = AsyncHTTP<'a>;
pub type ThreadlocalAsyncHttp<'a> = ThreadlocalAsyncHTTP<'a>;
pub use HTTPClientResult as http_client_result;
pub use bun_http_types::FetchRedirect::FetchRedirect;
pub use bun_http_types::Method::Method;
pub use bun_picohttp as picohttp;
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Default)]
pub enum HTTPVerboseLevel {
#[default]
None,
Headers,
Curl,
}
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Default)]
pub enum Protocol {
#[default]
Http1_1,
Http2,
Http3,
}
pub use bun_http_types::Encoding::Encoding;
pub use header_value_iterator::HeaderValueIterator;
pub use init_error::InitError;
pub const EXTREMELY_VERBOSE: bool = false;
pub struct HTTPResponseMetadata {
pub url: bun_ptr::RawSlice<u8>,
pub owned_buf: Box<[u8]>,
pub response: bun_picohttp::Response<'static>,
}
impl Default for HTTPResponseMetadata {
fn default() -> Self {
Self {
url: bun_ptr::RawSlice::EMPTY,
owned_buf: Box::default(),
response: bun_picohttp::Response::default(),
}
}
}
impl Drop for HTTPResponseMetadata {
fn drop(&mut self) {
let list = self.response.headers.list;
if !list.is_empty() {
unsafe { bun_core::heap::destroy(core::ptr::from_ref(list).cast_mut()) };
}
self.response.headers = bun_picohttp::HeaderList::default();
self.response.status = b"";
}
}
pub use bun_http_types::{ETag, FetchCacheMode, FetchRequestMode, MimeType, URLPath};
use bun_core::MutableString;
use bun_http_types::FetchRedirect::CommonAbortReason;
use core::sync::atomic::{AtomicBool, AtomicU32, AtomicUsize, Ordering};
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Default)]
pub enum HTTPUpgradeState {
#[default]
None = 0,
Pending = 1,
Upgraded = 2,
}
#[derive(Clone, Copy)]
pub struct Flags {
pub disable_timeout: bool,
pub disable_keepalive: bool,
pub disable_decompression: bool,
pub did_have_handshaking_error: bool,
pub force_last_modified: bool,
pub redirected: bool,
pub proxy_tunneling: bool,
pub reject_unauthorized: bool,
pub is_preconnect_only: bool,
pub is_streaming_request_body: bool,
pub defer_fail_until_connecting_is_complete: bool,
pub upgrade_state: HTTPUpgradeState,
pub protocol: Protocol,
pub force_http2: bool,
pub force_http1: bool,
pub force_http3: bool,
pub h3_retried: bool,
pub h3_alt_svc_fallback: bool,
pub is_node_http_client: bool,
pub omit_connection_header: bool,
pub is_page_egress: bool,
}
impl Default for Flags {
fn default() -> Self {
Self {
disable_timeout: false,
disable_keepalive: false,
disable_decompression: false,
did_have_handshaking_error: false,
force_last_modified: false,
redirected: false,
proxy_tunneling: false,
reject_unauthorized: true,
is_preconnect_only: false,
is_streaming_request_body: false,
defer_fail_until_connecting_is_complete: false,
upgrade_state: HTTPUpgradeState::None,
protocol: Protocol::Http1_1,
force_http2: false,
force_http1: false,
force_http3: false,
h3_retried: false,
h3_alt_svc_fallback: false,
is_node_http_client: false,
is_page_egress: false,
omit_connection_header: false,
}
}
}
pub static ASYNC_HTTP_ID_MONOTONIC: AtomicU32 = AtomicU32::new(1);
pub static EXPERIMENTAL_HTTP3_CLIENT_FROM_CLI: AtomicBool = AtomicBool::new(false);
const MAX_REDIRECT_URL_LENGTH: usize = 128 * 1024;
#[unsafe(export_name = "BUN_DEFAULT_MAX_HTTP_HEADER_SIZE")]
pub static MAX_HTTP_HEADER_SIZE: AtomicUsize = AtomicUsize::new(16 * 1024);
#[inline]
pub fn max_http_header_size() -> usize {
MAX_HTTP_HEADER_SIZE.load(Ordering::Relaxed)
}
#[inline]
pub fn set_max_http_header_size(v: usize) {
MAX_HTTP_HEADER_SIZE.store(v, Ordering::Relaxed);
}
pub static OVERRIDDEN_DEFAULT_USER_AGENT: std::sync::OnceLock<&'static [u8]> =
std::sync::OnceLock::new();
pub static IDLE_TIMEOUT_SECONDS: AtomicU32 = AtomicU32::new(300);
#[inline]
pub fn idle_timeout_seconds() -> c_uint {
IDLE_TIMEOUT_SECONDS.load(Ordering::Relaxed)
}
pub const END_OF_CHUNKED_HTTP1_1_ENCODING_RESPONSE_BODY: &[u8] = b"0\r\n\r\n";
pub static TEMP_HOSTNAME: bun_core::RacyCell<[u8; 8192]> = bun_core::RacyCell::new([0; 8192]);
const MAX_TLS_RECORD_SIZE: usize = 16 * 1024;
pub const MAX_H2_RETRIES: u8 = 5;
const PREALLOCATE_MAX: usize = 1024 * 1024 * 256;
#[inline]
pub fn cleanup(_force: bool) {
}
pub fn h3_alt_svc_enabled() -> bool {
let cli = EXPERIMENTAL_HTTP3_CLIENT_FROM_CLI.load(Ordering::Relaxed);
cli || bun_core::env_var::feature_flag::BUN_FEATURE_FLAG_EXPERIMENTAL_HTTP3_CLIENT
.get()
.unwrap_or(false)
}
pub fn strip_port_from_host(host: &[u8]) -> &[u8] {
if host.is_empty() {
return host;
}
if host[0] == b'[' {
if let Some(bracket) = host.iter().rposition(|&b| b == b']') {
return &host[0..bracket + 1];
}
return host;
}
if let Some(colon) = host.iter().rposition(|&b| b == b':') {
return &host[0..colon];
}
host
}
#[derive(Copy, Clone, PartialEq, Eq)]
pub enum ShouldContinue {
ContinueStreaming,
Finished,
}
#[derive(Copy, Clone, Eq, PartialEq)]
pub(crate) enum HeaderResult {
HasBody,
Finished,
}
impl HTTPClient<'_> {
#[inline]
pub(crate) fn apply_multiplexed_headers(
&mut self,
status_code: u32,
headers: &[picohttp::Header],
) -> Result<HeaderResult, bun_core::Error> {
let mut response = picohttp::Response {
minor_version: 0,
status_code,
status: b"",
headers: picohttp::HeaderList { list: headers },
bytes_read: 0,
};
self.state.pending_response = Some(unsafe { response.detach_lifetime() });
let should_continue = self.handle_response_metadata(&mut response)?;
self.state.pending_response = Some(unsafe { response.detach_lifetime() });
self.state.transfer_encoding = Encoding::Identity;
if self.state.response_stage == ResponseStage::BodyChunk {
self.state.response_stage = ResponseStage::Body;
}
self.state.flags.allow_keepalive = true;
Ok(if should_continue == ShouldContinue::Finished {
HeaderResult::Finished
} else {
HeaderResult::HasBody
})
}
}
#[derive(Default, Copy, Clone)]
pub enum BodySize {
TotalReceived(usize),
ContentLength(usize),
#[default]
Unknown,
}
#[derive(Default)]
pub struct HTTPClientResult<'a> {
pub body: Option<&'a mut MutableString>,
pub has_more: bool,
pub redirected: bool,
pub can_stream: bool,
pub is_http2: bool,
pub fail: Option<bun_core::Error>,
pub metadata: Option<HTTPResponseMetadata>,
pub body_size: BodySize,
pub certificate_info: Option<CertificateInfo>,
pub tls_info: Option<BunTlsInfo>,
}
impl<'a> HTTPClientResult<'a> {
pub fn abort_reason(&self) -> Option<CommonAbortReason> {
if self.is_timeout() {
return Some(CommonAbortReason::Timeout);
}
if self.is_abort() {
return Some(CommonAbortReason::UserAbort);
}
None
}
pub fn is_success(&self) -> bool {
self.fail.is_none()
}
pub fn is_timeout(&self) -> bool {
matches!(self.fail, Some(e) if e == bun_core::err!("Timeout"))
}
pub fn is_abort(&self) -> bool {
matches!(self.fail, Some(e) if e == bun_core::err!("Aborted") || e == bun_core::err!("AbortedBeforeConnecting"))
}
#[inline]
pub unsafe fn detach_lifetime(self) -> HTTPClientResult<'static> {
HTTPClientResult {
body: self
.body
.map(|b| unsafe { &mut *core::ptr::from_mut::<MutableString>(b) }),
has_more: self.has_more,
redirected: self.redirected,
can_stream: self.can_stream,
is_http2: self.is_http2,
fail: self.fail,
metadata: self.metadata,
body_size: self.body_size,
certificate_info: self.certificate_info,
tls_info: self.tls_info,
}
}
}
pub type HTTPClientResultCallbackFunction =
fn(*mut (), *mut AsyncHTTP<'static>, HTTPClientResult<'_>);
#[derive(Copy, Clone)]
pub struct HTTPClientResultCallback {
pub ctx: *mut (),
pub function: HTTPClientResultCallbackFunction,
pub release_at_shutdown: Option<unsafe fn(*mut ())>,
}
impl HTTPClientResultCallback {
pub fn run(self, async_http: *mut AsyncHTTP<'static>, result: HTTPClientResult<'_>) {
(self.function)(self.ctx, async_http, result);
}
pub fn new<T>(
this: *mut T,
callback: fn(*mut T, *mut AsyncHTTP<'static>, HTTPClientResult<'_>),
) -> Self {
Self {
ctx: this.cast::<()>(),
function: unsafe {
bun_ptr::cast_fn_ptr::<
fn(*mut T, *mut AsyncHTTP<'static>, HTTPClientResult<'_>),
HTTPClientResultCallbackFunction,
>(callback)
},
release_at_shutdown: None,
}
}
pub fn new_with_release<T>(
this: *mut T,
callback: fn(*mut T, *mut AsyncHTTP<'static>, HTTPClientResult<'_>),
release: unsafe fn(*mut ()),
) -> Self {
let mut cb = Self::new(this, callback);
cb.release_at_shutdown = Some(release);
cb
}
}
pub struct ThreadlocalAsyncHTTP<'a> {
pub async_http: AsyncHTTP<'a>,
}
impl<'a> ThreadlocalAsyncHTTP<'a> {
pub fn new(async_http: AsyncHTTP<'a>) -> Box<Self> {
Box::new(Self { async_http })
}
}
pub trait SocketTimeout {
fn timeout(&self, seconds: core::ffi::c_uint);
fn set_timeout_minutes(&self, minutes: core::ffi::c_uint);
fn set_timeout(&self, seconds: core::ffi::c_uint);
}
pub fn hash_header_name(name: &[u8]) -> u64 {
bun_wyhash::hash_ascii_lowercase(0, name)
}
use bun_core::ZigStringSlice;
use bun_url::URL;
use core::ptr::NonNull;
pub fn intern_extension_method(name: &[u8]) -> &'static [u8] {
use std::sync::Mutex;
static EXTENSION_METHODS: Mutex<Vec<&'static [u8]>> = Mutex::new(Vec::new());
let mut table = EXTENSION_METHODS.lock().unwrap();
if let Some(found) = table.iter().find(|token| **token == name) {
return found;
}
let leaked: &'static [u8] = Box::leak(name.to_vec().into_boxed_slice());
table.push(leaked);
leaked
}
pub struct HTTPClient<'a> {
pub method: Method,
pub extension_method: Option<&'static [u8]>,
pub header_entries: headers::EntryList,
pub header_buf: &'a [u8],
pub url: URL<'a>,
pub connected_url: URL<'a>,
pub verbose: HTTPVerboseLevel,
pub remaining_redirect_count: i8,
pub allow_retry: bool,
pub h2_retries: u8,
pub redirect_type: FetchRedirect,
pub redirect: Vec<u8>,
pub prev_redirect: Vec<u8>,
pub progress_node: Option<NonNull<bun_core::Progress::Node>>,
pub flags: Flags,
pub state: InternalState<'a>,
pub tls_props: Option<ssl_config::SharedPtr>,
pub custom_ssl_ctx: Option<http_context::HTTPContextRc<true>>,
pub result_callback: HTTPClientResultCallback,
pub if_modified_since: &'a [u8],
pub request_content_len_buf: [u8; b"-4294967295".len()],
pub http_proxy: Option<URL<'a>>,
pub proxy_headers: Option<Headers>,
pub proxy_authorization: Option<Vec<u8>>,
pub proxy_tunnel: Option<proxy_tunnel::RefPtr>,
pub h2: Option<NonNull<h2::Stream>>,
pub h3: Option<NonNull<h3::Stream>>,
pub pending_h2: Option<NonNull<h2::PendingConnect>>,
pub signals: Signals,
pub async_http_id: u32,
pub hostname: Option<&'a [u8]>,
pub unix_socket_path: ZigStringSlice,
pub tls_info: Option<BunTlsInfo>,
}
impl<'a> HTTPClient<'a> {
#[inline(always)]
pub fn as_erased_ptr(&self) -> NonNull<HTTPClient<'static>> {
NonNull::from(self).cast::<HTTPClient<'static>>()
}
#[inline(always)]
pub(crate) fn from_erased_backref<'b>(
p: NonNull<HTTPClient<'static>>,
) -> &'b mut HTTPClient<'static> {
unsafe { &mut *p.as_ptr() }
}
}
impl Drop for HTTPClient<'_> {
fn drop(&mut self) {
self.close_proxy_tunnel(false);
debug_assert!(self.h2.is_none());
if let Some(ctx) = self.custom_ssl_ctx.take() {
ctx.deref();
}
self.unix_socket_path = ZigStringSlice::EMPTY;
}
}
pub static HTTP_THREAD: bun_core::ThreadCell<core::mem::MaybeUninit<HTTPThread>> =
bun_core::ThreadCell::new(core::mem::MaybeUninit::uninit());
pub(crate) static HTTP_THREAD_INIT: core::sync::atomic::AtomicBool =
core::sync::atomic::AtomicBool::new(false);
#[inline]
pub fn http_thread() -> &'static mut HTTPThread {
assert!(
HTTP_THREAD_INIT.load(core::sync::atomic::Ordering::Acquire),
"http_thread() called before HTTPThread::init()"
);
unsafe { (*HTTP_THREAD.get()).assume_init_mut() }
}
#[inline]
pub fn http_thread_mut() -> &'static mut HTTPThread {
http_thread()
}
pub static SOCKET_ASYNC_HTTP_ABORT_TRACKER: bun_core::RacyCell<
Option<bun_collections::ArrayHashMap<u32, bun_uws::AnySocket>>,
> = bun_core::RacyCell::new(None);
use core::ffi::c_uint;
use bstr::BStr;
use bun_boringssl as boringssl;
use bun_collections::{ArrayHashMap, VecExt};
use bun_core::StringBuilder;
use bun_core::{FeatureFlags, Global, Output, err};
use bun_core::{OwnedString, String as BunString, Tag as BunStringTag, immutable as strings};
use bun_uws as uws;
use bun_http_types::ETag::StringPointer;
use bun_wyhash::Wyhash11 as Wyhash;
use crate::http_context::HTTPSocket as HttpSocket;
use crate::internal_state::{RequestStage, ResponseStage, Stage};
bun_core::declare_scope!(fetch, visible);
pub type GenHttpContext<const SSL: bool> = http_context::HTTPContext<SSL>;
const HOST_HEADER_NAME: &[u8] = b"Host";
const CONTENT_LENGTH_HEADER_NAME: &[u8] = b"Content-Length";
const CHUNKED_ENCODED_HEADER: picohttp::Header =
picohttp::Header::new(b"Transfer-Encoding", b"chunked");
pub const END_OF_CHUNKED_HTTP1_1_ENCODING_BODY: &[u8] = b"0\r\n\r\n";
const CONNECTION_HEADER: picohttp::Header = picohttp::Header::new(b"Connection", b"keep-alive");
const ACCEPT_HEADER: picohttp::Header = picohttp::Header::new(b"Accept", b"*/*");
const ACCEPT_ENCODING_NO_COMPRESSION: &[u8] = b"identity";
const ACCEPT_ENCODING_COMPRESSION: &[u8] = b"gzip, deflate, br, zstd";
const ACCEPT_ENCODING_HEADER_COMPRESSION: picohttp::Header =
picohttp::Header::new(b"Accept-Encoding", ACCEPT_ENCODING_COMPRESSION);
const ACCEPT_ENCODING_HEADER_NO_COMPRESSION: picohttp::Header =
picohttp::Header::new(b"Accept-Encoding", ACCEPT_ENCODING_NO_COMPRESSION);
const ACCEPT_ENCODING_HEADER: picohttp::Header = if FeatureFlags::DISABLE_COMPRESSION_IN_HTTP_CLIENT
{
ACCEPT_ENCODING_HEADER_NO_COMPRESSION
} else {
ACCEPT_ENCODING_HEADER_COMPRESSION
};
fn get_user_agent_header() -> picohttp::Header {
let ua = OVERRIDDEN_DEFAULT_USER_AGENT.get().copied().unwrap_or(b"");
picohttp::Header::new(
b"User-Agent",
if !ua.is_empty() {
ua
} else {
Global::user_agent.as_bytes()
},
)
}
#[inline(always)]
fn hash_header_const(name: &[u8]) -> u64 {
hash_header_name(name)
}
static AUTHORIZATION_HEADER_HASH: std::sync::LazyLock<u64> =
std::sync::LazyLock::new(|| hash_header_name(b"Authorization"));
static PROXY_AUTHORIZATION_HEADER_HASH: std::sync::LazyLock<u64> =
std::sync::LazyLock::new(|| hash_header_name(b"Proxy-Authorization"));
static COOKIE_HEADER_HASH: std::sync::LazyLock<u64> =
std::sync::LazyLock::new(|| hash_header_name(b"Cookie"));
const PRINT_EVERY: usize = 0;
static PRINT_EVERY_I: AtomicUsize = AtomicUsize::new(0);
const MAX_REQUEST_HEADERS: usize = 256;
static SHARED_REQUEST_HEADERS_BUF: bun_core::RacyCell<[picohttp::Header; MAX_REQUEST_HEADERS]> =
bun_core::RacyCell::new([picohttp::Header::ZERO; MAX_REQUEST_HEADERS]);
static SHARED_RESPONSE_HEADERS_BUF: bun_core::RacyCell<[picohttp::Header; 256]> =
bun_core::RacyCell::new([picohttp::Header::ZERO; 256]);
static SINGLE_PACKET_SMALL_BUFFER: bun_core::RacyCell<[u8; 16 * 1024]> =
bun_core::RacyCell::new([0; 16 * 1024]);
mod scratch {
use super::*;
#[inline]
pub(super) fn request_headers() -> &'static mut [picohttp::Header; MAX_REQUEST_HEADERS] {
unsafe { &mut *SHARED_REQUEST_HEADERS_BUF.get() }
}
#[inline]
pub(super) fn response_headers() -> &'static mut [picohttp::Header; 256] {
unsafe { &mut *SHARED_RESPONSE_HEADERS_BUF.get() }
}
#[inline]
pub(super) fn single_packet_small_buffer() -> &'static mut [u8; 16 * 1024] {
unsafe { &mut *SINGLE_PACKET_SMALL_BUFFER.get() }
}
#[inline]
pub fn temp_hostname() -> &'static mut [u8; 8192] {
unsafe { &mut *TEMP_HOSTNAME.get() }
}
}
pub use scratch::temp_hostname;
#[derive(Copy, Clone, PartialEq, Eq)]
pub enum AlpnOffer {
H1,
H2Only,
H1OrH2,
}
#[allow(clippy::not_unsafe_ptr_arg_deref)]
pub fn configure_http_client_with_alpn(
ssl: &mut boringssl::c::SSL,
hostname: *const core::ffi::c_char,
offer: AlpnOffer,
tls_props: Option<&ssl_config::SSLConfig>,
) {
unsafe {
if !hostname.is_null() && *hostname != 0 {
boringssl::c::SSL_set_tlsext_host_name(ssl, hostname);
}
boringssl::c::SSL_clear_options(ssl, boringssl::c::SSL_OP_LEGACY_SERVER_CONNECT);
boringssl::c::SSL_set_options(ssl, boringssl::c::SSL_OP_LEGACY_SERVER_CONNECT);
const ALPN_H1: &[u8] = &[8, b'h', b't', b't', b'p', b'/', b'1', b'.', b'1'];
const ALPN_H2: &[u8] = &[2, b'h', b'2'];
const ALPN_H2_H1: &[u8] = &[
2, b'h', b'2', 8, b'h', b't', b't', b'p', b'/', b'1', b'.', b'1',
];
let alpns: &'static [u8] = match offer {
AlpnOffer::H1 => ALPN_H1,
AlpnOffer::H1OrH2 => ALPN_H2_H1,
AlpnOffer::H2Only => ALPN_H2,
};
let rc = boringssl::c::SSL_set_alpn_protos(ssl, alpns.as_ptr(), alpns.len());
debug_assert_eq!(rc, 0);
boringssl::c::SSL_enable_signed_cert_timestamps(ssl);
boringssl::c::SSL_enable_ocsp_stapling(ssl);
if let Some(props) = tls_props {
if !props.tls12_cipher_list.is_null() {
let _ = boringssl::c::SSL_set_cipher_list(ssl, props.tls12_cipher_list);
}
if !props.tls13_cipher_suites.is_null() {
let _ = boringssl::c::SSL_set_ciphersuites(ssl, props.tls13_cipher_suites);
}
if !props.tls_curves_list.is_null() {
let _ = boringssl::c::SSL_set1_curves_list(ssl, props.tls_curves_list);
}
if !props.tls_sigalgs_list.is_null() {
let _ = boringssl::c::SSL_set1_sigalgs_list(ssl, props.tls_sigalgs_list);
}
if let Some(ders) = &props.ca_certs_der {
apply_ca_certs_der(ssl, ders);
}
}
}
}
unsafe fn apply_ca_certs_der(ssl: &mut boringssl::c::SSL, ders: &[Box<[u8]>]) {
unsafe extern "C" {
fn X509_STORE_new() -> *mut boringssl::c::X509_STORE;
fn X509_STORE_free(store: *mut boringssl::c::X509_STORE);
fn X509_STORE_add_cert(
store: *mut boringssl::c::X509_STORE,
x509: *mut boringssl::c::X509,
) -> core::ffi::c_int;
}
let store = unsafe { X509_STORE_new() };
if store.is_null() {
return;
}
for der in ders {
let mut p = der.as_ptr();
let x509 = unsafe {
boringssl::c::d2i_X509(core::ptr::null_mut(), &mut p, der.len() as core::ffi::c_long)
};
if x509.is_null() {
continue;
}
let ok = unsafe { X509_STORE_add_cert(store, x509) };
unsafe { boringssl::c::X509_free(x509) };
if ok != 1 {
unsafe { X509_STORE_free(store) };
return;
}
}
if unsafe { boringssl::c::SSL_set0_verify_cert_store(ssl, store) } != 1 {
unsafe { X509_STORE_free(store) };
}
}
use bun_http_types::ETag::HeaderEntryColumns;
impl<const SSL: bool> SocketTimeout for HttpSocket<SSL> {
fn timeout(&self, seconds: c_uint) {
uws::NewSocketHandler::<SSL>::timeout(self, seconds)
}
fn set_timeout_minutes(&self, minutes: c_uint) {
uws::NewSocketHandler::<SSL>::set_timeout_minutes(self, minutes)
}
fn set_timeout(&self, seconds: c_uint) {
uws::NewSocketHandler::<SSL>::set_timeout(self, seconds)
}
}
#[inline]
pub(crate) fn abort_tracker() -> &'static mut ArrayHashMap<u32, uws::AnySocket> {
unsafe { (*SOCKET_ASYNC_HTTP_ABORT_TRACKER.get()).get_or_insert_with(ArrayHashMap::new) }
}
fn get_tls_hostname<'c>(client: &'c HTTPClient<'_>, allow_proxy_url: bool) -> &'c [u8] {
if allow_proxy_url {
if let Some(proxy) = &client.http_proxy {
return proxy.hostname;
}
}
if let Some(props) = &client.tls_props {
let sn = props.get().server_name;
if !sn.is_null() {
let sn_slice = unsafe { bun_core::ffi::cstr(sn) }.to_bytes();
if !sn_slice.is_empty() {
return sn_slice;
}
}
}
if let Some(host) = &client.hostname {
return strip_port_from_host(host);
}
client.url.hostname
}
#[derive(Clone, Copy)]
enum PendingH2Resolution {
H2(h2::SessionPtr),
H1,
LeaderFailed,
}
struct InitialRequestPayloadResult {
has_sent_headers: bool,
has_sent_body: bool,
try_sending_more_data: bool,
}
fn write_proxy_auth_and_headers(writer: &mut Vec<u8>, client: &HTTPClient) {
let user_provided_proxy_auth = client
.proxy_headers
.as_ref()
.map(|hdrs| hdrs.get(b"proxy-authorization").is_some())
.unwrap_or(false);
if let Some(auth) = &client.proxy_authorization {
if !user_provided_proxy_auth {
writer.extend_from_slice(b"Proxy-Authorization: ");
writer.extend_from_slice(auth);
writer.extend_from_slice(b"\r\n");
}
}
if let Some(hdrs) = &client.proxy_headers {
let slice = hdrs.entries.slice();
let names = slice.items_name();
let values = slice.items_value();
for (idx, name_ptr) in names.iter().enumerate() {
writer.extend_from_slice(hdrs.as_str(*name_ptr));
writer.extend_from_slice(b": ");
writer.extend_from_slice(hdrs.as_str(values[idx]));
writer.extend_from_slice(b"\r\n");
}
}
}
fn write_proxy_connect(writer: &mut Vec<u8>, client: &HTTPClient) -> Result<(), bun_core::Error> {
let port: &[u8] = if client.url.get_port().is_some() {
client.url.port
} else if client.url.is_https() {
b"443"
} else {
b"80"
};
writer.extend_from_slice(b"CONNECT ");
writer.extend_from_slice(client.url.hostname);
writer.extend_from_slice(b":");
writer.extend_from_slice(port);
writer.extend_from_slice(b" HTTP/1.1\r\n");
writer.extend_from_slice(b"Host: ");
writer.extend_from_slice(client.url.hostname);
writer.extend_from_slice(b":");
writer.extend_from_slice(port);
writer.extend_from_slice(b"\r\nProxy-Connection: Keep-Alive\r\n");
write_proxy_auth_and_headers(writer, client);
writer.extend_from_slice(b"\r\n");
Ok(())
}
fn write_proxy_request(
writer: &mut Vec<u8>,
request: &picohttp::Request<'_>,
client: &HTTPClient,
) -> Result<(), bun_core::Error> {
writer.extend_from_slice(request.method);
writer.extend_from_slice(b" http://");
writer.extend_from_slice(client.url.hostname);
if client.url.get_port().is_some() {
writer.extend_from_slice(b":");
writer.extend_from_slice(client.url.port);
}
writer.extend_from_slice(request.path);
writer.extend_from_slice(b" HTTP/1.1\r\nProxy-Connection: Keep-Alive\r\n");
write_proxy_auth_and_headers(writer, client);
for header in request.headers {
writer.extend_from_slice(header.name());
writer.extend_from_slice(b": ");
writer.extend_from_slice(header.value());
writer.extend_from_slice(b"\r\n");
}
writer.extend_from_slice(b"\r\n");
Ok(())
}
fn write_request(
writer: &mut Vec<u8>,
request: &picohttp::Request<'_>,
) -> Result<(), bun_core::Error> {
writer.extend_from_slice(request.method);
writer.extend_from_slice(b" ");
writer.extend_from_slice(request.path);
writer.extend_from_slice(b" HTTP/1.1\r\n");
for header in request.headers {
writer.extend_from_slice(header.name());
writer.extend_from_slice(b": ");
writer.extend_from_slice(header.value());
writer.extend_from_slice(b"\r\n");
}
writer.extend_from_slice(b"\r\n");
Ok(())
}
#[cold]
pub fn print_request(
protocol: Protocol,
request: &picohttp::Request<'_>,
url: &[u8],
ignore_insecure: bool,
body: &[u8],
curl: bool,
) {
if curl {
let request_ = picohttp::Request {
method: request.method,
path: url,
minor_version: request.minor_version,
headers: request.headers,
bytes_read: request.bytes_read,
};
Output::pretty_errorln(format_args!("{}", request_.curl(ignore_insecure, body)));
}
let ver: &str = match protocol {
Protocol::Http1_1 => "HTTP/1.1",
Protocol::Http2 => "HTTP/2",
Protocol::Http3 => "HTTP/3",
};
Output::pretty_errorln(format_args!(
"> {} {} {}",
ver,
BStr::new(request.method),
BStr::new(url),
));
for header in request.headers {
Output::pretty_errorln(format_args!("> {}", header));
}
Output::flush();
}
#[cold]
fn print_response(response: &picohttp::Response<'_>) {
Output::pretty_errorln(format_args!("{}", response));
Output::flush();
}
fn write_to_socket<const IS_SSL: bool>(
socket: HttpSocket<IS_SSL>,
data: &[u8],
) -> Result<usize, bun_core::Error> {
let mut remaining = data;
let mut total_written: usize = 0;
while !remaining.is_empty() {
let amount = socket.write(remaining);
if amount < 0 {
return Err(err!(WriteFailed));
}
let wrote = usize::try_from(amount).expect("int cast");
total_written += wrote;
remaining = &remaining[wrote..];
if wrote == 0 {
break;
}
}
Ok(total_written)
}
fn write_to_socket_with_buffer_fallback<const IS_SSL: bool>(
socket: HttpSocket<IS_SSL>,
buffer: &mut bun_io::StreamBuffer,
data: &[u8],
) -> Result<usize, bun_core::Error> {
let amount = write_to_socket::<IS_SSL>(socket, data)?;
if amount < data.len() {
let _ = buffer.write(&data[amount..]);
}
Ok(amount)
}
pub(crate) fn get_cert_error_from_no(error_no: i32) -> bun_core::Error {
let name: &'static str = match error_no {
0 => "OK", 2 => "UNABLE_TO_GET_ISSUER_CERT",
3 => "UNABLE_TO_GET_CRL",
4 => "UNABLE_TO_DECRYPT_CERT_SIGNATURE",
5 => "UNABLE_TO_DECRYPT_CRL_SIGNATURE",
6 => "UNABLE_TO_DECODE_ISSUER_PUBLIC_KEY",
7 => "CERT_SIGNATURE_FAILURE",
8 => "CRL_SIGNATURE_FAILURE",
9 => "CERT_NOT_YET_VALID",
10 => "CERT_HAS_EXPIRED",
11 => "CRL_NOT_YET_VALID",
12 => "CRL_HAS_EXPIRED",
13 => "ERROR_IN_CERT_NOT_BEFORE_FIELD",
14 => "ERROR_IN_CERT_NOT_AFTER_FIELD",
15 => "ERROR_IN_CRL_LAST_UPDATE_FIELD",
16 => "ERROR_IN_CRL_NEXT_UPDATE_FIELD",
17 => "OUT_OF_MEM",
18 => "DEPTH_ZERO_SELF_SIGNED_CERT",
19 => "SELF_SIGNED_CERT_IN_CHAIN",
20 => "UNABLE_TO_GET_ISSUER_CERT_LOCALLY",
21 => "UNABLE_TO_VERIFY_LEAF_SIGNATURE",
22 => "CERT_CHAIN_TOO_LONG",
23 => "CERT_REVOKED",
24 => "INVALID_CA",
25 => "PATH_LENGTH_EXCEEDED",
26 => "INVALID_PURPOSE",
27 => "CERT_UNTRUSTED",
28 => "CERT_REJECTED",
29 => "SUBJECT_ISSUER_MISMATCH",
30 => "AKID_SKID_MISMATCH",
31 => "AKID_ISSUER_SERIAL_MISMATCH",
32 => "KEYUSAGE_NO_CERTSIGN",
33 => "UNABLE_TO_GET_CRL_ISSUER",
34 => "UNHANDLED_CRITICAL_EXTENSION",
35 => "KEYUSAGE_NO_CRL_SIGN",
36 => "UNHANDLED_CRITICAL_CRL_EXTENSION",
37 => "INVALID_NON_CA",
38 => "PROXY_PATH_LENGTH_EXCEEDED",
39 => "KEYUSAGE_NO_DIGITAL_SIGNATURE",
40 => "PROXY_CERTIFICATES_NOT_ALLOWED",
41 => "INVALID_EXTENSION",
42 => "INVALID_POLICY_EXTENSION",
43 => "NO_EXPLICIT_POLICY",
44 => "DIFFERENT_CRL_SCOPE",
45 => "UNSUPPORTED_EXTENSION_FEATURE",
46 => "UNNESTED_RESOURCE",
47 => "PERMITTED_VIOLATION",
48 => "EXCLUDED_VIOLATION",
49 => "SUBTREE_MINMAX",
50 => "APPLICATION_VERIFICATION",
51 => "UNSUPPORTED_CONSTRAINT_TYPE",
52 => "UNSUPPORTED_CONSTRAINT_SYNTAX",
53 => "UNSUPPORTED_NAME_SYNTAX",
54 => "CRL_PATH_VALIDATION_ERROR",
56 => "SUITE_B_INVALID_VERSION",
57 => "SUITE_B_INVALID_ALGORITHM",
58 => "SUITE_B_INVALID_CURVE",
59 => "SUITE_B_INVALID_SIGNATURE_ALGORITHM",
60 => "SUITE_B_LOS_NOT_ALLOWED",
61 => "SUITE_B_CANNOT_SIGN_P_384_WITH_P_256",
62 => "HOSTNAME_MISMATCH",
63 => "EMAIL_MISMATCH",
64 => "IP_ADDRESS_MISMATCH",
65 => "INVALID_CALL",
66 => "STORE_LOOKUP",
67 => "NAME_CONSTRAINTS_WITHOUT_SANS",
_ => "UNKNOWN_CERTIFICATE_VERIFICATION_ERROR",
};
bun_core::Error::from_name(name)
}
impl<'a> HTTPClient<'a> {
#[inline]
fn request_body(&self) -> &[u8] {
self.state.request_body.slice()
}
#[inline]
fn body_out_str(&self) -> Option<&MutableString> {
body_out::opt_mut(self.state.body_out_str).map(|b| &*b)
}
#[inline]
fn proxy_tunnel_ptr(&self) -> Option<NonNull<ProxyTunnel>> {
self.proxy_tunnel.as_ref().map(|p| p.data)
}
#[inline]
fn close_proxy_tunnel(&mut self, shutdown: bool) {
if let Some(t) = self.proxy_tunnel.take() {
let tunnel = proxy_tunnel::raw_as_mut(t.as_ptr());
if shutdown {
tunnel.shutdown();
}
tunnel.detach_socket();
t.deref();
}
}
fn dispatch_result_and_reset(&mut self, clear_proxy_tunneling: bool) {
let callback = self.result_callback;
let result = unsafe { self.to_result().detach_lifetime() };
self.state.reset();
if clear_proxy_tunneling {
self.flags.proxy_tunneling = false;
}
callback.run(self.parent_async_http(), result);
}
#[inline]
fn progress_node_mut(&mut self) -> Option<&mut bun_core::Progress::Node> {
self.progress_node.map(|mut p| unsafe { p.as_mut() })
}
fn report_progress(&mut self, completed: usize) {
if let Some(progress) = self.progress_node_mut() {
progress.activate();
progress.set_completed_items(completed);
unsafe { (*progress.context_ptr()).maybe_refresh() };
}
}
}
pub(crate) mod body_out {
use super::{MutableString, NonNull};
#[inline]
pub(crate) fn as_mut<'a>(mut p: NonNull<MutableString>) -> &'a mut MutableString {
unsafe { p.as_mut() }
}
#[inline]
pub(super) fn opt_mut<'a>(p: Option<NonNull<MutableString>>) -> Option<&'a mut MutableString> {
p.map(as_mut)
}
#[inline]
pub(super) fn take_list(p: Option<NonNull<MutableString>>) -> Option<Vec<u8>> {
p.map(|p| core::mem::take(&mut as_mut(p).list))
}
#[inline]
pub(super) fn restore_list(p: Option<NonNull<MutableString>>, v: Option<Vec<u8>>) {
if let (Some(p), Some(v)) = (p, v) {
as_mut(p).list = v;
}
}
}
impl<'a> HTTPClient<'a> {
pub fn check_server_identity<const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
cert_error: HTTPCertError,
ssl: &mut boringssl::c::SSL,
allow_proxy_url: bool,
) -> bool {
if self.flags.reject_unauthorized {
let cert_chain = unsafe { boringssl::c::SSL_get_peer_cert_chain(ssl) };
if !cert_chain.is_null() {
let x509 = unsafe { boringssl::c::sk_X509_value(cert_chain, 0) };
if !x509.is_null() {
let hostname = get_tls_hostname(self, allow_proxy_url);
let is_proxy_certificate = allow_proxy_url && self.http_proxy.is_some();
if !is_proxy_certificate && self.signals.get(signals::Field::CertErrors) {
let cert_size =
unsafe { boringssl::c::i2d_X509(x509, core::ptr::null_mut()) };
let mut cert = vec![0u8; usize::try_from(cert_size).expect("int cast")]
.into_boxed_slice();
let mut cert_ptr = cert.as_mut_ptr();
let result_size =
unsafe { boringssl::c::i2d_X509(x509, &raw mut cert_ptr) };
debug_assert!(result_size == cert_size);
self.state.certificate_info = Some(CertificateInfo {
cert,
hostname: Box::<[u8]>::from(hostname),
cert_error,
});
self.state.flags.is_waiting_for_cert_check = true;
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
return true;
} else {
if boringssl::check_x509_server_identity(unsafe { &mut *x509 }, hostname) {
return true;
}
}
}
}
self.close_and_fail::<IS_SSL>(err!(ERR_TLS_CERT_ALTNAME_INVALID), socket);
return false;
}
true
}
pub fn register_abort_tracker<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
if self.signals.aborted.is_some() {
let any = if IS_SSL {
uws::AnySocket::SocketTls(uws::SocketTLS::from_any(socket.socket))
} else {
uws::AnySocket::SocketTcp(uws::SocketTCP::from_any(socket.socket))
};
let _ = abort_tracker().put(self.async_http_id, any);
}
}
pub fn unregister_abort_tracker(&mut self) {
if self.signals.aborted.is_some() {
let _ = abort_tracker().swap_remove(&self.async_http_id);
}
}
pub fn on_open<const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
) -> Result<(), bun_core::Error> {
if cfg!(debug_assertions) {
if let Some(proxy) = &self.http_proxy {
debug_assert!(IS_SSL == proxy.is_https());
} else {
debug_assert!(IS_SSL == self.url.is_https());
}
}
self.register_abort_tracker::<IS_SSL>(socket);
bun_core::scoped_log!(fetch, "Connected {} \n", BStr::new(self.url.href));
self.set_timeout(&socket);
if !self.flags.disable_keepalive {
let _ = socket.set_keep_alive(true, 60);
}
if self.signals.get(signals::Field::Aborted) {
self.close_and_abort::<IS_SSL>(socket);
return Err(err!(ClientAborted));
}
if self.state.request_stage == RequestStage::Pending {
self.state.request_stage = RequestStage::Opened;
}
if IS_SSL {
let ssl_ptr: *mut boringssl::c::SSL = socket
.get_native_handle()
.map(|p| p.cast())
.unwrap_or(core::ptr::null_mut());
if !ssl_ptr.is_null() && unsafe { boringssl::c::SSL_is_init_finished(ssl_ptr) } == 0 {
let raw_hostname = get_tls_hostname(self, self.http_proxy.is_some());
let mut owned: Vec<u8>; let host_z: *const core::ffi::c_char = if !strings::is_ip_address(raw_hostname) {
let temp = scratch::temp_hostname();
if raw_hostname.len() < temp.len() {
temp[..raw_hostname.len()].copy_from_slice(raw_hostname);
temp[raw_hostname.len()] = 0;
temp.as_ptr().cast::<core::ffi::c_char>()
} else {
owned = Vec::with_capacity(raw_hostname.len() + 1);
owned.extend_from_slice(raw_hostname);
owned.push(0);
owned.as_ptr().cast::<core::ffi::c_char>()
}
} else {
core::ptr::null()
};
configure_http_client_with_alpn(
unsafe { &mut *ssl_ptr },
host_z,
self.alpn_offer(),
self.tls_props.as_deref(),
);
if let Some(host_str) = core::str::from_utf8(raw_hostname).ok() {
let port = if let Some(proxy) = &self.http_proxy {
proxy.get_port_auto()
} else {
self.url.get_port_auto()
};
let profile_salt = self
.tls_props
.as_deref()
.map(|props| props.content_hash())
.unwrap_or(0);
bao_boringssl_bridge::session_cache::offer_session(
ssl_ptr,
host_str,
port,
profile_salt,
);
}
}
} else {
self.first_call::<IS_SSL>(socket);
}
Ok(())
}
pub fn can_offer_h2(&self) -> bool {
if self.signals.get(signals::Field::CertErrors) {
return false;
}
if self.flags.force_http1 {
return false;
}
if self.http_proxy.is_some() {
return false;
}
if self.flags.is_preconnect_only {
return false;
}
if self.unix_socket_path.slice().len() > 0 {
return false;
}
if matches!(
self.state.original_request_body,
HTTPRequestBody::Sendfile(_)
) {
return false;
}
if self.flags.is_page_egress {
return true;
}
self.flags.force_http2
|| bun_core::env_var::feature_flag::BUN_FEATURE_FLAG_EXPERIMENTAL_HTTP2_CLIENT
.get()
.unwrap_or(false)
}
pub fn alpn_offer(&self) -> AlpnOffer {
if !self.can_offer_h2() {
return AlpnOffer::H1;
}
if self.flags.force_http2 {
AlpnOffer::H2Only
} else {
AlpnOffer::H1OrH2
}
}
pub fn can_try_h3_alt_svc(&self) -> bool {
if self.flags.h3_alt_svc_fallback {
return false;
}
if self.signals.get(signals::Field::CertErrors) {
return false;
}
if self.flags.force_http1 || self.flags.force_http2 {
return false;
}
if self.http_proxy.is_some() {
return false;
}
if self.flags.is_preconnect_only {
return false;
}
if self.unix_socket_path.slice().len() > 0 {
return false;
}
if matches!(
self.state.original_request_body,
HTTPRequestBody::Sendfile(_)
) {
return false;
}
if self.has_tls_options_unsupported_by_h3() {
return false;
}
h3_alt_svc_enabled()
}
fn has_tls_options_unsupported_by_h3(&self) -> bool {
self.signals.get(signals::Field::CertErrors)
|| self
.tls_props
.as_ref()
.is_some_and(|tls| tls.get().requires_custom_request_ctx)
}
pub fn first_call<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
if FeatureFlags::IS_FETCH_PRECONNECT_SUPPORTED {
if self.flags.is_preconnect_only {
self.on_preconnect::<IS_SSL>(socket);
return;
}
}
if IS_SSL {
let ssl_ptr: *mut boringssl::c::SSL = socket
.get_native_handle()
.map(|p| p.cast())
.unwrap_or(core::ptr::null_mut());
if !ssl_ptr.is_null() {
self.tls_info = Some(unsafe { BunTlsInfo::from_ssl(ssl_ptr) });
}
let mut proto: *const u8 = core::ptr::null();
let mut proto_len: c_uint = 0;
let alpn = unsafe {
boringssl::c::SSL_get0_alpn_selected(ssl_ptr, &raw mut proto, &raw mut proto_len);
bun_core::ffi::slice(proto, proto_len as usize)
};
if alpn == b"h2" {
bun_core::scoped_log!(fetch, "ALPN negotiated h2 {}", BStr::new(self.url.href));
let tls_socket = uws::SocketTLS::from_any(socket.socket);
let ctx = self.get_ssl_ctx::<true>();
let session = h2::ClientSession::create(ctx, tls_socket, self);
GenHttpContext::<true>::tag_as_h2(tls_socket, session.as_ptr());
self.resolve_pending_h2(PendingH2Resolution::H2(session));
h2::ClientSession::attach_leader(session, self);
return;
}
self.flags.protocol = Protocol::Http1_1;
self.resolve_pending_h2(PendingH2Resolution::H1);
if self.flags.force_http2 {
self.close_and_fail::<IS_SSL>(err!(HTTP2Unsupported), socket);
return;
}
}
match self.state.request_stage {
RequestStage::Opened | RequestStage::Pending => {
self.on_writable::<true, IS_SSL>(socket);
}
_ => {}
}
}
pub fn retry_after_h2_coalesce(&mut self) {
self.start_::<true>();
}
pub fn retry_from_h2(&mut self) {
debug_assert!(self.h2.is_none());
self.unregister_abort_tracker();
self.flags.protocol = Protocol::Http1_1;
self.h2_retries += 1;
let body = core::mem::replace(
&mut self.state.original_request_body,
HTTPRequestBody::Bytes(b""),
);
let body_out = self.state.body_out_str.take().unwrap();
self.state.reset();
self.start(body, body_out::as_mut(body_out));
}
pub fn retry_from_h3_to_tcp(&mut self) {
debug_assert!(self.h3.is_none());
debug_assert!(!self.flags.force_http3);
self.unregister_abort_tracker();
self.flags.protocol = Protocol::Http1_1;
self.flags.h3_alt_svc_fallback = true;
let body = core::mem::replace(
&mut self.state.original_request_body,
HTTPRequestBody::Bytes(b""),
);
let body_out = self.state.body_out_str.take().unwrap();
self.state.reset();
self.start(body, body_out::as_mut(body_out));
}
pub fn fail_from_h2(&mut self, err: bun_core::Error) {
debug_assert!(self.h2.is_none());
debug_assert!(self.h3.is_none());
self.unregister_abort_tracker();
if self.state.stage != Stage::Done && self.state.stage != Stage::Fail {
self.state.request_stage = RequestStage::Fail;
self.state.response_stage = ResponseStage::Fail;
self.state.fail = Some(err);
self.state.stage = Stage::Fail;
if self.flags.defer_fail_until_connecting_is_complete {
return;
}
self.dispatch_result_and_reset(false);
}
}
pub fn on_close<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
bun_core::scoped_log!(fetch, "Closed {}\n", BStr::new(self.url.href));
self.unregister_abort_tracker();
if self.signals.get(signals::Field::Aborted) {
self.fail(err!(Aborted));
return;
}
self.close_proxy_tunnel(true);
let in_progress = self.state.stage != Stage::Done
&& self.state.stage != Stage::Fail
&& !self.state.flags.is_redirect_pending;
if self.state.flags.is_redirect_pending {
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.do_redirect::<IS_SSL>(ctx, socket);
return;
}
if in_progress {
if self.state.is_chunked_encoding() {
if matches!(self.state.chunked_decoder._state, bun_picohttp::ChunkedState::TrailerLineHead | bun_picohttp::ChunkedState::TrailerLineMiddle) {
self.state.flags.received_last_chunk = true;
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
} else if self.state.content_length.is_none()
&& self.state.response_stage == ResponseStage::Body
{
self.state.flags.received_last_chunk = true;
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
}
if self.allow_retry
&& self.method.is_idempotent()
&& self.state.response_stage != ResponseStage::Body
&& self.state.response_stage != ResponseStage::BodyChunk
{
self.allow_retry = false;
self.state.response_message_buffer = MutableString::default();
let body = core::mem::replace(
&mut self.state.original_request_body,
HTTPRequestBody::Bytes(b""),
);
let body_out = self.state.body_out_str.take().unwrap();
self.start(body, body_out::as_mut(body_out));
return;
}
if in_progress {
self.fail(err!(ConnectionClosed));
}
}
pub fn on_timeout<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
if self.flags.disable_timeout {
return;
}
bun_core::scoped_log!(fetch, "Timeout {}\n", BStr::new(self.url.href));
self.fail(err!(Timeout));
GenHttpContext::<IS_SSL>::terminate_socket(socket);
}
pub fn on_connect_error(&mut self) {
bun_core::scoped_log!(fetch, "onConnectError {}\n", BStr::new(self.url.href));
self.fail(err!(ConnectionRefused));
}
#[inline]
fn get_request_body_send_buffer(&self) -> http_thread::RequestBodyBuffer {
let actual_estimated_size =
self.request_body().len() + self.estimated_request_header_byte_length();
let estimated_size = if HTTPClient::is_https(self) {
actual_estimated_size.min(MAX_TLS_RECORD_SIZE)
} else {
actual_estimated_size * 2
};
http_thread().get_request_body_send_buffer(estimated_size)
}
pub fn is_keep_alive_possible(&self) -> bool {
if FeatureFlags::ENABLE_KEEPALIVE {
if self.unix_socket_path.slice().len() > 0 {
return false;
}
if self.state.flags.allow_keepalive && !self.flags.disable_keepalive {
return true;
}
}
false
}
pub fn proxy_auth_hash(&self) -> u64 {
let mut combined: u64 = 0;
let mut any = false;
let mut name_lower_buf = [0u8; 256];
if let Some(sni_raw) = &self.hostname {
let sni = strip_port_from_host(sni_raw);
if !strings::eql_case_insensitive_ascii(sni, self.url.hostname, true) {
let sni_lower: &[u8] = if sni.len() <= name_lower_buf.len() {
strings::copy_lowercase(sni, &mut name_lower_buf[0..sni.len()])
} else {
sni
};
combined = combined.wrapping_add(bun_wyhash::hash(sni_lower));
any = true;
}
}
let mut user_provided_auth = false;
if let Some(hdrs) = &self.proxy_headers {
let slice = hdrs.entries.slice();
let names = slice.items_name();
let values = slice.items_value();
for (idx, name_ptr) in names.iter().enumerate() {
let name = hdrs.as_str(*name_ptr);
let value = hdrs.as_str(values[idx]);
let name_lower: &[u8] = if name.len() <= name_lower_buf.len() {
strings::copy_lowercase(name, &mut name_lower_buf[0..name.len()])
} else {
name
};
let mut h = Wyhash::init(0);
h.update(name_lower);
h.update(b":");
h.update(value);
combined = combined.wrapping_add(h.final_());
any = true;
if strings::eql_case_insensitive_ascii(name, b"proxy-authorization", true) {
user_provided_auth = true;
}
}
}
if !user_provided_auth {
if let Some(auth) = &self.proxy_authorization {
let mut h = Wyhash::init(0);
h.update(b"proxy-authorization:");
h.update(auth);
combined = combined.wrapping_add(h.final_());
any = true;
}
}
if any { combined } else { 0 }
}
pub fn get_ssl_ctx<const IS_SSL: bool>(&self) -> *mut GenHttpContext<IS_SSL> {
if IS_SSL {
if let Some(ctx) = self.custom_ssl_ctx.as_ref() {
return ctx.as_ptr().cast::<GenHttpContext<IS_SSL>>();
}
(&raw mut http_thread().https_context).cast::<GenHttpContext<IS_SSL>>()
} else {
(&raw mut http_thread().http_context).cast::<GenHttpContext<IS_SSL>>()
}
}
#[inline]
fn ssl_ctx_mut<'c, const IS_SSL: bool>(
ctx: *mut GenHttpContext<IS_SSL>,
) -> &'c mut GenHttpContext<IS_SSL> {
unsafe { &mut *ctx }
}
pub fn set_custom_ssl_ctx(&mut self, ctx: NonNull<HttpsContext>) {
let new_ref = unsafe { http_context::HTTPContextRc::<true>::init_ref(ctx.as_ptr()) };
if let Some(old) = self.custom_ssl_ctx.replace(new_ref) {
old.deref();
}
}
pub fn header_str(&self, ptr: StringPointer) -> &'a [u8] {
let buf: &'a [u8] = self.header_buf;
&buf[ptr.offset as usize..][..ptr.length as usize]
}
pub fn build_request(&mut self, body_len: usize) -> picohttp::Request<'static> {
let mut header_count: usize = 0;
let header_entries = self.header_entries.slice();
let header_names = header_entries.items_name();
let header_values = header_entries.items_value();
let request_headers_buf = scratch::request_headers();
let mut override_accept_encoding = false;
let mut override_accept_header = false;
let mut override_host_header = false;
let mut override_connection_header = false;
let mut override_user_agent = false;
let mut add_transfer_encoding = true;
let mut original_content_length: Option<&[u8]> = None;
const MAX_DEFAULT_HEADERS: usize = 6;
const MAX_USER_HEADERS: usize = MAX_REQUEST_HEADERS - MAX_DEFAULT_HEADERS;
for (i, head) in header_names.iter().enumerate() {
let name = self.header_str(*head);
let hash = hash_header_name(name);
let will_append = header_count < MAX_USER_HEADERS;
match hash {
h if h == hash_header_const(b"Content-Length") => {
original_content_length = Some(self.header_str(header_values[i]));
continue;
}
h if h == hash_header_const(b"Connection") => {
if will_append {
override_connection_header = true;
let connection_value = self.header_str(header_values[i]);
if bun_core::strings::eql_case_insensitive_ascii_check_length(
connection_value,
b"close",
) {
self.flags.disable_keepalive = true;
} else if bun_core::strings::eql_case_insensitive_ascii_check_length(
connection_value,
b"keep-alive",
) {
self.flags.disable_keepalive = false;
}
}
}
h if h == hash_header_const(b"if-modified-since") => {
if self.flags.force_last_modified && self.if_modified_since.is_empty() {
self.if_modified_since =
unsafe { bun_ptr::detach_lifetime(self.header_str(header_values[i])) };
}
}
h if h == hash_header_const(HOST_HEADER_NAME) => {
if will_append {
override_host_header = true;
}
}
h if h == hash_header_const(b"Accept") => {
if will_append {
override_accept_header = true;
}
}
h if h == hash_header_const(b"User-Agent") => {
if will_append {
override_user_agent = true;
}
}
h if h == hash_header_const(b"Accept-Encoding") => {
if will_append {
override_accept_encoding = true;
}
}
h if h == hash_header_const(b"Upgrade") => {
if will_append {
let value = self.header_str(header_values[i]);
if !bun_core::strings::eql_any_case_insensitive_ascii(
value,
&[b"h2", b"h2c"],
) {
self.flags.upgrade_state = HTTPUpgradeState::Pending;
}
}
}
h if h == hash_header_const(CHUNKED_ENCODED_HEADER.name()) => {
if !self.flags.is_streaming_request_body {
continue;
}
if will_append {
add_transfer_encoding = false;
}
}
_ => {}
}
if !will_append {
continue;
}
request_headers_buf[header_count] =
picohttp::Header::new(name, self.header_str(header_values[i]));
header_count += 1;
}
if !override_connection_header &&
!self.flags.disable_keepalive &&
!self.flags.omit_connection_header
{
request_headers_buf[header_count] = CONNECTION_HEADER;
header_count += 1;
}
if !override_user_agent {
request_headers_buf[header_count] = get_user_agent_header();
header_count += 1;
}
if !override_accept_header {
request_headers_buf[header_count] = ACCEPT_HEADER;
header_count += 1;
}
if !override_host_header {
request_headers_buf[header_count] =
picohttp::Header::new(HOST_HEADER_NAME, self.url.host);
header_count += 1;
}
if !override_accept_encoding && !self.flags.disable_decompression {
request_headers_buf[header_count] = ACCEPT_ENCODING_HEADER;
header_count += 1;
}
if body_len > 0 || self.method.has_request_body() {
if self.flags.is_streaming_request_body {
if let Some(content_length) = original_content_length {
if add_transfer_encoding {
request_headers_buf[header_count] =
picohttp::Header::new(CONTENT_LENGTH_HEADER_NAME, content_length);
header_count += 1;
}
} else if add_transfer_encoding
&& self.flags.upgrade_state == HTTPUpgradeState::None
{
request_headers_buf[header_count] = CHUNKED_ENCODED_HEADER;
header_count += 1;
}
} else {
let value: &[u8] = match bun_core::fmt::buf_print(
&mut self.request_content_len_buf,
format_args!("{body_len}"),
) {
Ok(s) => unsafe { bun_ptr::detach_lifetime(s) },
Err(_) => b"0",
};
request_headers_buf[header_count] =
picohttp::Header::new(CONTENT_LENGTH_HEADER_NAME, value);
header_count += 1;
}
} else if let Some(content_length) = original_content_length
&& (self.flags.is_node_http_client
|| matches!(bun_core::parse_unsigned::<usize>(content_length, 10), Ok(0)))
{
request_headers_buf[header_count] =
picohttp::Header::new(CONTENT_LENGTH_HEADER_NAME, content_length);
header_count += 1;
}
picohttp::Request {
method: self
.extension_method
.unwrap_or(self.method.as_str().as_bytes()),
path: unsafe { bun_ptr::detach_lifetime(self.url.pathname) },
minor_version: 1,
headers: unsafe { bun_ptr::detach_lifetime(&request_headers_buf[0..header_count]) },
bytes_read: 0,
}
}
pub fn do_redirect<const IS_SSL: bool>(
&mut self,
ctx: *mut GenHttpContext<IS_SSL>,
socket: HttpSocket<IS_SSL>,
) {
if self.flags.protocol != Protocol::Http1_1 {
return self.do_redirect_multiplexed();
}
bun_core::scoped_log!(fetch, "doRedirect");
if matches!(self.state.original_request_body, HTTPRequestBody::Stream(_)) {
self.flags.is_streaming_request_body = false;
}
let _ = core::mem::ManuallyDrop::new(core::mem::take(&mut self.unix_socket_path));
let request_body: &[u8] = if self.state.flags.resend_request_body_on_redirect
&& matches!(self.state.original_request_body, HTTPRequestBody::Bytes(_))
{
match &self.state.original_request_body {
HTTPRequestBody::Bytes(b) => b,
_ => unreachable!(),
}
} else {
b""
};
self.state.response_message_buffer = MutableString::default();
let body_out_str = self.state.body_out_str.unwrap();
self.remaining_redirect_count = self.remaining_redirect_count.saturating_sub(1);
self.flags.redirected = true;
debug_assert!(self.redirect_type == FetchRedirect::Follow);
self.unregister_abort_tracker();
if self.proxy_tunnel.is_some() {
bun_core::scoped_log!(fetch, "close the tunnel");
self.close_proxy_tunnel(true);
GenHttpContext::<IS_SSL>::close_socket(socket);
} else if self.state.request_stage == RequestStage::Done
&& self.is_keep_alive_possible()
&& !socket.is_closed_or_has_error()
&& (!IS_SSL || self.http_proxy.is_some() || self.hostname.is_none())
{
bun_core::scoped_log!(fetch, "Keep-Alive release in redirect");
debug_assert!(!self.connected_url.hostname.is_empty());
Self::ssl_ctx_mut(ctx).release_socket(
socket,
self.flags.did_have_handshaking_error && !self.flags.reject_unauthorized,
self.flags.reject_unauthorized,
self.connected_url.hostname,
self.connected_url.get_port_auto(),
self.tls_props.as_ref(),
None,
b"",
0,
0,
None,
);
} else {
GenHttpContext::<IS_SSL>::close_socket(socket);
}
self.connected_url = URL::default();
self.prev_redirect = Vec::new();
if self.state.flags.clear_hostname_on_redirect {
self.state.flags.clear_hostname_on_redirect = false;
self.hostname = None;
}
if self.remaining_redirect_count == 0 {
self.fail(err!(TooManyRedirects));
return;
}
self.state.reset();
bun_core::scoped_log!(fetch, "doRedirect state reset");
self.flags.proxy_tunneling = false;
self.close_proxy_tunnel(false);
self.flags.protocol = Protocol::Http1_1;
self.start(
HTTPRequestBody::Bytes(request_body),
body_out::as_mut(body_out_str),
);
}
pub fn is_https(&self) -> bool {
if let Some(proxy) = &self.http_proxy {
return proxy.is_https();
}
self.url.is_https()
}
pub fn start(&mut self, body: HTTPRequestBody<'a>, body_out_str: &mut MutableString) {
body_out_str.reset();
debug_assert!(self.state.response_message_buffer.list.capacity() == 0);
self.state = InternalState::init(body, body_out_str);
if self.is_https() {
self.start_::<true>();
} else {
self.start_::<false>();
}
}
fn start_<const IS_SSL: bool>(&mut self) {
self.unregister_abort_tracker();
self.flags.defer_fail_until_connecting_is_complete = true;
if self.signals.get(signals::Field::Aborted) {
self.fail(err!(AbortedBeforeConnecting));
self.complete_connecting_process();
return;
}
if !IS_SSL {
if self.flags.force_http2 {
self.fail(err!(HTTP2Unsupported));
self.complete_connecting_process();
return;
}
}
if IS_SSL {
if !self.flags.force_http3 && self.can_try_h3_alt_svc() {
if let Some(alt_port) =
h3::alt_svc::lookup(self.url.hostname, self.url.get_port_auto())
{
let h3_ctx = h3::ClientContext::get_or_create(unsafe {
NonNull::new_unchecked(http_thread().uws_loop)
});
if let Some(ctx) = h3_ctx {
if !h3::ClientContext::as_mut(ctx).connect(
self,
self.url.hostname,
alt_port,
) {
self.fail(err!(ConnectionRefused));
}
self.complete_connecting_process();
return;
}
}
}
}
if self.flags.force_http2 && self.signals.get(signals::Field::CertErrors) {
self.fail(err!(HTTP2Unsupported));
self.complete_connecting_process();
return;
}
if self.flags.force_http3 {
if self.signals.get(signals::Field::CertErrors) {
self.fail(err!(HTTP3Unsupported));
self.complete_connecting_process();
return;
}
if !IS_SSL {
self.fail(err!(HTTP3Unsupported));
self.complete_connecting_process();
return;
}
if self.http_proxy.is_some() || self.unix_socket_path.slice().len() > 0 {
self.fail(err!(HTTP3Unsupported));
self.complete_connecting_process();
return;
}
if self.has_tls_options_unsupported_by_h3() {
self.fail(err!(HTTP3Unsupported));
self.complete_connecting_process();
return;
}
let Some(ctx) = h3::ClientContext::get_or_create(unsafe {
NonNull::new_unchecked(http_thread().uws_loop)
}) else {
self.fail(err!(HTTP3Unsupported));
self.complete_connecting_process();
return;
};
if !h3::ClientContext::as_mut(ctx).connect(
self,
self.url.hostname,
self.url.get_port_auto(),
) {
self.fail(err!(ConnectionRefused));
}
self.complete_connecting_process();
return;
}
let socket = match http_thread().connect::<IS_SSL>(self) {
Ok(Some(s)) => s,
Ok(None) => {
self.complete_connecting_process();
return;
}
Err(err) => {
self.fail(err);
self.complete_connecting_process();
return;
}
};
if socket.is_closed()
&& (self.state.response_stage != ResponseStage::Done
&& self.state.response_stage != ResponseStage::Fail)
{
GenHttpContext::<IS_SSL>::mark_socket_as_dead(socket);
self.fail(err!(ConnectionClosed));
self.complete_connecting_process();
return;
}
if self.state.request_stage == RequestStage::Pending {
self.register_abort_tracker::<IS_SSL>(socket);
}
self.complete_connecting_process();
}
fn estimated_request_header_byte_length(&self) -> usize {
let sliced = self.header_entries.slice();
let mut count: usize = 0;
for head in sliced.items_name() {
count += head.length as usize;
}
for value in sliced.items_value() {
count += value.length as usize;
}
count
}
#[inline(never)]
fn send_initial_request_payload<const IS_FIRST_CALL: bool, const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
) -> Result<InitialRequestPayloadResult, bun_core::Error> {
let mut request_body_buffer = self.get_request_body_send_buffer();
let mut temporary_send_buffer = request_body_buffer.to_array_list();
let writer = &mut temporary_send_buffer;
let request = self.build_request(self.state.original_request_body.len());
if self.http_proxy.is_some() {
if self.url.is_https() {
bun_core::scoped_log!(fetch, "start proxy tunneling (https proxy)");
self.flags.proxy_tunneling = true;
write_proxy_connect(writer, self)?;
} else {
bun_core::scoped_log!(fetch, "start proxy request (http proxy)");
write_proxy_request(writer, &request, self)?;
}
} else {
bun_core::scoped_log!(fetch, "normal request");
write_request(writer, &request)?;
}
let headers_len = temporary_send_buffer.len();
if !self.request_body().is_empty()
&& temporary_send_buffer.capacity() - temporary_send_buffer.len() > 0
&& !self.flags.proxy_tunneling
{
let spare = temporary_send_buffer.capacity() - temporary_send_buffer.len();
let wrote = spare.min(self.request_body().len());
debug_assert!(wrote > 0);
temporary_send_buffer.extend_from_slice(&self.request_body()[0..wrote]);
}
let to_send = &temporary_send_buffer[self.state.request_sent_len..];
if cfg!(debug_assertions) {
debug_assert!(!socket.is_shutdown());
debug_assert!(!socket.is_closed());
}
let amount = write_to_socket::<IS_SSL>(socket, to_send)?;
if IS_FIRST_CALL {
if amount == 0 {
return Ok(InitialRequestPayloadResult {
has_sent_headers: self.state.request_sent_len >= headers_len,
has_sent_body: false,
try_sending_more_data: false,
});
}
}
self.state.request_sent_len += amount;
let has_sent_headers = self.state.request_sent_len >= headers_len;
if has_sent_headers && self.verbose != HTTPVerboseLevel::None {
print_request(
Protocol::Http1_1,
&request,
self.url.href,
!self.flags.reject_unauthorized,
self.request_body(),
self.verbose == HTTPVerboseLevel::Curl,
);
}
if has_sent_headers && !self.request_body().is_empty() {
self.state.request_body = bun_ptr::RawSlice::new(
&self.state.request_body.slice()[self.state.request_sent_len - headers_len..],
);
}
let has_sent_body = if matches!(self.state.original_request_body, HTTPRequestBody::Bytes(_))
{
self.request_body().is_empty()
} else {
false
};
Ok(InitialRequestPayloadResult {
has_sent_headers,
has_sent_body,
try_sending_more_data: amount == to_send.len() && (!has_sent_body || !has_sent_headers),
})
}
pub fn flush_stream<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
self.write_to_stream::<IS_SSL>(socket, b"");
}
fn write_to_stream_using_buffer<const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
buffer: &mut bun_io::StreamBuffer,
data: &[u8],
) -> Result<bool, bun_core::Error> {
let to_send_len = buffer.slice().len();
if to_send_len > 0 {
let amount = write_to_socket::<IS_SSL>(socket, buffer.slice())?;
self.state.request_sent_len += amount;
buffer.cursor += amount;
if amount < to_send_len {
if !data.is_empty() {
let _ = buffer.write(data); }
return Ok(true);
}
if buffer.is_empty() {
buffer.reset();
}
}
if !data.is_empty() {
let sent = write_to_socket_with_buffer_fallback::<IS_SSL>(socket, buffer, data)?;
self.state.request_sent_len += sent;
return Ok(sent < data.len());
}
Ok(false)
}
pub fn write_to_stream<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>, data: &[u8]) {
bun_core::scoped_log!(fetch, "flushStream");
let upgrade_state = self.flags.upgrade_state;
let (stream_buffer_ptr, ended) = {
let HTTPRequestBody::Stream(stream) = &mut self.state.original_request_body else {
return;
};
let Some(buf) = stream.buffer else { return };
(buf, stream.ended)
};
let stream_buffer = ThreadSafeStreamBuffer::from_attached(stream_buffer_ptr);
if upgrade_state == HTTPUpgradeState::Pending {
return;
}
let buffer = stream_buffer.acquire();
let was_empty = buffer.is_empty() && data.is_empty();
if was_empty && ended {
stream_buffer.release();
self.request_stream_detach();
if upgrade_state == HTTPUpgradeState::Upgraded {
socket.shutdown();
}
return;
}
let has_backpressure =
match self.write_to_stream_using_buffer::<IS_SSL>(socket, buffer, data) {
Ok(b) => b,
Err(err) => {
stream_buffer.release();
self.request_stream_detach();
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
if has_backpressure {
stream_buffer.release();
} else {
if ended {
self.state.request_stage = RequestStage::Done;
stream_buffer.release();
self.request_stream_detach();
if upgrade_state == HTTPUpgradeState::Upgraded {
socket.shutdown();
}
} else {
if !was_empty {
stream_buffer.report_drain();
}
stream_buffer.release();
}
}
}
#[inline]
fn request_stream_detach(&mut self) {
if let HTTPRequestBody::Stream(stream) = &mut self.state.original_request_body {
stream.detach();
}
}
pub fn on_writable<const IS_FIRST_CALL: bool, const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
) {
if self.signals.get(signals::Field::Aborted) {
self.close_and_abort::<IS_SSL>(socket);
return;
}
if FeatureFlags::IS_FETCH_PRECONNECT_SUPPORTED {
if self.flags.is_preconnect_only {
self.on_preconnect::<IS_SSL>(socket);
return;
}
}
if let Some(proxy) = self.proxy_tunnel_ptr() {
ProxyTunnel::on_writable::<IS_SSL>(proxy, socket);
}
if self.state.flags.is_waiting_for_cert_check {
return;
}
match self.state.request_stage {
RequestStage::Pending | RequestStage::Headers | RequestStage::Opened => {
bun_core::scoped_log!(fetch, "sendInitialRequestPayload");
self.set_timeout(&socket);
let result =
match self.send_initial_request_payload::<IS_FIRST_CALL, IS_SSL>(socket) {
Ok(r) => r,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
let has_sent_headers = result.has_sent_headers;
let has_sent_body = result.has_sent_body;
let try_sending_more_data = result.try_sending_more_data;
if has_sent_headers && has_sent_body {
if self.flags.proxy_tunneling {
self.state.request_stage = RequestStage::ProxyHandshake;
} else {
self.state.request_stage = RequestStage::Body;
if self.flags.is_streaming_request_body {
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
}
}
return;
}
if has_sent_headers {
if self.flags.proxy_tunneling {
self.state.request_stage = RequestStage::ProxyHandshake;
} else {
self.state.request_stage = RequestStage::Body;
if self.flags.is_streaming_request_body {
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
}
}
debug_assert!(
(matches!(self.state.original_request_body, HTTPRequestBody::Bytes(_))
&& !self.request_body().is_empty())
|| matches!(
self.state.original_request_body,
HTTPRequestBody::Sendfile(_) | HTTPRequestBody::Stream(_)
)
);
if try_sending_more_data {
self.on_writable::<false, IS_SSL>(socket);
}
} else {
self.state.request_stage = RequestStage::Headers;
}
}
RequestStage::Body => {
bun_core::scoped_log!(fetch, "send body");
self.set_timeout(&socket);
match &mut self.state.original_request_body {
HTTPRequestBody::Bytes(_) => {
let to_send = self.request_body();
if !to_send.is_empty() {
let sent = match write_to_socket::<IS_SSL>(socket, to_send) {
Ok(s) => s,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
self.state.request_sent_len += sent;
self.state.request_body =
bun_ptr::RawSlice::new(&self.state.request_body.slice()[sent..]);
}
if self.request_body().is_empty() {
self.state.request_stage = RequestStage::Done;
return;
}
}
HTTPRequestBody::Stream(_) => {
self.flush_stream::<IS_SSL>(socket);
}
HTTPRequestBody::Sendfile(sendfile) => {
if IS_SSL {
panic!(
"sendfile is only supported without SSL. This code should never have been reached!"
);
}
match sendfile.write(socket.fd()) {
crate::send_file::Status::Done => {
self.state.request_stage = RequestStage::Done;
return;
}
crate::send_file::Status::Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
crate::send_file::Status::Again => {
uws::SocketTCP::from_any(socket.socket)
.mark_needs_more_for_sendfile();
}
}
}
}
}
RequestStage::ProxyBody => {
bun_core::scoped_log!(fetch, "send proxy body");
if let Some(proxy_ptr) = self.proxy_tunnel.as_ref().map(|p| p.as_ptr()) {
let proxy = proxy_tunnel::raw_as_mut(proxy_ptr);
match &self.state.original_request_body {
HTTPRequestBody::Bytes(_) => {
self.set_timeout(&socket);
let to_send = self.request_body();
let Ok(sent) = ProxyTunnel::write(proxy, to_send) else {
return;
};
self.state.request_sent_len += sent;
self.state.request_body =
bun_ptr::RawSlice::new(&self.state.request_body.slice()[sent..]);
if self.request_body().is_empty() {
self.state.request_stage = RequestStage::Done;
return;
}
}
HTTPRequestBody::Stream(_) => {
self.flush_stream::<IS_SSL>(socket);
}
HTTPRequestBody::Sendfile(_) => {
panic!(
"sendfile is only supported without SSL. This code should never have been reached!"
);
}
}
}
}
RequestStage::ProxyHeaders => {
bun_core::scoped_log!(fetch, "send proxy headers");
if let Some(proxy_ptr) = self.proxy_tunnel.as_ref().map(|p| p.as_ptr()) {
let proxy = proxy_tunnel::raw_as_mut(proxy_ptr);
self.set_timeout(&socket);
let mut temporary_send_buffer: Vec<u8> = Vec::with_capacity(16 * 1024);
let writer = &mut temporary_send_buffer;
let request = self.build_request(self.request_body().len());
if write_request(writer, &request).is_err() {
self.close_and_fail::<IS_SSL>(err!(OutOfMemory), socket);
return;
}
let headers_len = temporary_send_buffer.len();
if !self.request_body().is_empty()
&& temporary_send_buffer.capacity() - temporary_send_buffer.len() > 0
{
let spare = temporary_send_buffer.capacity() - temporary_send_buffer.len();
let wrote = spare.min(self.request_body().len());
debug_assert!(wrote > 0);
temporary_send_buffer.extend_from_slice(&self.request_body()[0..wrote]);
}
let to_send = &temporary_send_buffer[self.state.request_sent_len..];
if cfg!(debug_assertions) {
debug_assert!(!socket.is_shutdown());
debug_assert!(!socket.is_closed());
}
let Ok(amount) = ProxyTunnel::write(proxy, to_send) else {
return;
};
if IS_FIRST_CALL {
if amount == 0 {
bun_core::scoped_log!(fetch, "is_first_call and amount == 0");
return;
}
}
self.state.request_sent_len += amount;
let has_sent_headers = self.state.request_sent_len >= headers_len;
if has_sent_headers && !self.request_body().is_empty() {
self.state.request_body = bun_ptr::RawSlice::new(
&self.state.request_body.slice()
[self.state.request_sent_len - headers_len..],
);
}
let has_sent_body = self.request_body().is_empty();
if has_sent_headers && has_sent_body {
self.state.request_stage = RequestStage::Done;
return;
}
if has_sent_headers {
self.state.request_stage = RequestStage::ProxyBody;
if self.flags.is_streaming_request_body {
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.progress_update::<IS_SSL>(ctx, socket);
}
debug_assert!(!self.request_body().is_empty());
if amount == to_send.len() {
self.on_writable::<false, IS_SSL>(socket);
}
} else {
self.state.request_stage = RequestStage::ProxyHeaders;
}
}
}
_ => {}
}
}
pub fn resume_after_cert_check<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
if !self.state.flags.is_waiting_for_cert_check {
return;
}
bun_core::scoped_log!(fetch, "resumeAfterCertCheck");
self.state.flags.is_waiting_for_cert_check = false;
self.on_writable::<true, IS_SSL>(socket);
}
pub fn close_and_fail<const IS_SSL: bool>(
&mut self,
err: bun_core::Error,
socket: HttpSocket<IS_SSL>,
) {
bun_core::scoped_log!(fetch, "closeAndFail: {:?}", err);
GenHttpContext::<IS_SSL>::terminate_socket(socket);
self.fail(err);
}
fn start_proxy_handshake<const IS_SSL: bool>(
&mut self,
socket: HttpSocket<IS_SSL>,
start_payload: &[u8],
) {
bun_core::scoped_log!(fetch, "startProxyHandshake");
let ssl_options = if let Some(tls) = &self.tls_props {
tls.get().clone()
} else {
crate::ssl_config::SSLConfig::ZERO
};
ProxyTunnel::start::<IS_SSL>(self, socket, &ssl_options, start_payload);
}
#[inline]
fn handle_short_read<const IS_SSL: bool>(
&mut self,
incoming_data: &[u8],
socket: HttpSocket<IS_SSL>,
needs_move: bool,
) {
if needs_move {
let to_copy = incoming_data;
if !to_copy.is_empty() {
let _ = self.state.response_message_buffer.append(to_copy); }
}
self.set_timeout(&socket);
}
pub fn handle_on_data_headers<const IS_SSL: bool>(
&mut self,
incoming_data: &[u8],
ctx: *mut GenHttpContext<IS_SSL>,
socket: HttpSocket<IS_SSL>,
) {
bun_core::scoped_log!(
fetch,
"handleOnDataHeader data: {}",
BStr::new(incoming_data)
);
let mut to_read = bun_ptr::RawSlice::new(incoming_data);
macro_rules! to_read {
() => {
to_read.slice()
};
}
let mut needs_move = true;
if !self.state.response_message_buffer.list.is_empty() {
let _ = self
.state
.response_message_buffer
.append_slice_exact(incoming_data);
to_read = bun_ptr::RawSlice::new(self.state.response_message_buffer.list.as_slice());
needs_move = false;
}
loop {
let mut amount_read: usize = 0;
self.state.pending_response = Some(picohttp::Response::default());
if to_read!().len() < 16 {
bun_core::scoped_log!(fetch, "handleShortRead");
self.handle_short_read::<IS_SSL>(to_read!(), socket, needs_move);
return;
}
let shared_resp = scratch::response_headers();
let response = match picohttp::Response::parse_parts(
to_read!(),
shared_resp,
Some(&mut amount_read),
) {
Ok(r) => r,
Err(picohttp::ParseResponseError::ShortRead) => {
const MAX_RESPONSE_HEADER_BUFFER: usize = 1024 * 1024;
if to_read!().len() > MAX_RESPONSE_HEADER_BUFFER {
self.close_and_fail::<IS_SSL>(err!(ResponseHeadersTooLarge), socket);
return;
}
self.handle_short_read::<IS_SSL>(to_read!(), socket, needs_move);
return;
}
Err(e) => {
self.close_and_fail::<IS_SSL>(e.into(), socket);
return;
}
};
let response = unsafe { response.detach_lifetime() };
self.state.pending_response = Some(response);
let bytes_read =
(usize::try_from(response.bytes_read).expect("int cast")).min(to_read.len());
to_read = bun_ptr::RawSlice::new(&to_read.slice()[bytes_read..]);
if response.status_code == 101 {
if self.flags.upgrade_state == HTTPUpgradeState::None {
self.close_and_fail::<IS_SSL>(err!(UnrequestedUpgrade), socket);
return;
}
self.flags.upgrade_state = HTTPUpgradeState::Upgraded;
self.signals
.store(signals::Field::Upgraded, true, Ordering::Relaxed);
self.flush_stream::<IS_SSL>(socket);
break;
}
if response.status_code >= 100 && response.status_code < 200 {
bun_core::scoped_log!(fetch, "information headers");
self.state.pending_response = None;
if !needs_move {
let remaining = to_read!().len();
let buffer = &mut self.state.response_message_buffer.list;
let consumed = buffer.len().saturating_sub(remaining);
buffer.drain_front(consumed);
to_read = bun_ptr::RawSlice::new(buffer.as_slice());
}
if to_read!().is_empty() {
return;
}
continue;
}
break;
}
let mut response: picohttp::Response<'static> = self.state.pending_response.unwrap();
let should_continue = match self.handle_response_metadata(&mut response) {
Ok(s) => s,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
self.state.pending_response = Some(response);
if (self.state.content_encoding_i as usize) < response.headers.list.len()
&& !self.state.flags.did_set_content_encoding
{
self.state.flags.did_set_content_encoding = true;
self.state.content_encoding_i = u8::MAX;
self.state.pending_response = Some(response);
}
if should_continue == ShouldContinue::Finished {
if !to_read!().is_empty() {
self.state.flags.allow_keepalive = false;
}
if self.state.flags.is_redirect_pending {
self.do_redirect::<IS_SSL>(ctx, socket);
return;
}
self.clone_metadata();
self.state.flags.received_last_chunk = true;
self.state.content_length = Some(0);
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
if self.flags.proxy_tunneling && self.proxy_tunnel.is_none() {
self.start_proxy_handshake::<IS_SSL>(socket, to_read!());
return;
}
self.clone_metadata();
if to_read!().is_empty() {
if self.signals.get(signals::Field::HeaderProgress) {
self.progress_update::<IS_SSL>(ctx, socket);
}
return;
}
if self.state.response_stage == ResponseStage::Body {
let report_progress = match self.handle_response_body(to_read!(), true) {
Ok(b) => b,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
if report_progress {
if !self.state.is_done()
&& !self.signals.get(signals::Field::ResponseBodyStreaming)
{
return;
}
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
} else if self.state.response_stage == ResponseStage::BodyChunk {
self.set_timeout(&socket);
let report_progress = match self.handle_response_body_chunked_encoding(to_read!()) {
Ok(b) => b,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
if report_progress {
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
}
if self.signals.get(signals::Field::HeaderProgress) {
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
}
pub fn on_data<const IS_SSL: bool>(
&mut self,
incoming_data: &[u8],
ctx: *mut GenHttpContext<IS_SSL>,
socket: HttpSocket<IS_SSL>,
) {
bun_core::scoped_log!(fetch, "onData {}", incoming_data.len());
if self.signals.get(signals::Field::Aborted) {
self.close_and_abort::<IS_SSL>(socket);
return;
}
if let Some(proxy) = self.proxy_tunnel_ptr() {
self.set_timeout(&socket);
ProxyTunnel::receive(proxy, incoming_data);
return;
}
if self.state.flags.is_waiting_for_cert_check {
self.state.pending_response = None;
self.close_and_fail::<IS_SSL>(err!(UnexpectedData), socket);
return;
}
match self.state.response_stage {
ResponseStage::Pending | ResponseStage::Headers => {
self.handle_on_data_headers::<IS_SSL>(incoming_data, ctx, socket);
}
ResponseStage::Body => {
self.set_timeout(&socket);
let report_progress = match self.handle_response_body(incoming_data, false) {
Ok(b) => b,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
if report_progress {
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
}
ResponseStage::BodyChunk => {
self.set_timeout(&socket);
let report_progress =
match self.handle_response_body_chunked_encoding(incoming_data) {
Ok(b) => b,
Err(err) => {
self.close_and_fail::<IS_SSL>(err, socket);
return;
}
};
if report_progress {
self.progress_update::<IS_SSL>(ctx, socket);
return;
}
}
ResponseStage::Fail => {}
_ => {
self.state.pending_response = None;
self.close_and_fail::<IS_SSL>(err!(UnexpectedData), socket);
return;
}
}
}
pub fn close_and_abort<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
self.close_and_fail::<IS_SSL>(err!(Aborted), socket);
}
fn complete_connecting_process(&mut self) {
if self.flags.defer_fail_until_connecting_is_complete {
self.flags.defer_fail_until_connecting_is_complete = false;
if self.state.stage == Stage::Fail {
self.dispatch_result_and_reset(true);
}
}
}
fn resolve_pending_h2(&mut self, resolution: PendingH2Resolution) {
let Some(pc_ptr) = self.pending_h2.take() else {
return;
};
let Some(pc) = h2::PendingConnect::unregister_from(
pc_ptr.as_ptr(),
Self::ssl_ctx_mut(self.get_ssl_ctx::<true>()),
) else {
return;
};
for waiter_ptr in pc.waiters.iter().copied() {
let waiter = h2::PendingConnect::waiter_mut(waiter_ptr);
if waiter.signals.get(signals::Field::Aborted) {
waiter.fail(err!(Aborted));
continue;
}
match resolution {
PendingH2Resolution::H2(session) => h2::ClientSession::enqueue(session, waiter),
PendingH2Resolution::H1 => {
if waiter.flags.force_http2 {
waiter.fail(err!(HTTP2Unsupported));
continue;
}
waiter.flags.force_http1 = true;
waiter.start_::<true>();
}
PendingH2Resolution::LeaderFailed => waiter.start_::<true>(),
}
}
}
fn fail(&mut self, err: bun_core::Error) {
self.unregister_abort_tracker();
self.resolve_pending_h2(PendingH2Resolution::LeaderFailed);
self.close_proxy_tunnel(true);
if self.state.stage != Stage::Done && self.state.stage != Stage::Fail {
self.state.request_stage = RequestStage::Fail;
self.state.response_stage = ResponseStage::Fail;
self.state.fail = Some(err);
self.state.stage = Stage::Fail;
if !self.flags.defer_fail_until_connecting_is_complete {
self.dispatch_result_and_reset(true);
}
}
}
pub fn clone_metadata(&mut self) {
debug_assert!(self.state.pending_response.is_some());
if let Some(response) = self.state.pending_response {
if let Some(old) = self.state.cloned_metadata.take() {
drop(old); }
let mut builder = picohttp::StringBuilder::default();
response.count(&mut builder);
builder.count(self.url.href);
let _ = builder.allocate();
let headers_buf = bun_core::heap::release(
vec![picohttp::Header::ZERO; response.headers.list.len()].into_boxed_slice(),
);
let cloned_response = response.clone(headers_buf, &mut builder);
self.state.pending_response = None;
let href = bun_ptr::RawSlice::new(unsafe { builder.append_raw(self.url.href) });
let owned_buf = builder.move_to_slice();
self.state.cloned_metadata = Some(HTTPResponseMetadata {
owned_buf,
response: cloned_response,
url: href,
});
} else {
self.state.cloned_metadata = Some(HTTPResponseMetadata::default());
}
}
pub fn set_timeout<S: SocketTimeout>(&self, socket: &S) {
if self.flags.disable_timeout || idle_timeout_seconds() == 0 {
socket.set_timeout(0);
return;
}
socket.set_timeout(idle_timeout_seconds());
}
pub fn drain_response_body<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
match self.state.stage {
Stage::Done | Stage::Fail => return,
_ => {}
}
if self.state.fail.is_some() {
return;
}
if self.state.flags.is_redirect_pending {
return;
}
let Some(body_out_str) = self.body_out_str() else {
return;
};
if body_out_str.list.is_empty() {
return;
}
let ctx = self.get_ssl_ctx::<IS_SSL>();
self.send_progress_update_without_stage_check::<IS_SSL>(ctx, socket);
}
fn send_progress_update_without_stage_check<const IS_SSL: bool>(
&mut self,
ctx: *mut GenHttpContext<IS_SSL>,
socket: HttpSocket<IS_SSL>,
) {
if IS_SSL && self.tls_info.is_none() {
let ssl_ptr: *mut boringssl::c::SSL = socket
.get_native_handle()
.map(|p| p.cast())
.unwrap_or(core::ptr::null_mut());
if !ssl_ptr.is_null()
&& unsafe { boringssl::c::SSL_is_init_finished(ssl_ptr) } == 1
{
self.tls_info = Some(unsafe { BunTlsInfo::from_ssl(ssl_ptr) });
}
}
if self.flags.protocol != Protocol::Http1_1 {
return self.send_progress_update_multiplexed();
}
let body = self.state.body_out_str;
let body_snapshot = body_out::take_list(body);
let callback = self.result_callback;
let (
has_more,
redirected,
can_stream,
is_http2,
fail,
metadata,
body_size,
certificate_info,
tls_info,
) = {
let r = self.to_result();
(
r.has_more,
r.redirected,
r.can_stream,
r.is_http2,
r.fail,
r.metadata,
r.body_size,
r.certificate_info,
r.tls_info,
)
}; let is_done = !has_more;
bun_core::scoped_log!(fetch, "progressUpdate {}", is_done);
if is_done {
self.unregister_abort_tracker();
let tunnel_poolable = if let Some(t) = self.proxy_tunnel.as_deref() {
self.state.request_stage == RequestStage::Done
&& t.write_buffer.is_empty()
&& t.wrapper
.as_ref()
.map(|w| !w.is_shutdown())
.unwrap_or(false)
} else {
true
};
let request_side_drained = match &self.state.original_request_body {
HTTPRequestBody::Bytes(_) => self.state.request_body.is_empty(),
_ => true,
};
if self.is_keep_alive_possible()
&& !socket.is_closed_or_has_error()
&& tunnel_poolable
&& request_side_drained
{
bun_core::scoped_log!(fetch, "release socket");
let tunnel = self.proxy_tunnel.take();
if let Some(t) = &tunnel {
proxy_tunnel::raw_as_mut(t.as_ptr()).detach_owner(&*self);
}
let had_tunnel = tunnel.is_some();
Self::ssl_ctx_mut(ctx).release_socket(
socket,
self.flags.did_have_handshaking_error && !self.flags.reject_unauthorized,
self.flags.reject_unauthorized,
self.connected_url.hostname,
self.connected_url.get_port_auto(),
self.tls_props.as_ref(),
tunnel,
if had_tunnel { self.url.hostname } else { b"" },
if had_tunnel {
self.url.get_port_auto()
} else {
0
},
if had_tunnel || (IS_SSL && self.http_proxy.is_none()) {
self.proxy_auth_hash()
} else {
0
},
None,
);
} else {
if self.proxy_tunnel.is_some() {
bun_core::scoped_log!(fetch, "close the tunnel");
self.close_proxy_tunnel(true);
}
GenHttpContext::<IS_SSL>::close_socket(socket);
}
self.state.reset();
self.state.response_stage = ResponseStage::Done;
self.state.request_stage = RequestStage::Done;
self.state.stage = Stage::Done;
self.flags.proxy_tunneling = false;
bun_core::scoped_log!(fetch, "done");
}
body_out::restore_list(body, body_snapshot);
let async_http = self.parent_async_http();
let result = HTTPClientResult {
body: body_out::opt_mut(body),
has_more,
redirected,
can_stream,
is_http2,
fail,
metadata,
body_size,
certificate_info,
tls_info,
};
callback.run(async_http, result);
if PRINT_EVERY != 0 {
let i = PRINT_EVERY_I.fetch_add(1, Ordering::Relaxed) + 1;
if i.is_multiple_of(PRINT_EVERY) {
Output::prettyln(format_args!("Heap stats for HTTP thread\n"));
Output::flush();
PRINT_EVERY_I.store(0, Ordering::Relaxed);
}
}
}
fn send_progress_update_multiplexed(&mut self) {
debug_assert!(self.flags.protocol != Protocol::Http1_1);
let body = self.state.body_out_str;
let body_snapshot = body_out::take_list(body);
let callback = self.result_callback;
let (
has_more,
redirected,
can_stream,
is_http2,
fail,
metadata,
body_size,
certificate_info,
tls_info,
) = {
let r = self.to_result();
(
r.has_more,
r.redirected,
r.can_stream,
r.is_http2,
r.fail,
r.metadata,
r.body_size,
r.certificate_info,
r.tls_info,
)
}; let is_done = !has_more;
bun_core::scoped_log!(fetch, "progressUpdate {}", is_done);
if is_done {
self.unregister_abort_tracker();
self.state.reset();
self.state.response_stage = ResponseStage::Done;
self.state.request_stage = RequestStage::Done;
self.state.stage = Stage::Done;
self.flags.proxy_tunneling = false;
}
body_out::restore_list(body, body_snapshot);
let async_http = self.parent_async_http();
let result = HTTPClientResult {
body: body_out::opt_mut(body),
has_more,
redirected,
can_stream,
is_http2,
fail,
metadata,
body_size,
certificate_info,
tls_info,
};
callback.run(async_http, result);
}
fn do_redirect_multiplexed(&mut self) {
debug_assert!(self.flags.protocol != Protocol::Http1_1);
bun_core::scoped_log!(fetch, "doRedirectMultiplexed");
if self.state.flags.clear_hostname_on_redirect {
self.state.flags.clear_hostname_on_redirect = false;
self.hostname = None;
}
if matches!(self.state.original_request_body, HTTPRequestBody::Stream(_)) {
self.flags.is_streaming_request_body = false;
}
let _ = core::mem::ManuallyDrop::new(core::mem::take(&mut self.unix_socket_path));
let request_body: &[u8] = if self.state.flags.resend_request_body_on_redirect
&& matches!(self.state.original_request_body, HTTPRequestBody::Bytes(_))
{
match &self.state.original_request_body {
HTTPRequestBody::Bytes(b) => b,
_ => unreachable!(),
}
} else {
b""
};
self.state.response_message_buffer = MutableString::default();
let body_out_str = self.state.body_out_str.unwrap();
self.remaining_redirect_count = self.remaining_redirect_count.saturating_sub(1);
self.flags.redirected = true;
debug_assert!(self.redirect_type == FetchRedirect::Follow);
self.unregister_abort_tracker();
self.connected_url = URL::default();
self.prev_redirect = Vec::new();
if self.remaining_redirect_count == 0 {
self.fail(err!(TooManyRedirects));
return;
}
self.state.reset();
self.flags.proxy_tunneling = false;
self.flags.protocol = Protocol::Http1_1;
self.start(
HTTPRequestBody::Bytes(request_body),
body_out::as_mut(body_out_str),
);
}
pub fn progress_update_h3(&mut self) {
debug_assert!(self.flags.protocol == Protocol::Http3);
if self.state.stage == Stage::Done || self.state.stage == Stage::Fail {
return;
}
if self.state.flags.is_redirect_pending && self.state.fail.is_none() {
if self.state.is_done() {
self.do_redirect_multiplexed();
}
return;
}
self.send_progress_update_multiplexed();
}
pub fn do_redirect_h3(&mut self) {
debug_assert!(self.flags.protocol == Protocol::Http3);
self.do_redirect_multiplexed();
}
pub fn progress_update<const IS_SSL: bool>(
&mut self,
ctx: *mut GenHttpContext<IS_SSL>,
socket: HttpSocket<IS_SSL>,
) {
if self.state.stage != Stage::Done && self.state.stage != Stage::Fail {
if self.state.flags.is_redirect_pending && self.state.fail.is_none() {
if self.state.is_done() {
self.do_redirect::<IS_SSL>(ctx, socket);
}
return;
}
self.send_progress_update_without_stage_check::<IS_SSL>(ctx, socket);
}
}
pub fn on_preconnect<const IS_SSL: bool>(&mut self, socket: HttpSocket<IS_SSL>) {
bun_core::scoped_log!(fetch, "onPreconnect({})", BStr::new(self.url.href));
self.unregister_abort_tracker();
let ctx = self.get_ssl_ctx::<IS_SSL>();
Self::ssl_ctx_mut(ctx).release_socket(
socket,
self.flags.did_have_handshaking_error && !self.flags.reject_unauthorized,
self.flags.reject_unauthorized,
self.url.hostname,
self.url.get_port_auto(),
self.tls_props.as_ref(),
None,
b"",
0,
0,
None,
);
self.state.reset();
self.state.response_stage = ResponseStage::Done;
self.state.request_stage = RequestStage::Done;
self.state.stage = Stage::Done;
self.flags.proxy_tunneling = false;
self.result_callback.run(
self.parent_async_http(),
HTTPClientResult {
fail: None,
metadata: None,
has_more: false,
..Default::default()
},
);
}
#[inline]
fn parent_async_http(&mut self) -> *mut AsyncHTTP<'static> {
unsafe {
bun_core::from_field_ptr!(AsyncHTTP<'static>, client, std::ptr::from_mut::<Self>(self))
}
}
pub fn to_result(&mut self) -> HTTPClientResult<'_> {
let body_size: BodySize = if self.state.is_chunked_encoding() {
BodySize::TotalReceived(self.state.total_body_received)
} else if let Some(content_length) = self.state.content_length {
BodySize::ContentLength(content_length)
} else {
BodySize::Unknown
};
let mut certificate_info: Option<CertificateInfo> = None;
if let Some(info) = self.state.certificate_info.take() {
certificate_info = Some(info);
} else if let Some(metadata) = self.state.cloned_metadata.take() {
return HTTPClientResult {
metadata: Some(metadata),
body: body_out::opt_mut(self.state.body_out_str),
redirected: self.flags.redirected,
fail: self.state.fail,
has_more: certificate_info.is_some()
|| (self.state.fail.is_none() && !self.state.is_done()),
body_size,
certificate_info: None,
can_stream: (self.state.request_stage == RequestStage::Body
|| self.state.request_stage == RequestStage::ProxyBody)
&& self.flags.is_streaming_request_body,
is_http2: self.flags.protocol != Protocol::Http1_1,
tls_info: self.tls_info.clone(),
};
}
HTTPClientResult {
body: body_out::opt_mut(self.state.body_out_str),
metadata: None,
redirected: self.flags.redirected,
fail: self.state.fail,
has_more: certificate_info.is_some()
|| (self.state.fail.is_none() && !self.state.is_done()),
body_size,
certificate_info,
can_stream: (self.state.request_stage == RequestStage::Body
|| self.state.request_stage == RequestStage::ProxyBody)
&& self.flags.is_streaming_request_body,
is_http2: self.flags.protocol != Protocol::Http1_1,
tls_info: self.tls_info.clone(),
}
}
pub fn handle_response_body(
&mut self,
incoming_data: &[u8],
is_only_buffer: bool,
) -> Result<bool, bun_core::Error> {
debug_assert!(self.state.transfer_encoding == Encoding::Identity);
let content_length = self.state.content_length;
if let Some(len) = content_length
&& incoming_data.len() > len.saturating_sub(self.state.total_body_received)
{
self.state.flags.allow_keepalive = false;
}
if is_only_buffer
&& let Some(len) = content_length
&& incoming_data.len() >= len
{
self.handle_response_body_from_single_packet(&incoming_data[0..len])?;
Ok(true)
} else {
self.handle_response_body_from_multiple_packets(incoming_data)
}
}
fn handle_response_body_from_single_packet(
&mut self,
incoming_data: &[u8],
) -> Result<(), bun_core::Error> {
if !self.state.is_chunked_encoding() {
self.state.total_body_received += incoming_data.len();
bun_core::scoped_log!(
fetch,
"handleResponseBodyFromSinglePacket {}",
self.state.total_body_received
);
}
if !self.state.flags.is_redirect_pending {
if self.state.encoding.is_compressed() {
let body_out = self.state.body_out_str.unwrap();
self.state
.decompress_bytes(incoming_data, body_out::as_mut(body_out), true)?;
} else {
self.state
.get_body_buffer()
.append_slice_exact(incoming_data)?;
}
if self.state.response_message_buffer.owns(incoming_data) {
if cfg!(debug_assertions) {
debug_assert!(
self.state.get_body_buffer().list.as_ptr()
!= self.state.response_message_buffer.list.as_ptr()
);
}
self.state.response_message_buffer = MutableString::default();
}
}
self.report_progress(incoming_data.len());
Ok(())
}
fn handle_response_body_from_multiple_packets(
&mut self,
incoming_data: &[u8],
) -> Result<bool, bun_core::Error> {
let content_length = self.state.content_length;
let remainder: &[u8] = if let Some(cl) = content_length {
let remaining_content_length = cl.saturating_sub(self.state.total_body_received);
&incoming_data[0..incoming_data.len().min(remaining_content_length)]
} else {
incoming_data
};
if !self.state.flags.is_redirect_pending {
let buffer = self.state.get_body_buffer();
if buffer.list.is_empty() && incoming_data.len() < PREALLOCATE_MAX {
let _ = buffer.list.try_reserve_exact(incoming_data.len());
}
let _ = buffer.write(remainder)?;
}
self.state.total_body_received += remainder.len();
bun_core::scoped_log!(
fetch,
"handleResponseBodyFromMultiplePackets {}",
self.state.total_body_received
);
let total_received = self.state.total_body_received;
self.report_progress(total_received);
let is_done =
content_length.is_some() && self.state.total_body_received >= content_length.unwrap();
if is_done
|| self.signals.get(signals::Field::ResponseBodyStreaming)
|| content_length.is_none()
{
let is_final_chunk = is_done;
let buffer_snap = core::mem::take(&mut self.state.get_body_buffer().list);
let processed = self
.state
.process_body_buffer(buffer_snap, is_final_chunk)?;
self.state.flags.is_libdeflate_fast_path_disabled = true;
let total_received = self.state.total_body_received;
self.report_progress(total_received);
return Ok(is_done || processed);
}
Ok(false)
}
pub fn handle_response_body_chunked_encoding(
&mut self,
incoming_data: &[u8],
) -> Result<bool, bun_core::Error> {
let small_len = 16 * 1024usize;
if incoming_data.len() <= small_len && self.state.get_body_buffer().list.is_empty() {
self.handle_response_body_chunked_encoding_from_single_packet(incoming_data)
} else {
self.handle_response_body_chunked_encoding_from_multiple_packets(incoming_data)
}
}
fn handle_response_body_chunked_encoding_from_multiple_packets(
&mut self,
incoming_data: &[u8],
) -> Result<bool, bun_core::Error> {
let (decoder, body_buf) = self.state.chunked_decoder_and_body_buffer();
body_buf.append_slice(incoming_data)?;
decoder.consume_trailer = 1;
let mut bytes_decoded = incoming_data.len();
let pret = unsafe {
picohttp::phr_decode_chunked(
&raw mut *decoder,
body_buf
.list
.as_mut_ptr()
.add(body_buf.list.len().saturating_sub(incoming_data.len())),
&raw mut bytes_decoded,
)
};
let new_len = body_buf
.list
.len()
.saturating_sub(incoming_data.len() - bytes_decoded);
body_buf.list.truncate(new_len);
let buffer_len = body_buf.list.len();
self.state.total_body_received += bytes_decoded;
bun_core::scoped_log!(
fetch,
"handleResponseBodyChunkedEncodingFromMultiplePackets {}",
self.state.total_body_received
);
match pret {
-1 => return Err(err!(InvalidHTTPResponse)),
-2 => {
self.report_progress(buffer_len);
if self.signals.get(signals::Field::ResponseBodyStreaming) {
self.state.flags.is_libdeflate_fast_path_disabled = true;
let buffer_snap = core::mem::take(&mut self.state.get_body_buffer().list);
return self.state.process_body_buffer(buffer_snap, false);
}
return Ok(false);
}
_ => {
self.state.flags.received_last_chunk = true;
let buffer_snap = core::mem::take(&mut self.state.get_body_buffer().list);
let _ = self.state.process_body_buffer(buffer_snap, true)?;
self.report_progress(buffer_len);
return Ok(true);
}
}
}
fn handle_response_body_chunked_encoding_from_single_packet(
&mut self,
incoming_data: &[u8],
) -> Result<bool, bun_core::Error> {
let small = scratch::single_packet_small_buffer();
debug_assert!(incoming_data.len() <= small.len());
self.state.chunked_decoder.consume_trailer = 1;
let in_len = incoming_data.len();
let buffer: &mut [u8] = if self.state.response_message_buffer.owns(incoming_data) {
let base = self.state.response_message_buffer.list.as_mut_ptr();
let off = incoming_data.as_ptr() as usize - base as usize;
unsafe { bun_core::ffi::slice_mut(base.add(off), in_len) }
} else {
small[0..in_len].copy_from_slice(incoming_data);
&mut small[0..in_len]
};
let mut bytes_decoded = in_len;
let pret = unsafe {
picohttp::phr_decode_chunked(
&raw mut self.state.chunked_decoder,
buffer.as_mut_ptr().add(buffer.len().saturating_sub(in_len)),
&raw mut bytes_decoded,
)
};
let new_len = buffer.len().saturating_sub(in_len - bytes_decoded);
let buffer = &mut buffer[..new_len];
self.state.total_body_received += bytes_decoded;
bun_core::scoped_log!(
fetch,
"handleResponseBodyChunkedEncodingFromSinglePacket {}",
self.state.total_body_received
);
match pret {
-1 => Err(err!(InvalidHTTPResponse)),
-2 => {
self.report_progress(buffer.len());
self.state.get_body_buffer().append_slice_exact(buffer)?;
if self.signals.get(signals::Field::ResponseBodyStreaming) {
self.state.flags.is_libdeflate_fast_path_disabled = true;
let buffer_snap = core::mem::take(&mut self.state.get_body_buffer().list);
return self.state.process_body_buffer(buffer_snap, true);
}
Ok(false)
}
_ => {
self.state.flags.received_last_chunk = true;
self.handle_response_body_from_single_packet(buffer)?;
debug_assert!(
self.body_out_str()
.map(|b| b.list.as_ptr())
.unwrap_or(core::ptr::null())
!= buffer.as_ptr()
);
self.report_progress(buffer.len());
Ok(true)
}
}
}
pub fn handle_response_metadata(
&mut self,
response: &mut picohttp::Response,
) -> Result<ShouldContinue, bun_core::Error> {
let mut location: &[u8] = b"";
let mut pretend_304 = false;
let mut is_server_sent_events = false;
for (header_i, header) in response.headers.list.iter().enumerate() {
match hash_header_name(header.name()) {
h if h == hash_header_const(b"Content-Length") => {
if self.flags.proxy_tunneling
&& self.proxy_tunnel.is_none()
&& response.status_code == 200
{
continue;
}
let Ok(content_length) = bun_core::parse_unsigned::<usize>(header.value(), 10)
else {
return Err(err!(InvalidContentLength));
};
if self.method.has_body() {
if self
.state
.content_length
.is_some_and(|prev| prev != content_length)
{
return Err(err!(InvalidContentLength));
}
self.state.content_length = Some(content_length);
} else {
self.state.content_length = Some(0);
}
}
h if h == hash_header_const(b"Content-Type") => {
if strings::index_of(header.value(), b"text/event-stream").is_some() {
is_server_sent_events = true;
}
}
h if h == hash_header_const(b"Content-Encoding") => {
if !self.flags.disable_decompression {
if header.value() == b"gzip" {
self.state.encoding = Encoding::Gzip;
self.state.content_encoding_i = header_i as u8;
} else if header.value() == b"deflate" {
self.state.encoding = Encoding::Deflate;
self.state.content_encoding_i = header_i as u8;
} else if header.value() == b"br" {
self.state.encoding = Encoding::Brotli;
self.state.content_encoding_i = header_i as u8;
} else if header.value() == b"zstd" {
self.state.encoding = Encoding::Zstd;
self.state.content_encoding_i = header_i as u8;
}
}
}
h if h == hash_header_const(b"Transfer-Encoding") => {
if self.flags.proxy_tunneling
&& self.proxy_tunnel.is_none()
&& response.status_code == 200
{
continue;
}
if header.value() == b"gzip" {
if !self.flags.disable_decompression {
self.state.transfer_encoding = Encoding::Gzip;
}
} else if header.value() == b"deflate" {
if !self.flags.disable_decompression {
self.state.transfer_encoding = Encoding::Deflate;
}
} else if header.value() == b"br" {
if !self.flags.disable_decompression {
self.state.transfer_encoding = Encoding::Brotli;
}
} else if header.value() == b"zstd" {
if !self.flags.disable_decompression {
self.state.transfer_encoding = Encoding::Zstd;
}
} else if header.value() == b"identity" {
self.state.transfer_encoding = Encoding::Identity;
} else if header.value() == b"chunked" {
self.state.transfer_encoding = Encoding::Chunked;
} else {
return Err(err!(UnsupportedTransferEncoding));
}
}
h if h == hash_header_const(b"Location") => {
location = header.value();
}
h if h == hash_header_const(b"Connection") => {
if response.status_code >= 200 && response.status_code <= 299 {
if bun_core::strings::eql_case_insensitive_ascii_check_length(
header.value(),
b"close",
) {
self.state.flags.allow_keepalive = false;
} else if bun_core::strings::eql_case_insensitive_ascii_check_length(
header.value(),
b"keep-alive",
) {
self.state.flags.allow_keepalive = true;
}
}
}
h if h == hash_header_const(b"Last-Modified") => {
pretend_304 = self.flags.force_last_modified
&& response.status_code > 199
&& response.status_code < 300
&& !self.if_modified_since.is_empty()
&& self.if_modified_since == header.value();
}
h if h == hash_header_const(b"Alt-Svc") => {
if self.is_https()
&& self.unix_socket_path.slice().len() == 0
&& h3_alt_svc_enabled()
{
h3::alt_svc::record(
self.url.hostname,
self.url.get_port_auto(),
header.value(),
);
}
}
_ => {}
}
}
if self.verbose != HTTPVerboseLevel::None {
print_response(response);
}
if pretend_304 {
response.status_code = 304;
}
if (response.status_code >= 100 && response.status_code < 200)
|| response.status_code == 204
|| response.status_code == 304
{
self.state.content_length = Some(0);
}
if !self.flags.proxy_tunneling {
if self.state.content_length.is_none()
&& self.state.transfer_encoding != Encoding::Chunked
{
self.state.flags.allow_keepalive = false;
}
}
let mut is_proxy_connect_failure = false;
if self.flags.proxy_tunneling && self.proxy_tunnel.is_none() {
if response.status_code == 200 {
return Ok(ShouldContinue::ContinueStreaming);
}
self.flags.proxy_tunneling = false;
self.flags.disable_keepalive = true;
is_proxy_connect_failure = true;
}
let status_code = response.status_code;
if status_code == 407 {
self.flags.disable_keepalive = true;
}
let is_redirect = status_code >= 300 && status_code <= 399;
if is_redirect {
if !is_proxy_connect_failure
&& self.redirect_type == FetchRedirect::Follow
&& !location.is_empty()
&& self.remaining_redirect_count > 0
{
match status_code {
302 | 301 | 307 | 308 | 303 => {
if status_code != 303
&& matches!(
self.state.original_request_body,
HTTPRequestBody::Stream(_)
)
{
return Err(err!(RequestBodyNotReusable));
}
let is_same_origin;
{
if let Some(i) = strings::index_of(location, b"://") {
let mut string_builder = StringBuilder::default();
let is_protocol_relative = i == 0;
let protocol_name: &[u8] = if is_protocol_relative {
self.url.display_protocol()
} else {
&location[0..i]
};
let is_http = protocol_name == b"http";
if is_http || protocol_name == b"https" {
} else {
return Err(err!(UnsupportedRedirectProtocol));
}
if (protocol_name.len() * usize::from(is_protocol_relative))
+ location.len()
> MAX_REDIRECT_URL_LENGTH
{
return Err(err!(RedirectURLTooLong));
}
string_builder.count(location);
if is_protocol_relative {
if is_http {
string_builder.count(b"http");
} else {
string_builder.count(b"https");
}
}
string_builder.allocate()?;
if is_protocol_relative {
if is_http {
let _ = string_builder.append(b"http");
} else {
let _ = string_builder.append(b"https");
}
}
let _ = string_builder.append(location);
if cfg!(debug_assertions) {
debug_assert!(string_builder.cap == string_builder.len);
}
let input =
BunString::borrow_utf8(string_builder.allocated_slice());
let normalized_url =
OwnedString::new(bun_url::href_from_string(&input));
if normalized_url.tag() == BunStringTag::Dead {
return Err(err!(RedirectURLInvalid));
}
let normalized_url_str = normalized_url.to_owned_slice();
let new_url: URL<'a> =
unsafe { URL::parse(&normalized_url_str).erase_lifetime() };
is_same_origin = strings::eql_case_insensitive_ascii(
strings::without_trailing_slash(new_url.origin),
strings::without_trailing_slash(self.url.origin),
true,
);
self.url = new_url;
debug_assert!(self.prev_redirect.is_empty());
self.prev_redirect =
core::mem::replace(&mut self.redirect, normalized_url_str);
} else if location.starts_with(b"//") {
let mut string_builder = StringBuilder::default();
let protocol_name = self.url.display_protocol();
if protocol_name.len() + 1 + location.len()
> MAX_REDIRECT_URL_LENGTH
{
return Err(err!(RedirectURLTooLong));
}
let is_http = protocol_name == b"http";
if is_http {
string_builder.count(b"http:");
} else {
string_builder.count(b"https:");
}
string_builder.count(location);
string_builder.allocate()?;
if is_http {
let _ = string_builder.append(b"http:");
} else {
let _ = string_builder.append(b"https:");
}
let _ = string_builder.append(location);
if cfg!(debug_assertions) {
debug_assert!(string_builder.cap == string_builder.len);
}
let input =
BunString::borrow_utf8(string_builder.allocated_slice());
let normalized_url =
OwnedString::new(bun_url::href_from_string(&input));
if normalized_url.tag() == BunStringTag::Dead {
return Err(err!(RedirectURLInvalid));
}
let normalized_url_str = normalized_url.to_owned_slice();
let new_url: URL<'a> =
unsafe { URL::parse(&normalized_url_str).erase_lifetime() };
is_same_origin = strings::eql_case_insensitive_ascii(
strings::without_trailing_slash(new_url.origin),
strings::without_trailing_slash(self.url.origin),
true,
);
self.url = new_url;
debug_assert!(self.prev_redirect.is_empty());
self.prev_redirect =
core::mem::replace(&mut self.redirect, normalized_url_str);
} else {
let original_url = self.url.clone();
let base = BunString::borrow_utf8(original_url.href);
let rel = BunString::borrow_utf8(location);
let new_url_ = OwnedString::new(bun_url::join(&base, &rel));
if new_url_.is_empty() {
return Err(err!(InvalidRedirectURL));
}
let new_url = new_url_.to_owned_slice();
self.url = unsafe { URL::parse(&new_url).erase_lifetime() };
is_same_origin = strings::eql_case_insensitive_ascii(
strings::without_trailing_slash(self.url.origin),
strings::without_trailing_slash(original_url.origin),
true,
);
debug_assert!(self.prev_redirect.is_empty());
self.prev_redirect =
core::mem::replace(&mut self.redirect, new_url);
}
}
if ((status_code == 301 || status_code == 302)
&& self.method == Method::POST)
|| (status_code == 303
&& self.method != Method::GET
&& self.method != Method::HEAD)
{
self.method = Method::GET;
if self.header_entries.len() > 0 {
const REQUEST_BODY_HEADER: [&[u8]; 3] = [
b"Content-Encoding",
b"Content-Language",
b"Content-Location",
];
let mut i: usize = 0;
let mut len = self.header_entries.len();
'outer: while i < len {
let names = self.header_entries.items_name();
let name = self.header_str(names[i]);
match name.len() {
l if l == b"Content-Type".len() => {
let hash = hash_header_name(name);
if hash == hash_header_const(b"Content-Type") {
let _ = self.header_entries.ordered_remove(i);
len = self.header_entries.len();
continue 'outer;
}
}
l if l == b"Content-Encoding".len() => {
let hash = hash_header_name(name);
for hash_value in REQUEST_BODY_HEADER {
if hash == hash_header_const(hash_value) {
let _ = self.header_entries.ordered_remove(i);
len = self.header_entries.len();
continue 'outer;
}
}
}
_ => {}
}
i += 1;
}
}
}
if !is_same_origin {
self.state.flags.clear_hostname_on_redirect = true;
}
if !is_same_origin && self.header_entries.len() > 0 {
struct H {
name: &'static [u8],
hash: u64,
}
let headers_to_remove: [H; 4] = [
H {
name: b"Authorization",
hash: *AUTHORIZATION_HEADER_HASH,
},
H {
name: b"Proxy-Authorization",
hash: *PROXY_AUTHORIZATION_HEADER_HASH,
},
H {
name: b"Cookie",
hash: *COOKIE_HEADER_HASH,
},
H {
name: HOST_HEADER_NAME,
hash: hash_header_const(HOST_HEADER_NAME),
},
];
for to_remove in headers_to_remove.iter() {
let mut i = 0;
while i < self.header_entries.len() {
let name = self.header_str(self.header_entries.items_name()[i]);
if name.len() == to_remove.name.len()
&& hash_header_name(name) == to_remove.hash
{
let _ = self.header_entries.ordered_remove(i);
} else {
i += 1;
}
}
}
}
self.state.flags.is_redirect_pending = true;
if self.method.has_request_body() {
self.state.flags.resend_request_body_on_redirect = true;
}
}
_ => {}
}
} else if !is_proxy_connect_failure && self.redirect_type == FetchRedirect::Error {
return Err(err!(UnexpectedRedirect));
}
}
self.state.response_stage = if self.state.transfer_encoding == Encoding::Chunked {
ResponseStage::BodyChunk
} else {
ResponseStage::Body
};
let content_length = self.state.content_length;
if let Some(length) = content_length {
bun_core::scoped_log!(
fetch,
"handleResponseMetadata: content_length is {} and transfer_encoding {:?}",
length,
self.state.transfer_encoding
);
} else {
bun_core::scoped_log!(
fetch,
"handleResponseMetadata: content_length is null and transfer_encoding {:?}",
self.state.transfer_encoding
);
}
if self.flags.upgrade_state == HTTPUpgradeState::Upgraded {
self.state.content_length = None;
self.state.flags.allow_keepalive = false;
return Ok(ShouldContinue::ContinueStreaming);
}
if self.method.has_body()
&& (content_length.is_none()
|| content_length.unwrap() > 0
|| !self.state.flags.allow_keepalive
|| self.state.transfer_encoding == Encoding::Chunked
|| is_server_sent_events)
{
Ok(ShouldContinue::ContinueStreaming)
} else {
Ok(ShouldContinue::Finished)
}
}
}