Skip to content

Commit 67f36f0

Browse files
committed
Add -sNODENET backend for real outgoing TCP via node:net
Adds a new NODENET setting that backs the POSIX sockets API directly with Node.js's node:net module, giving real, non-blocking outgoing (client) TCP sockets without WebSockets, an external proxy process, or pthreads. Unlike PROXY_POSIX_SOCKETS this is single-threaded and event-driven: socket readiness is delivered through the same emscripten_set_socket_*_callback hooks the default WebSocket backend uses, so it drops into existing readiness reactors unchanged. This initial backend supports outgoing TCP only: connect, send, recv and close, plus get/setsockopt (SO_ERROR, TCP_NODELAY, SO_KEEPALIVE and the TCP keep-alive tunables). There is no bind/listen/accept (server) support and no UDP yet; those land in follow-ups. - new nodenet_sock_ops in libsockfs.js implementing the sock_ops contract over net.createConnection - route get/setsockopt through the backend under -sNODENET, and compile out the weak __syscall_setsockopt stub via a libstubs variation so the JS symbol wins - test/sockets/test_nodenet.c: outgoing connect/send/recv against a loopback echo server started by the test harness
1 parent e11edd8 commit 67f36f0

7 files changed

Lines changed: 445 additions & 2 deletions

File tree

src/lib/libsockfs.js

Lines changed: 244 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,11 @@ addToLibrary({
88
$SOCKFS__postset: () => {
99
addAtInit('SOCKFS.root = FS.mount(SOCKFS, {}, null);');
1010
},
11-
$SOCKFS__deps: ['$FS'],
11+
$SOCKFS__deps: ['$FS',
12+
#if NODENET
13+
'$ERRNO_CODES',
14+
#endif
15+
],
1216
$SOCKFS: {
1317
#if expectToReceiveOnModule('websocket')
1418
websocketArgs: {},
@@ -69,6 +73,8 @@ addToLibrary({
6973
pending: [],
7074
recv_queue: [],
7175
#if SOCKET_WEBRTC
76+
#elif NODENET
77+
sock_ops: SOCKFS.nodenet_sock_ops
7278
#else
7379
sock_ops: SOCKFS.websocket_sock_ops
7480
#endif
@@ -726,7 +732,243 @@ addToLibrary({
726732

727733
return res;
728734
}
729-
}
735+
},
736+
#if NODENET
737+
// Outgoing TCP over node:net, using the same sock_ops contract and
738+
// SOCKFS.emit readiness callbacks as the WebSocket backend. Client TCP only
739+
// for now, with no bind, listen, accept or UDP.
740+
nodenet_sock_ops: {
741+
errnoForNode(e) {
742+
return (e && e.code && ERRNO_CODES[e.code]) || {{{ cDefs.ECONNREFUSED }}};
743+
},
744+
// Replay buffered opts once the socket is live.
745+
applyOptions(sock) {
746+
var conn = sock.connection;
747+
var o = sock.opts;
748+
if (!conn || !o) return;
749+
if (o.noDelay !== undefined) {
750+
try { conn.setNoDelay(!!o.noDelay); } catch (e) {}
751+
}
752+
SOCKFS.nodenet_sock_ops.applyKeepAlive(sock);
753+
},
754+
// The keepalive tunables arrive from C in seconds, but node wants
755+
// milliseconds, so we scale by 1000. A non-positive value keeps node's
756+
// default for that field.
757+
applyKeepAlive(sock) {
758+
var conn = sock.connection;
759+
var o = sock.opts;
760+
if (!conn || !o || o.keepAlive === undefined) return;
761+
try {
762+
conn.setKeepAlive(
763+
!!o.keepAlive,
764+
(o.keepAliveIdle || 0) * 1000,
765+
(o.keepAliveIntvl || 0) * 1000,
766+
o.keepAliveCnt || 0);
767+
} catch (e) {}
768+
},
769+
// Forward a connected node socket's events onto sock.
770+
wireConnection(sock, conn) {
771+
sock.connection = conn;
772+
conn.on('data', (buf) => {
773+
var data = new Uint8Array(buf.length);
774+
data.set(buf);
775+
sock.recv_queue.push({ addr: sock.daddr, port: sock.dport, data });
776+
sock.recv_bytes = (sock.recv_bytes || 0) + data.length;
777+
// If the peer outruns the reader, pause node and resume in recvmsg.
778+
if (sock.recv_bytes >= 262144 /* 256 KiB */) {
779+
try { conn.pause(); } catch (e) {}
780+
sock.paused = true;
781+
}
782+
SOCKFS.emit('message', sock.stream.fd);
783+
});
784+
// A peer FIN surfaces as EOF to the reader.
785+
conn.on('end', () => {
786+
sock.readClosed = true;
787+
SOCKFS.emit('message', sock.stream.fd);
788+
});
789+
conn.on('close', () => {
790+
sock.readClosed = true;
791+
sock.state = 'closed';
792+
SOCKFS.emit('close', sock.stream.fd);
793+
});
794+
// Backpressure relieved, so we are writable again.
795+
conn.on('drain', () => {
796+
sock.writeBlocked = false;
797+
SOCKFS.emit('open', sock.stream.fd);
798+
});
799+
conn.on('error', (e) => {
800+
sock.error = SOCKFS.nodenet_sock_ops.errnoForNode(e);
801+
// Let a failed connect resolve so SO_ERROR can be read.
802+
if (sock.state === 'connecting') sock.state = 'connected';
803+
SOCKFS.emit('error', [sock.stream.fd, sock.error, (e && e.message) || 'socket error']);
804+
});
805+
},
806+
poll(sock) {
807+
var mask = 0;
808+
if (sock.recv_queue.length || sock.readClosed || sock.error) {
809+
mask |= ({{{ cDefs.POLLRDNORM }}} | {{{ cDefs.POLLIN }}});
810+
}
811+
if (sock.error) {
812+
// Mark writable on error so SO_ERROR can be read.
813+
mask |= {{{ cDefs.POLLOUT }}};
814+
} else if (sock.connection && sock.state === 'connected' && !sock.writeBlocked) {
815+
mask |= {{{ cDefs.POLLOUT }}};
816+
}
817+
if (sock.readClosed) mask |= {{{ cDefs.POLLHUP }}};
818+
return mask;
819+
},
820+
ioctl(sock, request, arg) {
821+
switch (request) {
822+
case {{{ cDefs.FIONREAD }}}:
823+
var bytes = sock.recv_queue.length ? sock.recv_queue[0].data.length : 0;
824+
{{{ makeSetValue('arg', '0', 'bytes', 'i32') }}};
825+
return 0;
826+
case {{{ cDefs.FIONBIO }}}:
827+
var on = {{{ makeGetValue('arg', '0', 'i32') }}};
828+
if (on) sock.stream.flags |= {{{ cDefs.O_NONBLOCK }}};
829+
else sock.stream.flags &= ~{{{ cDefs.O_NONBLOCK }}};
830+
return 0;
831+
default:
832+
return {{{ cDefs.EINVAL }}};
833+
}
834+
},
835+
close(sock) {
836+
sock.state = 'closed';
837+
if (sock.connection) { try { sock.connection.destroy(); } catch (e) {} sock.connection = null; }
838+
return 0;
839+
},
840+
connect(sock, addr, port) {
841+
if (sock.connection) {
842+
throw new FS.ErrnoError(sock.state === 'connecting' ? {{{ cDefs.EALREADY }}} : {{{ cDefs.EISCONN }}});
843+
}
844+
sock.daddr = addr;
845+
sock.dport = port;
846+
sock.state = 'connecting';
847+
var conn = require('net').createConnection({ host: addr, port, allowHalfOpen: true });
848+
conn.once('connect', () => {
849+
sock.state = 'connected';
850+
sock.saddr = conn.localAddress;
851+
sock.sport = conn.localPort;
852+
sock.daddr = conn.remoteAddress || addr;
853+
sock.dport = conn.remotePort || port;
854+
SOCKFS.nodenet_sock_ops.applyOptions(sock);
855+
SOCKFS.emit('open', sock.stream.fd);
856+
});
857+
SOCKFS.nodenet_sock_ops.wireConnection(sock, conn);
858+
},
859+
sendmsg(sock, buffer, offset, length, addr, port) {
860+
if (!sock.connection || sock.state === 'closed') {
861+
throw new FS.ErrnoError({{{ cDefs.ENOTCONN }}});
862+
}
863+
if (ArrayBuffer.isView(buffer)) { offset += buffer.byteOffset; buffer = buffer.buffer; }
864+
var data = new Uint8Array(buffer.slice(offset, offset + length));
865+
var ok;
866+
try {
867+
ok = sock.connection.write(data);
868+
} catch (e) {
869+
throw new FS.ErrnoError(SOCKFS.nodenet_sock_ops.errnoForNode(e));
870+
}
871+
if (!ok) sock.writeBlocked = true; // cleared on 'drain'
872+
return length;
873+
},
874+
recvmsg(sock, length) {
875+
var queued = sock.recv_queue.shift();
876+
if (!queued) {
877+
if (sock.readClosed) return null; // EOF
878+
if (!sock.connection) {
879+
throw new FS.ErrnoError({{{ cDefs.ENOTCONN }}});
880+
}
881+
throw new FS.ErrnoError({{{ cDefs.EAGAIN }}});
882+
}
883+
var q = queued.data;
884+
var bytesRead = Math.min(length, q.length);
885+
var res = { buffer: q.subarray(0, bytesRead), addr: queued.addr, port: queued.port };
886+
if (bytesRead < q.length) {
887+
queued.data = q.subarray(bytesRead);
888+
sock.recv_queue.unshift(queued);
889+
}
890+
sock.recv_bytes = Math.max(0, (sock.recv_bytes || 0) - bytesRead);
891+
if (sock.paused && sock.recv_bytes < 262144 && sock.connection) {
892+
sock.paused = false;
893+
try { sock.connection.resume(); } catch (e) {}
894+
}
895+
return res;
896+
},
897+
setsockopt(sock, level, optname, optval, optlen) {
898+
sock.opts ||= {};
899+
var val = {{{ makeGetValue('optval', 0, 'i32') }}};
900+
if (level === {{{ cDefs.SOL_SOCKET }}}) {
901+
switch (optname) {
902+
case 9: // SO_KEEPALIVE
903+
sock.opts.keepAlive = !!val;
904+
SOCKFS.nodenet_sock_ops.applyKeepAlive(sock);
905+
return 0;
906+
case 8: // SO_RCVBUF. Node TCP cannot set this, so we just store it.
907+
sock.opts.recvBuf = val;
908+
return 0;
909+
case 7: // SO_SNDBUF. Node TCP cannot set this, so we just store it.
910+
sock.opts.sendBuf = val;
911+
return 0;
912+
case 2: // SO_REUSEADDR. Node manages reuse, so we just store it.
913+
sock.opts.reuseAddr = !!val;
914+
return 0;
915+
}
916+
} else if (level === {{{ cDefs.IPPROTO_TCP }}}) {
917+
switch (optname) {
918+
case 1: // TCP_NODELAY
919+
sock.opts.noDelay = !!val;
920+
if (sock.connection) { try { sock.connection.setNoDelay(!!val); } catch (e) {} }
921+
return 0;
922+
case 4: // TCP_KEEPIDLE (seconds)
923+
sock.opts.keepAliveIdle = val;
924+
SOCKFS.nodenet_sock_ops.applyKeepAlive(sock);
925+
return 0;
926+
case 5: // TCP_KEEPINTVL (seconds)
927+
sock.opts.keepAliveIntvl = val;
928+
SOCKFS.nodenet_sock_ops.applyKeepAlive(sock);
929+
return 0;
930+
case 6: // TCP_KEEPCNT (probe count)
931+
sock.opts.keepAliveCnt = val;
932+
SOCKFS.nodenet_sock_ops.applyKeepAlive(sock);
933+
return 0;
934+
}
935+
}
936+
// Accept unknown options silently, like a permissive stack.
937+
return 0;
938+
},
939+
getsockopt(sock, level, optname, optval, optlen) {
940+
sock.opts ||= {};
941+
var val;
942+
if (level === {{{ cDefs.SOL_SOCKET }}}) {
943+
switch (optname) {
944+
case {{{ cDefs.SO_ERROR }}}:
945+
{{{ makeSetValue('optval', 0, 'sock.error || 0', 'i32') }}};
946+
{{{ makeSetValue('optlen', 0, 4, 'i32') }}};
947+
sock.error = null; // SO_ERROR reads and clears
948+
return 0;
949+
case 9: val = sock.opts.keepAlive ? 1 : 0; break; // SO_KEEPALIVE
950+
case 8: val = sock.opts.recvBuf || 65536; break; // SO_RCVBUF
951+
case 7: val = sock.opts.sendBuf || 65536; break; // SO_SNDBUF
952+
case 2: val = sock.opts.reuseAddr ? 1 : 0; break; // SO_REUSEADDR
953+
default: return -{{{ cDefs.ENOPROTOOPT }}};
954+
}
955+
} else if (level === {{{ cDefs.IPPROTO_TCP }}}) {
956+
switch (optname) {
957+
case 1: val = sock.opts.noDelay ? 1 : 0; break; // TCP_NODELAY
958+
case 4: val = sock.opts.keepAliveIdle || 0; break; // TCP_KEEPIDLE
959+
case 5: val = sock.opts.keepAliveIntvl || 0; break;// TCP_KEEPINTVL
960+
case 6: val = sock.opts.keepAliveCnt || 0; break; // TCP_KEEPCNT
961+
default: return -{{{ cDefs.ENOPROTOOPT }}};
962+
}
963+
} else {
964+
return -{{{ cDefs.ENOPROTOOPT }}};
965+
}
966+
{{{ makeSetValue('optval', 0, 'val', 'i32') }}};
967+
{{{ makeSetValue('optlen', 0, 4, 'i32') }}};
968+
return 0;
969+
}
970+
},
971+
#endif
730972
},
731973

732974
/*

src/lib/libsyscall.js

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -445,6 +445,10 @@ var SyscallsLibrary = {
445445
__syscall_getsockopt__deps: ['$getSocketFromFD'],
446446
__syscall_getsockopt: (fd, level, optname, optval, optlen, d1) => {
447447
var sock = getSocketFromFD(fd);
448+
#if NODENET
449+
// The node:net backend handles all socket options.
450+
return sock.sock_ops.getsockopt(sock, level, optname, optval, optlen);
451+
#else
448452
// Minimal getsockopt aimed at resolving https://github.com/emscripten-core/emscripten/issues/2211
449453
// so only supports SOL_SOCKET with SO_ERROR.
450454
if (level === {{{ cDefs.SOL_SOCKET }}}) {
@@ -456,7 +460,17 @@ var SyscallsLibrary = {
456460
}
457461
}
458462
return -{{{ cDefs.ENOPROTOOPT }}}; // The option is unknown at the level indicated.
463+
#endif
459464
},
465+
#if NODENET
466+
// Only the node:net backend implements setsockopt. By default it stays
467+
// unimplemented.
468+
__syscall_setsockopt__deps: ['$getSocketFromFD'],
469+
__syscall_setsockopt: (fd, level, optname, optval, optlen, d1) => {
470+
var sock = getSocketFromFD(fd);
471+
return sock.sock_ops.setsockopt(sock, level, optname, optval, optlen);
472+
},
473+
#endif
460474
__syscall_sendmsg__deps: ['$getSocketFromFD', '$getSocketAddress'],
461475
__syscall_sendmsg: (fd, message, flags, d1, d2, d3) => {
462476
var sock = getSocketFromFD(fd);

src/settings.js

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -419,6 +419,20 @@ var WEBSOCKET_URL = 'ws://';
419419
// [link]
420420
var PROXY_POSIX_SOCKETS = false;
421421

422+
// If 1, the POSIX sockets API is backed by Node.js's ``node:net`` module,
423+
// giving real non-blocking outgoing TCP sockets with no WebSockets, proxy
424+
// process or pthreads. This only works under node and is ignored elsewhere.
425+
//
426+
// For now it supports outgoing TCP only. There is no bind, listen, accept or
427+
// UDP yet.
428+
//
429+
// It is single-threaded and event-driven. Socket readiness comes through the
430+
// same ``emscripten_set_socket_*_callback`` hooks the WebSocket backend uses,
431+
// so it works with existing readiness reactors. It cannot be combined with the
432+
// WebSocket emulation, PROXY_POSIX_SOCKETS or SOCKET_WEBRTC.
433+
// [link]
434+
var NODENET = false;
435+
422436
// A string containing a comma separated list of WebSocket subprotocols
423437
// as would be present in the Sec-WebSocket-Protocol header.
424438
// You can set 'null', if you don't want to specify it.

system/lib/libc/emscripten_syscall_stubs.c

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -248,10 +248,14 @@ weak int __syscall_prlimit64(int pid, int resource, intptr_t new_limit, intptr_t
248248
return 0;
249249
}
250250

251+
// Under NODENET the JS library provides setsockopt, so drop this weak stub to
252+
// leave the symbol undefined for it to fill.
253+
#ifndef __EMSCRIPTEN_NODENET__
251254
weak int __syscall_setsockopt(int sockfd, int level, int optname, intptr_t optval, size_t optlen, int dummy) {
252255
REPORT(setsockopt);
253256
return -ENOPROTOOPT; // The option is unknown at the level indicated.
254257
}
258+
#endif
255259

256260
UNIMPLEMENTED(acct, (intptr_t filename))
257261
UNIMPLEMENTED(mincore, (intptr_t addr, size_t length, intptr_t vec))

0 commit comments

Comments
 (0)