From d5b6ca70c8f220ae1733e88df22328342b16ec62 Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Wed, 25 Mar 2026 14:30:04 +0100 Subject: [PATCH 1/7] 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 --- pkg/api/api.go | 11 +++ pkg/api/chunk_stream.go | 184 +++++++++++++++++++++++++++++++++------- 2 files changed, 166 insertions(+), 29 deletions(-) diff --git a/pkg/api/api.go b/pkg/api/api.go index 2b713bc76d4..93595168cc3 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -845,7 +845,18 @@ func (s *Service) newStampedPutter(ctx context.Context, opts putterOptions, stam return nil, errInvalidPostageBatch } + return s.newStampedPutterWithBatch(ctx, opts, stamp, storedBatch) +} + +// newStampedPutterWithBatch creates a stamped putter using a pre-fetched batch. +// This avoids the database lookup when batch info is already cached. +func (s *Service) newStampedPutterWithBatch(ctx context.Context, opts putterOptions, stamp *postage.Stamp, storedBatch *postage.Batch) (storer.PutterSession, error) { + if !opts.Deferred && s.beeMode == DevMode { + return nil, errUnsupportedDevNodeOperation + } + var session storer.PutterSession + var err error if opts.Deferred || opts.Pin { session, err = s.storer.Upload(ctx, opts.Pin, opts.TagID) if err != nil { diff --git a/pkg/api/chunk_stream.go b/pkg/api/chunk_stream.go index 2f91939f21a..03399c4f6aa 100644 --- a/pkg/api/chunk_stream.go +++ b/pkg/api/chunk_stream.go @@ -8,6 +8,7 @@ import ( "context" "errors" "net/http" + "strconv" "time" "github.com/ethersphere/bee/v2/pkg/cac" @@ -28,7 +29,7 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques logger := s.logger.WithName("chunks_stream").Build() headers := struct { - BatchID []byte `map:"Swarm-Postage-Batch-Id" validate:"required"` + BatchID []byte `map:"Swarm-Postage-Batch-Id"` // Optional: omit if caller provides pre-signed stamps per chunk SwarmTag uint64 `map:"Swarm-Tag"` }{} if response := s.mapStructure(r.Header, &headers); response != nil { @@ -36,6 +37,15 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques return } + // Fallback: read tag from query parameter (browser WebSocket can't set headers) + if headers.SwarmTag == 0 { + if qTag := r.URL.Query().Get("swarm-tag"); qTag != "" { + if parsed, err := strconv.ParseUint(qTag, 10, 64); err == nil { + headers.SwarmTag = parsed + } + } + } + var ( tag uint64 err error @@ -55,29 +65,36 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques } } - // if tag not specified use direct upload - // Using context.Background here because the putter's lifetime extends beyond that of the HTTP request. - putter, err := s.newStamperPutter(context.Background(), putterOptions{ - BatchID: headers.BatchID, - TagID: tag, - Deferred: tag != 0, - }) - if err != nil { - logger.Debug("get putter failed", "error", err) - logger.Error(nil, "get putter failed") - switch { - case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable): - jsonhttp.UnprocessableEntity(w, "batch not usable yet or does not exist") - case errors.Is(err, postage.ErrNotFound): - jsonhttp.NotFound(w, "batch with id not found") - case errors.Is(err, errInvalidPostageBatch): - jsonhttp.BadRequest(w, "invalid batch id") - case errors.Is(err, errUnsupportedDevNodeOperation): - jsonhttp.BadRequest(w, errUnsupportedDevNodeOperation) - default: - jsonhttp.BadRequest(w, nil) + // Create connection-level putter only if BatchID is provided. + // If BatchID is not provided, the API caller is expected to provide + // pre-signed stamps with each chunk (and is also expected to keep + // track of stamp state over time). + var putter storer.PutterSession + if len(headers.BatchID) > 0 { + // if tag not specified use direct upload + // Using context.Background here because the putter's lifetime extends beyond that of the HTTP request. + putter, err = s.newStamperPutter(context.Background(), putterOptions{ + BatchID: headers.BatchID, + TagID: tag, + Deferred: tag != 0, + }) + if err != nil { + logger.Debug("get putter failed", "error", err) + logger.Error(nil, "get putter failed") + switch { + case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable): + jsonhttp.UnprocessableEntity(w, "batch not usable yet or does not exist") + case errors.Is(err, postage.ErrNotFound): + jsonhttp.NotFound(w, "batch with id not found") + case errors.Is(err, errInvalidPostageBatch): + jsonhttp.BadRequest(w, "invalid batch id") + case errors.Is(err, errUnsupportedDevNodeOperation): + jsonhttp.BadRequest(w, errUnsupportedDevNodeOperation) + default: + jsonhttp.BadRequest(w, nil) + } + return } - return } upgrader := websocket.Upgrader{ @@ -95,13 +112,46 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques } s.wsWg.Add(1) - go s.handleUploadStream(logger, wsConn, putter) + var decode chunkDecoder + if len(headers.BatchID) > 0 { + decode = decodeChunkWithoutStamp + } else { + decode = decodeChunkWithStamp + } + go s.handleUploadStream(logger, wsConn, putter, tag, decode) +} + +// chunkDecoder extracts chunk data and optionally a stamp from a websocket message. +// When BatchID is provided in headers, decodeChunkWithoutStamp is used (no stamp in message). +// When BatchID is not provided, decodeChunkWithStamp is used (stamp prepended to chunk data). +type chunkDecoder func(msg []byte) (chunkData []byte, stamp *postage.Stamp, err error) + +// decodeChunkWithoutStamp returns the message as-is (used when BatchID provided in headers). +func decodeChunkWithoutStamp(msg []byte) ([]byte, *postage.Stamp, error) { + return msg, nil, nil +} + +// decodeChunkWithStamp extracts a stamp from the first 113 bytes of the message. +// Returns an error if the message is too small or the stamp is invalid. +func decodeChunkWithStamp(msg []byte) ([]byte, *postage.Stamp, error) { + if len(msg) < postage.StampSize+swarm.SpanSize { + return nil, nil, errors.New("message too small for stamp + chunk") + } + + stamp := &postage.Stamp{} + if err := stamp.UnmarshalBinary(msg[:postage.StampSize]); err != nil { + return nil, nil, errors.New("invalid stamp") + } + + return msg[postage.StampSize:], stamp, nil } func (s *Service) handleUploadStream( logger log.Logger, conn *websocket.Conn, putter storer.PutterSession, + tag uint64, + decode chunkDecoder, ) { defer s.wsWg.Done() @@ -111,11 +161,23 @@ func (s *Service) handleUploadStream( gone = make(chan struct{}) err error ) + + // Cache for batch validation to avoid database lookups for every chunk + // Key: batch ID hex string, Value: stored batch info + // This avoids the expensive batchStore.Get() call for each chunk + batchCache := make(map[string]*postage.Batch) + defer func() { cancel() _ = conn.Close() - if err = putter.Done(swarm.ZeroAddress); err != nil { - logger.Error(err, "chunk upload stream: syncing chunks failed") + + // No cleanup needed for batch cache - it's just metadata + + // Only call Done on connection-level putter if it exists + if putter != nil { + if err = putter.Done(swarm.ZeroAddress); err != nil { + logger.Error(err, "chunk upload stream: syncing chunks failed") + } } }() @@ -190,14 +252,71 @@ func (s *Service) handleUploadStream( return } - chunk, err := cac.NewWithDataSpan(msg) + // Decode the message using the appropriate decoder + chunkData, stamp, err := decode(msg) if err != nil { - logger.Debug("chunk upload stream: create chunk failed", "error", err) + logger.Debug("chunk upload stream: decode failed", "error", err) + logger.Error(nil, "chunk upload stream: "+err.Error()) + sendErrorClose(websocket.CloseInternalServerErr, err.Error()) + return + } + + // Determine the putter to use + var ( + chunk swarm.Chunk + chunkPutter = putter + ) + + // If stamp was extracted, create a per-chunk putter + if stamp != nil { + batchID := stamp.BatchID() + batchIDHex := string(batchID) + + storedBatch, exists := batchCache[batchIDHex] + if !exists { + storedBatch, err = s.batchStore.Get(batchID) + if err != nil { + logger.Debug("chunk upload stream: batch validation failed", "error", err) + logger.Error(nil, "chunk upload stream: batch validation failed") + if errors.Is(err, storage.ErrNotFound) { + sendErrorClose(websocket.CloseInternalServerErr, "batch not found") + } else { + sendErrorClose(websocket.CloseInternalServerErr, "batch validation failed") + } + return + } + batchCache[batchIDHex] = storedBatch + } + + chunkPutter, err = s.newStampedPutterWithBatch(ctx, putterOptions{ + BatchID: batchID, + TagID: tag, + Deferred: tag != 0, + }, stamp, storedBatch) + if err != nil { + logger.Debug("chunk upload stream: failed to create stamped putter", "error", err) + logger.Error(nil, "chunk upload stream: failed to create stamped putter") + switch { + case errors.Is(err, errBatchUnusable) || errors.Is(err, postage.ErrNotUsable): + sendErrorClose(websocket.CloseInternalServerErr, "batch not usable") + case errors.Is(err, postage.ErrNotFound): + sendErrorClose(websocket.CloseInternalServerErr, "batch not found") + default: + sendErrorClose(websocket.CloseInternalServerErr, "stamped putter creation failed") + } + return + } + } + + chunk, err = cac.NewWithDataSpan(chunkData) + if err != nil { + logger.Debug("chunk upload stream: create chunk failed", "error", err, "chunk_size", len(chunkData)) logger.Error(nil, "chunk upload stream: create chunk failed") + sendErrorClose(websocket.CloseInternalServerErr, "invalid chunk data") return } - err = putter.Put(ctx, chunk) + err = chunkPutter.Put(ctx, chunk) if err != nil { logger.Debug("chunk upload stream: write chunk failed", "address", chunk.Address(), "error", err) logger.Error(nil, "chunk upload stream: write chunk failed") @@ -210,6 +329,13 @@ func (s *Service) handleUploadStream( return } + // Clean up per-chunk putter + if chunkPutter != putter { + if err := chunkPutter.Done(swarm.ZeroAddress); err != nil { + logger.Error(err, "chunk upload stream: failed to finalize per-chunk putter") + } + } + err = sendMsg(websocket.BinaryMessage, successWsMsg) if err != nil { s.logger.Debug("chunk upload stream: sending success message failed", "error", err) From 0430d51236cbebd81ce1beab558f5030877801c3 Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Tue, 31 Mar 2026 10:56:06 +0200 Subject: [PATCH 2/7] fix: openapi --- openapi/Swarm.yaml | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/openapi/Swarm.yaml b/openapi/Swarm.yaml index 4168e63f80e..6bf0e89fc3c 100644 --- a/openapi/Swarm.yaml +++ b/openapi/Swarm.yaml @@ -319,7 +319,18 @@ paths: - Chunk parameters: - $ref: "SwarmCommon.yaml#/components/parameters/SwarmTagParameter" - - $ref: "SwarmCommon.yaml#/components/parameters/SwarmPostageBatchId" + - in: query + name: swarm-tag + schema: + $ref: "SwarmCommon.yaml#/components/schemas/Uid" + required: false + description: "Associate upload with an existing Tag UID (use when WebSocket client cannot set custom headers)" + - in: header + name: swarm-postage-batch-id + description: "ID of Postage Batch that is used to upload data with. Optional when chunks include pre-signed postage stamps." + required: false + schema: + $ref: "SwarmCommon.yaml#/components/schemas/SwarmAddress" responses: "200": description: "Connection established" From d3676b43066c453817fc608234c26359c2eece65 Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Wed, 1 Apr 2026 12:09:45 +0200 Subject: [PATCH 3/7] fix: error handling --- pkg/api/chunk_stream.go | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/pkg/api/chunk_stream.go b/pkg/api/chunk_stream.go index 03399c4f6aa..07bebaf9476 100644 --- a/pkg/api/chunk_stream.go +++ b/pkg/api/chunk_stream.go @@ -40,9 +40,13 @@ func (s *Service) chunkUploadStreamHandler(w http.ResponseWriter, r *http.Reques // Fallback: read tag from query parameter (browser WebSocket can't set headers) if headers.SwarmTag == 0 { if qTag := r.URL.Query().Get("swarm-tag"); qTag != "" { - if parsed, err := strconv.ParseUint(qTag, 10, 64); err == nil { - headers.SwarmTag = parsed + parsed, err := strconv.ParseUint(qTag, 10, 64) + if err != nil { + logger.Debug("invalid swarm-tag query parameter", "value", qTag, "error", err) + jsonhttp.BadRequest(w, "invalid swarm-tag query parameter") + return } + headers.SwarmTag = parsed } } From 168a40e0a5404c731ae5e95f3d2a0ba2d7a23cc7 Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Wed, 1 Apr 2026 13:57:44 +0200 Subject: [PATCH 4/7] fix: add tests --- pkg/api/chunk_stream.go | 6 +- pkg/api/chunk_stream_test.go | 143 +++++++++++++++++++++++++++++++++++ 2 files changed, 146 insertions(+), 3 deletions(-) diff --git a/pkg/api/chunk_stream.go b/pkg/api/chunk_stream.go index 07bebaf9476..739044bc1f0 100644 --- a/pkg/api/chunk_stream.go +++ b/pkg/api/chunk_stream.go @@ -274,9 +274,9 @@ func (s *Service) handleUploadStream( // If stamp was extracted, create a per-chunk putter if stamp != nil { batchID := stamp.BatchID() - batchIDHex := string(batchID) + batchIDKey := string(batchID) - storedBatch, exists := batchCache[batchIDHex] + storedBatch, exists := batchCache[batchIDKey] if !exists { storedBatch, err = s.batchStore.Get(batchID) if err != nil { @@ -289,7 +289,7 @@ func (s *Service) handleUploadStream( } return } - batchCache[batchIDHex] = storedBatch + batchCache[batchIDKey] = storedBatch } chunkPutter, err = s.newStampedPutterWithBatch(ctx, putterOptions{ diff --git a/pkg/api/chunk_stream_test.go b/pkg/api/chunk_stream_test.go index 5ba62cdd76c..2982343d5c5 100644 --- a/pkg/api/chunk_stream_test.go +++ b/pkg/api/chunk_stream_test.go @@ -11,7 +11,11 @@ import ( "time" "github.com/ethersphere/bee/v2/pkg/api" + "github.com/ethersphere/bee/v2/pkg/crypto" + "github.com/ethersphere/bee/v2/pkg/postage" + mockbatchstore "github.com/ethersphere/bee/v2/pkg/postage/batchstore/mock" mockpost "github.com/ethersphere/bee/v2/pkg/postage/mock" + testingpostage "github.com/ethersphere/bee/v2/pkg/postage/testing" "github.com/ethersphere/bee/v2/pkg/spinlock" testingc "github.com/ethersphere/bee/v2/pkg/storage/testing" mockstorer "github.com/ethersphere/bee/v2/pkg/storer/mock" @@ -104,3 +108,142 @@ func TestChunkUploadStream(t *testing.T) { } }) } + +// nolint:paralleltest +func TestChunkUploadStreamWithStamp(t *testing.T) { + // Generate signer and batch for pre-signed stamps + key, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + signer := crypto.NewDefaultSigner(key) + owner, err := signer.EthereumAddress() + if err != nil { + t.Fatal(err) + } + + // Generate chunks and their pre-signed stamps + chunks := make([]swarm.Chunk, 5) + stampBytes := make([][]byte, 5) + + for i := range 5 { + chunks[i] = testingc.GenerateTestRandomChunk() + stamp := testingpostage.MustNewValidStamp(signer, chunks[i].Address()) + sb, err := stamp.MarshalBinary() + if err != nil { + t.Fatal(err) + } + stampBytes[i] = sb + } + + // Mock batch store: accept all batch IDs and return a batch with the correct owner + batchStore := mockbatchstore.New( + mockbatchstore.WithAcceptAllExistsFunc(), + mockbatchstore.WithBatch(&postage.Batch{ + Owner: owner.Bytes(), + }), + ) + + // No Swarm-Postage-Batch-Id header — triggers per-chunk stamp mode + wsHeaders := http.Header{} + wsHeaders.Set(api.ContentTypeHeader, "application/octet-stream") + + var ( + storerMock = mockstorer.New() + _, wsConn, _, chanStorer = newTestServer(t, testServerOptions{ + Storer: storerMock, + Post: mockpost.New(mockpost.WithAcceptAll()), + BatchStore: batchStore, + WsPath: "/chunks/stream", + WsHeaders: wsHeaders, + DirectUpload: true, + }) + ) + + t.Run("upload with pre-signed stamps", func(t *testing.T) { + for i := range 5 { + // Prepend stamp bytes to chunk data + msg := append(stampBytes[i], chunks[i].Data()...) + + err := wsConn.SetWriteDeadline(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + + err = wsConn.WriteMessage(websocket.BinaryMessage, msg) + if err != nil { + t.Fatal(err) + } + + err = wsConn.SetReadDeadline(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + + mt, msg, err := wsConn.ReadMessage() + if err != nil { + t.Fatal(err) + } + + if mt != websocket.BinaryMessage || !bytes.Equal(msg, api.SuccessWsMsg) { + t.Fatal("invalid response", mt, string(msg)) + } + } + + for _, c := range chunks { + err := spinlock.Wait(100*time.Millisecond, func() bool { return chanStorer.Has(c.Address()) }) + if err != nil { + t.Fatal(err) + } + } + }) +} + +// nolint:paralleltest +func TestChunkUploadStreamInvalidStamp(t *testing.T) { + // No Swarm-Postage-Batch-Id header — triggers per-chunk stamp mode + wsHeaders := http.Header{} + wsHeaders.Set(api.ContentTypeHeader, "application/octet-stream") + + var ( + storerMock = mockstorer.New() + _, wsConn, _, _ = newTestServer(t, testServerOptions{ + Storer: storerMock, + Post: mockpost.New(mockpost.WithAcceptAll()), + WsPath: "/chunks/stream", + WsHeaders: wsHeaders, + DirectUpload: true, + }) + ) + + t.Run("message too small for stamp", func(t *testing.T) { + // Send a message smaller than StampSize + SpanSize + tooSmall := make([]byte, postage.StampSize) + + err := wsConn.SetWriteDeadline(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + + err = wsConn.WriteMessage(websocket.BinaryMessage, tooSmall) + if err != nil { + t.Fatal(err) + } + + err = wsConn.SetReadDeadline(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + + _, _, err = wsConn.ReadMessage() + if err == nil { + t.Fatal("expected failure on read") + } + // nolint:errorlint + if cerr, ok := err.(*websocket.CloseError); !ok { + t.Fatal("invalid error on read") + } else if cerr.Text != "message too small for stamp + chunk" { + t.Fatalf("incorrect response on error, exp: (message too small for stamp + chunk) got (%s)", cerr.Text) + } + }) +} From 3b2bd9a9a5fedbb42ebb9bdfb0ffcab3017928cb Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Wed, 1 Apr 2026 14:15:56 +0200 Subject: [PATCH 5/7] fix: linter --- pkg/api/chunk_stream_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/api/chunk_stream_test.go b/pkg/api/chunk_stream_test.go index 2982343d5c5..359931f9d23 100644 --- a/pkg/api/chunk_stream_test.go +++ b/pkg/api/chunk_stream_test.go @@ -206,7 +206,7 @@ func TestChunkUploadStreamInvalidStamp(t *testing.T) { wsHeaders.Set(api.ContentTypeHeader, "application/octet-stream") var ( - storerMock = mockstorer.New() + storerMock = mockstorer.New() _, wsConn, _, _ = newTestServer(t, testServerOptions{ Storer: storerMock, Post: mockpost.New(mockpost.WithAcceptAll()), From 811edc04d6486fc0d90ad52f3d5b18ce7d018412 Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Thu, 2 Apr 2026 11:42:01 +0200 Subject: [PATCH 6/7] fix: chunkPutter cleanup --- pkg/api/chunk_stream.go | 6 ++++++ pkg/api/chunk_stream_test.go | 2 +- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/pkg/api/chunk_stream.go b/pkg/api/chunk_stream.go index 739044bc1f0..732392a66e6 100644 --- a/pkg/api/chunk_stream.go +++ b/pkg/api/chunk_stream.go @@ -316,6 +316,9 @@ func (s *Service) handleUploadStream( if err != nil { logger.Debug("chunk upload stream: create chunk failed", "error", err, "chunk_size", len(chunkData)) logger.Error(nil, "chunk upload stream: create chunk failed") + if chunkPutter != putter { + _ = chunkPutter.Cleanup() + } sendErrorClose(websocket.CloseInternalServerErr, "invalid chunk data") return } @@ -324,6 +327,9 @@ func (s *Service) handleUploadStream( if err != nil { logger.Debug("chunk upload stream: write chunk failed", "address", chunk.Address(), "error", err) logger.Error(nil, "chunk upload stream: write chunk failed") + if chunkPutter != putter { + _ = chunkPutter.Cleanup() + } switch { case errors.Is(err, postage.ErrBucketFull): sendErrorClose(websocket.CloseInternalServerErr, "batch is overissued") diff --git a/pkg/api/chunk_stream_test.go b/pkg/api/chunk_stream_test.go index 359931f9d23..add7165fa89 100644 --- a/pkg/api/chunk_stream_test.go +++ b/pkg/api/chunk_stream_test.go @@ -41,7 +41,7 @@ func TestChunkUploadStream(t *testing.T) { ) t.Run("upload and verify", func(t *testing.T) { - chsToGet := []swarm.Chunk{} + chsToGet := make([]swarm.Chunk, 0, 5) for range 5 { ch := testingc.GenerateTestRandomChunk() From c860af6c2862cefe6491e918ef996b882dec95fe Mon Sep 17 00:00:00 2001 From: Attila Gazso <230163+agazso@users.noreply.github.com> Date: Thu, 2 Apr 2026 14:10:21 +0200 Subject: [PATCH 7/7] fix: comment --- pkg/api/chunk_stream.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/api/chunk_stream.go b/pkg/api/chunk_stream.go index 732392a66e6..be648b964be 100644 --- a/pkg/api/chunk_stream.go +++ b/pkg/api/chunk_stream.go @@ -167,7 +167,7 @@ func (s *Service) handleUploadStream( ) // Cache for batch validation to avoid database lookups for every chunk - // Key: batch ID hex string, Value: stored batch info + // Key: batch ID, Value: stored batch info // This avoids the expensive batchStore.Get() call for each chunk batchCache := make(map[string]*postage.Batch)