fix(events): isolate subscription overflow instead of killing the shared channel
douyin-release-gate / verify (push) Failing after 21m22s
douyin-release-gate / verify (push) Failing after 21m22s
One slow account's full subscription buffer no longer tears down the WebSocket shared by all accounts on the gateway. The overflowing subscription is retired with an explicit error and replays unacked deliveries on resubscribe; the buffer grows from 2 to 64 messages.
This commit is contained in:
@@ -14,6 +14,7 @@ import (
|
||||
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/environment"
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
type channelMessage struct {
|
||||
@@ -26,6 +27,8 @@ type eventSubscription struct {
|
||||
hub *gatewayEventChannel
|
||||
id, alias string
|
||||
messages chan channelMessage
|
||||
done chan struct{}
|
||||
failErr error
|
||||
}
|
||||
type gatewayEventChannel struct {
|
||||
target environment.Gateway
|
||||
@@ -61,7 +64,7 @@ func acquireEventSubscription(ctx context.Context, target environment.Gateway, p
|
||||
eventChannels.items[key] = hub
|
||||
}
|
||||
alias, _ := params["alias"].(string)
|
||||
sub := &eventSubscription{hub: hub, id: fmt.Sprint(subscriptionSequence.Add(1)), alias: alias, messages: make(chan channelMessage, 2)}
|
||||
sub := &eventSubscription{hub: hub, id: fmt.Sprint(subscriptionSequence.Add(1)), alias: alias, messages: make(chan channelMessage, 64), done: make(chan struct{})}
|
||||
// Reserve the subscription before releasing the pool lock, so closing the
|
||||
// previous last account cannot retire a channel another account is joining.
|
||||
hub.mu.Lock()
|
||||
@@ -170,8 +173,12 @@ func (h *gatewayEventChannel) receive() {
|
||||
select {
|
||||
case sub.messages <- message:
|
||||
default:
|
||||
h.fail(errors.New("gateway WS subscription queue full"))
|
||||
return
|
||||
// One slow account must not tear down the shared channel: retire only
|
||||
// this subscription. The gateway keeps unacknowledged deliveries and
|
||||
// replays them when the account session resubscribes.
|
||||
spec := fmt.Errorf("event subscription queue overflow; the account session will resubscribe and replay unacknowledged deliveries")
|
||||
logrus.WithFields(logrus.Fields{"alias": sub.alias, "subscription": sub.id}).Error("event subscription queue overflow; isolating subscription from shared channel")
|
||||
sub.fail(spec)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -191,6 +198,20 @@ func (h *gatewayEventChannel) fail(err error) {
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// fail retires one subscription without touching the shared channel: the
|
||||
// account session restarts and the gateway replays unacknowledged deliveries.
|
||||
// Only the receive goroutine calls this, so close(done) happens exactly once.
|
||||
func (s *eventSubscription) fail(err error) {
|
||||
if err == nil {
|
||||
err = errors.New("event subscription failed")
|
||||
}
|
||||
s.hub.mu.Lock()
|
||||
delete(s.hub.subs, s.id)
|
||||
s.hub.mu.Unlock()
|
||||
s.failErr = err
|
||||
close(s.done)
|
||||
}
|
||||
func (s *eventSubscription) poll(ctx context.Context, acks []string) ([]creator.ListenerDelivery, error) {
|
||||
if len(acks) > 0 {
|
||||
if err := s.hub.send(map[string]any{"type": "ack", "alias": s.alias, "subscription": s.id, "delivery_ids": acks}); err != nil {
|
||||
@@ -208,6 +229,8 @@ func (s *eventSubscription) poll(ctx context.Context, acks []string) ([]creator.
|
||||
}
|
||||
}
|
||||
return message.Deliveries, nil
|
||||
case <-s.done:
|
||||
return nil, s.failErr
|
||||
case <-s.hub.done:
|
||||
return nil, errors.New("gateway WS disconnected")
|
||||
case <-ctx.Done():
|
||||
|
||||
@@ -3,10 +3,13 @@ package api
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/environment"
|
||||
"github.com/gorilla/websocket"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -133,6 +136,78 @@ func TestEventChannelReportsDisconnect(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventChannelSubscriptionOverflowIsolatesAccounts(t *testing.T) {
|
||||
var connections atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
c, err := (&websocket.Upgrader{}).Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
defer c.Close()
|
||||
connections.Add(1)
|
||||
var write sync.Mutex
|
||||
send := func(v map[string]any) {
|
||||
write.Lock()
|
||||
defer write.Unlock()
|
||||
_ = c.WriteJSON(v)
|
||||
}
|
||||
for {
|
||||
var v map[string]any
|
||||
if c.ReadJSON(&v) != nil {
|
||||
return
|
||||
}
|
||||
if v["type"] != "subscribe" {
|
||||
continue
|
||||
}
|
||||
send(map[string]any{"type": "subscribed", "subscription": v["subscription"]})
|
||||
if v["alias"] == "flooded" {
|
||||
for i := 0; i < 80; i++ {
|
||||
send(map[string]any{"type": "deliveries", "subscription": v["subscription"],
|
||||
"deliveries": []any{map[string]any{"kind": "open", "delivery_id": fmt.Sprintf("flood-%d", i)}}})
|
||||
}
|
||||
}
|
||||
if v["alias"] == "late" {
|
||||
send(map[string]any{"type": "deliveries", "subscription": v["subscription"],
|
||||
"deliveries": []any{map[string]any{"kind": "open", "delivery_id": "late-1"}}})
|
||||
}
|
||||
}
|
||||
}))
|
||||
defer server.Close()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
target := environment.Gateway{Endpoint: server.URL, Token: "test"}
|
||||
flooded, err := acquireEventSubscription(ctx, target, map[string]any{"alias": "flooded"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer flooded.close()
|
||||
late, err := acquireEventSubscription(ctx, target, map[string]any{"alias": "late"})
|
||||
if err != nil {
|
||||
t.Fatalf("shared channel refused a joining account after another account overflowed: %v", err)
|
||||
}
|
||||
defer late.close()
|
||||
items, err := late.poll(ctx, nil)
|
||||
if err != nil || items[0].DeliveryID != "late-1" {
|
||||
t.Fatalf("other account deliveries broke after one account overflowed: %+v %v", items, err)
|
||||
}
|
||||
if connections.Load() != 1 {
|
||||
t.Fatalf("overflow tore down the shared channel, connections=%d", connections.Load())
|
||||
}
|
||||
// Buffered deliveries sent before the overflow fired remain valid work; the
|
||||
// subscription keeps reporting them until the buffer drains, then surfaces
|
||||
// the explicit overflow error.
|
||||
for {
|
||||
_, err = flooded.poll(ctx, nil)
|
||||
if err != nil {
|
||||
if !strings.Contains(err.Error(), "overflow") {
|
||||
t.Fatalf("overflowed subscription must fail with an explicit overflow error, got %v", err)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventChannelClosingLastAccountDoesNotRetireJoiningAccount(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
c, err := (&websocket.Upgrader{}).Upgrade(w, r, nil)
|
||||
|
||||
Reference in New Issue
Block a user