From e4c0186f96fb39de19e56bcfc3bfe793c088e331 Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 8 Oct 2026 21:57:12 +0800 Subject: [PATCH] fix(events): isolate subscription overflow instead of killing the shared channel 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. --- docs/realtime-event-push-plan.md | 1 + internal/controlplane/api/event_channel.go | 29 ++++++- .../controlplane/api/event_channel_test.go | 75 +++++++++++++++++++ 3 files changed, 102 insertions(+), 3 deletions(-) diff --git a/docs/realtime-event-push-plan.md b/docs/realtime-event-push-plan.md index e2b64c9..2134a32 100644 --- a/docs/realtime-event-push-plan.md +++ b/docs/realtime-event-push-plan.md @@ -67,5 +67,6 @@ SSE 是列表失效通知而非唯一数据存储;服务重启或页面重连 - Go 全项目测试通过;PostgreSQL 事件/监听相关测试含 race 检查通过。新增通道业务函数覆盖 65%–100%,封面工作线程 97.7%,监听会话 83.9%,封面结果更新 77.3%。这是受影响业务范围的覆盖,不是整个 Go 项目的覆盖率。 - 前端及网页信号 Node 单元测试:19 项通过;事件页覆盖 98.83%;antd lint 无问题,web 构建通过;没有编写或运行 E2E。 - 已检查同代次旧会话停止、共享通道最后订阅关闭与新订阅加入、单账号启动/错误隔离、历史与实时并行、封面过载不额外阻塞确认等边界。 +- 订阅缓冲溢出只隔离溢出订阅并显式报错(该账号会话重启后由网关补送未确认投递),共享通道与其他账号不受影响。 - 部署需同时更新后台与网关,旧事件协议不保留兼容。保留浏览器 profile、指纹与 Cookie,不通过重建环境或清 Cookie 迁移;暂停账号不会因此恢复运行。 - 尚未进行真实抖音通知到数据库、页面的全链路延迟验收;前述 P95 目标仍待人工采样核对,不能以模拟测试、连接成功或通知列表读取成功替代。 diff --git a/internal/controlplane/api/event_channel.go b/internal/controlplane/api/event_channel.go index 71a85eb..9fecf28 100644 --- a/internal/controlplane/api/event_channel.go +++ b/internal/controlplane/api/event_channel.go @@ -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(): diff --git a/internal/controlplane/api/event_channel_test.go b/internal/controlplane/api/event_channel_test.go index 5efa453..a3c7249 100644 --- a/internal/controlplane/api/event_channel_test.go +++ b/internal/controlplane/api/event_channel_test.go @@ -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)