nx_net_chan_test.nx source
↩ module page · 140 lines · 4834 B
1// nx_net_chan_test.nx -- prove typed i64 channel works over TCP
2// loopback with multi-message exchange.
3//
4// Server thread: accept connection, recv N messages, verify they
5// equal 0..N-1, send back doubled value. Client thread: connect,
6// send 0..N-1, recv N doubled values, verify each == 2*i.
7//
8// This is the distributed-MIMD building block: nx_net_chan_send /
9// recv is the same shape as nx_chan_send / recv, so existing
10// pipeline / pool / DAG code can use net_chan wherever they use
11// chan, transparently distributing across hosts when run on a
12// cluster.
13
14import "nx_kernel_v2.nx"
15import "nx_log.nx"
16import "nx_atom.nx"
17import "nx_thread.nx"
18import "nx_socket.nx"
19import "nx_net_chan.nx"
20
21const NX_PORT: i64 = 38919
22const N_MSGS: i64 = 100
23
24struct NetTestState {
25 done_flag: i64, // atomic; client sets 1 on success, -1 on fail
26 diag: i64, // diagnostic code when failure
27}
28
29func _sleep_ms(ms: i64) -> i64 {
30 let ts_raw: *u8 = sys_mmap(16)
31 let ts: *i64 = ts_raw as *i64
32 ts[0] = ms / 1000
33 ts[1] = (ms % 1000) * 1000000
34 return __syscall(SYS_CLOCK_NANOSLEEP, 1, 0, ts as i64, 0, 0, 0)
35}
36
37// Client: connect, send 0..N-1, recv N doubled values.
38func client_main(arg: *u8) -> i64 {
39 let s: *NetTestState = arg as *NetTestState
40 let done_addr: *i64 = (arg as *i64)
41
42 _sleep_ms(150) // let server reach accept()
43
44 let fd: i64 = nx_sock_socket(NX_AF_INET, NX_SOCK_STREAM, NX_IPPROTO_TCP)
45 if fd < 0 { s.diag = 10; nx_atom_store_i64(done_addr, -1, NX_MO_SEQ_CST); return 0 }
46 let addr: *u8 = sys_mmap(16)
47 nx_sock_sin_init(addr, 0x0100007F, NX_PORT) // 127.0.0.1
48 let cr: i64 = nx_sock_connect(fd, addr, 16)
49 if cr < 0 { s.diag = 20; nx_atom_store_i64(done_addr, -1, NX_MO_SEQ_CST); return 0 }
50
51 let nc: *NxNetChan = nx_net_chan_from_fd(fd)
52
53 var i: i64 = 0
54 while i < N_MSGS {
55 if nx_net_chan_send(nc, i) != 0 {
56 s.diag = 30; nx_atom_store_i64(done_addr, -1, NX_MO_SEQ_CST); return 0
57 }
58 i = i + 1
59 }
60
61 let out_raw: *u8 = sys_mmap(16)
62 let out: *i64 = out_raw as *i64
63 var j: i64 = 0
64 while j < N_MSGS {
65 if nx_net_chan_recv(nc, out) != 0 {
66 s.diag = 40; nx_atom_store_i64(done_addr, -1, NX_MO_SEQ_CST); return 0
67 }
68 if *out != j * 2 {
69 s.diag = 50; nx_atom_store_i64(done_addr, -1, NX_MO_SEQ_CST); return 0
70 }
71 j = j + 1
72 }
73 nx_net_chan_close(nc)
74 s.diag = 1
75 nx_atom_store_i64(done_addr, 1, NX_MO_SEQ_CST)
76 return 0
77}
78
79func main() -> nx_exit {
80 println("=== nx_net_chan typed-i64 over TCP smoke ===" as *u8)
81 println("port:" as *u8); print_i64(NX_PORT); println("" as *u8)
82 println("msgs:" as *u8); print_i64(N_MSGS); println("" as *u8)
83
84 let state_raw: *u8 = sys_mmap(64)
85 let s: *NetTestState = state_raw as *NetTestState
86 s.done_flag = 0
87 s.diag = 0
88
89 // Server listener BEFORE spawning client.
90 let lfd: i64 = nx_sock_socket(NX_AF_INET, NX_SOCK_STREAM, NX_IPPROTO_TCP)
91 if lfd < 0 { println("FAIL: socket" as *u8); return 1 }
92 nx_sock_reuseaddr(lfd)
93 let laddr: *u8 = sys_mmap(16)
94 nx_sock_sin_init(laddr, 0, NX_PORT)
95 if nx_sock_bind(lfd, laddr, 16) < 0 {
96 println("FAIL: bind" as *u8); return 2
97 }
98 if nx_sock_listen(lfd, 1) < 0 { println("FAIL: listen" as *u8); return 3 }
99
100 let tid: i64 = nx_thread_spawn_fn(client_main, state_raw, 65536)
101 if tid <= 0 { println("FAIL: spawn" as *u8); return 4 }
102
103 let cfd: i64 = nx_sock_accept(lfd, 0 as *u8, 0 as *i64)
104 if cfd < 0 { println("FAIL: accept" as *u8); return 5 }
105 let server_nc: *NxNetChan = nx_net_chan_from_fd(cfd)
106
107 let in_raw: *u8 = sys_mmap(16)
108 let in_v: *i64 = in_raw as *i64
109 var k: i64 = 0
110 while k < N_MSGS {
111 if nx_net_chan_recv(server_nc, in_v) != 0 {
112 println("FAIL: server recv" as *u8); return 6
113 }
114 if *in_v != k {
115 println("FAIL: server got unexpected" as *u8); return 7
116 }
117 if nx_net_chan_send(server_nc, *in_v * 2) != 0 {
118 println("FAIL: server send" as *u8); return 8
119 }
120 k = k + 1
121 }
122 nx_net_chan_close(server_nc)
123 sys_close(lfd)
124
125 // Wait for client verdict.
126 let done_addr: *i64 = (state_raw as *i64)
127 var spins: i64 = 0
128 while nx_atom_load_i64(done_addr, NX_MO_SEQ_CST) == 0 {
129 nx_thread_yield()
130 spins = spins + 1
131 if spins > 100000000 { println("FAIL: client timeout" as *u8); return 9 }
132 }
133 if nx_atom_load_i64(done_addr, NX_MO_SEQ_CST) != 1 {
134 println("FAIL: client diag=" as *u8); print_i64(s.diag); println("" as *u8)
135 return 10
136 }
137 println("PASS: 100 i64 messages round-tripped over TCP loopback." as *u8)
138 println("nx_net_chan API identical to nx_chan -- transparent distribution." as *u8)
139 return 0
140}