Skip to content

feat: gsoc grained api - #5614

Open
nugaon wants to merge 30 commits into
masterfrom
feat/gsoc-grained-api
Open

nugaon wants to merge 30 commits into
masterfrom
feat/gsoc-grained-api

Conversation

@nugaon

@nugaon nugaon commented Sep 15, 2026

Copy link
Copy Markdown
Member

merging #5497 and #5593 PRs into one.

I took out the slow client handling part because the incoming updates frequency are unpredictable.
The feature is still notifying the user in case of piled up messages with warning log.

after the 1st master merge you can see the new code.


The following changes were made by Claude Sonnet 5 while chasing down the CI failures on this branch:

  • pkg/api/gsoc.go: replaced the fixed-size buffered channel + non-blocking default-drop ("slow consumer") logic with an unbounded FIFO queue (gsocQueue). The old default branch could silently drop a GSOC update on any goroutine-scheduling jitter (not just a genuinely slow client) and, after an earlier refactor, could also permanently hang Handle() — called synchronously by push/pull-sync — since nothing ever unblocked it for a client that never sends a close frame. The queue guarantees in-order delivery with no cap and no dropped messages, while a truly dead connection still times out via the existing per-message write deadline.
  • pkg/api/gsoc.go: moved Subscribe() to run synchronously in gsocWsHandler, before the connection is handed off to its own goroutine, closing a window where a GSOC update could arrive before the handler was registered and be silently missed.
  • pkg/api/gsoc_test.go: updated TestGsocWebsocketSlowConsumer to assert the new unbounded-queue behavior (the full backlog is delivered, in order, once the consumer catches up, instead of being disconnected), and added a synchronization fix specific to this test's in-memory net.Pipe transport — Dial() returning is not a happens-before guarantee that the server has reached Subscribe(), so the test now waits before sending and confirms with a throwaway round-trip before building the unread backlog it asserts on.
  • pkg/api/gsoc_test.go: removed the now-unnecessary fixed sleep between updates in TestGsocWebsocketMessageOrdering — ordering is guaranteed by the fix above rather than by pacing.

🤖 Generated with Claude Code

Add `Swarm-Soc-Fields` header to allow clients to request specific SOC fields (address, recoveredpubkey, identifier, signature, wrappedaddress, span, payload) in GSOC WebSocket messages. Add `Swarm-Cache-Wrapped-Chunk` header to enable caching of wrapped chunks on the node. Update GSOC handler to pass full SOC object instead of just payload, enabling access to all chunk properties. Adjust WebSocket buffer sizes to accommodate maximum SOC fields message size.
Add tests for `Swarm-Soc-Fields` header to verify requesting specific SOC fields (identifier, wrappedAddress, payload) and full wrapped chunk data (span + payload). Add test for `Swarm-Cache-Wrapped-Chunk` header to verify wrapped chunks are cached and retrievable. Update test helpers to support custom headers and return storer instance. Update gsoc handler signature to accept full SOC object instead of payload bytes.
gsoc.Handle spawned a goroutine per subscriber handler, so message
delivery order to a subscriber was not guaranteed. Call handlers
synchronously in registration order instead.
@nugaon
nugaon force-pushed the feat/gsoc-grained-api branch from f8f40ac to 8c0ba37 Compare September 15, 2026 08:31
@nugaon
nugaon marked this pull request as ready for review September 15, 2026 14:11
Comment thread pkg/api/gsoc.go
Comment thread pkg/api/gsoc.go Outdated
Comment thread pkg/api/gsoc.go Outdated
Comment thread pkg/api/api_test.go Outdated
Comment thread pkg/api/gsoc_test.go Outdated
Comment thread pkg/api/gsoc_test.go
Comment thread pkg/api/gsoc.go
Comment thread pkg/api/gsoc.go
Comment thread pkg/gsoc/gsoc.go

@martinconic martinconic left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changes look good, checked them manually also, but discovered with AI some things, please take a look

Comment thread pkg/api/gsoc.go
// Caching is a node-local side effect independent of this
// subscriber's connection, so it must not be aborted just
// because the websocket closes mid-write.
if err := s.storer.Cache().Put(context.Background(), c.WrappedChunk()); err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This callback runs inline on the pushsync (pushsync.go:255) and pullsync (pullsync.go:365) stream goroutines since 125e62e. So Cache().Put — a Sharky write plus a LevelDB index update — now runs synchronously inside the protocol handler, in pushsync's case before the chunk is stored and before the receipt goes back to the waiting peer. Note the CAC branch at pushsync.go:246 deliberately uses safe.Go for exactly this reason. Two further issues: with N subscribers on one address this runs N times sequentially for the same chunk, and context.Background() bypasses both the stream timeout and node shutdown.

Could we carry the chunk on the queue instead and do the Put on the writer goroutine, with a context derived from s.quit? That does change semantics — caching would stop when the connection closes, and a chunk could be dropped along with an evicted message — so the comment on lines 186-188 needs updating either way, since it currently argues the opposite.

Comment thread pkg/api/gsoc.go
}

headers := struct {
SocFields string `map:"Swarm-Soc-Fields"`

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Right now, a browser dApp cannot send Swarm-Soc-Fields or Swarm-Cache-Wrapped-Chunk?
Should we add query parameter fallback so ws://node/gsoc/subscribe/{addr}?swarm-soc-fields=payload,signature works from any browse ?

Comment thread pkg/api/gsoc.go
s.logger.Debug("gsoc ws: set write deadline failed", "error", err)
return
case <-wake:
for {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This loop only exits when the queue empties, which never happens under sustained backlog — and the yield below covers s.quit and gone but not ticker.C. Measured: 2 pings in 2s at a 100ms period, ~20 expected. Clients then drop the connection on their own pong timeout.

Simplest fix is to delete this for and write one message per outer-select iteration, re-arming wake after each pop(). Shorter than the current code, and fixes the ping and the drop log together.

Comment thread pkg/api/gsoc.go
if err != nil {
s.logger.Debug("gsoc ws: write message failed", "error", err)
return
if dropped := queue.droppedCount(); dropped > 0 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same gate as the ping: this only runs when the queue drains to empty, which doesn't happen while messages are being dropped. So the warning fires after the problem is over, never during.

Moving it into case <-ticker.C: fixes that and rate-limits it to one line per ping period. Worth adding the SOC address as a log key — right now it doesn't say which subscription is behind.

Comment thread pkg/gsoc/gsoc.go
go func(hh Handler) {
hh(c.WrappedChunk().Data()[swarm.SpanSize:])
}(*hh)
(*hh)(c)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not sure but should we add this instead, wdyt @janos ?

safe.Run(l.logger, "gsoc-handler", func() {
        (*hh)(c)
    })

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants