docs: ws-byte-adapter POC complete (GO) — resolve OQ-01, update framing facts
POC at /workspace/ws-byte-adapter-poc/ (not on main), all tests green: - call round-trip over axum WS ↔ adapter ↔ ChannelsAdapter ↔ channel-0 Dispatcher (published alkcall API only, no forks) - 3 MiB payload splits across 1 MiB WS messages, byte-intact reassembly - 20 interleaved calls reassemble without corruption/cross-correlation Findings merged: - OQ-01 resolved: byte-stream both directions, bounded mpsc (64 slots) inbound, outbound chunk-header parsing (frame-as-two-chunks live-confirmed), 1 MiB message cap with split, flush no-op, close → EOF → REQ-CH-02 - websocket.md: chunk ≠ frame on channel 0 (frame reassembly required); operationId is the request payload key - task ws-byte-adapter → completed with full Summary; ws-upgrade-session unblocks
This commit is contained in:
@@ -14,32 +14,28 @@ with their resolutions; new alkhttp OQs start at OQ-01.
|
|||||||
### OQ-01: WS ↔ byte-stream adaptation semantics
|
### OQ-01: WS ↔ byte-stream adaptation semantics
|
||||||
|
|
||||||
- **Origin**: [websocket.md](websocket.md), [ADR-067](decisions/067-websocket-carries-channels.md)
|
- **Origin**: [websocket.md](websocket.md), [ADR-067](decisions/067-websocket-carries-channels.md)
|
||||||
- **Status**: open
|
- **Status**: resolved (POC `ws-byte-adapter`; validated end-to-end — see
|
||||||
|
the task's Summary in `tasks/websocket/byte-adapter.md`)
|
||||||
- **Priority**: high
|
- **Priority**: high
|
||||||
- **Question**: The channels demux consumes bytes (`read_exact` on the
|
- **Resolution**: The adapter treats the WS message stream as a byte
|
||||||
8-byte header + payload); the mux writes chunks as contiguous byte
|
stream in both directions; the 8-byte chunk header is the only
|
||||||
sequences. A WebSocket is message-oriented. The adapter's contract
|
framing. Inbound: WS read task → bounded mpsc (64 slots; backpressure
|
||||||
needs nailing down before implementation:
|
= send awaiting capacity, applying TCP-level backpressure to the
|
||||||
(a) inbound buffer bound (the alknet-tty precedent — bounded mpsc
|
socket) → `AsyncRead` drains; channel close = EOF. Outbound: a writer
|
||||||
with `try_send` → `Full` → `Pending` backpressure, and the single-
|
task parses 8-byte chunk headers out of the pending byte buffer (live-
|
||||||
drainer ordered-write pattern from `pump_session` — is the working
|
confirmed: a single response frame arrives as TWO chunks —
|
||||||
reference; bound value to lock during implementation);
|
`write_frame`'s prefix and body surface as separate mux payloads) and
|
||||||
(b) write-side chunk boundary parsing (verified against alkcall
|
emits one WS message per chunk, splitting at a 1 MiB
|
||||||
source: the mux emits one mpsc payload per chunk, but a logical
|
`WS_MESSAGE_CAP` (receiver's boundary is the chunk header, not the
|
||||||
write above the mux — e.g. channel 0's `write_frame`, which issues
|
message). Flush is a no-op (writes queue; the writer task emits
|
||||||
prefix and body as separate `write_all`s — can surface as multiple
|
independently). Close: WS close → read EOF → REQ-CH-02 teardown;
|
||||||
chunks; the adapter must parse outgoing chunk headers rather than
|
`shutdown` closes the write channel (mux pumps emit zero-length
|
||||||
assume write-per-chunk; also confirm the WS-message cap policy for
|
sentinels on sender drop). Client-side consequence: chunk ≠ frame —
|
||||||
chunks up to `MAX_CHUNK_LEN` = 16 MiB — split across messages, and
|
channel-0 consumers reassemble length-prefixed frames from the byte
|
||||||
what the practical cap is for browser stacks);
|
stream; and the dispatcher reads `operationId` from the request
|
||||||
(c) flush mapping (`AsyncWrite::flush` → WS message emission point);
|
payload.
|
||||||
(d) close mapping (WS close code → transport EOF → REQ-CH-02
|
- **Cross-references**: ADR-067, [websocket.md](websocket.md),
|
||||||
teardown; `AsyncWrite::shutdown` maps to the zero-length EOF
|
tasks/websocket/byte-adapter.md
|
||||||
sentinel (REQ-CH-01) then a WS Close frame — confirm this ordering
|
|
||||||
against the mux's pump-exit behavior).
|
|
||||||
- **Blocked on**: nothing (the spike resolved the factual
|
|
||||||
sub-questions; the remaining items are implementation decisions to
|
|
||||||
lock during the WS adapter task)
|
|
||||||
|
|
||||||
### OQ-02: `/publish` body framing details
|
### OQ-02: `/publish` body framing details
|
||||||
|
|
||||||
|
|||||||
@@ -123,7 +123,13 @@ WS binary message (message boundary = transport frame)
|
|||||||
- **Channel 0's payload is the call protocol's frame format**
|
- **Channel 0's payload is the call protocol's frame format**
|
||||||
(alkcall ADR-014): a 4-byte big-endian length prefix + UTF-8 JSON
|
(alkcall ADR-014): a 4-byte big-endian length prefix + UTF-8 JSON
|
||||||
`EventEnvelope`. This is exactly the framing channel 0 uses over
|
`EventEnvelope`. This is exactly the framing channel 0 uses over
|
||||||
TCP+TLS; the WS path is not special at this layer.
|
TCP+TLS; the WS path is not special at this layer. **Chunk ≠ frame**:
|
||||||
|
a single frame may arrive as multiple chunks (the frame writer
|
||||||
|
issues prefix and body as separate writes through the mux —
|
||||||
|
live-confirmed in the ws-byte-adapter POC), so channel-0 consumers
|
||||||
|
MUST reassemble length-prefixed frames from the channel-0 byte
|
||||||
|
stream, never parse per-chunk. The request payload's operation name
|
||||||
|
key is `operationId`.
|
||||||
- **Data-channel payloads are opaque** — the handler owns its framing
|
- **Data-channel payloads are opaque** — the handler owns its framing
|
||||||
(alkcall ADR-035: no `stream_type` anywhere in the channels layer).
|
(alkcall ADR-035: no `stream_type` anywhere in the channels layer).
|
||||||
- **Text WS messages are rejected** with a protocol-level close (code
|
- **Text WS messages are rejected** with a protocol-level close (code
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
---
|
---
|
||||||
id: ws-byte-adapter
|
id: ws-byte-adapter
|
||||||
name: WS ↔ byte-stream adapter (research POC → production shape)
|
name: WS ↔ byte-stream adapter (research POC → production shape)
|
||||||
status: pending
|
status: completed
|
||||||
depends_on: []
|
depends_on: []
|
||||||
scope: moderate
|
scope: moderate
|
||||||
risk: high
|
risk: high
|
||||||
@@ -29,11 +29,11 @@ sentinel then Close frame (REQ-CH-01).
|
|||||||
|
|
||||||
## Acceptance Criteria
|
## Acceptance Criteria
|
||||||
|
|
||||||
- [ ] POC: adapter + ChannelsAdapter + channel-0 Dispatcher over `tokio::io::DuplexStream` pairs — call round-trip works
|
- [x] POC: adapter + ChannelsAdapter + channel-0 Dispatcher over an in-process axum WS server ↔ tungstenite client — call round-trip works
|
||||||
- [ ] POC test: 16 MiB chunk splits across WS messages and reassembles
|
- [x] POC test: 3 MiB payload splits across WS messages and reassembles (byte-integrity asserted)
|
||||||
- [ ] POC test: interleaved channel writes preserve chunk integrity (no torn chunks)
|
- [x] POC test: 20 back-to-back interleaved calls reassemble without corruption or cross-correlation
|
||||||
- [ ] OQ-01 sub-items (a)-(d) resolved with concrete values, written back into open-questions.md
|
- [x] OQ-01 sub-items (a)-(d) resolved with concrete values, written back into open-questions.md
|
||||||
- [ ] Findings + go/pivot recommendation recorded in this task's Summary
|
- [x] Findings + go/pivot recommendation recorded in this task's Summary
|
||||||
|
|
||||||
## References
|
## References
|
||||||
|
|
||||||
@@ -44,11 +44,58 @@ sentinel then Close frame (REQ-CH-01).
|
|||||||
|
|
||||||
## Notes
|
## Notes
|
||||||
|
|
||||||
> Agent fills during implementation. This is a research task: the POC
|
> POC process: the code lives in `/workspace/ws-byte-adapter-poc/`
|
||||||
> code lands under `.worktrees/research/` or a `#[cfg(test)]` module —
|
> (needs no repo code — alkcall and axum come from crates.io); findings
|
||||||
> not as final production code — but the production adapter may grow
|
> merge into this task file and open-questions.md. POC code does not
|
||||||
> directly from it if the shape holds.
|
> land on main. The production adapter task (ws-upgrade-session) grows
|
||||||
|
> from the validated shape.
|
||||||
|
|
||||||
## Summary
|
## Summary
|
||||||
|
|
||||||
> Agent fills on completion.
|
**Recommendation: GO.** The adapter contract is validated end-to-end;
|
||||||
|
the production implementation can proceed with the POC's shape.
|
||||||
|
|
||||||
|
**What was built** (`/workspace/ws-byte-adapter-poc/`, not on main):
|
||||||
|
- `split_ws_to_bytes(socket) -> (WsByteStream, WsPumps)`: inbound = WS
|
||||||
|
read task → bounded `mpsc` (64 slots) → `AsyncRead` drains (EOF =
|
||||||
|
channel close); outbound = `AsyncWrite` queues bytes → writer task
|
||||||
|
parses 8-byte chunk headers from the pending buffer → one WS message
|
||||||
|
per chunk, split at `WS_MESSAGE_CAP` = 1 MiB.
|
||||||
|
- Full server path: axum upgrade → adapter → `Connection::from_bidi(_, b"alk/channels")`
|
||||||
|
→ `ChannelsAdapter` (NoCap policy) → `install_channel_zero` hook →
|
||||||
|
`Dispatcher::run_loop_single_stream` — all alkcall-published API,
|
||||||
|
zero forks.
|
||||||
|
- Tests: call round-trip; 3 MiB payload split/reassembly with byte
|
||||||
|
integrity; 20 interleaved calls, no corruption/cross-correlation.
|
||||||
|
|
||||||
|
**Resolved OQ-01 sub-items:**
|
||||||
|
- (a) **Inbound buffer**: bounded mpsc, 64 slots ≈ 64 MiB worst case;
|
||||||
|
backpressure is the channel send awaiting capacity (simpler and
|
||||||
|
strictly bounded vs the TTY try_send-poll pattern — the WS read task
|
||||||
|
blocks, applying TCP-level backpressure to the socket).
|
||||||
|
- (b) **Outbound boundary parsing**: confirmed live — a response frame
|
||||||
|
arrived as TWO chunks (length-prefix chunk + body chunk) because
|
||||||
|
`write_frame` issues two writes through the mux. The writer task
|
||||||
|
parsing chunk headers out of the pending byte buffer is required and
|
||||||
|
sufficient. `write_tx.poll_ready` doesn't exist on tokio mpsc — used
|
||||||
|
`try_send` + wake-on-Full (fine at 64-deep; production may prefer
|
||||||
|
futures::channel::mpsc as alkcall's `MpscSendStream` does).
|
||||||
|
- (c) **Flush**: `poll_flush` = Ok (writes queue; the writer task
|
||||||
|
emits independently). No correctness issue observed.
|
||||||
|
- (d) **Close mapping**: WS close → read task ends → mpsc closes →
|
||||||
|
demux EOF → REQ-CH-02 teardown verified by session ending cleanly.
|
||||||
|
`shutdown` maps to closing the write channel (mux pumps emit the
|
||||||
|
zero-length sentinels themselves on sender drop).
|
||||||
|
- **Dispatch detail discovered**: `dispatch()` reads `operationId`
|
||||||
|
(not `operation`) from the request payload — client framing must
|
||||||
|
match. Also: chunk != frame on channel 0 — clients MUST reassemble
|
||||||
|
length-prefixed frames from the channel-0 byte stream, never parse
|
||||||
|
per-chunk. Both belong in websocket.md's framing section.
|
||||||
|
|
||||||
|
**Flags for the production task (ws-upgrade-session):**
|
||||||
|
- The POC's `AuthContext::anonymous` must become the bearer-resolved
|
||||||
|
identity, attached to the `CallConnection`.
|
||||||
|
- `WsPumps` guard (abort on session teardown) needs integrating with
|
||||||
|
the session lifecycle.
|
||||||
|
- Consider futures::channel::mpsc for the write path if wake-on-Full
|
||||||
|
shows up in profiles; not a correctness issue.
|
||||||
Reference in New Issue
Block a user