Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
140 changes: 73 additions & 67 deletions src/js/node/dgram.ts
Original file line number Diff line number Diff line change
Expand Up @@ -299,67 +299,82 @@ Socket.prototype.bind = function (port_, address_ /* , callback */) {
return;
}

let flags = uSockets.LISTEN_DISALLOW_REUSE_PORT_FAILURE;

if (state.reuseAddr) {
flags |= uSockets.LISTEN_REUSE_ADDR;
}

if (state.ipv6Only) {
flags |= uSockets.SOCKET_IPV6_ONLY;
}

if (state.reusePort) {
flags |= uSockets.LISTEN_REUSE_PORT;
}

// TODO flags
const family = this.type === "udp4" ? "IPv4" : "IPv6";
try {
Bun.udpSocket({
hostname: ip,
port: port || 0,
flags,
socket: {
data: (_socket, data, port, address) => {
this.emit("message", data, {
port: port,
address: address,
size: data.length,
// TODO check if this is correct
family,
});
},
error: error => {
this.emit("error", error);
},
},
}).$then(
socket => {
if (state.unrefOnBind) {
socket.unref();
state.unrefOnBind = false;
}
state.handle.socket = socket;
state.receiving = true;
state.bindState = BIND_STATE_BOUND;

this.emit("listening");
},
err => {
state.bindState = BIND_STATE_UNBOUND;
this.emit("error", err);
},
);
attachSocket(this, ip, port);
} catch (err) {
state.bindState = BIND_STATE_UNBOUND;
this.emit("error", err);
return;
}

this.emit("listening");
Comment thread
robobun marked this conversation as resolved.
});

return this;
};

const createSocketFn = $newRustFunction("udp_socket.rs", "UDPSocket.jsCreate", 1);

// Synchronously creates + binds the native socket for `self`, attaches it to
// the handle, and marks the socket BOUND. On failure, throws without touching
// `bindState`; callers own error reporting and the `listening` event.
function attachSocket(self, ip, port) {
const state = self[kStateSymbol];

let flags = uSockets.LISTEN_DISALLOW_REUSE_PORT_FAILURE;

if (state.reuseAddr) {
flags |= uSockets.LISTEN_REUSE_ADDR;
}

if (state.ipv6Only) {
flags |= uSockets.SOCKET_IPV6_ONLY;
}

if (state.reusePort) {
flags |= uSockets.LISTEN_REUSE_PORT;
}

const family = self.type === "udp4" ? "IPv4" : "IPv6";
const socket = createSocketFn({
hostname: ip,
port: port || 0,
flags,
socket: {
data: (_socket, data, port, address) => {
self.emit("message", data, {
port: port,
address: address,
size: data.length,
family,
});
},
error: error => {
self.emit("error", error);
},
},
});

if (state.unrefOnBind) {
socket.unref();
state.unrefOnBind = false;
}
state.handle.socket = socket;
state.receiving = true;
state.bindState = BIND_STATE_BOUND;
}

// node binds an unbound socket to a random port before a membership operation
// (libuv's `uv__udp_maybe_deferred_bind`), synchronously and without emitting
// "listening". A closed or already-binding socket is "not running".
function implicitBind(self) {
const state = self[kStateSymbol];
if (!state.handle || state.bindState !== BIND_STATE_UNBOUND) {
throw $ERR_SOCKET_DGRAM_NOT_RUNNING();
}
attachSocket(self, self.type === "udp4" ? "0.0.0.0" : "::", 0);
}

Socket.prototype.connect = function (port, address, callback) {
port = validatePort(port, "Port", false);
if (typeof address === "function") {
Expand Down Expand Up @@ -815,15 +830,12 @@ Socket.prototype.addMembership = function (multicastAddress, interfaceAddress) {
if (typeof interfaceAddress !== "undefined") {
validateString(interfaceAddress, "interfaceAddress");
}
const { handle, bindState } = this[kStateSymbol];
const { handle } = this[kStateSymbol];
if (!handle?.socket) {
if (!isIP(multicastAddress)) {
throw EINVAL("addMembership");
}
throw $ERR_SOCKET_DGRAM_NOT_RUNNING();
}
if (bindState === BIND_STATE_UNBOUND) {
this.bind({ port: 0, exclusive: true }, null);
implicitBind(this);
}
return handle.socket.addMembership(multicastAddress, interfaceAddress);
};
Expand All @@ -841,7 +853,7 @@ Socket.prototype.dropMembership = function (multicastAddress, interfaceAddress)
if (!isIP(multicastAddress)) {
throw EINVAL("dropMembership");
}
throw $ERR_SOCKET_DGRAM_NOT_RUNNING();
implicitBind(this);
}
return handle.socket.dropMembership(multicastAddress, interfaceAddress);
};
Expand All @@ -853,15 +865,12 @@ Socket.prototype.addSourceSpecificMembership = function (sourceAddress, groupAdd
validateString(interfaceAddress, "interfaceAddress");
}

const { handle, bindState } = this[kStateSymbol];
const { handle } = this[kStateSymbol];
if (!handle?.socket) {
if (!isIP(sourceAddress) || !isIP(groupAddress)) {
throw EINVAL("addSourceSpecificMembership");
}
throw $ERR_SOCKET_DGRAM_NOT_RUNNING();
}
if (bindState === BIND_STATE_UNBOUND) {
this.bind(0);
implicitBind(this);
}
return handle.socket.addSourceSpecificMembership(sourceAddress, groupAddress, interfaceAddress);
};
Expand All @@ -873,15 +882,12 @@ Socket.prototype.dropSourceSpecificMembership = function (sourceAddress, groupAd
validateString(interfaceAddress, "interfaceAddress");
}

const { handle, bindState } = this[kStateSymbol];
const { handle } = this[kStateSymbol];
if (!handle?.socket) {
if (!isIP(sourceAddress) || !isIP(groupAddress)) {
throw EINVAL("dropSourceSpecificMembership");
}
throw $ERR_SOCKET_DGRAM_NOT_RUNNING();
}
if (bindState === BIND_STATE_UNBOUND) {
this.bind(0);
implicitBind(this);
}
return handle.socket.dropSourceSpecificMembership(sourceAddress, groupAddress, interfaceAddress);
};
Expand Down
29 changes: 23 additions & 6 deletions src/js/node/dns.ts
Original file line number Diff line number Diff line change
Expand Up @@ -317,11 +317,15 @@ function lookup(hostname, options, callback) {
res.sort((a, b) => b.family - a.family);
}

if (options?.all) {
callback(null, res.map(mapLookupAll));
} else {
const [{ address, family }] = res;
callback(null, address, family);
try {
if (options?.all) {
callback(null, res.map(mapLookupAll));
} else {
const [{ address, family }] = res;
callback(null, address, family);
}
} catch (err) {
rethrowAsUncaughtException(err);
}
})
.catch(err => {
Expand All @@ -333,10 +337,23 @@ function lookup(hostname, options, callback) {
err.hostname = hostname;
err.message = `${syscall} ${err.code} ${hostname}`;
}
callback(err, undefined, undefined);
try {
callback(err, undefined, undefined);
} catch (err) {
rethrowAsUncaughtException(err);
}
});
}

// Re-raises a throw from a user-supplied lookup() callback as an
// uncaughtException, matching node. Without this the chained `.catch` above
// would receive it and invoke the callback a second time.
function rethrowAsUncaughtException(err) {
queueMicrotask(() => {
throw err;
});
Comment thread
claude[bot] marked this conversation as resolved.
}

function lookupService(address, port, callback) {
if (arguments.length < 3) {
throw $ERR_MISSING_ARGS("address", "port", "callback");
Expand Down
26 changes: 21 additions & 5 deletions src/runtime/socket/udp_socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -483,7 +483,7 @@ impl UDPSocket {
/// Recover `&UDPSocket` from the uws user-data slot. Centralises the
/// `unsafe { &*(*socket).user().cast() }` back-ref deref shared by every
/// `extern "C"` callback below — the user pointer was set to the
/// heap-allocated `UDPSocket` in [`udp_socket`] via
/// heap-allocated `UDPSocket` in [`Self::create`] via
/// `uws::udp::Socket::create(.., user_data = this_ptr)` and remains live
/// until `on_close` (uws guarantees no callback after close). All mutated
/// fields are `Cell`/`JsCell`, so a shared borrow is sufficient (R-2).
Expand All @@ -496,7 +496,26 @@ impl UDPSocket {
unsafe { &*user.cast::<UDPSocket>() }
}

/// `Bun.udpSocket(options)`. The bind happens synchronously inside
/// [`Self::create`]; the Promise is only the public API shape.
pub fn udp_socket(global_this: &JSGlobalObject, options: JSValue) -> JsResult<JSValue> {
let this_value = Self::create(global_this, options)?;
Ok(bun_jsc::JSPromise::resolved_promise_value(
global_this,
this_value,
))
}

// See `js_connect` — codegen `JsClass` derive owns the link name.
// `node:dgram` implicitly binds before a membership operation (libuv's
// deferred bind), so it needs the socket without a promise round-trip.
pub fn js_create(global_this: &JSGlobalObject, call_frame: &CallFrame) -> JsResult<JSValue> {
Self::create(global_this, call_frame.argument(0))
}

/// Creates, binds, and (optionally) connects the socket. Returns the JS
/// wrapper directly; every failure is thrown synchronously.
fn create(global_this: &JSGlobalObject, options: JSValue) -> JsResult<JSValue> {
bun_output::scoped_log!(UdpSocket, "udpSocket");

let vm = global_this.bun_vm_ptr();
Expand Down Expand Up @@ -643,10 +662,7 @@ impl UDPSocket {
scopeguard::ScopeGuard::into_inner(guard);

this.poll_ref.with_mut(|p| p.ref_(bun_io::js_vm_ctx()));
Ok(bun_jsc::JSPromise::resolved_promise_value(
global_this,
this_value,
))
Ok(this_value)
}

pub fn call_error_handler(&self, this_value_: JSValue, err: JSValue) {
Expand Down
Loading
Loading