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}