From b696c9a5a37030e30415489c0fa6944d52c605ab Mon Sep 17 00:00:00 2001 From: mirkobrombin Date: Sat, 5 Sep 2026 17:57:37 +0200 Subject: [PATCH] fix[closes #83]: harden StreamBus release blockers --- .github/workflows/cli.yaml | 2 +- .github/workflows/pallas.yml | 2 +- .github/workflows/test-build-release.yml | 6 +- CHANGELOG.md | 2 + CONTRIBUTING.md | 2 +- README.md | 2 +- bench/todomvc/rfw/go.mod | 2 +- cmd/rfw/initproj/init.go | 2 +- .../guide/getting-started-from-node.md | 10 +-- .../guide/realtime-dashboard-tutorial.md | 2 +- docs/articles/guide/transports.md | 3 +- examples/counter/go.mod | 2 +- examples/dashboard/go.mod | 2 +- examples/dynamic-list/go.mod | 2 +- go.mod | 4 +- go.sum | 4 +- host/streambus_delivery_test.go | 90 +++++++++++++++++++ host/streambus_server.go | 67 +++++++++++--- 18 files changed, 171 insertions(+), 35 deletions(-) create mode 100644 host/streambus_delivery_test.go diff --git a/.github/workflows/cli.yaml b/.github/workflows/cli.yaml index 1102fe19..57814f67 100644 --- a/.github/workflows/cli.yaml +++ b/.github/workflows/cli.yaml @@ -32,7 +32,7 @@ jobs: - name: Set up Go uses: actions/setup-go@v5 with: - go-version: 1.25 + go-version: 1.26 - name: Test run: go test -v ./... diff --git a/.github/workflows/pallas.yml b/.github/workflows/pallas.yml index f91cbc3c..edcf686f 100644 --- a/.github/workflows/pallas.yml +++ b/.github/workflows/pallas.yml @@ -16,7 +16,7 @@ jobs: - name: Set up Go uses: actions/setup-go@v5 with: - go-version: 1.25 + go-version: 1.26 - name: Install pallas run: go install github.com/vanilla-os/pallas/cmd/pallas@a07d6a955f40871655158d66178f90e69bd417a2 diff --git a/.github/workflows/test-build-release.yml b/.github/workflows/test-build-release.yml index c641306c..26dd9ac5 100644 --- a/.github/workflows/test-build-release.yml +++ b/.github/workflows/test-build-release.yml @@ -24,7 +24,7 @@ jobs: - name: Set up Go uses: actions/setup-go@v5 with: - go-version: 1.25 + go-version: 1.26 - name: Run Unit Tests run: go test -v ./... @@ -163,7 +163,7 @@ jobs: - name: Set up Go uses: actions/setup-go@v5 with: - go-version: 1.25 + go-version: 1.26 - name: Build release assets run: | @@ -228,7 +228,7 @@ jobs: - uses: actions/checkout@v4 - uses: actions/setup-go@v5 with: - go-version: "1.25" + go-version: "1.26" - uses: golangci/golangci-lint-action@v7 with: version: v2.12.2 diff --git a/CHANGELOG.md b/CHANGELOG.md index 5913b05a..2898d45e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ WebSocket as an automatic compatibility fallback. - Keep SSC components, actions, hydration, sessions, replay and broadcasts transport-independent. +- Keep StreamBus outbound failures from consuming SSC sequence numbers before + the connection can accept a frame. All notable changes to this project are documented in this file. diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index d4f3a640..2a14469c 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -5,7 +5,7 @@ pull request. ## Development setup -You need Go **1.25** or later (see `go.mod`). +You need Go **1.26** or later (see `go.mod`). Clone the repository and verify everything builds for both targets: diff --git a/README.md b/README.md index 918a4329..832d2d1a 100644 --- a/README.md +++ b/README.md @@ -39,7 +39,7 @@ rfw 2.1 includes: ## Getting Started -rfw requires Go 1.25 or newer. Coming from Node and never installed Go? See +rfw requires Go 1.26 or newer. Coming from Node and never installed Go? See [Getting started from Node](./docs/articles/guide/getting-started-from-node.md). ```bash diff --git a/bench/todomvc/rfw/go.mod b/bench/todomvc/rfw/go.mod index 31b57c7c..c170e01c 100644 --- a/bench/todomvc/rfw/go.mod +++ b/bench/todomvc/rfw/go.mod @@ -1,6 +1,6 @@ module github.com/rfwlab/bench/todomvc -go 1.25.0 +go 1.26.0 require github.com/rfwlab/rfw/v2 v2.0.0-beta.19 diff --git a/cmd/rfw/initproj/init.go b/cmd/rfw/initproj/init.go index 6a647972..0de1ad48 100644 --- a/cmd/rfw/initproj/init.go +++ b/cmd/rfw/initproj/init.go @@ -14,7 +14,7 @@ import ( "strings" ) -const scaffoldGoVersion = "1.25.0" +const scaffoldGoVersion = "1.26.0" // InitProject creates a new rfw project from the embedded template. func InitProject(projectName string, skipTidy bool) (err error) { diff --git a/docs/articles/guide/getting-started-from-node.md b/docs/articles/guide/getting-started-from-node.md index bbd24ed0..b7263b57 100644 --- a/docs/articles/guide/getting-started-from-node.md +++ b/docs/articles/guide/getting-started-from-node.md @@ -4,7 +4,7 @@ You know npm, `package.json`, and `npm run dev`. You have never installed Go. This guide gets you from zero to a running rfw app without assuming any Go background. -rfw requires **Go 1.25 or newer**. +rfw requires **Go 1.26 or newer**. ## 1. Install Go @@ -13,9 +13,9 @@ rfw requires **Go 1.25 or newer**. Most distro packages lag behind. Prefer the official tarball: ```bash -curl -LO https://go.dev/dl/go1.25.0.linux-amd64.tar.gz +curl -LO https://go.dev/dl/go1.26.0.linux-amd64.tar.gz sudo rm -rf /usr/local/go -sudo tar -C /usr/local -xzf go1.25.0.linux-amd64.tar.gz +sudo tar -C /usr/local -xzf go1.26.0.linux-amd64.tar.gz ``` Then add Go to your PATH (in `~/.bashrc` or `~/.zshrc`): @@ -24,7 +24,7 @@ Then add Go to your PATH (in `~/.bashrc` or `~/.zshrc`): export PATH=$PATH:/usr/local/go/bin ``` -If you prefer your package manager, check the version first: you need 1.25+. +If you prefer your package manager, check the version first: you need 1.26+. On Arch `pacman -S go` is current; on Debian/Ubuntu the `golang` package is often too old, use the tarball instead. @@ -50,7 +50,7 @@ winget install GoLang.Go ```bash go version -# go version go1.25.0 linux/amd64 +# go version go1.26.0 linux/amd64 ``` ## 2. The PATH gotcha diff --git a/docs/articles/guide/realtime-dashboard-tutorial.md b/docs/articles/guide/realtime-dashboard-tutorial.md index 11d50c60..c975c675 100644 --- a/docs/articles/guide/realtime-dashboard-tutorial.md +++ b/docs/articles/guide/realtime-dashboard-tutorial.md @@ -5,7 +5,7 @@ metrics, a component that renders them, a simulated data feed driven by a goroutine and a `time.Ticker`, a list rendered with `@for`, and a pause button wired with `@on:click`. No JavaScript is written at any point. -Prerequisites: Go 1.25+ and the rfw CLI. If you have neither, start with +Prerequisites: Go 1.26+ and the rfw CLI. If you have neither, start with [Getting started from Node](getting-started-from-node.md). ## 1. Scaffold the project diff --git a/docs/articles/guide/transports.md b/docs/articles/guide/transports.md index cfdd0988..360f7f51 100644 --- a/docs/articles/guide/transports.md +++ b/docs/articles/guide/transports.md @@ -39,7 +39,8 @@ Accepted values are: StreamBus runs on WebTransport over HTTP/3 at `/streambus`. RFW preserves its existing JSON SSC protocol and length-prefixes messages on a reliable QUIC stream. On the server, Warp StreamBus provides bounded queues, priorities, -replay storage and explicit backpressure before frames reach the network. +and explicit backpressure before frames reach the network. SSC keeps its own +sequenced replay history above the transport. The WebSocket endpoint remains mounted as a compatibility fallback. No application code or component API changes when the selected transport changes. diff --git a/examples/counter/go.mod b/examples/counter/go.mod index 3f66951e..4f03cca9 100644 --- a/examples/counter/go.mod +++ b/examples/counter/go.mod @@ -1,6 +1,6 @@ module github.com/rfwlab/examples/counter -go 1.25.0 +go 1.26.0 require github.com/rfwlab/rfw/v2 v2.0.0-beta.19 diff --git a/examples/dashboard/go.mod b/examples/dashboard/go.mod index bacbf988..08e3b989 100644 --- a/examples/dashboard/go.mod +++ b/examples/dashboard/go.mod @@ -1,6 +1,6 @@ module github.com/rfwlab/examples/dashboard -go 1.25.0 +go 1.26.0 require github.com/rfwlab/rfw/v2 v2.0.0-beta.19 diff --git a/examples/dynamic-list/go.mod b/examples/dynamic-list/go.mod index 9488cf1b..5f075c57 100644 --- a/examples/dynamic-list/go.mod +++ b/examples/dynamic-list/go.mod @@ -1,6 +1,6 @@ module github.com/rfwlab/examples/dynamic-list -go 1.25.0 +go 1.26.0 require github.com/rfwlab/rfw/v2 v2.0.0-beta.19 diff --git a/go.mod b/go.mod index 6244fbd9..9e5490e5 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module github.com/rfwlab/rfw/v2 -go 1.25.0 +go 1.26.0 require ( github.com/andybalholm/brotli v1.2.1 @@ -26,7 +26,7 @@ require ( github.com/mattn/go-isatty v0.0.20 // indirect github.com/quic-go/qpack v0.6.0 // indirect github.com/tdewolff/parse/v2 v2.8.3 // indirect - golang.org/x/crypto v0.55.0 // indirect + golang.org/x/crypto v0.56.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect gopkg.in/natefinch/lumberjack.v2 v2.2.1 // indirect diff --git a/go.sum b/go.sum index 3ecf71cb..f20621a8 100644 --- a/go.sum +++ b/go.sum @@ -50,8 +50,8 @@ github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZ github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= go.uber.org/mock v0.5.2 h1:LbtPTcP8A5k9WPXj54PPPbjcI4Y6lhyOZXn+VS7wNko= go.uber.org/mock v0.5.2/go.mod h1:wLlUxC2vVTPTaE3UD51E0BGOAElKrILxhVSDYQLld5o= -golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= -golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/crypto v0.56.0 h1:GUh5Ii4J5jtcseSMiRqr1jXCNHoxjeV9Fmekc2oLy6Y= +golang.org/x/crypto v0.56.0/go.mod h1:OMW5y6CY9l38uPLmxU6l6pwcXp1obtLo3e6gT7gQR2I= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= diff --git a/host/streambus_delivery_test.go b/host/streambus_delivery_test.go new file mode 100644 index 00000000..fd4d3d84 --- /dev/null +++ b/host/streambus_delivery_test.go @@ -0,0 +1,90 @@ +package host + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/mirkobrombin/go-warp/v2/streambus" +) + +func TestStreamBusInvalidPayloadDoesNotConsumeSequence(t *testing.T) { + connection := &streamBusConnection{} + session := newSession("stream-invalid", sessionOptions{replayLimit: 4}) + session.streamConnection = connection + + sendStreamBusSession(connection, session, Outbound{Payload: make(chan int)}) + + if session.outboundSeq != 0 { + t.Fatalf("sequence = %d, want 0", session.outboundSeq) + } + if len(session.replay) != 0 { + t.Fatalf("replay = %#v, want empty", session.replay) + } +} + +func TestStreamBusBackpressureClosesBeforeSequence(t *testing.T) { + bus := streambus.NewInMemory(streambus.Config{DefaultBuffer: 1, MaxBuffer: 1}) + t.Cleanup(func() { _ = bus.Close() }) + subscription, err := bus.Subscribe(context.Background(), streambus.SubscribeOptions{ + Topic: "blocked", Buffer: 1, Overflow: streambus.Block, + }) + if err != nil { + t.Fatalf("subscribe: %v", err) + } + connection := &streamBusConnection{bus: bus, topic: "blocked", subscription: subscription, buffer: 1} + session := newSession("stream-full", sessionOptions{replayLimit: 4}) + session.streamConnection = connection + + if _, err := bus.Publish(context.Background(), streambus.Frame{Topic: "blocked", Payload: []byte("one")}); err != nil { + t.Fatalf("publish first: %v", err) + } + waitForStreamBus(t, func() bool { return subscription.Stats().Queued == 0 }) + if _, err := bus.Publish(context.Background(), streambus.Frame{Topic: "blocked", Payload: []byte("two")}); err != nil { + t.Fatalf("publish second: %v", err) + } + waitForStreamBus(t, func() bool { return subscription.Stats().Queued == 1 }) + + sendStreamBusSession(connection, session, Outbound{Payload: "late"}) + + if session.outboundSeq != 0 { + t.Fatalf("sequence = %d, want 0", session.outboundSeq) + } + select { + case <-subscription.Done(): + default: + t.Fatal("subscription stayed open") + } +} + +func TestStreamBusEndpointDisablesTransportReplay(t *testing.T) { + endpoint := newStreamBusEndpoint(NewWSRuntime()) + t.Cleanup(func() { _ = endpoint.bus.Close() }) + topic := "rfw/connection/replay" + sequence, err := endpoint.bus.Publish(context.Background(), streambus.Frame{Topic: topic, Payload: []byte("{}")}) + if err != nil { + t.Fatalf("publish: %v", err) + } + subscription, err := endpoint.bus.Subscribe(context.Background(), streambus.SubscribeOptions{ + Topic: topic, Buffer: 1, Since: sequence, + }) + if subscription != nil { + _ = subscription.Close() + } + if !errors.Is(err, streambus.ErrReplayUnavailable) { + t.Fatalf("subscribe since = %v, want %v", err, streambus.ErrReplayUnavailable) + } +} + +func waitForStreamBus(t *testing.T, f func() bool) { + t.Helper() + deadline := time.Now().Add(time.Second) + for time.Now().Before(deadline) { + if f() { + return + } + time.Sleep(time.Millisecond) + } + t.Fatal("condition was not reached") +} diff --git a/host/streambus_server.go b/host/streambus_server.go index f5a42e2e..ec400b09 100644 --- a/host/streambus_server.go +++ b/host/streambus_server.go @@ -24,8 +24,11 @@ import ( ) const streamBusPath = "/streambus" +const streamBusConnectionBuffer = 256 var ( + errStreamBusBackpressure = errors.New("streambus: outbound queue is full") + streamEndpoints sync.Map streamClientsMu sync.RWMutex streamClients = make(map[string]map[*streamBusConnection]*Session) @@ -49,7 +52,7 @@ func newStreamBusEndpoint(runtime *WSRuntime) *streamBusEndpoint { return &streamBusEndpoint{ runtime: runtime, bus: streambus.NewInMemory(streambus.Config{ - DefaultBuffer: 256, MaxBuffer: 4096, ReplayCapacity: 256, + DefaultBuffer: streamBusConnectionBuffer, MaxBuffer: 4096, MaxPayloadBytes: maximum, }), } @@ -117,6 +120,7 @@ type streamBusConnection struct { bus *streambus.InMemory topic string subscription *streambus.Subscription + buffer int maximum int writeMu sync.Mutex closeOnce sync.Once @@ -125,13 +129,14 @@ type streamBusConnection struct { func newStreamBusConnection(ctx context.Context, request *http.Request, transport *wt.Session, stream *wt.Stream, bus *streambus.InMemory, maximum int) (*streamBusConnection, error) { topic := fmt.Sprintf("rfw/connection/%d", streamID.Add(1)) - subscription, err := bus.Subscribe(ctx, streambus.SubscribeOptions{Topic: topic, Buffer: 256, Overflow: streambus.Block}) + subscription, err := bus.Subscribe(ctx, streambus.SubscribeOptions{Topic: topic, Buffer: streamBusConnectionBuffer, Overflow: streambus.Block}) if err != nil { return nil, err } connection := &streamBusConnection{ request: request, transport: transport, stream: stream, reader: bufio.NewReader(stream), - bus: bus, topic: topic, subscription: subscription, maximum: maximum, done: make(chan struct{}), + bus: bus, topic: topic, subscription: subscription, buffer: streamBusConnectionBuffer, + maximum: maximum, done: make(chan struct{}), } go connection.writeLoop() return connection, nil @@ -150,6 +155,9 @@ func (c *streamBusConnection) Send(out Outbound) error { if err != nil { return err } + if !c.canAcceptOutbound() { + return errStreamBusBackpressure + } ctx, cancel := context.WithTimeout(c.transport.Context(), 5*time.Second) defer cancel() _, err = c.bus.Publish(ctx, streambus.Frame{ @@ -159,6 +167,17 @@ func (c *streamBusConnection) Send(out Outbound) error { return err } +func (c *streamBusConnection) canAcceptOutbound() bool { + if c == nil || c.subscription == nil { + return false + } + buffer := c.buffer + if buffer <= 0 { + buffer = streamBusConnectionBuffer + } + return c.subscription.Stats().Queued < buffer +} + func (c *streamBusConnection) writeLoop() { defer close(c.done) for frame := range c.subscription.Frames() { @@ -174,11 +193,17 @@ func (c *streamBusConnection) writeLoop() { func (c *streamBusConnection) Close() { c.closeOnce.Do(func() { - _ = c.transport.CloseWithError(0, "") - _ = c.subscription.Close() - c.writeMu.Lock() - _ = c.stream.Close() - c.writeMu.Unlock() + if c.transport != nil { + _ = c.transport.CloseWithError(0, "") + } + if c.subscription != nil { + _ = c.subscription.Close() + } + if c.stream != nil { + c.writeMu.Lock() + _ = c.stream.Close() + c.writeMu.Unlock() + } }) } @@ -330,7 +355,11 @@ func bindStreamBusConnection(connection *streamBusConnection, session *Session, } func sendStreamBusSession(connection *streamBusConnection, session *Session, out Outbound) { - if session == nil { + if connection == nil || session == nil { + return + } + if _, err := json.Marshal(out); err != nil { + log.Printf("streambus marshal outbound payload: %v", err) return } session.outboundMu.Lock() @@ -338,7 +367,14 @@ func sendStreamBusSession(connection *streamBusConnection, session *Session, out if session.streamConnection != connection { return } - _ = connection.Send(session.PrepareOutbound(out)) + if !connection.canAcceptOutbound() { + connection.Close() + return + } + if err := connection.Send(session.PrepareOutbound(out)); err != nil { + log.Printf("streambus send: %v", err) + connection.Close() + } } func replayStreamBus(connection *streamBusConnection, session *Session, acknowledged uint64) { @@ -349,11 +385,18 @@ func replayStreamBus(connection *streamBusConnection, session *Session, acknowle } messages, err := session.ReplayAfter(acknowledged) if err != nil { - _ = connection.Send(session.PrepareOutbound(Outbound{Error: NewActionError("resync_required", "message history is no longer available")})) + if err := connection.Send(session.PrepareOutbound(Outbound{Error: NewActionError("resync_required", "message history is no longer available")})); err != nil { + log.Printf("streambus replay: %v", err) + connection.Close() + } return } for _, message := range messages { - _ = connection.Send(message) + if err := connection.Send(message); err != nil { + log.Printf("streambus replay: %v", err) + connection.Close() + return + } } }