fix: restore web widget realtime replies (HH-556) (#136)

* fix: restore web widget realtime replies (HH-556)

* fix: make realtime delivery durable (HH-556)

* fix: make realtime message outbox atomic (HH-556)

---------

Co-authored-by: Rogee <rogee@ipao.vip>
This commit is contained in:
Rogee
2026-08-23 22:35:10 +08:00
committed by GitHub
co-authored by rogee
parent 17244bcc9d
commit 60a6ac4785
13 changed files with 987 additions and 136 deletions
+173 -26
View File
@@ -6,12 +6,31 @@ package ws
import (
"context"
"crypto/sha256"
"encoding/json"
"fmt"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/gorm"
)
const taskTypeRealtimeEventPublish = "realtime:event_publish"
const (
realtimeTargetAccount = "account"
realtimeTargetToken = "pubsub_token"
)
type realtimeEventPublishJob struct {
AccountID uint `json:"account_id"`
PubsubToken string `json:"pubsub_token,omitempty"`
Target string `json:"target,omitempty"`
EventType string `json:"event_type"`
Payload interface{} `json:"payload"`
}
// EventPublisher routes business events to both WebSocket and SSE clients.
// It is the central point that services call after mutations to push real-time
// updates, matching Chatwoot's pattern where controllers broadcast events
@@ -29,9 +48,10 @@ import (
// entry point, with Hub and SSERegistry as local delivery targets, and Redis
// Pub/Sub relay for multi-instance fan-out.
type EventPublisher struct {
hub MessageHandler // WebSocket Hub (local delivery)
sse *SSERegistry // SSE registry (local delivery)
relay *BroadcastRelay // Redis Pub/Sub relay (cross-instance delivery)
hub MessageHandler // WebSocket Hub (local delivery)
sse *SSERegistry // SSE registry (local delivery)
relay *BroadcastRelay // Redis Pub/Sub relay (cross-instance delivery)
worker *worker.WorkerPool
}
// NewEventPublisher creates an EventPublisher with all delivery targets.
@@ -52,6 +72,15 @@ func NewEventPublisherLocal(hub MessageHandler, sse *SSERegistry) *EventPublishe
}
}
// SetWorkerPool makes realtime delivery durable. The request path only stores
// the job; Redis failures are retried and remain visible in background_jobs.
func (p *EventPublisher) SetWorkerPool(wp *worker.WorkerPool) {
p.worker = wp
if wp != nil {
wp.Register(taskTypeRealtimeEventPublish, p.performPublishJob)
}
}
// PublishEvent publishes a real-time event to all delivery targets.
// This is the primary API that services call after business mutations.
//
@@ -67,7 +96,14 @@ func NewEventPublisherLocal(hub MessageHandler, sse *SSERegistry) *EventPublishe
//
// Reference: Chatwoot controllers call broadcast_event after mutations,
// which triggers Wisper → ActionCable → Redis Pub/Sub relay.
func (p *EventPublisher) PublishEvent(accountID uint, eventType string, payload interface{}) {
func (p *EventPublisher) PublishEvent(accountID uint, eventType string, payload interface{}) error {
if p.worker != nil {
return p.enqueuePublishJobs(context.Background(), realtimePublishJobs(accountID, "", eventType, payload))
}
return p.publishEvent(accountID, eventType, payload)
}
func (p *EventPublisher) publishEvent(accountID uint, eventType string, payload interface{}) error {
// Build the wire-format message (matching Chatwoot ActionCable event format)
wsMsg := &WSMessage{
Event: eventType,
@@ -77,8 +113,7 @@ func (p *EventPublisher) PublishEvent(accountID uint, eventType string, payload
data, err := json.Marshal(wsMsg)
if err != nil {
applogger.L().Warnf("event publisher: failed to marshal event %s: %v", eventType, err)
return
return fmt.Errorf("marshal event %s: %w", eventType, err)
}
// 1. Deliver to WebSocket Hub (local clients)
@@ -98,9 +133,10 @@ func (p *EventPublisher) PublishEvent(accountID uint, eventType string, payload
if p.relay != nil {
room := accountRoomNameHelper(accountID)
if err := p.relay.Publish(context.Background(), room, wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for %s: %v", eventType, err)
return fmt.Errorf("publish event %s to account room: %w", eventType, err)
}
}
return nil
}
// PublishConversationEvent publishes a real-time event scoped to a specific
@@ -150,7 +186,24 @@ func (p *EventPublisher) PublishConversationEvent(accountID uint, conversationID
// contact's pubsub_token room. Chatwoot's widget ActionCable connector
// subscribes to RoomChannel with pubsub_token, so widget-visible events must be
// available on that token-scoped room in addition to the dashboard account room.
func (p *EventPublisher) PublishWidgetEvent(accountID uint, pubsubToken string, eventType string, payload interface{}) {
func (p *EventPublisher) PublishWidgetEvent(accountID uint, pubsubToken string, eventType string, payload interface{}) error {
if p.worker != nil {
return p.enqueuePublishJobs(context.Background(), realtimePublishJobs(accountID, pubsubToken, eventType, payload))
}
return p.publishWidgetEvent(accountID, pubsubToken, eventType, payload)
}
func (p *EventPublisher) publishWidgetEvent(accountID uint, pubsubToken string, eventType string, payload interface{}) error {
if err := p.publishEvent(accountID, eventType, payload); err != nil {
return err
}
if pubsubToken == "" {
return nil
}
return p.publishTokenEvent(accountID, pubsubToken, eventType, payload)
}
func (p *EventPublisher) publishTokenEvent(accountID uint, pubsubToken string, eventType string, payload interface{}) error {
wsMsg := &WSMessage{
Event: eventType,
Data: payload,
@@ -159,32 +212,126 @@ func (p *EventPublisher) PublishWidgetEvent(accountID uint, pubsubToken string,
data, err := json.Marshal(wsMsg)
if err != nil {
applogger.L().Warnf("event publisher: failed to marshal widget event %s: %v", eventType, err)
return
return fmt.Errorf("marshal widget event %s: %w", eventType, err)
}
if p.hub != nil {
p.hub.SendToAccount(accountID, data)
if pubsubToken != "" {
p.hub.SendToRoom(pubsubTokenRoomNameHelper(pubsubToken), data)
}
}
if p.sse != nil {
p.sse.SendToAccount(accountID, SSEEvent{Type: eventType, Payload: payload})
p.hub.SendToRoom(pubsubTokenRoomNameHelper(pubsubToken), data)
}
if p.relay != nil {
room := accountRoomNameHelper(accountID)
if err := p.relay.Publish(context.Background(), room, wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for %s: %v", eventType, err)
}
if pubsubToken != "" {
if err := p.relay.Publish(context.Background(), pubsubTokenRoomNameHelper(pubsubToken), wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for widget %s: %v", eventType, err)
}
if err := p.relay.Publish(context.Background(), pubsubTokenRoomNameHelper(pubsubToken), wsMsg); err != nil {
return fmt.Errorf("publish widget event %s to token room: %w", eventType, err)
}
}
return nil
}
func realtimePublishJobs(accountID uint, pubsubToken, eventType string, payload interface{}) []realtimeEventPublishJob {
jobs := []realtimeEventPublishJob{{AccountID: accountID, Target: realtimeTargetAccount, EventType: eventType, Payload: payload}}
if pubsubToken != "" {
jobs = append(jobs, realtimeEventPublishJob{AccountID: accountID, PubsubToken: pubsubToken, Target: realtimeTargetToken, EventType: eventType, Payload: payload})
}
return jobs
}
func realtimePublishJobOptions(job realtimeEventPublishJob) ([]worker.EnqueueOption, error) {
opts := []worker.EnqueueOption{worker.WithQueue("events"), worker.WithMaxAttempts(3)}
if job.EventType == EventMessageCreated {
encoded, err := json.Marshal(job)
if err != nil {
return nil, fmt.Errorf("marshal realtime publish job: %w", err)
}
opts = append(opts, worker.WithIdempotencyKey(fmt.Sprintf("realtime:message.created:%s:%x", job.Target, sha256.Sum256(encoded))))
}
return opts, nil
}
func (p *EventPublisher) enqueuePublishJobs(ctx context.Context, jobs []realtimeEventPublishJob) error {
for _, job := range jobs {
opts, err := realtimePublishJobOptions(job)
if err != nil {
return err
}
if _, err := p.worker.Enqueue(ctx, taskTypeRealtimeEventPublish, job, opts...); err != nil {
applogger.L().Errorf("enqueue realtime event %s failed: %v", job.EventType, err)
return err
}
}
return nil
}
// EnqueueInTransaction atomically persists every target channel with the
// caller's business mutation. The returned jobs are safe to publish only after
// the transaction commits.
func (p *EventPublisher) EnqueueInTransaction(ctx context.Context, tx *gorm.DB, accountID uint, pubsubToken, eventType string, payload interface{}) ([]*model.BackgroundJob, error) {
if p.worker == nil {
return nil, worker.ErrWorkerDatabaseRequired
}
var createdJobs []*model.BackgroundJob
for _, job := range realtimePublishJobs(accountID, pubsubToken, eventType, payload) {
opts, err := realtimePublishJobOptions(job)
if err != nil {
return nil, err
}
backgroundJob, created, err := p.worker.EnqueueInTransaction(ctx, tx, taskTypeRealtimeEventPublish, job, opts...)
if err != nil {
return nil, err
}
if created {
createdJobs = append(createdJobs, backgroundJob)
}
}
return createdJobs, nil
}
// PublishEnqueued triggers Redis after the surrounding transaction commits.
// Database sweep recovery remains authoritative if this process exits here.
func (p *EventPublisher) PublishEnqueued(ctx context.Context, jobs []*model.BackgroundJob) {
for _, job := range jobs {
p.worker.Publish(ctx, job)
}
}
func (p *EventPublisher) performPublishJob(ctx context.Context, backgroundJob *model.BackgroundJob) error {
var job realtimeEventPublishJob
if err := json.Unmarshal(backgroundJob.Payload, &job); err != nil {
return worker.Permanent(fmt.Errorf("unmarshal realtime publish job: %w", err))
}
if job.AccountID == 0 || job.EventType == "" {
return worker.Permanent(fmt.Errorf("invalid realtime publish job: account_id=%d event_type=%q", job.AccountID, job.EventType))
}
var err error
switch job.Target {
case realtimeTargetToken:
err = p.publishTokenJob(job.AccountID, job.PubsubToken, job.EventType, job.Payload)
default:
err = p.publishAccountJob(job.AccountID, job.EventType, job.Payload)
}
if err != nil {
applogger.L().Errorf("realtime event job %d failed for %s: %v", backgroundJob.ID, job.EventType, err)
}
return err
}
func (p *EventPublisher) publishAccountJob(accountID uint, eventType string, payload interface{}) error {
if p.relay == nil {
return p.publishEvent(accountID, eventType, payload)
}
if err := p.relay.Publish(context.Background(), accountRoomNameHelper(accountID), &WSMessage{Event: eventType, Data: payload, AccountID: accountID}); err != nil {
return err
}
if p.sse != nil {
p.sse.SendToAccount(accountID, SSEEvent{Type: eventType, Payload: payload})
}
return nil
}
func (p *EventPublisher) publishTokenJob(accountID uint, pubsubToken, eventType string, payload interface{}) error {
if p.relay == nil {
return p.publishTokenEvent(accountID, pubsubToken, eventType, payload)
}
return p.relay.Publish(context.Background(), pubsubTokenRoomNameHelper(pubsubToken), &WSMessage{Event: eventType, Data: payload, AccountID: accountID})
}
// accountRoomNameHelper generates the room name for an account channel.
+100
View File
@@ -1,15 +1,63 @@
package ws
import (
"context"
"encoding/json"
"errors"
"fmt"
"sync"
"testing"
"time"
"github.com/alicebob/miniredis/v2"
"github.com/redis/go-redis/v9"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
type failPublishChannelOnceHook struct {
mu sync.Mutex
channel string
failed bool
attempts map[string]int
succeeded map[string]int
}
func (h *failPublishChannelOnceHook) DialHook(next redis.DialHook) redis.DialHook { return next }
func (h *failPublishChannelOnceHook) ProcessHook(next redis.ProcessHook) redis.ProcessHook {
return func(ctx context.Context, cmd redis.Cmder) error {
if cmd.Name() != "publish" || len(cmd.Args()) < 2 {
return next(ctx, cmd)
}
channel := fmt.Sprint(cmd.Args()[1])
h.mu.Lock()
h.attempts[channel]++
if channel == h.channel && !h.failed {
h.failed = true
h.mu.Unlock()
return errors.New("token room unavailable")
}
h.mu.Unlock()
err := next(ctx, cmd)
if err == nil {
h.mu.Lock()
h.succeeded[channel]++
h.mu.Unlock()
}
return err
}
}
func (h *failPublishChannelOnceHook) ProcessPipelineHook(next redis.ProcessPipelineHook) redis.ProcessPipelineHook {
return next
}
// === EventPublisher Tests ===
// Uses mockHandler (from broadcast_test.go) as the MessageHandler implementation,
// plus SSERegistry for SSE delivery verification.
@@ -709,3 +757,55 @@ func TestEventPublisher_WidgetEvent_PubsubTokenRoomDelivery(t *testing.T) {
assert.Equal(t, EventMessageCreated, roomMsg.Event)
assert.Equal(t, uint(1), roomMsg.AccountID)
}
func TestEventPublisher_DurableWidgetPublishRetriesPartialFailure(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{})
require.NoError(t, err)
require.NoError(t, db.AutoMigrate(&model.BackgroundJob{}))
mini := miniredis.RunT(t)
rdb := redis.NewClient(&redis.Options{Addr: mini.Addr()})
t.Cleanup(func() { require.NoError(t, rdb.Close()) })
tokenChannel := RedisPrefixRoom + "pubsub_token_visitor"
hook := &failPublishChannelOnceHook{
channel: tokenChannel, attempts: map[string]int{}, succeeded: map[string]int{},
}
rdb.AddHook(hook)
pool := worker.NewWorkerPoolWithOptions(db, worker.WithBackoff(func(int) time.Duration { return 0 }))
publisher := NewEventPublisher(nil, nil, NewBroadcastRelay(rdb, nil))
publisher.SetWorkerPool(pool)
payload := map[string]interface{}{"id": 7, "content": "hello widget"}
require.NoError(t, publisher.PublishWidgetEvent(1, "visitor", EventMessageCreated, payload))
require.NoError(t, publisher.PublishWidgetEvent(1, "visitor", EventMessageCreated, payload))
var count int64
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", taskTypeRealtimeEventPublish).Count(&count).Error)
require.Equal(t, int64(2), count, "account and token jobs must each be idempotent")
processed, err := pool.ProcessOne(context.Background())
require.True(t, processed)
require.NoError(t, err)
processed, err = pool.ProcessOne(context.Background())
require.True(t, processed)
require.ErrorContains(t, err, "token room unavailable")
var job model.BackgroundJob
require.NoError(t, db.Where("idempotency_key LIKE ?", "realtime:message.created:pubsub_token:%").First(&job).Error)
require.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
require.Contains(t, job.LastError, "token room unavailable")
processed, err = pool.ProcessOne(context.Background())
require.True(t, processed)
require.NoError(t, err)
require.NoError(t, db.First(&job, job.ID).Error)
require.Equal(t, model.BackgroundJobStatusCompleted, job.Status)
accountChannel := RedisPrefixRoom + "account_1"
hook.mu.Lock()
defer hook.mu.Unlock()
require.Equal(t, 1, hook.attempts[accountChannel])
require.Equal(t, 1, hook.succeeded[accountChannel])
require.Equal(t, 2, hook.attempts[tokenChannel])
require.Equal(t, 1, hook.succeeded[tokenChannel])
}