Read one length-delimited protobuf body off ``reader``. Reads the varint length one byte at a time (high bit = continue), then the exact body. ``readexactly`` has no buffer-size cap, so there is no 64 KB readline trap (report 13 §1.4 — this is *why* proto framing replaced NDJSON). R
(reader: asyncio.StreamReader)
| 60 | |
| 61 | |
| 62 | def _fail_stream(queue: asyncio.Queue[object], error: Exception) -> None: |
| 63 | """Discard buffered frames and deliver one terminal error without blocking.""" |
| 64 | while True: |
| 65 | try: |
| 66 | queue.get_nowait() |
| 67 | except asyncio.QueueEmpty: |
| 68 | break |
| 69 | queue.put_nowait(error) |
| 70 | queue.put_nowait(_STREAM_END) |
| 71 | |
| 72 | |
| 73 | class Transport: |
| 74 | def __init__(self, config: ClientConfig, binary_path: str) -> None: |
| 75 | self._config = config |
| 76 | self._binary_path = binary_path |
| 77 | self._process: asyncio.subprocess.Process | None = None |
| 78 | self._reader_task: asyncio.Task[None] | None = None |
| 79 | self._stderr_task: asyncio.Task[None] | None = None |
| 80 | self._pending: dict[int, _UnaryPending | _StreamPending] = {} |
no outgoing calls