nx_mcp_sse_pool_gate_t316.nx source
↩ module page · 59 lines · 3194 B
1import "nx_mcp_pending.nx"
2func pg_len(s:*u8)->i64{var n:i64=0;while s[n]!=0 as u8{n=n+1};return n}
3func pg_equal(a:*u8,n:i64,b:*u8)->i64{if n!=pg_len(b){return 0};var i:i64=0;while i<n{if a[i]!=b[i]{return 0};i=i+1};return 1}
4func pg_has(a:*u8,n:i64,b:*u8)->i64{let bn:i64=pg_len(b);var i:i64=0;while i+bn<=n{var j:i64=0;while j<bn{if a[i+j]!=b[j]{break};j=j+1};if j==bn{return 1};i=i+1};return 0}
5func pg_run(failure:i64)->i64{
6 let request:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"pool-sse\",\"method\":\"tools/call\",\"params\":{\"name\":\"nx_fs\"}}"
7 let notification:*u8="{\"jsonrpc\":\"2.0\",\"method\":\"notifications/tools/list_changed\"}\n"
8 let response:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"pool-sse\",\"result\":{}}\n"
9 let expected:*u8="{\"jsonrpc\":\"2.0\",\"method\":\"notifications/tools/list_changed\"}\n{\"jsonrpc\":\"2.0\",\"id\":\"pool-sse\",\"result\":{}}\n"
10 let budget:i64=pg_len(expected)+pg_len(request)+1024
11 let pool:*NxRequestPool=rp_new(1,budget)
12 let pending:*u8=sys_mmap(__size_of(NxMcpPending))
13 let control:*NxMcpControl=sys_mmap(__size_of(NxMcpControl)) as *NxMcpControl
14 if mct_read(request,pg_len(request),control)!=0{return 1}
15 var output:i64=0;var release:i64=0
16 if sys_pipe2(&output,0)!=0||sys_pipe2(&release,0)!=0{return 2}
17 let readfd:i64=output&RP_FD_MASK;let writefd:i64=(output>>32)&RP_FD_MASK
18 let releasefd:i64=release&RP_FD_MASK;let signalfd:i64=(release>>32)&RP_FD_MASK
19 let pid:i64=mp_spawn(pool,pending,request,control)
20 if pid==0{
21 sys_close(readfd);sys_close(writefd);sys_close(signalfd)
22 rp_write(pool.child_fd,notification,pg_len(notification))
23 var signal:i64=0
24 if sys_read(releasefd,(&signal) as *u8,1)!=1{sys_exit(9);return 9}
25 if failure==1{sys_exit(7);return 7}
26 rp_write(pool.child_fd,response,pg_len(response));sys_exit(0);return 0
27 }
28 if pid<0{return 3}
29 sys_close(releasefd)
30 while rp_worker(pool,0).used==0{
31 if mp_recover_pump_timeout(pool,pending,0,writefd,5,1000)<0{return 4}
32 }
33 // Notification bytes are buffered, but neither output nor pending ownership may finish early.
34 if pool.active!=1||mp_at(pending,0).pid!=pid{return 5}
35 if rp_worker(pool,0).used!=pg_len(notification){return 6}
36 var pollfd:i64=readfd|(RP_POLLIN<<32)
37 if sys_poll((&pollfd) as *u8,1,0)!=0{return 7}
38 var signal:i64=1
39 if sys_write(signalfd,(&signal) as *u8,1)!=1{return 8}
40 sys_close(signalfd)
41 while pool.active>0{if mp_recover_pump_timeout(pool,pending,0,writefd,5,1000)<0{return 9}}
42 if mp_at(pending,0).pid!=0{return 10}
43 sys_close(writefd)
44 let received:*u8=sys_mmap(budget)
45 let n:i64=sys_read(readfd,received,budget);sys_close(readfd)
46 if failure==0{if pg_equal(received,n,expected)!=1{return 11}}else{
47 if pg_has(received,n,"notifications/tools/list_changed")!=0{return 12}
48 if pg_has(received,n,"\"id\":\"pool-sse\"")!=1{return 13}
49 if pg_has(received,n,"\"error\":")!=1{return 14}
50 }
51 return 0
52}
53func main()->i64{
54 sys_alarm(20)
55 var rc:i64=pg_run(0);if rc!=0{return rc}
56 rc=pg_run(1);if rc!=0{return rc}
57 let message:*u8="PASS SSE parent pool: notifications buffered until worker exit; ordered final response; failed worker discards partial notification and returns correlated error\n"
58 sys_write(1,message,pg_len(message));return 0
59}