code wiki / (root) / nx_rpc.nx

nx_rpc.nx source

↩ module page · 150 lines · 5784 B

1// nx_rpc.nx -- typed RPC over nx_net_chan_blob. 2// 3// Builds on nx_net_chan_blob's variable-length frames to give a 4// request/response pattern: caller sends (method_id, request_blob), 5// server dispatches to a registered handler indexed by method_id, 6// returns (response_blob). Wire format per exchange: 7// 8// client -> server: 9// 8 bytes: method_id (i64 big-endian via nx_net_chan_send) 10// N bytes: request blob (length-prefix framed via 11// nx_net_chan_send_blob) 12// server -> client: 13// N bytes: response blob (length-prefix framed) 14// 15// Sentinel method_id NX_RPC_SHUTDOWN signals server to exit cleanly. 16// 17// Handler table: a fixed-size array indexed by method_id. Each 18// entry is a typed fn-ptr: func(*u8, n_req, *u8, max_resp) -> n_resp. 19// Unregistered methods get a -1 response. 20// 21// Composes against: [[nx_net_chan_tcp_i64]] (method_id transport), 22// [[nx_net_chan_blob_framing]] (request/response payload), 23// [[fn_ptr_indirect_call]] (typed dispatch). 24 25// nx_safety_envelope: 26// intended_use: AUTO_APPLIED -- primitive-specific tuning queued 27// sil_target: SIL1 28// evidence: [bulk_applied_2026-05-16, see-file-comment-for-detail] 29// verdict: NOT_YET_EVALUATED 30 31import "nx_syscalls.nx" 32import "nx_atom.nx" 33import "nx_net_chan.nx" 34 35const NX_RPC_SHUTDOWN: i64 = -1 36const NX_RPC_MAX_METHODS: i64 = 64 37 38struct NxRpcHandlerEntry { 39 fn_ptr: func(*u8, i64, *u8, i64) -> i64, 40 used: i64, 41} 42 43const NX_RPC_ENTRY_BYTES: i64 = 16 44 45struct NxRpcServer { 46 nc_addr: i64, // *NxNetChan 47 handlers_base: i64, // *NxRpcHandlerEntry[NX_RPC_MAX_METHODS] 48 n_calls: i64, // atomic; bumped per dispatch 49 shutdown_flag: i64, // atomic; 1 when shutdown signalled 50} 51 52func _nx_rpc_handler_at(s: *NxRpcServer, idx: i64) -> *NxRpcHandlerEntry { 53 return ((s.handlers_base) + idx * NX_RPC_ENTRY_BYTES) as *NxRpcHandlerEntry 54} 55 56// Create a server bound to an existing NetChan. Caller is responsible 57// for the listening socket + accept; pass the resulting connected 58// NetChan in. 59func nx_rpc_server_new(nc: *NxNetChan) -> *NxRpcServer { 60 let raw: *u8 = sys_mmap(128) 61 let s: *NxRpcServer = raw as *NxRpcServer 62 s.nc_addr = nc as i64 63 let table: *u8 = sys_mmap(NX_RPC_MAX_METHODS * NX_RPC_ENTRY_BYTES) 64 s.handlers_base = table as i64 65 s.n_calls = 0 66 s.shutdown_flag = 0 67 return s 68} 69 70// Register a handler for a method_id. Returns 0 on success, -1 on 71// out-of-range id or slot collision. 72func nx_rpc_register(s: *NxRpcServer, method_id: i64, 73 handler: func(*u8, i64, *u8, i64) -> i64) -> i64 { 74 if method_id < 0 { return -1 } 75 if method_id >= NX_RPC_MAX_METHODS { return -1 } 76 let entry: *NxRpcHandlerEntry = _nx_rpc_handler_at(s, method_id) 77 if entry.used == 1 { return -1 } 78 entry.fn_ptr = handler 79 entry.used = 1 80 return 0 81} 82 83// Server main loop: recv method_id, recv request blob, dispatch 84// to handler, send response blob. Returns when shutdown sentinel 85// received or NetChan closes. 86func nx_rpc_server_run(s: *NxRpcServer, scratch_buf: *u8, scratch_cap: i64) -> i64 { 87 let nc: *NxNetChan = s.nc_addr as *NxNetChan 88 let method_id_raw: *u8 = sys_mmap(16) 89 let method_id_p: *i64 = method_id_raw as *i64 90 let resp_buf: *u8 = sys_mmap(scratch_cap) 91 92 var running: i64 = 1 93 while running == 1 { 94 // 1. recv method_id 95 if nx_net_chan_recv(nc, method_id_p) != 0 { running = 0 } 96 else { 97 let id: i64 = *method_id_p 98 if id == NX_RPC_SHUTDOWN { 99 let sf: *i64 = ((s as i64) + 24) as *i64 100 nx_atom_store_i64(sf, 1, NX_MO_SEQ_CST) 101 running = 0 102 } else { 103 // 2. recv request blob 104 let req_n: i64 = nx_net_chan_recv_blob(nc, scratch_buf, scratch_cap) 105 if req_n < 0 { running = 0 } 106 else { 107 // 3. dispatch 108 var resp_n: i64 = -1 109 if id >= 0 { 110 if id < NX_RPC_MAX_METHODS { 111 let entry: *NxRpcHandlerEntry = _nx_rpc_handler_at(s, id) 112 if entry.used == 1 { 113 let fp: func(*u8, i64, *u8, i64) -> i64 = entry.fn_ptr 114 resp_n = fp(scratch_buf, req_n, resp_buf, scratch_cap) 115 } 116 } 117 } 118 // 4. send response (or empty blob if no handler) 119 if resp_n < 0 { resp_n = 0 } 120 if nx_net_chan_send_blob(nc, resp_buf, resp_n) != 0 { running = 0 } 121 let nc_addr: *i64 = ((s as i64) + 16) as *i64 122 nx_atom_faa_i64(nc_addr, 1, NX_MO_SEQ_CST) 123 } 124 } 125 } 126 } 127 return 0 128} 129 130// Client-side: send (method_id, request) and receive response. 131// Returns response length on success, -1 on send/recv error, 132// -2 if response exceeds caller's max_resp buffer. 133func nx_rpc_call(nc: *NxNetChan, method_id: i64, 134 req_buf: *u8, req_n: i64, 135 resp_buf: *u8, max_resp: i64) -> i64 { 136 if nx_net_chan_send(nc, method_id) != 0 { return -1 } 137 if nx_net_chan_send_blob(nc, req_buf, req_n) != 0 { return -1 } 138 return nx_net_chan_recv_blob(nc, resp_buf, max_resp) 139} 140 141// Client-side: send shutdown sentinel. 142func nx_rpc_shutdown_call(nc: *NxNetChan) -> i64 { 143 return nx_net_chan_send(nc, NX_RPC_SHUTDOWN) 144} 145 146// Stats / introspection. 147func nx_rpc_server_calls(s: *NxRpcServer) -> i64 { 148 let nc_addr: *i64 = ((s as i64) + 16) as *i64 149 return nx_atom_load_i64(nc_addr, NX_MO_SEQ_CST) 150}