use core::ptr::NonNull;
use core::sync::atomic::Ordering;
use bun_threading::Guarded;
use bun_uws as uws;
use bun_uws::quic;
use super::ClientSession;
use super::client_session::session_mut;
pub struct PendingConnect {
session: *mut ClientSession,
pc: *mut quic::PendingConnect,
loop_ptr: *mut uws::Loop,
}
impl Drop for PendingConnect {
fn drop(&mut self) {
unsafe { ClientSession::deref(self.session) };
}
}
impl PendingConnect {
#[inline]
fn pc_mut<'a>(&self) -> &'a mut quic::PendingConnect {
unsafe { &mut *self.pc }
}
pub fn register(session: *mut ClientSession, pc: *mut quic::PendingConnect, l: *mut uws::Loop) {
session_mut(session).ref_();
let self_ = Box::new(PendingConnect {
session,
pc,
loop_ptr: l,
});
let addrinfo = self_.pc_mut().addrinfo();
let self_ = bun_core::heap::into_raw(self_);
unsafe { bun_dns::internal::register_quic(addrinfo, self_.cast(), dns_resolved_thunk) };
}
pub fn r#loop(&self) -> *mut uws::Loop {
self.loop_ptr
}
pub unsafe fn on_dns_resolved(this: *mut PendingConnect) {
let this = unsafe { bun_core::heap::take(this) };
let session = this.session;
let s = session_mut(session);
if s.closed || s.pending.is_empty() {
this.pc_mut().cancel();
if !s.closed {
Self::fail_session(session, bun_core::err!("Aborted"));
}
return;
}
let Some(qs) = this.pc_mut().resolved() else {
Self::fail_session(session, bun_core::err!("DNSResolutionFailed"));
return;
};
s.qsocket = Some(NonNull::from(&mut *qs));
*qs.ext::<ClientSession>() = NonNull::new(session);
}
pub unsafe fn on_dns_resolved_threadsafe(this: *mut PendingConnect) {
let loop_ptr = bun_ptr::ParentRef::from(
NonNull::new(this).expect("on_dns_resolved_threadsafe: non-null"),
)
.r#loop();
RESOLVED.lock().push(Resolved(this));
unsafe { (*loop_ptr).wakeup() };
}
pub fn drain_resolved() {
let batch = core::mem::take(&mut *RESOLVED.lock());
for Resolved(head) in batch {
unsafe { PendingConnect::on_dns_resolved(head) };
}
}
pub fn fail_session(session: *mut ClientSession, err: bun_core::Error) {
let s = session_mut(session);
s.closed = true;
if let Some(ctx) = super::client_context::ClientContext::get() {
super::client_context::ClientContext::as_mut(ctx).unregister(s);
}
while !s.pending.is_empty() {
let stream = s.pending[0];
let cl = super::client_session::stream_ref(stream).client;
s.detach(stream);
if let Some(cl) = cl {
super::client_session::client_mut(cl).fail_from_h2(err);
}
}
let _ = super::LIVE_SESSIONS.fetch_sub(1, Ordering::Relaxed);
unsafe { ClientSession::deref(s) };
}
}
#[repr(transparent)]
struct Resolved(*mut PendingConnect);
unsafe impl Send for Resolved {}
static RESOLVED: Guarded<Vec<Resolved>> = Guarded::new(Vec::new());
unsafe extern "C" fn dns_resolved_thunk(pc: *mut core::ffi::c_void) {
unsafe { PendingConnect::on_dns_resolved_threadsafe(pc.cast()) };
}