nx_mcp_sse.nx source
↩ module page · 127 lines · 5393 B
1// MCP response decoding belongs to the native stdio transport, before protocol output.
2// No partial notification/response batch escapes when framing or correlation fails.
3import "nx_mcp_route.nx"
4import "nx_http_response_parse.nx"
5const MSS_BAD_BODY:i64=0-1
6const MSS_BAD_ID:i64=0-2
7const MSS_INCOMPLETE:i64=0-3
8const MSS_OVERFLOW:i64=0-4
9struct NxMcpStream { request:*u8, request_n:i64, output:*u8, capacity:i64, used:i64, responses:i64 }
10
11func mss_message_kind(request:*u8,request_n:i64,data:*u8,n:i64)->i64 {
12 var verdict:i64=0
13 let root:*NxValue=nx_value_parse_json(data,n,&verdict)
14 if verdict!=NX_VAL_PARSE_OK{return MSS_BAD_BODY}
15 var bad:i64=0
16 let version:*NxValue=mr_member(root,"jsonrpc",&bad)
17 let id:*NxValue=mr_member(root,"id",&bad)
18 let method:*NxValue=mr_member(root,"method",&bad)
19 let result:*NxValue=mr_member(root,"result",&bad)
20 let error:*NxValue=mr_member(root,"error",&bad)
21 if bad!=0||version==(0 as *NxValue){return MSS_BAD_BODY}
22 if version.kind!=NX_VAL_STRING{return MSS_BAD_BODY}
23 if mr_equal(version.str_ptr,version.str_len,"2.0")!=1{return MSS_BAD_BODY}
24 if method!=(0 as *NxValue){
25 if method.kind!=NX_VAL_STRING||method.str_len==0{return MSS_BAD_BODY}
26 if result!=(0 as *NxValue)||error!=(0 as *NxValue){return MSS_BAD_BODY}
27 // Bidirectional server requests need a separately owned response route; never mistake one for a notification.
28 if id!=(0 as *NxValue){return MSS_BAD_BODY}
29 return 0
30 }
31 if result==(0 as *NxValue)&&error==(0 as *NxValue){return MSS_BAD_BODY}
32 if result!=(0 as *NxValue)&&error!=(0 as *NxValue){return MSS_BAD_BODY}
33 if error!=(0 as *NxValue){
34 let code:*NxValue=mr_member(error,"code",&bad)
35 let message:*NxValue=mr_member(error,"message",&bad)
36 if bad!=0||code==(0 as *NxValue)||message==(0 as *NxValue){return MSS_BAD_BODY}
37 if code.kind!=NX_VAL_INT||message.kind!=NX_VAL_STRING{return MSS_BAD_BODY}
38 }
39 if mr_response_matches(request,request_n,data,n)!=1{return MSS_BAD_ID}
40 return 1
41}
42
43// SSE multiline data can contain formatting LF. Stdio uses one JSON value per line.
44// Remove only JSON whitespace outside strings after the complete value has been validated.
45func mss_append(s:*NxMcpStream,data:*u8,n:i64)->i64 {
46 var quoted:i64=0;var escaped:i64=0;var i:i64=0
47 while i<n{
48 let c:i64=data[i] as i64
49 var keep:i64=1
50 if quoted==0{
51 if c==32||c==9||c==10||c==13{keep=0}
52 if c==34{quoted=1}
53 }else{
54 if escaped==1{escaped=0}else{if c==92{escaped=1}else{if c==34{quoted=0}}}
55 }
56 if keep==1{if s.used>=s.capacity{return MSS_OVERFLOW};s.output[s.used]=data[i];s.used=s.used+1}
57 i=i+1
58 }
59 if s.used>=s.capacity{return MSS_OVERFLOW}
60 s.output[s.used]=10 as u8;s.used=s.used+1
61 return 0
62}
63func mss_accept(s:*NxMcpStream,data:*u8,n:i64)->i64 {
64 let kind:i64=mss_message_kind(s.request,s.request_n,data,n)
65 if kind<0{return kind}
66 if kind==1{if s.responses!=0{return MSS_BAD_ID};s.responses=1}
67 return mss_append(s,data,n)
68}
69
70// SSE field names are case-sensitive; CR, LF and CRLF delimit lines. A blank line
71// dispatches data joined with LF. EOF never implicitly dispatches an unfinished event.
72func mss_events(s:*NxMcpStream,src:*u8,n:i64)->i64 {
73 let data:*u8=sys_mmap(n+1)
74 var dn:i64=0;var pos:i64=0;var message:i64=1
75 if n>=3{if src[0]==239 as u8&&src[1]==187 as u8&&src[2]==191 as u8{pos=3}}
76 while pos<n{
77 let start:i64=pos
78 while pos<n{if src[pos]==10 as u8||src[pos]==13 as u8{break};pos=pos+1}
79 if pos==n{return MSS_INCOMPLETE}
80 let end:i64=pos;let delimiter:i64=src[pos] as i64;pos=pos+1
81 if delimiter==13&&pos<n{if src[pos]==10 as u8{pos=pos+1}}
82 if end==start{
83 if dn>0&&message==1{let rc:i64=mss_accept(s,data,dn-1);if rc<0{return rc}}
84 dn=0;message=1
85 }else{
86 var colon:i64=start
87 while colon<end{if src[colon]==58 as u8{break};colon=colon+1}
88 var value:i64=colon
89 if value<end{value=value+1;if value<end{if src[value]==32 as u8{value=value+1}}}
90 if mr_equal(src+start,colon-start,"data")==1{
91 let length:i64=end-value
92 if length>n-dn{return MSS_OVERFLOW}
93 var i:i64=0;while i<length{data[dn+i]=src[value+i];i=i+1}
94 dn=dn+length;data[dn]=10 as u8;dn=dn+1
95 }
96 if mr_equal(src+start,colon-start,"event")==1{
97 message=0
98 if value==end||mr_equal(src+value,end-value,"message")==1{message=1}
99 }
100 // Comments, event IDs, retry hints and unknown fields carry no JSON-RPC value.
101 }
102 }
103 if dn!=0{return MSS_INCOMPLETE}
104 if s.responses!=1{return MSS_INCOMPLETE}
105 return s.used
106}
107
108// Caller owns both buffers. Output is valid only for a positive return value.
109func mss_decode(request:*u8,request_n:i64,src:*u8,n:i64,sse:i64,out:*u8,capacity:i64)->i64 {
110 if n<=0||capacity<=0{return MSS_BAD_BODY}
111 let s:*NxMcpStream=sys_mmap(__size_of(NxMcpStream)) as *NxMcpStream
112 s.request=request;s.request_n=request_n;s.output=out;s.capacity=capacity;s.used=0;s.responses=0
113 if sse==1{return mss_events(s,src,n)}
114 let rc:i64=mss_accept(s,src,n)
115 if rc<0{return rc}
116 if s.responses!=1{return MSS_BAD_ID}
117 return s.used
118}
119func mss_media_type(src:*u8,n:i64,parsed:*i64)->i64 {
120 var off:i64=0;var length:i64=0
121 if nx_http_response_header(src,n,parsed,"Content-Type",12,&off,&length)!=1{return MSS_BAD_BODY}
122 var end:i64=0
123 while end<length{let c:i64=src[off+end] as i64;if c==59||c==32||c==9{break};end=end+1}
124 if nx_http_resp_cieq(src+off,end,"application/json",16)==1{return 0}
125 if nx_http_resp_cieq(src+off,end,"text/event-stream",17)==1{return 1}
126 return MSS_BAD_BODY
127}