server.mojo (12610B)
1 from std.collections import List, Optional 2 from std.io.io import _fdopen 3 from std.sys import stdin 4 5 from json import Value 6 7 from hyf_runtime.diagnostics import ( 8 append_internal_diagnostic as append_internal_diagnostic_to_dir, 9 effective_diagnostics_dir_for_runtime_paths, 10 ) 11 from hyf_runtime.startup import ( 12 RuntimeStartupContext, 13 resolve_startup_context_from_process, 14 ) 15 from hyf_core.capabilities.registry import ( 16 canonical_business_capability, 17 execute_gated_operation, 18 is_gated_operation, 19 ) 20 from hyf_runtime.config import operation_enabled 21 from hyf_core.errors import ( 22 CapabilityFailure, 23 CapabilityResult, 24 CapabilitySuccess, 25 ) 26 from hyf_core.metadata import hyf_protocol_version 27 from hyf_stdio.codec import ( 28 decode_request, 29 encode_error, 30 encode_success, 31 extract_request_correlation, 32 ) 33 from hyf_stdio.dispatch_observer import ( 34 DispatchAttemptObserver, 35 NoopDispatchAttemptObserver, 36 ) 37 from hyf_stdio.control.capabilities import ( 38 build_capabilities_output_with_runtime_context, 39 ) 40 from hyf_stdio.control.status import ( 41 build_status_output_with_runtime_context, 42 ) 43 from hyf_stdio.envelope import ( 44 WireErrorResponse, 45 WireRequest, 46 WireSuccessResponse, 47 ) 48 from hyf_stdio.errors import ( 49 WireError, 50 capability_disabled_error, 51 capability_unavailable_error, 52 internal_error, 53 invalid_request_error, 54 unsupported_capability_error, 55 ) 56 from hyf_stdio.meta import serialize_core_response_meta 57 from hyf_stdio.provider_execution import ( 58 execute_runtime_aware_business_capability, 59 ) 60 61 62 def _read_request_line() raises -> String: 63 with _fdopen["r"](stdin) as input_file: 64 return input_file.readline() 65 66 67 def _unsupported_response(request: WireRequest) -> WireErrorResponse: 68 return WireErrorResponse( 69 version=hyf_protocol_version(), 70 request_id=String(request.request_id), 71 trace_id=request.trace_id, 72 error=unsupported_capability_error(String(request.capability)), 73 ) 74 75 76 def _disabled_response(request: WireRequest) -> WireErrorResponse: 77 return WireErrorResponse( 78 version=hyf_protocol_version(), 79 request_id=String(request.request_id), 80 trace_id=request.trace_id, 81 error=capability_disabled_error(String(request.capability)), 82 ) 83 84 85 def _unavailable_response(request: WireRequest) -> WireErrorResponse: 86 return WireErrorResponse( 87 version=hyf_protocol_version(), 88 request_id=String(request.request_id), 89 trace_id=request.trace_id, 90 error=capability_unavailable_error(String(request.capability)), 91 ) 92 93 94 def _write_error(response: WireErrorResponse) raises: 95 print(encode_error(response)) 96 97 98 def _write_success(response: WireSuccessResponse) raises: 99 print(encode_success(response)) 100 101 102 def _diagnostic_value(value: String) -> String: 103 return value.replace("\n", "\\n").replace("\r", "\\r") 104 105 106 def _diagnostic_trace_id(trace_id: Optional[String]) -> String: 107 if trace_id: 108 return _diagnostic_value(String(trace_id.value())) 109 return "" 110 111 112 def _emit_internal_diagnostic( 113 request_id: String, 114 trace_id: Optional[String], 115 capability: String, 116 detail: String, 117 diagnostics_dir: String, 118 ): 119 append_internal_diagnostic_to_dir( 120 'hyf_internal_error request_id="' 121 + _diagnostic_value(request_id) 122 + '" trace_id="' 123 + _diagnostic_trace_id(trace_id) 124 + '" capability="' 125 + _diagnostic_value(capability) 126 + '" detail="' 127 + _diagnostic_value(detail) 128 + '"\n', 129 diagnostics_dir, 130 ) 131 132 133 def _wire_error_from_core_failure( 134 request_id: String, 135 trace_id: Optional[String], 136 failure: CapabilityFailure, 137 ) -> WireErrorResponse: 138 var code = String(failure.error.code) 139 if code == "invalid_input": 140 code = "invalid_request" 141 return WireErrorResponse( 142 version=hyf_protocol_version(), 143 request_id=request_id, 144 trace_id=trace_id, 145 error=WireError(code=code, message=String(failure.error.message)), 146 ) 147 148 149 def _wire_success_from_core_success( 150 request_id: String, 151 trace_id: Optional[String], 152 success: CapabilitySuccess, 153 ) raises -> WireSuccessResponse: 154 var meta: Optional[Value] = None 155 if success.meta: 156 meta = serialize_core_response_meta(success.meta.value()) 157 return WireSuccessResponse( 158 version=hyf_protocol_version(), 159 request_id=request_id, 160 trace_id=trace_id, 161 output=success.output.clone(), 162 meta=meta^, 163 ) 164 165 166 def _dispatch_capability_result( 167 request_id: String, 168 trace_id: Optional[String], 169 result: CapabilityResult, 170 ) raises -> String: 171 if result.failure: 172 return encode_error( 173 _wire_error_from_core_failure( 174 request_id, trace_id, result.failure.value() 175 ) 176 ) 177 return encode_success( 178 _wire_success_from_core_success( 179 request_id, trace_id, result.success.value() 180 ) 181 ) 182 183 184 def _dispatch_business_capability( 185 request: WireRequest, 186 request_id: String, 187 runtime_context: RuntimeStartupContext, 188 ) raises -> String: 189 var result = execute_runtime_aware_business_capability( 190 request.capability, 191 request.input.clone(), 192 request.context.copy(), 193 runtime_context, 194 ) 195 return _dispatch_capability_result(request_id, request.trace_id, result) 196 197 198 def _route_business_capability[ 199 O: DispatchAttemptObserver 200 ]( 201 mut observer: O, 202 request: WireRequest, 203 request_id: String, 204 runtime_context: RuntimeStartupContext, 205 ) raises -> String: 206 # ADR-0025 D45 CB01 pre-activation guard: a recognized hyf_ops_v2 request has 207 # already passed the strict capability-context parse, but activation belongs 208 # to C042-C046, so it must not reach the legacy shortcut handlers even when 209 # their legacy enable flag is set. C008/C009 retain this guard. 210 if request.operation_context: 211 return encode_error(_unavailable_response(request)) 212 213 # ADR-0027 D47 BP02: every path below is a real pre-dispatch boundary. The 214 # compile-time no-op production observer records nothing; a test-owned 215 # in-memory observer proves the guard above short-circuits before this point. 216 observer.record_dispatch_attempt(String(request.capability)) 217 218 if is_gated_operation(request.capability): 219 if not operation_enabled(runtime_context.config, request.capability): 220 return encode_error(_disabled_response(request)) 221 var gated = execute_gated_operation( 222 request.capability, request.input.clone() 223 ) 224 return _dispatch_capability_result(request_id, request.trace_id, gated) 225 226 var capability = canonical_business_capability(request.capability) 227 if not capability: 228 return encode_error(_unsupported_response(request)) 229 230 var descriptor = capability.value().copy() 231 if not descriptor.deterministic_enabled: 232 return encode_error(_disabled_response(request)) 233 234 if descriptor.implemented and descriptor.callable: 235 return _dispatch_business_capability( 236 request, request_id, runtime_context 237 ) 238 239 return encode_error(_unavailable_response(request)) 240 241 242 def handle_request(request: WireRequest) raises -> String: 243 return handle_request_with_runtime_context( 244 request, resolve_startup_context_from_process() 245 ) 246 247 248 def handle_request_with_runtime_context( 249 request: WireRequest, runtime_context: RuntimeStartupContext 250 ) raises -> String: 251 var observer = NoopDispatchAttemptObserver() 252 return handle_request_with_runtime_context_and_observer( 253 request, runtime_context, observer 254 ) 255 256 257 def handle_request_with_runtime_context_and_observer[ 258 O: DispatchAttemptObserver 259 ]( 260 request: WireRequest, 261 runtime_context: RuntimeStartupContext, 262 mut observer: O, 263 ) raises -> String: 264 var request_id = String(request.request_id) 265 var trace_id = request.trace_id 266 var diagnostics_dir = effective_diagnostics_dir_for_runtime_paths( 267 runtime_context.paths 268 ) 269 try: 270 if request.capability == "sys.status": 271 return encode_success( 272 WireSuccessResponse( 273 version=hyf_protocol_version(), 274 request_id=request_id, 275 trace_id=trace_id, 276 output=build_status_output_with_runtime_context( 277 runtime_context 278 ), 279 meta=None, 280 ) 281 ) 282 if request.capability == "sys.capabilities": 283 return encode_success( 284 WireSuccessResponse( 285 version=hyf_protocol_version(), 286 request_id=request_id, 287 trace_id=trace_id, 288 output=build_capabilities_output_with_runtime_context( 289 runtime_context 290 ), 291 meta=None, 292 ) 293 ) 294 return _route_business_capability( 295 observer, request.copy(), request_id, runtime_context 296 ) 297 except e: 298 _emit_internal_diagnostic( 299 request_id, 300 trace_id, 301 String(request.capability), 302 String(e), 303 diagnostics_dir, 304 ) 305 return encode_error( 306 WireErrorResponse( 307 version=hyf_protocol_version(), 308 request_id=request_id, 309 trace_id=trace_id, 310 error=internal_error(), 311 ) 312 ) 313 314 315 def handle_request_line(line: String) raises -> String: 316 try: 317 var request = decode_request(line) 318 return handle_request(request^) 319 except e: 320 var correlation = extract_request_correlation(line) 321 return encode_error( 322 WireErrorResponse( 323 version=hyf_protocol_version(), 324 request_id=correlation.request_id, 325 trace_id=correlation.trace_id, 326 error=invalid_request_error(String(e)), 327 ) 328 ) 329 330 331 def handle_request_line_with_runtime_context( 332 line: String, runtime_context: RuntimeStartupContext 333 ) raises -> String: 334 var observer = NoopDispatchAttemptObserver() 335 return handle_request_line_with_runtime_context_and_observer( 336 line, runtime_context, observer 337 ) 338 339 340 def handle_request_line_with_runtime_context_and_observer[ 341 O: DispatchAttemptObserver 342 ]( 343 line: String, runtime_context: RuntimeStartupContext, mut observer: O 344 ) raises -> String: 345 try: 346 var request = decode_request(line) 347 return handle_request_with_runtime_context_and_observer( 348 request^, runtime_context, observer 349 ) 350 except e: 351 var correlation = extract_request_correlation(line) 352 return encode_error( 353 WireErrorResponse( 354 version=hyf_protocol_version(), 355 request_id=correlation.request_id, 356 trace_id=correlation.trace_id, 357 error=invalid_request_error(String(e)), 358 ) 359 ) 360 361 362 def run_stdio_server() raises: 363 run_stdio_server_with_runtime_context( 364 resolve_startup_context_from_process() 365 ) 366 367 368 def handle_frame( 369 frame: String, runtime_context: RuntimeStartupContext 370 ) raises -> String: 371 if frame_too_large(frame): 372 return encode_error( 373 WireErrorResponse( 374 version=hyf_protocol_version(), 375 request_id="", 376 trace_id=None, 377 error=invalid_request_error( 378 "request frame exceeds the size limit" 379 ), 380 ) 381 ) 382 return handle_request_line_with_runtime_context(frame, runtime_context) 383 384 385 def run_stdio_session( 386 frames: List[String], runtime_context: RuntimeStartupContext 387 ) raises -> List[String]: 388 var responses = List[String]() 389 for frame in frames: 390 responses.append(handle_frame(frame, runtime_context)) 391 return responses^ 392 393 394 def run_stdio_server_with_runtime_context( 395 runtime_context: RuntimeStartupContext, 396 ) raises: 397 if stdin.isatty(): 398 return 399 400 with _fdopen["r"](stdin) as input_file: 401 while True: 402 var line = String() 403 try: 404 line = input_file.readline() 405 except e: 406 if String(e) == "EOF": 407 break 408 raise e^ 409 if line == "": 410 break 411 print(handle_frame(line, runtime_context)) 412 413 414 comptime MAX_FRAME_BYTES = 1048576 415 416 417 def frame_too_large(line: String) -> Bool: 418 return line.byte_length() > MAX_FRAME_BYTES