Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
e5a7b0a
feat: add SOC fields headers and wrapped chunk caching
nugaon Jun 8, 2026
7bbf91e
test: add SOC fields and wrapped chunk caching
nugaon Jun 8, 2026
9f63b76
docs: openapi with minor version bump
nugaon Jun 8, 2026
8b4d996
fix: linting
nugaon Jun 9, 2026
d38b6e8
refactor: warning log instead of debug
nugaon Jun 9, 2026
e411683
refactor: increase msg buffer to 2 and slow connection handling
nugaon Jun 9, 2026
5ef413a
docs: openapi random access desctiption
nugaon Jun 10, 2026
54df2e3
docs: default payload
nugaon Aug 13, 2026
a9921e3
fix: race issue on close
nugaon Aug 19, 2026
be9f2ac
fix: deduplication values in the header values
nugaon Aug 19, 2026
a504b5f
fix: context cancellation
nugaon Aug 19, 2026
35b3533
test: remaining
nugaon Aug 19, 2026
125e62e
fix: handle gsoc messages sequentially
nugaon Sep 2, 2026
d03e0e4
Merge branch 'master' into feat/gsoc-fine-grained
nugaon Sep 15, 2026
5e38e05
refactor: remove slow client handling
nugaon Sep 15, 2026
8c0ba37
refactor: remove slow consumer handling
nugaon Sep 15, 2026
9f03baa
refactor: implement message queuing for slow consumers
nugaon Sep 15, 2026
ccb847e
refactor: dedicated queue structure on api
nugaon Sep 15, 2026
525d2b4
test: improve synchronization in websocket
nugaon Sep 15, 2026
220bb02
refactor: implement bounded queue
nugaon Sep 22, 2026
d79f162
refactor: replace validSocFields slice
nugaon Sep 22, 2026
0105b92
refactor: implement GsocQueue with release functionality
nugaon Sep 22, 2026
943b885
refactor: test server initialization in API tests to include addition…
nugaon Sep 22, 2026
0451a8d
test: enhance gsoc pipe test for subscription handling and message qu…
nugaon Sep 22, 2026
351824c
test: enhance GSOC websocket tests for SOC fields and stalled consume…
nugaon Sep 22, 2026
1a57d6c
refactor: update test timeout settings and adjust batch depth in api …
nugaon Sep 22, 2026
043d794
test: optimize chunk building in GSOC websocket tests for performance
nugaon Sep 22, 2026
485b65c
refactor: improve websocket shutdown handling
nugaon Sep 22, 2026
6831a2c
test: add concurrent subscription handling
nugaon Sep 22, 2026
2ea522b
refactor: improve handler cleanup logic in subscription
nugaon Sep 22, 2026
154b4d2
refactor: enhance background context handling for websocket caching
nugaon Sep 23, 2026
f98a1b5
refactor: improve subscription handling
nugaon Sep 23, 2026
a6fa28f
refactor: add query parameter support
nugaon Sep 23, 2026
933b324
refactor: improve websocket message handling
nugaon Sep 23, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@ REACHABILITY_OVERRIDE_PUBLIC ?= false
BATCHFACTOR_OVERRIDE_PUBLIC ?= 5
BEE_IMAGE ?= ethersphere/bee:latest
PLATFORM ?= linux/amd64
# go test defaults to 10m per package binary, which the pkg/api suite exceeds
# under the race detector on the slower CI runners.
TEST_TIMEOUT ?= 30m

BEE_API_VERSION ?= "$(shell grep '^ version:' openapi/Swarm.yaml | awk '{print $$2}')"

Expand Down Expand Up @@ -99,9 +102,9 @@ check-whitespace:
.PHONY: test-race
test-race:
ifdef cover
$(GO) test -race -failfast -coverprofile=cover.out -v ./...
$(GO) test -race -failfast -timeout $(TEST_TIMEOUT) -coverprofile=cover.out -v ./...
else
$(GO) test -race -failfast -v ./...
$(GO) test -race -failfast -timeout $(TEST_TIMEOUT) -v ./...
endif

.PHONY: test-integration
Expand All @@ -127,9 +130,9 @@ endif
.PHONY: test-ci-race
test-ci-race:
ifdef cover
$(GO) test -race -coverprofile=cover.out ./...
$(GO) test -race -timeout $(TEST_TIMEOUT) -coverprofile=cover.out ./...
else
$(GO) test -race ./...
$(GO) test -race -timeout $(TEST_TIMEOUT) ./...
endif

.PHONY: build
Expand Down
22 changes: 18 additions & 4 deletions openapi/Swarm.yaml
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
openapi: 3.0.3

info:
version: 8.1.1
version: 8.2.0
title: Bee API
description: "API endpoints for interacting with the Swarm network, supporting file operations, messaging, and node management"

Expand Down Expand Up @@ -967,9 +967,23 @@ paths:
$ref: "SwarmCommon.yaml#/components/schemas/SwarmAddress"
required: true
description: "Single Owner Chunk address (which may have multiple payloads)"
responses:
"200":
description: Establishes a WebSocket subscription for incoming messages on the Single Owner Chunk address
- $ref: "SwarmCommon.yaml#/components/parameters/SwarmSocFieldsParameter"
- $ref: "SwarmCommon.yaml#/components/parameters/SwarmCacheWrappedChunkParameter"
- $ref: "SwarmCommon.yaml#/components/parameters/SwarmSocFieldsQueryParameter"
- $ref: "SwarmCommon.yaml#/components/parameters/SwarmCacheWrappedChunkQueryParameter"
responses:
"200":
description: >
Establishes a WebSocket subscription for incoming messages on the
Single Owner Chunk address. Each message is the binary serialization
of the Single Owner Chunk fields requested through the
swarm-soc-fields header or query parameter (defaults to the wrapped
chunk payload).
Pending messages are buffered per subscription up to a fixed limit;
a client that does not keep up with the incoming rate misses the
oldest of them, as they are discarded in favor of the newest ones.
"400":
$ref: "SwarmCommon.yaml#/components/responses/400"
"500":
$ref: "SwarmCommon.yaml#/components/responses/500"
default:
Expand Down
48 changes: 48 additions & 0 deletions openapi/SwarmCommon.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -1180,6 +1180,54 @@ components:
required: false
description: Associate upload with an existing Tag UID

SwarmSocFieldsParameter:
in: header
name: swarm-soc-fields
schema:
type: string
default: "payload"
required: false
description: >
Comma separated list of Single Owner Chunk fields to be serialized and
channeled on every incoming GSOC message, in the given order. Allowed
values are: address, recoveredPubKey, identifier, signature,
wrappedAddress, span, payload. When omitted it defaults to "payload".
In order to have random access on the response bytes define payload
as the last field in the list since it has variable length.

SwarmCacheWrappedChunkParameter:
in: header
name: swarm-cache-wrapped-chunk
schema:
type: boolean
required: false
description: >
Indicates whether the wrapped chunk of every incoming GSOC message should
be cached locally so that it can be resolved through the bytes endpoint
(useful when the single owner chunk wraps a root chunk larger than 4KB).

SwarmSocFieldsQueryParameter:
in: query
name: swarm-soc-fields
schema:
type: string
required: false
description: >
Same as the swarm-soc-fields header, for clients that cannot set request
headers (e.g. browser WebSocket). Takes precedence over the header when
both are given.

SwarmCacheWrappedChunkQueryParameter:
in: query
name: swarm-cache-wrapped-chunk
schema:
type: boolean
required: false
description: >
Same as the swarm-cache-wrapped-chunk header, for clients that cannot set
request headers (e.g. browser WebSocket). Takes precedence over the header
when both are given.

SwarmPinParameter:
in: header
name: swarm-pin
Expand Down
50 changes: 25 additions & 25 deletions pkg/api/accesscontrol_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,7 @@ func TestAccessLogicEachEndpointWithAct(t *testing.T) {
upTestOpts = append(upTestOpts, jsonhttptest.WithRequestHeader(api.SwarmCollectionHeader, "True"))
}
t.Run(v.name, func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -213,7 +213,7 @@ func TestAccessLogicWithoutAct(t *testing.T) {
)

t.Run("upload-w/-act-then-download-w/o-act", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -246,7 +246,7 @@ func TestAccessLogicWithoutAct(t *testing.T) {
})

t.Run("upload-w/o-act-then-download-w/-act", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -310,7 +310,7 @@ func TestAccessLogicInvalidPath(t *testing.T) {
)

t.Run("invalid-path-params", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -361,7 +361,7 @@ func TestAccessLogicHistory(t *testing.T) {
)

t.Run("empty-history-upload-then-download-and-check-data", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -397,7 +397,7 @@ func TestAccessLogicHistory(t *testing.T) {
})

t.Run("with-history-upload-then-download-and-check-data", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -443,7 +443,7 @@ func TestAccessLogicHistory(t *testing.T) {
})

t.Run("upload-then-download-wrong-history", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -479,7 +479,7 @@ func TestAccessLogicHistory(t *testing.T) {
})

t.Run("upload-wrong-history", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -502,7 +502,7 @@ func TestAccessLogicHistory(t *testing.T) {
})

t.Run("download-w/o-history", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -538,7 +538,7 @@ func TestAccessLogicTimestamp(t *testing.T) {
fileName = "sample.html"
)
t.Run("upload-then-download-with-timestamp-and-check-data", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -586,7 +586,7 @@ func TestAccessLogicTimestamp(t *testing.T) {

t.Run("download-w/o-timestamp", func(t *testing.T) {
encryptedRef := "a5df670544eaea29e61b19d8739faa4573b19e4426e58a173e51ed0b5e7e2ade"
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -601,7 +601,7 @@ func TestAccessLogicTimestamp(t *testing.T) {
)
})
t.Run("download-w/-invalid-timestamp", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -649,7 +649,7 @@ func TestAccessLogicPublisher(t *testing.T) {
)

t.Run("upload-then-download-w/-publisher-and-check-data", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -695,7 +695,7 @@ func TestAccessLogicPublisher(t *testing.T) {
})

t.Run("upload-then-download-invalid-publickey", func(t *testing.T) {
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down Expand Up @@ -753,7 +753,7 @@ func TestAccessLogicPublisher(t *testing.T) {
downloader = "03c712a7e29bc792ac8d8ae49793d28d5bda27ed70f0d90697b2fb456c0a168bd2"
encryptedRef = "a5df670544eaea29e61b19d8739faa4573b19e4426e58a173e51ed0b5e7e2ade"
)
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -778,7 +778,7 @@ func TestAccessLogicPublisher(t *testing.T) {
downloader = "03c712a7e29bc792ac8d8ae49793d28d5bda27ed70f0d90697b2fb456c0a168bd2"
testfile = "testfile1"
)
downloaderClient, _, _, _ := newTestServer(t, testServerOptions{
downloaderClient, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -800,7 +800,7 @@ func TestAccessLogicPublisher(t *testing.T) {

t.Run("download-w/o-publisher", func(t *testing.T) {
encryptedRef := "a5df670544eaea29e61b19d8739faa4573b19e4426e58a173e51ed0b5e7e2ade"
client, _, _, _ := newTestServer(t, testServerOptions{
client, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -820,13 +820,13 @@ func TestAccessLogicPublisher(t *testing.T) {
func TestAccessLogicGrantees(t *testing.T) {
t.Parallel()
var (
spk, _ = hex.DecodeString("a786dd84b61485de12146fd9c4c02d87e8fd95f0542765cb7fc3d2e428c0bcfa")
pk, _ = crypto.DecodeSecp256k1PrivateKey(spk)
storerMock = mockstorer.New()
h, fixtureHref = prepareHistoryFixture(storerMock)
logger = log.Noop
addr = swarm.RandAddress(t)
client, _, _, _ = newTestServer(t, testServerOptions{
spk, _ = hex.DecodeString("a786dd84b61485de12146fd9c4c02d87e8fd95f0542765cb7fc3d2e428c0bcfa")
pk, _ = crypto.DecodeSecp256k1PrivateKey(spk)
storerMock = mockstorer.New()
h, fixtureHref = prepareHistoryFixture(storerMock)
logger = log.Noop
addr = swarm.RandAddress(t)
client, _, _, _, _ = newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand All @@ -839,7 +839,7 @@ func TestAccessLogicGrantees(t *testing.T) {
publicKeyBytes = crypto.EncodeSecp256k1PublicKey(&pk.PublicKey)
publisher = hex.EncodeToString(publicKeyBytes)
)
clientwihtpublisher, _, _, _ := newTestServer(t, testServerOptions{
clientwihtpublisher, _, _, _, _ := newTestServer(t, testServerOptions{
Storer: storerMock,
Logger: logger,
Post: mockpost.New(mockpost.WithAcceptAll()),
Expand Down
4 changes: 2 additions & 2 deletions pkg/api/accounting_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ func TestAccountingInfo(t *testing.T) {
return ret, nil
}

testServer, _, _, _ := newTestServer(t, testServerOptions{
testServer, _, _, _, _ := newTestServer(t, testServerOptions{
AccountingOpts: []mock.Option{mock.WithPeerAccountingFunc(accountingFunc)},
})

Expand Down Expand Up @@ -109,7 +109,7 @@ func TestAccountingInfoError(t *testing.T) {
accountingFunc := func() (map[string]accounting.PeerInfo, error) {
return nil, wantErr
}
testServer, _, _, _ := newTestServer(t, testServerOptions{
testServer, _, _, _, _ := newTestServer(t, testServerOptions{
AccountingOpts: []mock.Option{mock.WithPeerAccountingFunc(accountingFunc)},
})

Expand Down
17 changes: 17 additions & 0 deletions pkg/api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,8 @@ const (
SwarmActTimestampHeader = "Swarm-Act-Timestamp"
SwarmActPublisherHeader = "Swarm-Act-Publisher"
SwarmActHistoryAddressHeader = "Swarm-Act-History-Address"
SwarmSocFieldsHeader = "Swarm-Soc-Fields"
SwarmCacheWrappedChunkHeader = "Swarm-Cache-Wrapped-Chunk"

ImmutableHeader = "Immutable"
GasPriceHeader = "Gas-Price"
Expand Down Expand Up @@ -180,6 +182,17 @@ type Service struct {
wsWg sync.WaitGroup // wait for all websockets to close on exit
quit chan struct{}

// bgCtx is the context form of quit: it is canceled by Close and bounds
// node-local work that a handler starts and that has to outlive the
// request or connection which triggered it.
bgCtx context.Context
bgCancel context.CancelFunc

// gsocCacheSubs holds the caching subscriptions shared by the GSOC
// websocket subscribers, keyed by GSOC address, see gsoc.go.
gsocCacheMu sync.Mutex
gsocCacheSubs map[string]*gsocCacheSub

overlay *swarm.Address
publicKey ecdsa.PublicKey
pssPublicKey ecdsa.PublicKey
Expand Down Expand Up @@ -340,6 +353,8 @@ func (s *Service) Configure(signer crypto.Signer, tracer *tracing.Tracer, o Opti
s.metrics = newMetrics()

s.quit = make(chan struct{})
s.bgCtx, s.bgCancel = context.WithCancel(context.Background())
s.gsocCacheSubs = make(map[string]*gsocCacheSub)

s.storer = e.Storer
s.resolver = e.Resolver
Expand Down Expand Up @@ -400,6 +415,7 @@ func (s *Service) SetIsWarmingUp(v bool) {
func (s *Service) Close() error {
s.logger.Info("api shutting down")
close(s.quit)
s.bgCancel()

done := make(chan struct{})
go func() {
Expand Down Expand Up @@ -607,6 +623,7 @@ func (s *Service) corsHandler(h http.Handler) http.Handler {
SwarmRedundancyStrategyHeader, SwarmRedundancyFallbackModeHeader, SwarmChunkRetrievalTimeoutHeader, SwarmLookAheadBufferSizeHeader,
SwarmFeedIndexHeader, SwarmFeedIndexNextHeader, SwarmSocSignatureHeader, SwarmOnlyRootChunk, GasPriceHeader, GasLimitHeader, ImmutableHeader,
SwarmActHeader, SwarmActTimestampHeader, SwarmActPublisherHeader, SwarmActHistoryAddressHeader,
SwarmSocFieldsHeader, SwarmCacheWrappedChunkHeader,
}
allowedHeadersStr := strings.Join(allowedHeaders, ", ")

Expand Down
Loading
Loading