feat: download mode for the chunk stream endpoint - #5618
martinconic wants to merge 5 commits into
Conversation
| s.metrics.ChunkStreamDeliveryCount.WithLabelValues("success").Inc() | ||
|
|
||
| chunkData := chunk.Data() | ||
| resp := make([]byte, 1+swarm.HashSize+len(chunkData)) |
There was a problem hiding this comment.
not sure about this custom serialization by hand thing... it is just specced out in the openapi spec and assumed that implementers should implement it by hand. it is fragile and breakable. why not use some sort of standard serialization format to both decode the request and encode the response?
aloknerurkar
left a comment
There was a problem hiding this comment.
Since we are making breaking changes and also building on the serialization thread — I'd suggest going further than reworking the download framing: use one connection for both directions, with request/response semantics and a request id.
Two things this fixes beyond tidiness:
-
Upload is round-trip bound. The upload loop is strictly sequential — read → put → ack → read, one chunk in flight. Each chunk costs a full RTT, so on a 50ms link that's ~20 chunks/sec regardless of bandwidth. Download already has a 16-worker pool; upload has none. stamper.Stamp already takes issuer.mtx (pkg/postage/stamper.go:43), so concurrent stamping is safe — the ceiling is the protocol, not the storage layer.
-
Neither direction is recoverable. Download replies carry no id, and the upload ack is successWsMsg = []byte{} — an empty frame with no address at all. Clients correlate purely by ordering. On any error both paths call sendErrorClose and drop the connection, so a client that had N requests outstanding cannot tell which completed. For bulk upload that means restarting from zero.
A shared envelope — [type][8-byte request-id][payload] for requests, [type][8-byte request-id][status][payload] for responses — would give:
- one read loop feeding one typed job channel, with a worker pool serving both Get and Put (removes the duplicated read/deadline/close-handler logic between the two handlers)
- pipelined uploads instead of one-at-a-time
- per-request error status instead of connection teardown, so a single bad chunk no longer kills the stream
- correlation for recovery after an error close
I'd keep this as raw binary rather than JSON-RPC or protobuf. JSON-RPC costs more on the wire with base64 and ~15x encode/decode — self-defeating for an endpoint that exists to cut per-chunk overhead. Protobuf is bee's p2p convention but has never appeared in pkg/api; requiring a schema compiler would be a new burden on bee-js. The binary envelope stays a four-line DataView parse in the browser with no dependency.
The upload stream endpoint was not written with care. I was going through the code and there is a lot of scope to cleanup. Multiple putters defined, deferred/direct upload semantics, stamped and unstamped chunks etc. I feel that this would make it much more usable for clients and also close to what elad mentioned about having grpc style chunk get/put API.
| } | ||
|
|
||
| // fetchAndSendChunk retrieves a single chunk and writes exactly one response | ||
| // frame for it. That one-frame-per-requested-address invariant is what lets a |
There was a problem hiding this comment.
Exactly one reply per requested address, and "a dropped frame is indistinguishable from a slow one" because replies carry no request id. But the protocol-error paths in the read loop (sendErrorClose + return on bad message type, bad length, unknown opcode, oversized batch) void every in-flight request silently.
A client that sends 256 addresses and then trips CloseMessageTooBig on its next frame has no way to determine which of the first 256 were answered. With no request id and no per-batch boundary marker in the reply stream, the only recovery is to discard everything and re-request.
The OpenAPI spec documents the close codes, but not that pending replies are lost when they fire. At minimum that consequence belongs in the spec. A framing with a request/batch id — which is largely what @acud point about a standard serialization format would give you for free — would make it recoverable instead.
Thank you for this, sounds like a good change. @acud do you also agree? |
|
@martinconic i generally encourage the initiative from @aloknerurkar, but i'm still not 100% convinced a custom encoding is right.
i think we can progress with this proposal, and make serialization a pluggable implementation. but in general i'd encourage to move away from hand written serialization at least on wire-formats (persistence is a different story). |
I am fine with protobuf. The only reason I didn't suggest it is that bee-js/javascript will be primary consumer if this is successful and I remember there were issues with javascript ecosystem when working with protobuf. I spoke to claude and I see that there is a new typescript compiler which makes things easier. So I am totally on board! |
Checklist
Description
Adds a request/response protocol to
GET /chunks/stream: a client opens one websocket and pipelines chunk downloads and uploads over it, instead of paying for an HTTP request per chunk.Closes #5417, closes #5599.
The first version of this PR used a hand-rolled
['D'][32-byte address]...framing. Following review (@acud, @aloknerurkar) it now uses protobuf with a request id per message, so each request gets its own status and a client can recover after an error. That earlier download mode never shipped and has been removed.Protocol
Negotiated with
Sec-WebSocket-Protocol: swarm-chunk-stream, or?mode=streamfor browser clients that cannot set headers. Messages are defined inpkg/api/pb/chunkstream.proto, one protobuf message per websocket binary message.Id, plus either aGetRequest(32-byteAddress, optionalCacheOption) or aPutRequest(Data, optionalStamp, requiredType).Id, aStatus, the chunkAddress,Datafor downloads, and a sanitisedError. Exactly one response per request; responses arrive in completion order, so clients match onId.OK,NOT_FOUND,ERROR,BAD_REQUEST, andBUSYwhen a queue is full (retry with backoff). Every enum has anUNSPECIFIED = 0value that is never sent.CHUNK_TYPE_CACorCHUNK_TYPE_SOC. SOCs are validated including their signature; an unspecified or unknown type is rejected withBAD_REQUEST.Swarm-Postage-Batch-Idon the connection stamps every chunk, or eachPutRequestcarries its own pre-signedStamp.Swarm-Tagmakes uploads deferred (stored locally, synced later) with either form of stamping.Swarm-Cache/?cache=set the connection default;CacheOptionoverrides it per request.What
OKmeans for an upload: for a direct upload, the chunk has been pushed to the network; for a tagged upload, it has been stored locally. It is never sent before that.Errors: a failing request gets its own status and does not close the connection. The connection is closed only for a message that cannot be decoded (
1003), a message over 64 KiB (1009), or node shutdown (1001). Responses received before a close are final; a client should resend any request whoseIdgot no response.The legacy upload stream (no subprotocol, or
?mode=upload) is unchanged and byte-for-byte identical to master.Notes for review
pb/convention. No new dependency.BUSYinstead of blocking the read loop.getter.DefaultFetchTimeout, an upload by 30s. The delivery write deadline is 5 minutes, so a client that pauses reading while its buffer drains doesn't get disconnected.topology.ErrNotFoundmaps toNOT_FOUND, matching howbzz.gotreats the same error fromstorer.Download.Testing
15 tests cover download and upload, interleaving, per-request errors, a failed push, fairness under blocked uploads, a full queue, per-chunk stamps, a feed-sized SOC, tags with per-chunk stamps, cache options, protocol violations, and shutdown — including with a client that has stopped reading. The fairness, queue-full and shutdown tests were checked against deliberately broken code to confirm they fail when they should.
make build,make lint,make vet,make testand the race detector all pass; the redocly problem count is unchanged from master.AI Disclosure