code wiki / (root) / nx_net_chan_test.nx

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}