code wiki / (root) / nx_net_chan.nx

nx_net_chan.nx source

↩ module page · 213 lines · 7586 B

1// nx_net_chan.nx -- typed i64 channel serialised over TCP. 2// 3// Same send/recv semantics as nx_chan (Vyukov MPMC), but the 4// "channel" is a connected TCP socket and each message is 8 bytes 5// big-endian on the wire. Composes with nx_thread_pool / nx_parallel 6// transparently -- a NetChan is a 64-bit handle (the socket fd) 7// that workers can send to or receive from. 8// 9// Wire format: 10// - Each message is a fixed 8-byte big-endian i64 frame. 11// - No length prefix, no checksum, no framing markers. TCP is 12// stream-oriented but with fixed-size frames the receiver can 13// always read exactly 8 bytes per message. 14// - For variable-length payloads (typed structs, byte arrays), 15// use nx_net_chan_blob (queued; length-prefix framed). 16// 17// Endianness: big-endian on the wire (network byte order convention) 18// even though both host and target are likely little-endian. Makes 19// cross-architecture wire compatibility trivial when the cluster 20// becomes heterogeneous. 21// 22// Composes against: [[nx_socket_tcp_loopback]] (transport), [[atomic_intrinsics_real_amo]] 23// (if multiple workers share a NetChan, the underlying socket is 24// the synchronisation point -- kernel handles the lock). 25 26// nx_safety_envelope: 27// intended_use: AUTO_APPLIED -- primitive-specific tuning queued 28// sil_target: SIL1 29// evidence: [bulk_applied_2026-05-16, see-file-comment-for-detail] 30// verdict: NOT_YET_EVALUATED 31 32import "nx_syscalls.nx" 33import "nx_atom.nx" 34import "nx_socket.nx" 35const NX_MAGIC_9223372036854775807: i64 = 9223372036854775807 36 37struct NxNetChan { 38 fd: i64, // connected TCP socket 39 closed: i64, // 1 after explicit close 40} 41 42const NX_NET_CHAN_BYTES: i64 = 16 43 44// Pack an i64 into 8 big-endian bytes at buf[0..7]. 45func _nx_net_pack_be64(buf: *u8, v: i64) -> i64 { 46 buf[0] = ((v >> 56) & 0xFF) as u8 47 buf[1] = ((v >> 48) & 0xFF) as u8 48 buf[2] = ((v >> 40) & 0xFF) as u8 49 buf[3] = ((v >> 32) & 0xFF) as u8 50 buf[4] = ((v >> 24) & 0xFF) as u8 51 buf[5] = ((v >> 16) & 0xFF) as u8 52 buf[6] = ((v >> 8) & 0xFF) as u8 53 buf[7] = (v & 0xFF) as u8 54 return 0 55} 56 57// Unpack 8 big-endian bytes at buf[0..7] into an i64. 58func _nx_net_unpack_be64(buf: *u8) -> i64 { 59 let b0: i64 = buf[0] 60 let b1: i64 = buf[1] 61 let b2: i64 = buf[2] 62 let b3: i64 = buf[3] 63 let b4: i64 = buf[4] 64 let b5: i64 = buf[5] 65 let b6: i64 = buf[6] 66 let b7: i64 = buf[7] 67 return (b0 << 56) | (b1 << 48) | (b2 << 40) | (b3 << 32) 68 | (b4 << 24) | (b5 << 16) | (b6 << 8) | b7 69} 70 71// Wrap an existing connected socket fd as a NetChan. Caller must 72// have already done socket+connect (or accept). The chan does NOT 73// take ownership of the fd until nx_net_chan_close. 74func nx_net_chan_from_fd(fd: i64) -> *NxNetChan { 75 let raw: *u8 = sys_mmap(NX_NET_CHAN_BYTES) 76 let c: *NxNetChan = raw as *NxNetChan 77 c.fd = fd 78 c.closed = 0 79 return c 80} 81 82// Blocking send of one i64 message. Returns 0 on success, -errno 83// on socket failure. Handles short writes by retrying. 84func nx_net_chan_send(c: *NxNetChan, v: i64) -> i64 { 85 if c.closed != 0 { return -1 } 86 let buf: *u8 = sys_mmap(16) 87 _nx_net_pack_be64(buf, v) 88 var sent: i64 = 0 89 while sent < 8 { 90 let remaining: i64 = 8 - sent 91 let p: *u8 = ((buf as i64) + sent) as *u8 92 let n: i64 = __syscall(NX_SYS_SENDTO, c.fd, p as i64, remaining, 0, 0, 0) 93 if n <= 0 { return n } 94 sent = sent + n 95 } 96 return 0 97} 98 99// Blocking receive of one i64 message. Returns 0 + writes value 100// to *out on success; returns -1 on socket close / failure. 101func nx_net_chan_recv(c: *NxNetChan, out: *i64) -> i64 { 102 if c.closed != 0 { return -1 } 103 let buf: *u8 = sys_mmap(16) 104 var got: i64 = 0 105 while got < 8 { 106 let remaining: i64 = 8 - got 107 let p: *u8 = ((buf as i64) + got) as *u8 108 let n: i64 = __syscall(NX_SYS_RECVFROM, c.fd, p as i64, remaining, 0, 0, 0) 109 if n <= 0 { return -1 } // 0 = peer closed; <0 = error 110 got = got + n 111 } 112 *out = _nx_net_unpack_be64(buf) 113 return 0 114} 115 116// ---- variable-length blob frames ------------------------------- 117// Wire format: 8-byte big-endian length followed by exactly `len` 118// payload bytes. No checksum or framing markers; relies on TCP 119// in-order delivery. Useful for typed-struct serialisation and 120// RPC payloads where the message size isn't fixed. 121 122const NX_NET_BLOB_MAX_LEN: i64 = 1073741824 // 1 GiB safety ceiling 123 124func _nx_net_write_n(fd: i64, buf: *u8, n: i64) -> i64 { 125 var sent: i64 = 0 126 while sent < n { 127 let p: *u8 = ((buf as i64) + sent) as *u8 128 let r: i64 = __syscall(NX_SYS_SENDTO, fd, p as i64, n - sent, 0, 0, 0) 129 if r <= 0 { return -1 } 130 sent = sent + r 131 } 132 return 0 133} 134 135func _nx_net_read_n(fd: i64, buf: *u8, n: i64) -> i64 { 136 var got: i64 = 0 137 while got < n { 138 let p: *u8 = ((buf as i64) + got) as *u8 139 let r: i64 = __syscall(NX_SYS_RECVFROM, fd, p as i64, n - got, 0, 0, 0) 140 if r <= 0 { return -1 } // 0 = peer closed; <0 = error 141 got = got + r 142 } 143 return 0 144} 145 146// Send a blob: writes 8-byte BE length header then `n` payload bytes. 147// Returns 0 on success, -1 on socket failure or oversize. 148func nx_net_chan_send_blob(c: *NxNetChan, buf: *u8, n: i64) -> i64 { 149 if c.closed != 0 { return -1 } 150 if n < 0 { return -1 } 151 if n > NX_NET_BLOB_MAX_LEN { return -1 } 152 let hdr: *u8 = sys_mmap(16) 153 _nx_net_pack_be64(hdr, n) 154 if _nx_net_write_n(c.fd, hdr, 8) != 0 { return -1 } 155 if n > 0 { 156 if _nx_net_write_n(c.fd, buf, n) != 0 { return -1 } 157 } 158 return 0 159} 160 161// Receive a blob: reads 8-byte length header then exactly len payload 162// bytes into `buf`. `max_len` is the caller's buffer cap; if the 163// incoming message exceeds it, returns -2 (caller can grow + retry). 164// Returns the actual payload length (>= 0) on success. 165func nx_net_chan_recv_blob(c: *NxNetChan, buf: *u8, max_len: i64) -> i64 { 166 if c.closed != 0 { return -1 } 167 let hdr: *u8 = sys_mmap(16) 168 if _nx_net_read_n(c.fd, hdr, 8) != 0 { return -1 } 169 let len: i64 = _nx_net_unpack_be64(hdr) 170 if len < 0 { return -1 } 171 if len > NX_NET_BLOB_MAX_LEN { return -1 } 172 if len > max_len { return -2 } 173 if len > 0 { 174 if _nx_net_read_n(c.fd, buf, len) != 0 { return -1 } 175 } 176 return len 177} 178 179// Close the underlying socket cleanly (RDWR shutdown + close fd). 180func nx_net_chan_close(c: *NxNetChan) -> i64 { 181 if c.closed != 0 { return 0 } 182 c.closed = 1 183 nx_sock_shutdown(c.fd, NX_SHUT_RDWR) 184 return sys_close(c.fd) 185} 186 187// ---- self-test (encoding round-trip; networking proved by smoke) - 188 189func main() -> i64 { 190 let buf: *u8 = sys_mmap(16) 191 // Round-trip 5 sentinel values. 192 let vals: *u8 = sys_mmap(64) 193 let v: *i64 = vals as *i64 194 v[0] = 0 195 v[1] = 1 196 v[2] = -1 197 v[3] = 0x0123456789ABCDEF 198 v[4] = 0 - NX_MAGIC_9223372036854775807 // close to i64::MIN 199 var i: i64 = 0 200 while i < 5 { 201 _nx_net_pack_be64(buf, v[i]) 202 let r: i64 = _nx_net_unpack_be64(buf) 203 if r != v[i] { return __syscall(93, 10 + i, 0, 0, 0, 0, 0) } 204 i = i + 1 205 } 206 // Verify endianness explicitly: pack 0x0102030405060708 should 207 // give bytes 01 02 03 04 05 06 07 08 (big-endian). 208 _nx_net_pack_be64(buf, 0x0102030405060708) 209 if buf[0] != 0x01 { return __syscall(93, 1, 0, 0, 0, 0, 0) } 210 if buf[1] != 0x02 { return __syscall(93, 2, 0, 0, 0, 0, 0) } 211 if buf[7] != 0x08 { return __syscall(93, 8, 0, 0, 0, 0, 0) } 212 return 0 213}