code wiki / (root) / nx_mcp_sse.nx

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}