degenbot.operator.operator_channel¶
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 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; 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¶
- degenbot.operator.operator_channel.OperatorHandler¶
- class degenbot.operator.operator_channel.StepSpec¶
A hop descriptor in the discovery-item shape the pipeline consumes.
- type¶
The typed pool family for the hop (V2/V3/V4).
- address¶
The pool address (
0x+ 40 hex).
- hash¶
The V4 pool id (
0x+ 64 hex) whentypeisPoolKind.V4.
- type: degenbot.pathfinding.PoolKind¶
- degenbot.operator.operator_channel.step_from_wire(step: dict[str, Any]) StepSpec¶
Translate a wire
stepsentry into aStepSpec.- Parameters:
step –
{"family": "V2|V3|V4", "address": "0x..", "hash": "0x.."?}.- Returns:
A
StepSpecmappingfamilyto its typed pool kind.- Raises:
ValueError – if
familyis not V2/V3/V4 oraddressis absent.
- degenbot.operator.operator_channel.FLEET_POSTURE_THRESHOLD_KEYS¶
- degenbot.operator.operator_channel.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).
- Parameters:
op – The command op (set_fleet_posture or get_fleet_posture).
payload – The command payload — a partial patch over
FLEET_POSTURE_THRESHOLD_KEYSfor 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 typedPostureRetuneErrorsurfaces verbatim).- Return type:
A handler-shaped response
- degenbot.operator.operator_channel.wrap_handler(handler: OperatorHandler) OperatorHandler¶
Wrap a host
(op, payload) -> responsehandler 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.- Parameters:
handler – The raw host
(op, payload) -> responsecoroutine.- Returns:
The wrapped handler the
OperatorServerdispatches to.
- class degenbot.operator.operator_channel.OperatorServer(handler: OperatorHandler, *, socket_path: str, request_timeout: float = 60.0)¶
A Unix-domain-socket command server for a live bot.
Run
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 bywrap_handler()so a failing command replies with{"ok": false}instead of raising into the host.- Parameters:
handler – async
(op, payload) -> dict(seeOperatorHandler).socket_path – filesystem path for the Unix domain socket.
request_timeout – seconds to wait for a request line before dropping the connection (default 60).
- async serve() None¶
Start the socket server and accept connections until closed.
Runs forever (cancellable); the caller starts it with
asyncio.create_taskand cancels it on shutdown. Theserve_foreverloop runs as its own future soclose()can cancel it (without that,wait_closedinclosewould block on the still-pending loop).
- async wait_ready(timeout_s: float = 5.0) None¶
Wait until
serve()has bound the socket and is listening.- Parameters:
timeout_s – seconds to wait before raising
TimeoutError.
- async close() None¶
Stop accepting connections and remove the socket file.
Self-sufficient: cancels the in-flight
serve_foreverloop (sowait_closed()cannot block on it) then closes the server and unlinks the socket. Safe to call whetherserve()is still running, was cancelled by the host, or is running on another thread’s event loop.
- async degenbot.operator.operator_channel.send_command(socket_path: str, op: str, payload: dict[str, Any]) dict[str, Any]¶
Connect to a running
OperatorServerand send one command.- Parameters:
socket_path – the Unix domain socket path the server listens on.
op – command op (
add_path/discover).payload – the command payload.
- Returns:
The server’s response dict (
{"ok": bool, "detail"/"error": ...}).- Raises:
RuntimeError – if the server sends no response line.