Skip to content

Commit d5b6ca7

Browse files
committed
feat: per-chunk stamped putter for WebSocket chunk uploads
- Support per-chunk postage stamps in WebSocket chunk stream endpoint - When BatchID header is omitted, each chunk message must include a 113-byte stamp prefix - Add chunkDecoder abstraction (decodeChunkWithStamp / decodeChunkWithoutStamp) - Extract newStampedPutterWithBatch to avoid repeated batch DB lookups - Cache batch metadata per WebSocket connection for performance - Add swarm-tag query parameter fallback for browser WebSocket compatibility
1 parent fec3ecd commit d5b6ca7

2 files changed

Lines changed: 166 additions & 29 deletions

File tree

pkg/api/api.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -845,7 +845,18 @@ func (s *Service) newStampedPutter(ctx context.Context, opts putterOptions, stam
845845
return nil, errInvalidPostageBatch
846846
}
847847

848+
return s.newStampedPutterWithBatch(ctx, opts, stamp, storedBatch)
849+
}
850+
851+
// newStampedPutterWithBatch creates a stamped putter using a pre-fetched batch.
852+
// This avoids the database lookup when batch info is already cached.
853+
func (s *Service) newStampedPutterWithBatch(ctx context.Context, opts putterOptions, stamp *postage.Stamp, storedBatch *postage.Batch) (storer.PutterSession, error) {
854+
if !opts.Deferred && s.beeMode == DevMode {
855+
return nil, errUnsupportedDevNodeOperation
856+
}
857+
848858
var session storer.PutterSession
859+
var err error
849860
if opts.Deferred || opts.Pin {
850861
session, err = s.storer.Upload(ctx, opts.Pin, opts.TagID)
851862
if err != nil {

pkg/api/chunk_stream.go

Lines changed: 155 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import (
88
"context"
99
"errors"
1010
"net/http"
11+
"strconv"
1112
"time"
1213

1314
"github.com/ethersphere/bee/v2/pkg/cac"
@@ -28,14 +29,23 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques
2829
logger := s.logger.WithName("chunks_stream").Build()
2930

3031
headers := struct {
31-
BatchID []byte `map:"Swarm-Postage-Batch-Id" validate:"required"`
32+
BatchID []byte `map:"Swarm-Postage-Batch-Id"` // Optional: omit if caller provides pre-signed stamps per chunk
3233
SwarmTag uint64 `map:"Swarm-Tag"`
3334
}{}
3435
if response := s.mapStructure(r.Header, &headers); response != nil {
3536
response("invalid header params", logger, w)
3637
return
3738
}
3839

40+
// Fallback: read tag from query parameter (browser WebSocket can't set headers)
41+
if headers.SwarmTag == 0 {
42+
if qTag := r.URL.Query().Get("swarm-tag"); qTag != "" {
43+
if parsed, err := strconv.ParseUint(qTag, 10, 64); err == nil {
44+
headers.SwarmTag = parsed
45+
}
46+
}
47+
}
48+
3949
var (
4050
tag uint64
4151
err error
@@ -55,29 +65,36 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques
5565
}
5666
}
5767

58-
// if tag not specified use direct upload
59-
// Using context.Background here because the putter's lifetime extends beyond that of the HTTP request.
60-
putter, err := s.newStamperPutter(context.Background(), putterOptions{
61-
BatchID: headers.BatchID,
62-
TagID: tag,
63-
Deferred: tag != 0,
64-
})
65-
if err != nil {
66-
logger.Debug("get putter failed", "error", err)
67-
logger.Error(nil, "get putter failed")
68-
switch {
69-
case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable):
70-
jsonhttp.UnprocessableEntity(w, "batch not usable yet or does not exist")
71-
case errors.Is(err, postage.ErrNotFound):
72-
jsonhttp.NotFound(w, "batch with id not found")
73-
case errors.Is(err, errInvalidPostageBatch):
74-
jsonhttp.BadRequest(w, "invalid batch id")
75-
case errors.Is(err, errUnsupportedDevNodeOperation):
76-
jsonhttp.BadRequest(w, errUnsupportedDevNodeOperation)
77-
default:
78-
jsonhttp.BadRequest(w, nil)
68+
// Create connection-level putter only if BatchID is provided.
69+
// If BatchID is not provided, the API caller is expected to provide
70+
// pre-signed stamps with each chunk (and is also expected to keep
71+
// track of stamp state over time).
72+
var putter storer.PutterSession
73+
if len(headers.BatchID) > 0 {
74+
// if tag not specified use direct upload
75+
// Using context.Background here because the putter's lifetime extends beyond that of the HTTP request.
76+
putter, err = s.newStamperPutter(context.Background(), putterOptions{
77+
BatchID: headers.BatchID,
78+
TagID: tag,
79+
Deferred: tag != 0,
80+
})
81+
if err != nil {
82+
logger.Debug("get putter failed", "error", err)
83+
logger.Error(nil, "get putter failed")
84+
switch {
85+
case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable):
86+
jsonhttp.UnprocessableEntity(w, "batch not usable yet or does not exist")
87+
case errors.Is(err, postage.ErrNotFound):
88+
jsonhttp.NotFound(w, "batch with id not found")
89+
case errors.Is(err, errInvalidPostageBatch):
90+
jsonhttp.BadRequest(w, "invalid batch id")
91+
case errors.Is(err, errUnsupportedDevNodeOperation):
92+
jsonhttp.BadRequest(w, errUnsupportedDevNodeOperation)
93+
default:
94+
jsonhttp.BadRequest(w, nil)
95+
}
96+
return
7997
}
80-
return
8198
}
8299

83100
upgrader := websocket.Upgrader{
@@ -95,13 +112,46 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques
95112
}
96113

97114
s.wsWg.Add(1)
98-
go s.handleUploadStream(logger, wsConn, putter)
115+
var decode chunkDecoder
116+
if len(headers.BatchID) > 0 {
117+
decode = decodeChunkWithoutStamp
118+
} else {
119+
decode = decodeChunkWithStamp
120+
}
121+
go s.handleUploadStream(logger, wsConn, putter, tag, decode)
122+
}
123+
124+
// chunkDecoder extracts chunk data and optionally a stamp from a websocket message.
125+
// When BatchID is provided in headers, decodeChunkWithoutStamp is used (no stamp in message).
126+
// When BatchID is not provided, decodeChunkWithStamp is used (stamp prepended to chunk data).
127+
type chunkDecoder func(msg []byte) (chunkData []byte, stamp *postage.Stamp, err error)
128+
129+
// decodeChunkWithoutStamp returns the message as-is (used when BatchID provided in headers).
130+
func decodeChunkWithoutStamp(msg []byte) ([]byte, *postage.Stamp, error) {
131+
return msg, nil, nil
132+
}
133+
134+
// decodeChunkWithStamp extracts a stamp from the first 113 bytes of the message.
135+
// Returns an error if the message is too small or the stamp is invalid.
136+
func decodeChunkWithStamp(msg []byte) ([]byte, *postage.Stamp, error) {
137+
if len(msg) < postage.StampSize+swarm.SpanSize {
138+
return nil, nil, errors.New("message too small for stamp + chunk")
139+
}
140+
141+
stamp := &postage.Stamp{}
142+
if err := stamp.UnmarshalBinary(msg[:postage.StampSize]); err != nil {
143+
return nil, nil, errors.New("invalid stamp")
144+
}
145+
146+
return msg[postage.StampSize:], stamp, nil
99147
}
100148

101149
func (s *Service) handleUploadStream(
102150
logger log.Logger,
103151
conn *websocket.Conn,
104152
putter storer.PutterSession,
153+
tag uint64,
154+
decode chunkDecoder,
105155
) {
106156
defer s.wsWg.Done()
107157

@@ -111,11 +161,23 @@ func (s *Service) handleUploadStream(
111161
gone = make(chan struct{})
112162
err error
113163
)
164+
165+
// Cache for batch validation to avoid database lookups for every chunk
166+
// Key: batch ID hex string, Value: stored batch info
167+
// This avoids the expensive batchStore.Get() call for each chunk
168+
batchCache := make(map[string]*postage.Batch)
169+
114170
defer func() {
115171
cancel()
116172
_ = conn.Close()
117-
if err = putter.Done(swarm.ZeroAddress); err != nil {
118-
logger.Error(err, "chunk upload stream: syncing chunks failed")
173+
174+
// No cleanup needed for batch cache - it's just metadata
175+
176+
// Only call Done on connection-level putter if it exists
177+
if putter != nil {
178+
if err = putter.Done(swarm.ZeroAddress); err != nil {
179+
logger.Error(err, "chunk upload stream: syncing chunks failed")
180+
}
119181
}
120182
}()
121183

@@ -190,14 +252,71 @@ func (s *Service) handleUploadStream(
190252
return
191253
}
192254

193-
chunk, err := cac.NewWithDataSpan(msg)
255+
// Decode the message using the appropriate decoder
256+
chunkData, stamp, err := decode(msg)
194257
if err != nil {
195-
logger.Debug("chunk upload stream: create chunk failed", "error", err)
258+
logger.Debug("chunk upload stream: decode failed", "error", err)
259+
logger.Error(nil, "chunk upload stream: "+err.Error())
260+
sendErrorClose(websocket.CloseInternalServerErr, err.Error())
261+
return
262+
}
263+
264+
// Determine the putter to use
265+
var (
266+
chunk swarm.Chunk
267+
chunkPutter = putter
268+
)
269+
270+
// If stamp was extracted, create a per-chunk putter
271+
if stamp != nil {
272+
batchID := stamp.BatchID()
273+
batchIDHex := string(batchID)
274+
275+
storedBatch, exists := batchCache[batchIDHex]
276+
if !exists {
277+
storedBatch, err = s.batchStore.Get(batchID)
278+
if err != nil {
279+
logger.Debug("chunk upload stream: batch validation failed", "error", err)
280+
logger.Error(nil, "chunk upload stream: batch validation failed")
281+
if errors.Is(err, storage.ErrNotFound) {
282+
sendErrorClose(websocket.CloseInternalServerErr, "batch not found")
283+
} else {
284+
sendErrorClose(websocket.CloseInternalServerErr, "batch validation failed")
285+
}
286+
return
287+
}
288+
batchCache[batchIDHex] = storedBatch
289+
}
290+
291+
chunkPutter, err = s.newStampedPutterWithBatch(ctx, putterOptions{
292+
BatchID: batchID,
293+
TagID: tag,
294+
Deferred: tag != 0,
295+
}, stamp, storedBatch)
296+
if err != nil {
297+
logger.Debug("chunk upload stream: failed to create stamped putter", "error", err)
298+
logger.Error(nil, "chunk upload stream: failed to create stamped putter")
299+
switch {
300+
case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable):
301+
sendErrorClose(websocket.CloseInternalServerErr, "batch not usable")
302+
case errors.Is(err, postage.ErrNotFound):
303+
sendErrorClose(websocket.CloseInternalServerErr, "batch not found")
304+
default:
305+
sendErrorClose(websocket.CloseInternalServerErr, "stamped putter creation failed")
306+
}
307+
return
308+
}
309+
}
310+
311+
chunk, err = cac.NewWithDataSpan(chunkData)
312+
if err != nil {
313+
logger.Debug("chunk upload stream: create chunk failed", "error", err, "chunk_size", len(chunkData))
196314
logger.Error(nil, "chunk upload stream: create chunk failed")
315+
sendErrorClose(websocket.CloseInternalServerErr, "invalid chunk data")
197316
return
198317
}
199318

200-
err = putter.Put(ctx, chunk)
319+
err = chunkPutter.Put(ctx, chunk)
201320
if err != nil {
202321
logger.Debug("chunk upload stream: write chunk failed", "address", chunk.Address(), "error", err)
203322
logger.Error(nil, "chunk upload stream: write chunk failed")
@@ -210,6 +329,13 @@ func (s *Service) handleUploadStream(
210329
return
211330
}
212331

332+
// Clean up per-chunk putter
333+
if chunkPutter != putter {
334+
if err := chunkPutter.Done(swarm.ZeroAddress); err != nil {
335+
logger.Error(err, "chunk upload stream: failed to finalize per-chunk putter")
336+
}
337+
}
338+
213339
err = sendMsg(websocket.BinaryMessage, successWsMsg)
214340
if err != nil {
215341
s.logger.Debug("chunk upload stream: sending success message failed", "error", err)

0 commit comments

Comments
 (0)