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}