nx_mcp_queue_gate_t230.nx source
↩ module page · 75 lines · 4663 B
1// nx_mcp_queue_gate_t230.nx -- Manages JSON-RPC request queuing, cancellation, and worker recovery in a distributed system.
2import "nx_mcp_stdio_candidate_t230.nx"
3import "nx_gate_verdict.nx"
4func qg_add(q:*NxMcpQueue,p:*NxRequestPool,rows:*u8,s:*u8,now:i64,fd:i64)->i64 {
5 var c:NxMcpControl
6 if mct_read(s,mr_len(s),&c)!=0{return -1}
7 return mq_enqueue(q,p,rows,s,mr_len(s),&c,now,fd)
8}
9func qg_has(b:*u8,n:i64,s:*u8)->i64 {
10 let l:i64=mr_len(s);var i:i64=0
11 while i+l<=n{var j:i64=0;while j<l{if b[i+j]!=s[j]{break};j=j+1};if j==l{return 1};i=i+1};return 0
12}
13func main()->i64 {
14 let c:*i64=gv_ctr()
15 let plan:*i64=gv_plan_new("queue-beyond-executor-capacity\nqueued-id-duplicate\nrunning-id-duplicate\nqueued-cancel-no-dispatch\ncancel-releases-budget\nfifo-preserved\nbarrier-waits-without-drain\nworker-failure-recovers\nqueue-dispatch-after-recovery\nexpired-never-dispatched\nmemory-refusal-explicit\nfinite-poll-returns\n")
16 let p:*NxRequestPool=rp_new(1,8192)
17 let rows:*u8=sys_mmap(__size_of(NxMcpPending))
18 let q:*NxMcpQueue=mq_new(8192,100)
19 q.first=0
20 let fd:i64=sys_openat_wr("knowledge/gates/mcp-queue-protocol-t230.jsonl",MODE_0644)
21 if fd<0{return 2}
22 var pipe:i64=0;if sys_pipe2(&pipe,0)!=0{return 2}
23 let rd:i64=pipe&RP_FD_MASK;let wr:i64=(pipe>>32)&RP_FD_MASK
24 let active:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"active\",\"method\":\"tools/list\"}"
25 var ctl:NxMcpControl;mct_read(active,mr_len(active),&ctl)
26 let pid:i64=mp_spawn(p,rows,active,&ctl)
27 if pid==0{var x:i64=0;sys_read(rd,&x as *u8,1);sys_exit(7);return 7}
28 let a:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"a\",\"method\":\"tools/list\"}"
29 let b:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"b\",\"method\":\"tools/list\"}"
30 let d:*u8="{\"jsonrpc\":\"2.0\",\"id\":\"c\",\"method\":\"tools/list\"}"
31 let now:i64=sys_now_ms()
32 let added:i64=qg_add(q,p,rows,a,now,fd)+qg_add(q,p,rows,b,now,fd)+qg_add(q,p,rows,d,now,fd)
33 gv_plan_check(plan,"queue-beyond-executor-capacity",added==3&&q.count==3&&p.active==1,c)
34 gv_plan_check(plan,"queued-id-duplicate",qg_add(q,p,rows,b,now,fd)==0&&q.count==3,c)
35 gv_plan_check(plan,"running-id-duplicate",qg_add(q,p,rows,active,now,fd)==0&&q.count==3,c)
36 let held:i64=q.used
37 let original_budget:i64=q.budget;q.budget=q.used
38 let cancel:*u8="{\"jsonrpc\":\"2.0\",\"method\":\"notifications/cancelled\",\"params\":{\"requestId\":\"b\"}}"
39 mct_read(cancel,mr_len(cancel),&ctl)
40 gv_plan_check(plan,"queued-cancel-no-dispatch",mq_cancel(q,cancel,&ctl,fd)==1&&p.active==1&&q.count==2,c)
41 gv_plan_check(plan,"cancel-releases-budget",q.used<held,c)
42 q.budget=original_budget
43 gv_plan_check(plan,"fifo-preserved",mq_same(q.head,NX_JSON_STRING,"a",1)==1&&mq_same(q.tail,NX_JSON_STRING,"c",1)==1,c)
44 q.head.barrier=1
45 gv_plan_check(plan,"barrier-waits-without-drain",mq_ready(q,p)==0&&p.active==1,c)
46 gv_plan_check(plan,"finite-poll-returns",mp_recover_pump_timeout(p,rows,0,fd,1,0)>=0&&p.active==1,c)
47 var x:i64=1;sys_write(wr,&x as *u8,1);sys_close(wr);sys_close(rd)
48 var okay:i64=1
49 while p.active>0{if mp_recover_pump_timeout(p,rows,0,fd,1,1000)<0{okay=0;break}}
50 gv_plan_check(plan,"worker-failure-recovers",okay==1&&p.active==0&&mq_ready(q,p)==1,c)
51 // Dispatch actual request workers. Missing credential fails before HTTPS:
52 // validates queue -> worker -> correlated failure, not a network mock.
53 let av:*i64=sys_mmap(7*__size_of(i64))
54 av[0]="candidate" as i64;av[1]="https://nishifamily.com" as i64
55 av[2]="knowledge/gates/nonexistent-cap-t230" as i64
56 av[3]="8192" as i64;av[4]="8192" as i64;av[5]="8192" as i64;av[6]="1" as i64
57 var nd:*NxMcpQueued=q.head
58 while nd!=(0 as *NxMcpQueued){nd.deadline_ms=sys_now_ms()+10000;nd=nd.next}
59 var dispatched:i64=0
60 while q.count>0||p.active>0{
61 if ms_dispatch_queue(q,p,rows,fd,1,8192,7,av,7)<0{okay=0;break}
62 if p.active>0{dispatched=dispatched+1;if mp_recover_pump_timeout(p,rows,0,fd,1,1000)<0{okay=0;break}}
63 }
64 gv_plan_check(plan,"queue-dispatch-after-recovery",okay==1&&q.count==0&&q.used==0&&p.active==0&&dispatched>=2,c)
65 qg_add(q,p,rows,a,10,fd);mq_expire(q,110,fd)
66 gv_plan_check(plan,"expired-never-dispatched",q.count==0&&p.active==0&&q.used==0,c)
67 let budget:i64=q.budget;q.budget=__size_of(NxMcpQueued)+1
68 gv_plan_check(plan,"memory-refusal-explicit",qg_add(q,p,rows,a,now,fd)==0&&q.count==0,c);q.budget=budget
69 sys_close(fd)
70 var n:i64=0;let output:*u8=sys_read_file("knowledge/gates/mcp-queue-protocol-t230.jsonl",&n)
71 gv_kv("protocol_bytes",n)
72 if qg_has(output,n,"\"id\":\"a\"")==0||qg_has(output,n,"\"id\":\"c\"")==0{return 3}
73 gv_plan_finish(plan,c)
74 return gv_verdict("MCP-QUEUE",c,"Real occupied worker, queued cancellation, crash recovery, later actual workers, expiry and configured memory refusal.")
75}