From fd10607033686c2b47db7a76b84f9c2be29a123c Mon Sep 17 00:00:00 2001 From: Dobromir Popov Date: Sat, 1 Aug 2026 01:52:42 +0300 Subject: [PATCH] story: DGR-042 Carry native frames through direct and existing relay seams --- .../evidence/DGR-042/README.md | 73 +++++ ...through-direct-and-existing-relay-seams.md | 14 +- .scratch/distributed-gguf-runtime/prd.json | 3 +- .../meshnet_node/native_activation_seam.py | 298 ++++++++++++++++++ tests/test_native_activation_seam.py | 159 ++++++++++ 5 files changed, 539 insertions(+), 8 deletions(-) create mode 100644 .scratch/distributed-gguf-runtime/evidence/DGR-042/README.md create mode 100644 packages/node/meshnet_node/native_activation_seam.py create mode 100644 tests/test_native_activation_seam.py diff --git a/.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md b/.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md new file mode 100644 index 0000000..920f9dd --- /dev/null +++ b/.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md @@ -0,0 +1,73 @@ +# DGR-042 evidence — native frames through direct and relay seams + +**Date:** 2026-08-01 +**Authority:** `.scratch/distributed-gguf-runtime/prd.json`. + +## Implemented + +- Added `NativeActivationSeam`, a Route-Session-scoped adapter with exactly two + selectable transports. Direct traffic calls the generated + `ShardRuntimeStub.Session()` once and keeps its bidirectional gRPC stream + open for the session. Its request and response hand-off queues are bounded. +- Relay traffic calls the existing persistent relay request shape with + `POST /native/session`, `application/x-protobuf`, and the exact + `SessionRequest.SerializeToString()` body. It parses only the returned + `SessionResponse`; neither the adapter nor the relay contract rewrites a + protobuf frame. Relay failure is explicitly uncertain and is never retried. +- `NativeFrameContext` validates Route Session, epoch, work, and deadline + fields against the versioned protobuf request before either path sends it. + The unchanged existing relay header contract receives request/billing ID, + node attribution, route, work, and deadline copies for control-plane + telemetry/billing correlation. `NativeSeamTelemetry` reports per-node, + per-request seam byte/latency observations without interpreting frames. +- Deterministic fake-worker tests cover a single direct stream, byte-identical + relay request frames, relay disconnect/no replay, cancellation, correlation + headers, telemetry, and bounded direct buffering. + +## Changed files + +- `packages/node/meshnet_node/native_activation_seam.py` +- `tests/test_native_activation_seam.py` +- `.scratch/distributed-gguf-runtime/prd.json` +- `.scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md` +- `.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md` +- `.ralph-tui/progress.md` + +## Commands and results + +```bash +PYTHONPATH=packages/node:packages/tracker /home/popov/.hermes/hermes-agent/venv/bin/python -m pytest -q \ + tests/test_native_activation_seam.py tests/test_native_shard_protocol.py \ + tests/test_native_worker_supervisor.py tests/test_native_registration.py \ + tests/test_ralph_prd_schema.py +# 172 passed, 2 skipped in 2.01s + +/home/popov/.hermes/hermes-agent/venv/bin/python -m ruff check \ + packages/node/meshnet_node/native_activation_seam.py tests/test_native_activation_seam.py +# All checks passed! + +python3 -m compileall -q packages tests +# exit 0 + +git diff --check +# exit 0 +``` + +No model download, GPU, API credit, native worker build, or upstream patch was +required. Native CMake/CTest and patch-stack gates do not apply to this +Python-only transport adapter. + +## Limitations and dependency handoff + +- Relay is deliberately a sequence of opaque existing relay RPC bodies, not a + gRPC tunnel. The direct path alone is a long-lived gRPC stream; this avoids + changing relay behavior while preserving native frame bytes. +- This fixture lane uses an injected generated-stub-shaped fake worker and an + injected existing-relay-client-shaped callable. DGR-054/DGR-058 must use the + adapter with certified workers and add real route-loss/restart policy; they + must retain the no-replay rule after an uncertain relay send. +- DGR-024 supplied the versioned generated `Session` protocol and the prior + raw-frame identity proof. DGR-040 supplied the verified worker lifecycle; + its published native listen address is the direct endpoint for this seam. +- Existing Transformer HTTP routes and relay routing, load balancing, billing, + and peer behavior were not changed. diff --git a/.scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md b/.scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md index 816f8c7..db24e5c 100644 --- a/.scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md +++ b/.scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md @@ -1,7 +1,7 @@ # DGR-042: Carry native frames through direct and existing relay seams -- **Status / triage:** specification only; `ready-for-agent`; `passes: false` +- **Status / triage:** completed; `passes: true` - **Execution mode:** `AFK` - **Milestone:** `M2` - **Dependencies:** `DGR-024`, `DGR-040` @@ -18,11 +18,11 @@ Fresh Ralph session: read `.scratch/distributed-gguf-runtime/RALPH-CONTEXT.md`, ## Acceptance criteria -- [ ] Direct paths use the long-lived gRPC activation stream. -- [ ] Relayed paths carry byte-identical versioned protobuf frames through the existing relay contract. -- [ ] Request/work identity, cancellation, deadlines, telemetry, billing correlation, and per-node attribution survive both paths. -- [ ] Fake-worker tests cover direct, relay, disconnect, cancellation, and bounded buffering. -- [ ] Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff. +- [x] Direct paths use the long-lived gRPC activation stream. +- [x] Relayed paths carry byte-identical versioned protobuf frames through the existing relay contract. +- [x] Request/work identity, cancellation, deadlines, telemetry, billing correlation, and per-node attribution survive both paths. +- [x] Fake-worker tests cover direct, relay, disconnect, cancellation, and bounded buffering. +- [x] Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff. ## Shared quality gates @@ -36,4 +36,4 @@ Fresh Ralph session: read `.scratch/distributed-gguf-runtime/RALPH-CONTEXT.md`, ## Evidence handoff -Write and verify `.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md`. Until every criterion and applicable gate has real evidence, this story remains `passes: false`. Legacy evidence is provenance only, not completion credit. +Verified evidence: `.scratch/distributed-gguf-runtime/evidence/DGR-042/README.md`. Legacy evidence remains provenance only and grants no implementation completion credit. diff --git a/.scratch/distributed-gguf-runtime/prd.json b/.scratch/distributed-gguf-runtime/prd.json index d75d8e9..6baf022 100644 --- a/.scratch/distributed-gguf-runtime/prd.json +++ b/.scratch/distributed-gguf-runtime/prd.json @@ -1276,7 +1276,8 @@ "Fake-worker tests cover direct, relay, disconnect, cancellation, and bounded buffering.", "Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff." ], - "passes": false, + "passes": true, + "completionNotes": "Completed 2026-08-01. Added a bounded NativeActivationSeam that keeps direct traffic on one generated gRPC Session stream and uses the existing relay request contract to carry serialized SessionRequest/SessionResponse bytes without rewriting them. Fixture tests cover direct, opaque relay, disconnect/no replay, cancellation, correlation headers, telemetry, and bounded buffering. See evidence/DGR-042/README.md.", "notes": "Generated source issue: .scratch/distributed-gguf-runtime/issues/042-carry-native-frames-through-direct-and-existing-relay-seams.md; prd.json is authoritative.", "blocks": [ "DGR-054", diff --git a/packages/node/meshnet_node/native_activation_seam.py b/packages/node/meshnet_node/native_activation_seam.py new file mode 100644 index 0000000..bed50bb --- /dev/null +++ b/packages/node/meshnet_node/native_activation_seam.py @@ -0,0 +1,298 @@ +"""Native activation transport over direct gRPC or the existing relay RPC. + +This is deliberately a *seam adapter*, not a new relay protocol. Direct +peers use one generated ``ShardRuntime.Session`` bidi stream for the lifetime +of a Route Session. A relayed peer uses the relay's existing HTTP-shaped +binary-body contract: each body is exactly a serialized ``SessionRequest`` or +``SessionResponse``. The relay only routes those bytes and restores its own +request id; it does not deserialize a native frame. + +The correlation headers are duplicated outside the opaque frame solely for +the existing tracker/relay observability and billing path. The authoritative +work, route, epoch, deadline, and cancellation information remains in the +versioned protobuf frame and is validated before it is sent. +""" + +from __future__ import annotations + +from collections.abc import Callable, Iterator +from dataclasses import dataclass +from queue import Empty, Full, Queue +import threading +import time +from typing import Protocol + +from .native_protocol import pb + +NATIVE_RELAY_PATH = "/native/session" +NATIVE_FRAME_CONTENT_TYPE = "application/x-protobuf" + + +class NativeActivationSeamError(RuntimeError): + """The activation seam cannot safely continue this Route Session.""" + + +class NativeActivationBufferFull(NativeActivationSeamError): + """The caller exceeded the negotiated local hand-off buffer.""" + + +class NativeActivationDisconnected(NativeActivationSeamError): + """A direct or relay transport disconnected with an uncertain outcome.""" + + +class RelayRequest(Protocol): + """The existing ``_RelayHopClient.request`` shape, kept dependency-free.""" + + def __call__( + self, path: str, body: bytes, headers: dict[str, str] + ) -> tuple[int, dict[str, str], bytes]: ... + + +@dataclass(frozen=True) +class NativeFrameContext: + """Correlation owned by Meshnet around one opaque native frame.""" + + request_id: str + node_id: str + route_session_id: str + route_epoch: int + work_id: str = "" + deadline_unix_nanos: int = 0 + + def __post_init__(self) -> None: + if not self.request_id or not self.node_id or not self.route_session_id: + raise ValueError("request, node, and Route Session identities are required") + if self.route_epoch < 0 or self.deadline_unix_nanos < 0: + raise ValueError("route epoch and deadline must be non-negative") + + def headers(self) -> dict[str, str]: + """Headers retained by the existing relay/Tracker accounting path.""" + return { + "Content-Type": NATIVE_FRAME_CONTENT_TYPE, + "X-Meshnet-Native-Frame": "shard-runtime/v1", + "X-Meshnet-Request-Id": self.request_id, + "X-Meshnet-Node-Id": self.node_id, + "X-Meshnet-Session": self.route_session_id, + "X-Meshnet-Route-Epoch": str(self.route_epoch), + "X-Meshnet-Work-Id": self.work_id, + "X-Meshnet-Deadline-Unix-Nanos": str(self.deadline_unix_nanos), + # The relay request id is restored on reply and is intentionally + # distinct from the caller/billing request id above. + "X-Meshnet-Activation-Id": self.request_id, + } + + +@dataclass(frozen=True) +class NativeSeamTelemetry: + transport: str + request_id: str + node_id: str + work_id: str + request_bytes: int + response_bytes: int + elapsed_seconds: float + + +TelemetrySink = Callable[[NativeSeamTelemetry], None] + + +def _request_identity(request: pb.SessionRequest) -> tuple[str, int, str, int]: + kind = request.WhichOneof("kind") + if kind == "open": + return request.open.route_session_id, request.open.route_epoch, "", 0 + if kind == "chunk": + item = request.chunk.envelope + return item.route_session_id, item.route_epoch, item.work_id, item.deadline_unix_nanos + if kind == "decode": + # DecodeStep relies on the already opened Route Session, while work + # identity/deadline are carried on every decode frame. + return "", 0, request.decode.work_id, request.decode.deadline_unix_nanos + if kind in {"cancel", "release"}: + item = getattr(request, kind) + return item.route_session_id, item.route_epoch, item.work_id, 0 + if kind == "flow_control": + return "", 0, "", 0 + raise NativeActivationSeamError("native SessionRequest has no frame kind") + + +def _validate_request(request: pb.SessionRequest, context: NativeFrameContext) -> None: + if request.ByteSize() == 0: + raise NativeActivationSeamError("empty native SessionRequest is not a versioned frame") + route_session, epoch, work_id, deadline = _request_identity(request) + if route_session and route_session != context.route_session_id: + raise NativeActivationSeamError("native frame Route Session differs from seam context") + if route_session and epoch != context.route_epoch: + raise NativeActivationSeamError("native frame route epoch differs from seam context") + if context.work_id and work_id and work_id != context.work_id: + raise NativeActivationSeamError("native frame work identity differs from seam context") + if context.deadline_unix_nanos and deadline and deadline != context.deadline_unix_nanos: + raise NativeActivationSeamError("native frame deadline differs from seam context") + + +def _response_work_id(response: pb.SessionResponse) -> str: + kind = response.WhichOneof("kind") + if kind == "chunk": + return response.chunk.envelope.work_id + if kind == "ack": + return response.ack.work_id + if kind == "status": + return response.status.work_id + return "" + + +class NativeActivationSeam: + """One Route-Session-to-worker seam with bounded direct buffering. + + ``direct_stub`` is the generated ``ShardRuntimeStub`` and is selected when + it is available. ``relay_request`` has the exact signature of the + existing persistent relay client; no relay server or bridge API changes + are needed. Relay calls are intentionally not retried: a failed send may + already have mutated downstream Hot KV state. + """ + + def __init__( + self, + context: NativeFrameContext, + *, + direct_stub=None, + relay_request: RelayRequest | None = None, + max_buffered_frames: int = 8, + telemetry: TelemetrySink | None = None, + ) -> None: + if (direct_stub is None) == (relay_request is None): + raise ValueError("provide exactly one of direct_stub or relay_request") + if max_buffered_frames < 1: + raise ValueError("max_buffered_frames must be positive") + self.context = context + self._direct_stub = direct_stub + self._relay_request = relay_request + self._telemetry = telemetry + self._closed = False + self._failure: BaseException | None = None + self._responses: Queue[pb.SessionResponse | BaseException] = Queue(maxsize=max_buffered_frames) + self._requests: Queue[pb.SessionRequest | object] | None = None + self._thread: threading.Thread | None = None + self._stop = object() + if direct_stub is not None: + self._requests = Queue(maxsize=max_buffered_frames) + self._thread = threading.Thread(target=self._run_direct, daemon=True, name="native-activation-grpc") + self._thread.start() + + @property + def transport(self) -> str: + return "direct-grpc" if self._direct_stub is not None else "relay" + + def _direct_requests(self) -> Iterator[pb.SessionRequest]: + assert self._requests is not None + while True: + item = self._requests.get() + if item is self._stop: + return + assert isinstance(item, pb.SessionRequest) + yield item + + def _run_direct(self) -> None: + try: + assert self._direct_stub is not None + for response in self._direct_stub.Session(self._direct_requests()): + self._put_response(response) + except BaseException as exc: + self._failure = exc + self._put_response(exc) + + def _put_response(self, value: pb.SessionResponse | BaseException) -> None: + # A worker may finish while a caller is abandoning the session. Do not + # let an unconsumed response turn into an unbounded producer queue. + try: + self._responses.put(value, timeout=0.1) + except Full: + self._failure = NativeActivationBufferFull("native response buffer is full") + + def send(self, request: pb.SessionRequest) -> pb.SessionResponse | None: + """Send one already-versioned protobuf frame without rewriting it.""" + if self._closed: + raise NativeActivationDisconnected("native activation seam is closed") + if self._failure is not None: + raise NativeActivationDisconnected("native activation stream failed") from self._failure + _validate_request(request, self.context) + frame = request.SerializeToString() + if self._direct_stub is not None: + assert self._requests is not None + try: + self._requests.put_nowait(request) + except Full as exc: + raise NativeActivationBufferFull("native direct request buffer is full") from exc + return None + + assert self._relay_request is not None + started = time.monotonic() + try: + status, _, response_frame = self._relay_request(NATIVE_RELAY_PATH, frame, self.context.headers()) + except Exception as exc: + self._closed = True + raise NativeActivationDisconnected("relay outcome is uncertain; refusing replay") from exc + if status != 200: + self._closed = True + raise NativeActivationDisconnected(f"relay native frame returned HTTP {status}") + response = pb.SessionResponse() + try: + response.ParseFromString(response_frame) + except Exception as exc: + self._closed = True + raise NativeActivationSeamError("relay returned a malformed native response frame") from exc + self._validate_response(response) + self._record(len(frame), len(response_frame), started) + return response + + def receive(self, timeout: float | None = None) -> pb.SessionResponse: + """Receive the next response from the one long-lived direct stream.""" + if self._direct_stub is None: + raise NativeActivationSeamError("relay sends return their response synchronously") + try: + value = self._responses.get(timeout=timeout) + except Empty as exc: + raise TimeoutError("timed out waiting for native direct response") from exc + if isinstance(value, BaseException): + raise NativeActivationDisconnected("native direct stream disconnected") from value + self._validate_response(value) + # gRPC owns its framing, but this records the actual protobuf payload + # size at the seam for the same telemetry shape as relay. + self._record(0, len(value.SerializeToString()), time.monotonic()) + return value + + def cancel(self, reason: str = "cancelled") -> pb.SessionResponse | None: + """Propagate cancellation through the same path and correlation fields.""" + return self.send(pb.SessionRequest(cancel=pb.CancelSignal( + route_session_id=self.context.route_session_id, + route_epoch=self.context.route_epoch, + work_id=self.context.work_id, + reason=reason, + ))) + + def _validate_response(self, response: pb.SessionResponse) -> None: + work_id = _response_work_id(response) + if self.context.work_id and work_id and work_id != self.context.work_id: + raise NativeActivationSeamError("native response work identity differs from seam context") + + def _record(self, request_bytes: int, response_bytes: int, started: float) -> None: + if self._telemetry is not None: + self._telemetry(NativeSeamTelemetry( + transport=self.transport, request_id=self.context.request_id, + node_id=self.context.node_id, work_id=self.context.work_id, + request_bytes=request_bytes, response_bytes=response_bytes, + elapsed_seconds=max(0.0, time.monotonic() - started), + )) + + def close(self) -> None: + if self._closed: + return + self._closed = True + if self._requests is not None: + try: + self._requests.put_nowait(self._stop) + except Full: + # The bounded queue is intentionally never expanded during + # shutdown; the worker will observe process/session teardown. + pass + if self._thread is not None: + self._thread.join(timeout=1.0) diff --git a/tests/test_native_activation_seam.py b/tests/test_native_activation_seam.py new file mode 100644 index 0000000..1a9ef82 --- /dev/null +++ b/tests/test_native_activation_seam.py @@ -0,0 +1,159 @@ +"""DGR-042 seam tests with a deterministic fake generated worker.""" + +from __future__ import annotations + +from collections.abc import Iterator +import threading + +import pytest + +from meshnet_node.native_activation_seam import ( + NATIVE_RELAY_PATH, + NativeActivationBufferFull, + NativeActivationDisconnected, + NativeActivationSeam, + NativeFrameContext, +) +from meshnet_node.native_protocol import pb + + +def _context(**changes: object) -> NativeFrameContext: + values: dict[str, object] = dict( + request_id="billing-request-7", node_id="node-tail", route_session_id="route-9", + route_epoch=4, work_id="work-3", deadline_unix_nanos=987654321, + ) + values.update(changes) + return NativeFrameContext(**values) + + +def _open() -> pb.SessionRequest: + return pb.SessionRequest(open=pb.SessionOpen( + schema_version=pb.SCHEMA_VERSION_1, route_session_id="route-9", route_epoch=4, + )) + + +def _chunk() -> pb.SessionRequest: + return pb.SessionRequest(chunk=pb.ActivationChunk(envelope=pb.Envelope( + schema_version=pb.SCHEMA_VERSION_1, route_session_id="route-9", route_epoch=4, + work_id="work-3", deadline_unix_nanos=987654321, + ))) + + +def _ack(request: pb.SessionRequest) -> pb.SessionResponse: + if request.WhichOneof("kind") == "open": + return pb.SessionResponse(accepted=pb.SessionAccepted( + schema_version=pb.SCHEMA_VERSION_1, route_session_id="route-9", route_epoch=4, + )) + route, epoch, work, _ = ("route-9", 4, "work-3", 0) + del route, epoch + return pb.SessionResponse(ack=pb.Ack(work_id=work, idempotency_step=1)) + + +class _FakeGrpcWorker: + def __init__(self, *, block: bool = False) -> None: + self.calls = 0 + self.received: list[pb.SessionRequest] = [] + self.started = threading.Event() + self.consumed = threading.Event() + self.release = threading.Event() + self.block = block + + def Session(self, requests: Iterator[pb.SessionRequest]): + self.calls += 1 + self.started.set() + for request in requests: + self.received.append(request) + self.consumed.set() + if self.block: + self.release.wait(1) + yield _ack(request) + + +def test_direct_uses_one_long_lived_grpc_stream_and_preserves_correlation(): + worker = _FakeGrpcWorker() + telemetry = [] + seam = NativeActivationSeam(_context(), direct_stub=worker, telemetry=telemetry.append) + try: + seam.send(_open()) + seam.send(_chunk()) + assert seam.receive(1).WhichOneof("kind") == "accepted" + assert seam.receive(1).ack.work_id == "work-3" + assert worker.calls == 1 + assert [frame.SerializeToString() for frame in worker.received] == [ + _open().SerializeToString(), _chunk().SerializeToString() + ] + assert telemetry[-1].request_id == "billing-request-7" + assert telemetry[-1].node_id == "node-tail" + finally: + seam.close() + + +def test_relay_carries_byte_identical_protobuf_frames_and_all_correlation_headers(): + captured: list[tuple[str, bytes, dict[str, str]]] = [] + + def relay(path: str, body: bytes, headers: dict[str, str]): + captured.append((path, body, headers)) + request = pb.SessionRequest() + request.ParseFromString(body) + return 200, {}, _ack(request).SerializeToString() + + seam = NativeActivationSeam(_context(), relay_request=relay) + response = seam.send(_chunk()) + assert response is not None and response.ack.work_id == "work-3" + path, body, headers = captured[0] + assert path == NATIVE_RELAY_PATH + assert body == _chunk().SerializeToString() + assert headers == { + "Content-Type": "application/x-protobuf", "X-Meshnet-Native-Frame": "shard-runtime/v1", + "X-Meshnet-Request-Id": "billing-request-7", "X-Meshnet-Node-Id": "node-tail", + "X-Meshnet-Session": "route-9", "X-Meshnet-Route-Epoch": "4", + "X-Meshnet-Work-Id": "work-3", "X-Meshnet-Deadline-Unix-Nanos": "987654321", + "X-Meshnet-Activation-Id": "billing-request-7", + } + + +def test_relay_disconnect_is_uncertain_and_is_never_replayed(): + calls = 0 + + def disconnected(*_): + nonlocal calls + calls += 1 + raise OSError("relay vanished") + + seam = NativeActivationSeam(_context(), relay_request=disconnected) + with pytest.raises(NativeActivationDisconnected, match="uncertain"): + seam.send(_chunk()) + with pytest.raises(NativeActivationDisconnected): + seam.send(_chunk()) + assert calls == 1 + + +def test_cancellation_uses_the_same_opaque_relay_contract(): + received = [] + + def relay(_path, body, _headers): + request = pb.SessionRequest() + request.ParseFromString(body) + received.append(request) + return 200, {}, pb.SessionResponse(ack=pb.Ack(work_id="work-3")).SerializeToString() + + seam = NativeActivationSeam(_context(), relay_request=relay) + response = seam.cancel("client disconnected") + assert response is not None + assert received[0].cancel.work_id == "work-3" + assert received[0].cancel.reason == "client disconnected" + + +def test_direct_request_buffer_is_bounded(): + worker = _FakeGrpcWorker(block=True) + seam = NativeActivationSeam(_context(), direct_stub=worker, max_buffered_frames=1) + try: + assert worker.started.wait(1) + seam.send(_open()) + assert worker.consumed.wait(1) + seam.send(_chunk()) + with pytest.raises(NativeActivationBufferFull): + seam.send(_chunk()) + finally: + worker.release.set() + seam.close()