feat(dispatch): queue async events durably
This commit is contained in:
@@ -6,12 +6,13 @@ import (
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
// TaskTypeEventDispatch is the async task type for event dispatch.
|
||||
// Reference: Chatwoot EventDispatcherJob (Sidekiq worker)
|
||||
// TODO: asynq integration in P10 (async task processing)
|
||||
const TaskTypeEventDispatch = "event:dispatch_async"
|
||||
|
||||
// EventListener is the interface that all event listeners must implement.
|
||||
@@ -36,13 +37,30 @@ type EventListener interface {
|
||||
type Dispatcher struct {
|
||||
mu sync.RWMutex
|
||||
listeners map[string]EventListener // name → listener (prevents duplicates)
|
||||
worker *worker.WorkerPool
|
||||
}
|
||||
|
||||
// NewDispatcher creates a Dispatcher with empty listener registry.
|
||||
func NewDispatcher() *Dispatcher {
|
||||
return &Dispatcher{
|
||||
func NewDispatcher(workers ...*worker.WorkerPool) *Dispatcher {
|
||||
d := &Dispatcher{
|
||||
listeners: make(map[string]EventListener),
|
||||
}
|
||||
if len(workers) > 0 {
|
||||
d.SetWorkerPool(workers[0])
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
// SetWorkerPool wires DispatchAsync to the durable background job worker.
|
||||
func (d *Dispatcher) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
d.mu.Lock()
|
||||
d.worker = wp
|
||||
d.mu.Unlock()
|
||||
if wp != nil {
|
||||
wp.Register(TaskTypeEventDispatch, func(ctx context.Context, job *model.BackgroundJob) error {
|
||||
return d.handleAsyncJob(ctx, job.Payload)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Register adds a listener. Duplicate names overwrite.
|
||||
@@ -88,25 +106,42 @@ func (d *Dispatcher) Dispatch(ctx context.Context, event *ChannelEvent) error {
|
||||
}
|
||||
|
||||
// DispatchAsync enqueues an event for async processing.
|
||||
// TODO: Implement with asynq in P10. Currently falls back to synchronous dispatch.
|
||||
// When no durable worker is configured, it preserves the legacy synchronous fallback.
|
||||
func (d *Dispatcher) DispatchAsync(ctx context.Context, event *ChannelEvent) error {
|
||||
// For now, fall back to synchronous dispatch.
|
||||
// P10 will implement asynq task enqueue:
|
||||
// payload, _ := json.Marshal(event)
|
||||
// task := asynq.NewTask(TaskTypeEventDispatch, payload)
|
||||
// _, err = d.taskClient.Enqueue(task, asynq.Queue("events"), asynq.MaxRetry(3))
|
||||
wp := d.workerPool()
|
||||
if wp != nil {
|
||||
_, err := wp.Enqueue(ctx, TaskTypeEventDispatch, event, worker.WithQueue("events"), worker.WithMaxAttempts(3))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
applogger.L().Infof("async dispatch enqueued event %s", event.Type)
|
||||
return nil
|
||||
}
|
||||
applogger.L().Infof("async dispatch (sync fallback) event %s", event.Type)
|
||||
return d.Dispatch(ctx, event)
|
||||
}
|
||||
|
||||
// HandleAsyncTask processes an asynq task for async event dispatch.
|
||||
// TODO: Implement with asynq.Server handler in P10.
|
||||
func (d *Dispatcher) workerPool() *worker.WorkerPool {
|
||||
d.mu.RLock()
|
||||
defer d.mu.RUnlock()
|
||||
return d.worker
|
||||
}
|
||||
|
||||
func (d *Dispatcher) handleAsyncJob(ctx context.Context, payload []byte) error {
|
||||
var event ChannelEvent
|
||||
if err := json.Unmarshal(payload, &event); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal event payload: %w", err)
|
||||
}
|
||||
return d.Dispatch(ctx, &event)
|
||||
}
|
||||
|
||||
// HandleAsyncTask processes a durable async event payload for compatibility with
|
||||
// older callers that only need payload validation.
|
||||
func HandleAsyncTask(ctx context.Context, payload []byte) error {
|
||||
var event ChannelEvent
|
||||
if err := json.Unmarshal(payload, &event); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal event payload: %w", err)
|
||||
}
|
||||
applogger.L().Infof("handling async task for event %s", event.Type)
|
||||
// TODO: Call dispatcher.Dispatch(ctx, &event) after asynq server setup in P10
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
package channel
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
type workerDispatchListener struct {
|
||||
name string
|
||||
count atomic.Int32
|
||||
}
|
||||
|
||||
func (l *workerDispatchListener) Name() string { return l.name }
|
||||
|
||||
func (l *workerDispatchListener) OnEvent(ctx context.Context, event *ChannelEvent) error {
|
||||
l.count.Add(1)
|
||||
return nil
|
||||
}
|
||||
|
||||
func newChannelWorkerDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:channel-dispatcher-worker?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatalf("sqlite db handle: %v", err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err := db.AutoMigrate(&model.BackgroundJob{}); err != nil {
|
||||
t.Fatalf("migrate background jobs: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
db.Exec("DELETE FROM background_jobs")
|
||||
sqlDB.Close()
|
||||
})
|
||||
return db
|
||||
}
|
||||
|
||||
func TestDispatcherDispatchAsyncEnqueuesDurableJob(t *testing.T) {
|
||||
db := newChannelWorkerDB(t)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 11, 0, 0, 0, time.UTC) }))
|
||||
dispatcher := NewDispatcher(wp)
|
||||
listener := &workerDispatchListener{name: "capture"}
|
||||
dispatcher.Register(listener)
|
||||
|
||||
event := NewChannelEvent(EventConversationCreated, ChannelWebWidget, 1, 2)
|
||||
if err := dispatcher.DispatchAsync(context.Background(), event); err != nil {
|
||||
t.Fatalf("dispatch async: %v", err)
|
||||
}
|
||||
if listener.count.Load() != 0 {
|
||||
t.Fatalf("listener ran synchronously before worker processed job")
|
||||
}
|
||||
|
||||
var count int64
|
||||
if err := db.Model(&model.BackgroundJob{}).Where("job_type = ? AND status = ?", TaskTypeEventDispatch, model.BackgroundJobStatusQueued).Count(&count).Error; err != nil {
|
||||
t.Fatalf("count jobs: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("expected one queued dispatch job, got %d", count)
|
||||
}
|
||||
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
if err != nil || !processed {
|
||||
t.Fatalf("process dispatch job: processed=%v err=%v", processed, err)
|
||||
}
|
||||
if listener.count.Load() != 1 {
|
||||
t.Fatalf("listener was not called by durable dispatch job")
|
||||
}
|
||||
}
|
||||
@@ -2,12 +2,18 @@ package dispatch
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/channel"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
const TaskTypeEventListenerDispatch = "event:listener_dispatch"
|
||||
|
||||
// EventDispatcher wraps the channel.Dispatcher and adds:
|
||||
// - Sync/async split: sync listeners run immediately; async listeners are queued
|
||||
// - Event name routing: listeners subscribe to specific event names
|
||||
@@ -53,6 +59,7 @@ type EventDispatcher struct {
|
||||
channelDispatcher *channel.Dispatcher
|
||||
registry *ListenerRegistry
|
||||
entries map[string]listenerEntry // listener name → entry
|
||||
worker *worker.WorkerPool
|
||||
}
|
||||
|
||||
// NewEventDispatcher creates a new EventDispatcher wrapping the given
|
||||
@@ -65,6 +72,13 @@ func NewEventDispatcher(cd *channel.Dispatcher) *EventDispatcher {
|
||||
}
|
||||
}
|
||||
|
||||
func (ed *EventDispatcher) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
ed.worker = wp
|
||||
if wp != nil {
|
||||
wp.Register(TaskTypeEventListenerDispatch, ed.performListenerJob)
|
||||
}
|
||||
}
|
||||
|
||||
// RegisterSync adds a sync-mode listener for the given event names.
|
||||
// If eventNames is empty, the listener receives all events.
|
||||
func (ed *EventDispatcher) RegisterSync(listener channel.EventListener, eventNames ...string) {
|
||||
@@ -112,14 +126,7 @@ func (ed *EventDispatcher) Dispatch(ctx context.Context, event *channel.ChannelE
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// Async: run in background goroutine
|
||||
go func(l channel.EventListener, e *channel.ChannelEvent) {
|
||||
asyncCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
if err := l.OnEvent(asyncCtx, e); err != nil {
|
||||
applogger.L().Errorf("async listener %s error on event %s: %v", l.Name(), e.Type, err)
|
||||
}
|
||||
}(entry.listener, event)
|
||||
ed.dispatchListenerAsync(ctx, entry.listener.Name(), event, "async listener")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -136,16 +143,48 @@ func (ed *EventDispatcher) DispatchAsync(ctx context.Context, event *channel.Cha
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
go func(l channel.EventListener, e *channel.ChannelEvent) {
|
||||
asyncCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
if err := l.OnEvent(asyncCtx, e); err != nil {
|
||||
applogger.L().Errorf("async dispatch: listener %s error on event %s: %v", l.Name(), e.Type, err)
|
||||
}
|
||||
}(entry.listener, event)
|
||||
ed.dispatchListenerAsync(ctx, entry.listener.Name(), event, "async dispatch")
|
||||
}
|
||||
}
|
||||
|
||||
type listenerJobPayload struct {
|
||||
ListenerName string `json:"listener_name"`
|
||||
Event *channel.ChannelEvent `json:"event"`
|
||||
}
|
||||
|
||||
func (ed *EventDispatcher) dispatchListenerAsync(ctx context.Context, listenerName string, event *channel.ChannelEvent, logPrefix string) {
|
||||
if ed.worker != nil {
|
||||
payload := listenerJobPayload{ListenerName: listenerName, Event: event}
|
||||
if _, err := ed.worker.Enqueue(ctx, TaskTypeEventListenerDispatch, payload, worker.WithQueue("events"), worker.WithMaxAttempts(3)); err != nil {
|
||||
applogger.L().Errorf("%s %s enqueue error on event %s: %v", logPrefix, listenerName, event.Type, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
entry := ed.entries[listenerName]
|
||||
go func(l channel.EventListener, e *channel.ChannelEvent) {
|
||||
asyncCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
if err := l.OnEvent(asyncCtx, e); err != nil {
|
||||
applogger.L().Errorf("%s: listener %s error on event %s: %v", logPrefix, l.Name(), e.Type, err)
|
||||
}
|
||||
}(entry.listener, event)
|
||||
}
|
||||
|
||||
func (ed *EventDispatcher) performListenerJob(ctx context.Context, job *model.BackgroundJob) error {
|
||||
var payload listenerJobPayload
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
return fmt.Errorf("unmarshal listener dispatch payload: %w", err)
|
||||
}
|
||||
entry, ok := ed.entries[payload.ListenerName]
|
||||
if !ok {
|
||||
return fmt.Errorf("listener %q not registered", payload.ListenerName)
|
||||
}
|
||||
if payload.Event == nil {
|
||||
return fmt.Errorf("listener %q job missing event", payload.ListenerName)
|
||||
}
|
||||
return entry.listener.OnEvent(ctx, payload.Event)
|
||||
}
|
||||
|
||||
// ChannelDispatcher returns the underlying channel.Dispatcher for direct access
|
||||
// if needed (e.g. for channel-level dispatch without the enhanced routing).
|
||||
func (ed *EventDispatcher) ChannelDispatcher() *channel.Dispatcher {
|
||||
@@ -155,4 +194,4 @@ func (ed *EventDispatcher) ChannelDispatcher() *channel.Dispatcher {
|
||||
// Registry returns the listener registry for inspection/testing.
|
||||
func (ed *EventDispatcher) Registry() *ListenerRegistry {
|
||||
return ed.registry
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,107 @@
|
||||
package dispatch
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/channel"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
type dispatchWorkerListener struct {
|
||||
name string
|
||||
count atomic.Int32
|
||||
}
|
||||
|
||||
func (l *dispatchWorkerListener) Name() string { return l.name }
|
||||
|
||||
func (l *dispatchWorkerListener) OnEvent(ctx context.Context, event *channel.ChannelEvent) error {
|
||||
l.count.Add(1)
|
||||
return nil
|
||||
}
|
||||
|
||||
func newDispatchWorkerDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:dispatch-worker?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatalf("sqlite db handle: %v", err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err := db.AutoMigrate(&model.BackgroundJob{}); err != nil {
|
||||
t.Fatalf("migrate background jobs: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
db.Exec("DELETE FROM background_jobs")
|
||||
sqlDB.Close()
|
||||
})
|
||||
return db
|
||||
}
|
||||
|
||||
func TestEventDispatcherQueuesAsyncListenersDurably(t *testing.T) {
|
||||
db := newDispatchWorkerDB(t)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 11, 30, 0, 0, time.UTC) }))
|
||||
ed := NewEventDispatcher(channel.NewDispatcher())
|
||||
ed.SetWorkerPool(wp)
|
||||
syncListener := &dispatchWorkerListener{name: "sync-listener"}
|
||||
asyncListener := &dispatchWorkerListener{name: "async-listener"}
|
||||
ed.RegisterSync(syncListener, string(channel.EventConversationCreated))
|
||||
ed.RegisterAsync(asyncListener, string(channel.EventConversationCreated))
|
||||
|
||||
event := channel.NewChannelEvent(channel.EventConversationCreated, channel.ChannelWebWidget, 1, 2)
|
||||
if err := ed.Dispatch(context.Background(), event); err != nil {
|
||||
t.Fatalf("dispatch: %v", err)
|
||||
}
|
||||
if syncListener.count.Load() != 1 {
|
||||
t.Fatalf("sync listener should run immediately")
|
||||
}
|
||||
if asyncListener.count.Load() != 0 {
|
||||
t.Fatalf("async listener should wait for durable worker")
|
||||
}
|
||||
|
||||
var count int64
|
||||
if err := db.Model(&model.BackgroundJob{}).Where("job_type = ? AND status = ?", TaskTypeEventListenerDispatch, model.BackgroundJobStatusQueued).Count(&count).Error; err != nil {
|
||||
t.Fatalf("count jobs: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("expected one queued listener dispatch job, got %d", count)
|
||||
}
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
if err != nil || !processed {
|
||||
t.Fatalf("process listener job: processed=%v err=%v", processed, err)
|
||||
}
|
||||
if asyncListener.count.Load() != 1 {
|
||||
t.Fatalf("async listener was not called by worker")
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventDispatcherDispatchAsyncQueuesAllMatchingListeners(t *testing.T) {
|
||||
db := newDispatchWorkerDB(t)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 11, 45, 0, 0, time.UTC) }))
|
||||
ed := NewEventDispatcher(channel.NewDispatcher())
|
||||
ed.SetWorkerPool(wp)
|
||||
one := &dispatchWorkerListener{name: "one"}
|
||||
two := &dispatchWorkerListener{name: "two"}
|
||||
ed.RegisterSync(one, string(channel.EventMessageCreated))
|
||||
ed.RegisterAsync(two, string(channel.EventMessageCreated))
|
||||
|
||||
ed.DispatchAsync(context.Background(), channel.NewChannelEvent(channel.EventMessageCreated, channel.ChannelAPI, 1, 2))
|
||||
for i := 0; i < 2; i++ {
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
if err != nil || !processed {
|
||||
t.Fatalf("process listener job %d: processed=%v err=%v", i, processed, err)
|
||||
}
|
||||
}
|
||||
if one.count.Load() != 1 || two.count.Load() != 1 {
|
||||
t.Fatalf("expected both listeners via durable jobs, got one=%d two=%d", one.count.Load(), two.count.Load())
|
||||
}
|
||||
}
|
||||
@@ -48,7 +48,6 @@ type Option func(*WorkerPool)
|
||||
func NewWorkerPool(db ...*gorm.DB) *WorkerPool {
|
||||
wp := &WorkerPool{
|
||||
handlers: make(map[string]JobHandler),
|
||||
queues: []string{model.DefaultBackgroundJobQueue},
|
||||
workerID: fmt.Sprintf("worker-%d", time.Now().UnixNano()),
|
||||
workerCount: 1,
|
||||
pollInterval: 500 * time.Millisecond,
|
||||
|
||||
Reference in New Issue
Block a user