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)