degenbot.operator.operator_channel ================================== .. py:module:: degenbot.operator.operator_channel .. autoapi-nested-parse:: Operator command channel: JSON-lines over a Unix domain socket. Lets an operator steer a live bot without touching its process. The host bot runs an :class:`OperatorServer` (an asyncio task) bound to a Unix domain socket; the ``degenbot path add`` / ``degenbot path discover`` CLI (or any client) writes ONE JSON command line and reads ONE JSON response line. The wire format is minimal and versioned so a client built earlier still talks to a running host built later. Wire protocol (one JSON object per line, newline-terminated). Request lines:: { "op": "add_path", "steps": [ {"family": "V2|V3|V4", "address": "0x..", "hash": "0x..(32-byte v4 pool_id, v4 only)"}, ..., ], "directions": [true, false], } # optional; auto-resolved if absent {"op": "discover", "bound": 5} # bounded on-demand discovery { "op": "set_fleet_posture", "cordon_enter_events": 2, # ANY subset of the six cordon_* keys "cordon_enter_window_ms": 1000, # (the typed DEGENBOT_FLEET_CORDON_* "cordon_duty_percent": 2.0, # keys); at least one is required; "cordon_duty_window_ms": 5000, # absent keys keep the live value; "cordon_exit_clean_ms": 10000, # cordon_sim_intake_floor also takes "cordon_sim_intake_floor": 2, # null = restore half the slot cap } {"op": "get_fleet_posture"} # read the live thresholds + posture Response lines:: {"ok": true, "detail": "..."} {"ok": true, "detail": "", "effective": {...}} # fleet posture ops: # the six cordon_* values + "posture" {"ok": false, "error": "..."} The server maps a ``family`` string to the typed Rust-backed pool family the registration pipeline consumes, so no Python class object crosses the wire. The handler (supplied by the host) turns an ``(op, payload)`` pair into a response dict; the server guards the wire (JSON decode, op dispatch, exception -> ``{"ok": false}``) so a malformed or failing command never crashes the host. The fleet-posture ops (`set_fleet_posture` / `get_fleet_posture`) re-tune the LIVE cordon thresholds of the process posture owner through the `degenbot.fleet` mirror home; :func:`handle_fleet_posture_op` is the host-side helper both the runner's handler and the tests route them through. Validation of the six values lives ONCE in the Rust core (`PosturePolicyPatch::validate`) and is surfaced as the typed `degenbot.fleet.PostureRetuneError`; this layer adds the wire-hygiene checks (unknown key, empty patch) in front of it — defense-in-depth, not a second validator. Boot config stays the default source: the op re-tunes the live policy only. The protocol is a plain request/response: the host processes each command concurrently with the pump and replies once the registration pipeline has accepted it (add-path replies after the per-path registration body completes; discover replies after the bounded sweep is consumed). Module Contents --------------- .. py:data:: OperatorHandler .. py:class:: StepSpec A hop descriptor in the discovery-item shape the pipeline consumes. .. attribute:: type The typed pool family for the hop (V2/V3/V4). .. attribute:: address The pool address (``0x`` + 40 hex). .. attribute:: hash The V4 pool id (``0x`` + 64 hex) when ``type`` is ``PoolKind.V4``. .. py:attribute:: type :type: degenbot.pathfinding.PoolKind .. py:attribute:: address :type: str .. py:attribute:: hash :type: object | None :value: None .. py:function:: step_from_wire(step: dict[str, Any]) -> StepSpec Translate a wire ``steps`` entry into a :class:`StepSpec`. :param step: ``{"family": "V2|V3|V4", "address": "0x..", "hash": "0x.."?}``. :returns: A :class:`StepSpec` mapping ``family`` to its typed pool kind. :raises ValueError: if ``family`` is not V2/V3/V4 or ``address`` is absent. .. py:data:: FLEET_POSTURE_THRESHOLD_KEYS .. py:function:: handle_fleet_posture_op(op: str, payload: dict[str, Any]) -> dict[str, Any] Handle the two fleet-posture ops through the `degenbot.fleet` mirror. The host-side helper the runner's operator handler routes `set_fleet_posture` / `get_fleet_posture` through (and the round-trip tests exercise over a real socket). Sync by design — the channel body is two cheap FFI round-trips; the host's async handler calls it directly from its async body. The mirror home is imported lazily: the transport module stays independent of the compiled extension, and the home mints on the first op that actually runs (the ADR-013 barrier). :param op: The command op (`set_fleet_posture` or `get_fleet_posture`). :param payload: The command payload — a partial patch over :data:`FLEET_POSTURE_THRESHOLD_KEYS` for the set op (at least one key; `cordon_sim_intake_floor` also takes `None` = restore half the slot cap); ignored for the get op. :returns: ``{"effective": {...}}`` on success (the effective policy echoed — all six fields + the current posture), or ``{"error": "..."}`` on an unknown op/key, an empty patch, or a refused patch (the typed ``PostureRetuneError`` surfaces verbatim). :rtype: A handler-shaped response .. py:function:: wrap_handler(handler: OperatorHandler) -> OperatorHandler Wrap a host ``(op, payload) -> response`` handler with wire hygiene. Turns any raised exception into a ``{"ok": false, "error": ...}`` response (never crashes the host), and provides a convenient success/error spelling: the handler returns ``{"detail": ...}`` on success or ``{"error": ...}`` on failure, and the wrapper normalizes both to the wire shape. :param handler: The raw host ``(op, payload) -> response`` coroutine. :returns: The wrapped handler the :class:`OperatorServer` dispatches to. .. py:class:: OperatorServer(handler: OperatorHandler, *, socket_path: str, request_timeout: float = 60.0) A Unix-domain-socket command server for a live bot. Run :meth:`serve` as an asyncio task (e.g. a background task on the registration loop). It accepts one or more concurrent client connections, reads one JSON request line, routes it to the wrapped handler, and writes one JSON response line. The handler is wrapped by :func:`wrap_handler` so a failing command replies with ``{"ok": false}`` instead of raising into the host. :param handler: async ``(op, payload) -> dict`` (see :data:`OperatorHandler`). :param socket_path: filesystem path for the Unix domain socket. :param request_timeout: seconds to wait for a request line before dropping the connection (default 60). .. py:method:: serve() -> None :async: Start the socket server and accept connections until closed. Runs forever (cancellable); the caller starts it with ``asyncio.create_task`` and cancels it on shutdown. The ``serve_forever`` loop runs as its own future so :meth:`close` can cancel it (without that, ``wait_closed`` in ``close`` would block on the still-pending loop). .. py:method:: wait_ready(timeout_s: float = 5.0) -> None :async: Wait until :meth:`serve` has bound the socket and is listening. :param timeout_s: seconds to wait before raising :class:`TimeoutError`. .. py:method:: close() -> None :async: Stop accepting connections and remove the socket file. Self-sufficient: cancels the in-flight ``serve_forever`` loop (so ``wait_closed()`` cannot block on it) then closes the server and unlinks the socket. Safe to call whether ``serve()`` is still running, was cancelled by the host, or is running on another thread's event loop. .. py:function:: send_command(socket_path: str, op: str, payload: dict[str, Any]) -> dict[str, Any] :async: Connect to a running :class:`OperatorServer` and send one command. :param socket_path: the Unix domain socket path the server listens on. :param op: command op (``add_path`` / ``discover``). :param payload: the command payload. :returns: The server's response dict (``{"ok": bool, "detail"/"error": ...}``). :raises RuntimeError: if the server sends no response line.