hyf

Context-aware query service for Radroots
git clone https://radroots.dev/git/hyf.git
Log | Files | Refs | README | LICENSE

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