feat(automation): queue external action deliveries

This commit is contained in:
2026-06-05 16:24:59 +08:00
parent dc48b5cd1f
commit b3fb0b1267
12 changed files with 313 additions and 15 deletions
+9
View File
@@ -16,6 +16,7 @@ import (
channelmodel "github.com/gochat/gochat/internal/model/channel"
"github.com/gochat/gochat/internal/pubsub"
"github.com/gochat/gochat/internal/service"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/driver/postgres"
"gorm.io/gorm"
@@ -31,6 +32,7 @@ type App struct {
engine *gin.Engine
wsHub *ws.Hub
notificationDeliverySvc *service.NotificationDeliveryService
workerPool *worker.WorkerPool
}
// New creates and initializes the application.
@@ -80,6 +82,13 @@ func New(cfg *config.Config) (*App, error) {
// Run starts the notification delivery pipeline and the HTTP server.
func (a *App) Run() error {
if a.workerPool != nil {
if err := a.workerPool.Start(); err != nil {
return fmt.Errorf("failed to start background worker: %w", err)
}
applogger.L().Info("Background worker started")
}
// Start notification delivery service (Watermill router) in background
if a.notificationDeliverySvc != nil {
go func() {
+9 -4
View File
@@ -39,6 +39,7 @@ import (
"github.com/gochat/gochat/internal/search"
"github.com/gochat/gochat/internal/security"
"github.com/gochat/gochat/internal/service"
"github.com/gochat/gochat/internal/worker"
wspkg "github.com/gochat/gochat/internal/ws"
applogger "github.com/gochat/gochat/pkg/logger"
"github.com/spf13/viper"
@@ -103,6 +104,9 @@ func Bootstrap(env string) (*App, error) {
applogger.L().Info("Database migrations completed successfully")
}
workerPool := worker.NewWorkerPool(db)
automation.RegisterActionDeliveryJobs(workerPool, &dbProvider{db: db})
// Step 5: Connect to Redis (ref: Chatwoot config/cable.yml)
rdb, err := NewRedisClient(&cfg.Redis)
if err != nil {
@@ -319,7 +323,7 @@ func Bootstrap(env string) (*App, error) {
notionIntegrationService := service.NewNotionIntegrationService(integrationHookRepo)
// Create channel dispatcher for event-driven architecture (ref: Chatwoot Dispatcher)
channelDispatcher := channel.NewDispatcher()
channelDispatcher := channel.NewDispatcher(workerPool)
// P9: AgentBot rule engine services (declared early so listener can be registered on dispatcher)
botRuleService := automation.NewBotRuleService(&dbProvider{db: db})
@@ -388,7 +392,7 @@ func Bootstrap(env string) (*App, error) {
// P9: Register AgentBot rule listener on the dispatcher
channelDispatcher.Register(agentBotRuleListener)
// M6: Register automation rule listener on the dispatcher for event-triggered automation
automation.RegisterAutomationRuleListener(channelDispatcher, &dbProvider{db: db})
automation.RegisterAutomationRuleListenerWithWorker(channelDispatcher, &dbProvider{db: db}, workerPool)
// M6: Register CSAT survey listener on the dispatcher for conversation.resolved events
channelDispatcher.Register(automation.NewCsatSurveyListener(&dbProvider{db: db}))
@@ -604,8 +608,8 @@ func Bootstrap(env string) (*App, error) {
portalMemberService := service.NewPortalMemberService(portalMemberRepo)
// Automation services (P6 — Automation Rules + Macros + CSAT)
automationRuleService := automation.NewAutomationRuleService(&dbProvider{db: db})
macroService := automation.NewMacroService(&dbProvider{db: db})
automationRuleService := automation.NewAutomationRuleServiceWithWorker(&dbProvider{db: db}, workerPool)
macroService := automation.NewMacroServiceWithWorker(&dbProvider{db: db}, workerPool)
csatSurveyService := automation.NewCsatSurveyService(&dbProvider{db: db})
cannedResponseService := canned.NewCannedResponseService(&dbProvider{db: db})
@@ -879,6 +883,7 @@ func Bootstrap(env string) (*App, error) {
engine: engine,
wsHub: wsHub,
notificationDeliverySvc: notificationDeliverySvc,
workerPool: workerPool,
}, nil
}
+5 -2
View File
@@ -17,8 +17,10 @@ import (
)
const (
defaultActionDeliveryAttempts = 3
defaultActionDeliveryTimeout = 10 * time.Second
defaultActionDeliveryAttempts = 3
defaultActionDeliveryTimeout = 10 * time.Second
TaskTypeAutomationWebhookDelivery = "automation:webhook_delivery"
TaskTypeAutomationTranscriptDelivery = "automation:transcript_delivery"
)
// ActionDeliveryResult is copied into AutomationExecution.action_results so
@@ -30,6 +32,7 @@ type ActionDeliveryResult struct {
ResponseCode int
ResponseBody string
Retryable bool
Queued bool
}
type AutomationWebhookRequest struct {
@@ -0,0 +1,86 @@
package automation
import (
"context"
"encoding/json"
"fmt"
"sync"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
type automationWebhookDeliveryJob struct {
AccountID uint `json:"account_id"`
ConversationID uint `json:"conversation_id"`
EventName string `json:"event_name"`
URL string `json:"url"`
Payload map[string]interface{} `json:"payload"`
}
type automationTranscriptDeliveryJob struct {
AccountID uint `json:"account_id"`
ConversationID uint `json:"conversation_id"`
Recipient string `json:"recipient"`
}
var actionDeliveryRegistrations sync.Map
// RegisterActionDeliveryJobs wires Chatwoot-style automation external side
// effects into the durable worker. Callers may invoke this repeatedly; worker
// handlers are idempotently registered per WorkerPool instance.
func RegisterActionDeliveryJobs(wp *worker.WorkerPool, db DBProvider) {
if wp == nil || db == nil {
return
}
if _, loaded := actionDeliveryRegistrations.LoadOrStore(wp, struct{}{}); loaded {
return
}
runner := &actionDeliveryJobRunner{
db: db,
webhookDeliverer: defaultWebhookDelivererFactory(),
transcriptDeliverer: defaultTranscriptDelivererFactory(),
}
wp.Register(TaskTypeAutomationWebhookDelivery, runner.performWebhookDelivery)
wp.Register(TaskTypeAutomationTranscriptDelivery, runner.performTranscriptDelivery)
}
type actionDeliveryJobRunner struct {
db DBProvider
webhookDeliverer AutomationWebhookDeliverer
transcriptDeliverer AutomationTranscriptDeliverer
}
func (r *actionDeliveryJobRunner) performWebhookDelivery(ctx context.Context, job *model.BackgroundJob) error {
var payload automationWebhookDeliveryJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal automation webhook job: %w", err)
}
_, err := r.webhookDeliverer.DeliverWebhook(ctx, AutomationWebhookRequest{
AccountID: payload.AccountID,
ConversationID: payload.ConversationID,
EventName: payload.EventName,
URL: payload.URL,
Payload: payload.Payload,
})
return err
}
func (r *actionDeliveryJobRunner) performTranscriptDelivery(ctx context.Context, job *model.BackgroundJob) error {
var payload automationTranscriptDeliveryJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal automation transcript job: %w", err)
}
subject, body, err := NewActionService(r.db).buildTranscriptEmail(ctx, payload.AccountID, payload.ConversationID)
if err != nil {
return err
}
_, err = r.transcriptDeliverer.DeliverTranscript(ctx, AutomationTranscriptRequest{
AccountID: payload.AccountID,
ConversationID: payload.ConversationID,
Recipient: payload.Recipient,
Subject: subject,
Body: body,
})
return err
}
+41
View File
@@ -8,6 +8,7 @@ import (
"time"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
)
@@ -29,6 +30,7 @@ type ActionService struct {
db DBProvider
webhookDeliverer AutomationWebhookDeliverer
transcriptDeliverer AutomationTranscriptDeliverer
worker *worker.WorkerPool
}
var defaultWebhookDelivererFactory = func() AutomationWebhookDeliverer {
@@ -48,6 +50,17 @@ func NewActionService(db DBProvider) *ActionService {
}
}
func NewActionServiceWithWorker(db DBProvider, wp *worker.WorkerPool) *ActionService {
s := NewActionService(db)
s.SetWorkerPool(wp)
return s
}
func (s *ActionService) SetWorkerPool(wp *worker.WorkerPool) {
s.worker = wp
RegisterActionDeliveryJobs(wp, s.db)
}
func setAutomationActionDeliverersForTest(webhook AutomationWebhookDeliverer, transcript AutomationTranscriptDeliverer) func() {
originalWebhookFactory := defaultWebhookDelivererFactory
originalTranscriptFactory := defaultTranscriptDelivererFactory
@@ -310,6 +323,19 @@ func (s *ActionService) handleSendWebhookEvent(ctx context.Context, accountID, c
if err != nil {
return ActionDeliveryResult{DeliveryType: "webhook", Target: url}, err
}
if s.worker != nil {
_, err := s.worker.Enqueue(ctx, TaskTypeAutomationWebhookDelivery, automationWebhookDeliveryJob{
AccountID: accountID,
ConversationID: conversationID,
EventName: eventName,
URL: url,
Payload: payload,
}, worker.WithQueue("automation"), worker.WithMaxAttempts(defaultActionDeliveryAttempts))
if err != nil {
return ActionDeliveryResult{DeliveryType: "webhook", Target: url}, err
}
return ActionDeliveryResult{DeliveryType: "webhook", Target: url, ResponseBody: "queued", Queued: true}, nil
}
return s.webhookDeliverer.DeliverWebhook(ctx, AutomationWebhookRequest{
AccountID: accountID,
ConversationID: conversationID,
@@ -391,6 +417,20 @@ func (s *ActionService) handleSendEmailTranscript(ctx context.Context, accountID
return ActionDeliveryResult{DeliveryType: "email_transcript"}, fmt.Errorf("send_email_transcript action requires 'email' param")
}
if s.worker != nil {
for _, recipient := range recipients {
_, err := s.worker.Enqueue(ctx, TaskTypeAutomationTranscriptDelivery, automationTranscriptDeliveryJob{
AccountID: accountID,
ConversationID: conversationID,
Recipient: recipient,
}, worker.WithQueue("automation"), worker.WithMaxAttempts(defaultActionDeliveryAttempts))
if err != nil {
return ActionDeliveryResult{DeliveryType: "email_transcript", Target: strings.Join(recipients, ",")}, err
}
}
return ActionDeliveryResult{DeliveryType: "email_transcript", Target: strings.Join(recipients, ","), ResponseBody: "queued", Queued: true}, nil
}
subject, body, err := s.buildTranscriptEmail(ctx, accountID, conversationID)
if err != nil {
return ActionDeliveryResult{DeliveryType: "email_transcript", Target: strings.Join(recipients, ",")}, err
@@ -546,6 +586,7 @@ func applyDeliveryResult(result *ActionExecutionResult, delivery ActionDeliveryR
result.ResponseCode = delivery.ResponseCode
result.ResponseBody = delivery.ResponseBody
result.Retryable = delivery.Retryable
result.Queued = delivery.Queued
}
func firstStringParam(params map[string]interface{}, keys ...string) string {
@@ -11,6 +11,7 @@ import (
"time"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
type roundTripFunc func(*http.Request) (*http.Response, error)
@@ -143,6 +144,99 @@ func TestActionService_SendEmailTranscript_DeliversSplitRecipients(t *testing.T)
}
}
func TestActionService_SendWebhookEvent_QueuesDurableDelivery(t *testing.T) {
dbProvider := setupAutomationTestDBProvider(t)
db := dbProvider.DB()
accountID, _ := seedTestAccount(db, t)
inboxID := seedTestInbox(db, t, accountID)
contactID := seedTestContact(db, t, accountID)
conversationID := seedTestConversation(db, t, accountID, inboxID, contactID)
webhook := &recordingWebhookDeliverer{result: ActionDeliveryResult{DeliveryType: "webhook", Attempts: 1, ResponseCode: http.StatusOK}}
restore := setAutomationActionDeliverersForTest(webhook, &recordingTranscriptDeliverer{})
defer restore()
if err := db.Create(&model.Message{ConversationID: conversationID, AccountID: accountID, InboxID: inboxID, Content: "queued webhook", ContentType: "text", MessageType: "incoming"}).Error; err != nil {
t.Fatalf("seed message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 12, 30, 0, 0, time.UTC) }))
result, err := NewActionServiceWithWorker(dbProvider, wp).ExecuteWithResult(context.Background(), accountID, conversationID, Action{
ActionName: "send_webhook_event",
ActionParams: map[string]interface{}{"url": "https://hooks.example/queued", "_event_name": "conversation_created"},
}, ActionSourceAutomation, 99)
if err != nil {
t.Fatalf("queue webhook action: %v", err)
}
if len(webhook.requests) != 0 {
t.Fatalf("webhook should not deliver synchronously, got %d requests", len(webhook.requests))
}
if !result.Queued || result.DeliveryType != "webhook" || result.Target != "https://hooks.example/queued" {
t.Fatalf("unexpected queued webhook result: %#v", result)
}
var count int64
if err := db.Model(&model.BackgroundJob{}).Where("job_type = ? AND status = ?", TaskTypeAutomationWebhookDelivery, model.BackgroundJobStatusQueued).Count(&count).Error; err != nil {
t.Fatalf("count webhook jobs: %v", err)
}
if count != 1 {
t.Fatalf("expected one queued webhook job, got %d", count)
}
processed, err := wp.ProcessOne(context.Background())
if err != nil || !processed {
t.Fatalf("process webhook job: processed=%v err=%v", processed, err)
}
if len(webhook.requests) != 1 || webhook.requests[0].URL != "https://hooks.example/queued" || webhook.requests[0].EventName != "conversation_created" {
t.Fatalf("unexpected durable webhook request: %#v", webhook.requests)
}
}
func TestActionService_SendEmailTranscript_QueuesDurableDeliveries(t *testing.T) {
dbProvider := setupAutomationTestDBProvider(t)
db := dbProvider.DB()
accountID, _ := seedTestAccount(db, t)
inboxID := seedTestInbox(db, t, accountID)
contactID := seedTestContact(db, t, accountID)
conversationID := seedTestConversation(db, t, accountID, inboxID, contactID)
transcript := &recordingTranscriptDeliverer{result: ActionDeliveryResult{DeliveryType: "email_transcript", Attempts: 1}}
restore := setAutomationActionDeliverersForTest(&recordingWebhookDeliverer{}, transcript)
defer restore()
if err := db.Create(&model.Message{ConversationID: conversationID, AccountID: accountID, InboxID: inboxID, Content: "queued transcript", ContentType: "text", MessageType: "incoming"}).Error; err != nil {
t.Fatalf("seed message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 12, 45, 0, 0, time.UTC) }))
result, err := NewActionServiceWithWorker(dbProvider, wp).ExecuteWithResult(context.Background(), accountID, conversationID, Action{
ActionName: "send_email_transcript",
ActionParams: map[string]interface{}{"email": "first@example.com, second@example.com"},
}, ActionSourceAutomation, 99)
if err != nil {
t.Fatalf("queue transcript action: %v", err)
}
if len(transcript.requests) != 0 {
t.Fatalf("transcript should not deliver synchronously, got %d requests", len(transcript.requests))
}
if !result.Queued || result.Target != "first@example.com,second@example.com" {
t.Fatalf("unexpected queued transcript result: %#v", result)
}
for i := 0; i < 2; i++ {
processed, err := wp.ProcessOne(context.Background())
if err != nil || !processed {
t.Fatalf("process transcript job %d: processed=%v err=%v", i, processed, err)
}
}
if len(transcript.requests) != 2 {
t.Fatalf("expected two durable transcript deliveries, got %d", len(transcript.requests))
}
if transcript.requests[0].Recipient != "first@example.com" || transcript.requests[1].Recipient != "second@example.com" {
t.Fatalf("unexpected durable transcript recipients: %#v", transcript.requests)
}
if !strings.Contains(transcript.requests[0].Body, "queued transcript") {
t.Fatalf("expected durable transcript body to be rendered at job time: %#v", transcript.requests[0])
}
}
func TestAutomationRuleService_MatchAndExecute_RecordsEmailTranscriptFailureMetadata(t *testing.T) {
dbProvider := setupAutomationTestDBProvider(t)
db := dbProvider.DB()
@@ -45,6 +45,7 @@ func setupAutomationTestDB(t *testing.T) *gorm.DB {
&Macro{},
&MacroExecution{},
&AutomationExecution{},
&model.BackgroundJob{},
&CsatSurveyResponse{},
&ConversationLabel{},
&ConversationMute{},
@@ -29,6 +29,7 @@ type ActionExecutionResult struct {
ResponseCode int `json:"response_code,omitempty"`
ResponseBody string `json:"response_body,omitempty"`
Retryable bool `json:"retryable,omitempty"`
Queued bool `json:"queued,omitempty"`
}
// NewExecutionLogService creates a new ExecutionLogService.
+14
View File
@@ -8,6 +8,7 @@ import (
"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"
)
@@ -34,6 +35,13 @@ func NewAutomationRuleListener(db DBProvider) *AutomationRuleListener {
}
}
func NewAutomationRuleListenerWithWorker(db DBProvider, wp *worker.WorkerPool) *AutomationRuleListener {
return &AutomationRuleListener{
db: db,
ruleService: NewAutomationRuleServiceWithWorker(db, wp),
}
}
// Name returns the unique identifier for this listener.
func (l *AutomationRuleListener) Name() string {
return "automation_rule_listener"
@@ -288,6 +296,12 @@ func RegisterAutomationRuleListener(dispatcher *channel.Dispatcher, db DBProvide
applogger.L().Infof("registered automation rule listener with dispatcher")
}
func RegisterAutomationRuleListenerWithWorker(dispatcher *channel.Dispatcher, db DBProvider, wp *worker.WorkerPool) {
listener := NewAutomationRuleListenerWithWorker(db, wp)
dispatcher.Register(listener)
applogger.L().Infof("registered automation rule listener with durable worker")
}
// String helper for event name validation
func isValidAutomationEventName(name string) bool {
validNames := map[string]bool{
+16 -3
View File
@@ -6,6 +6,7 @@ import (
"strings"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/gorm"
)
@@ -14,7 +15,8 @@ import (
// Reference: Chatwoot Macros::ExecutionService — same action handler pattern as AutomationRule,
// but user-originated (supports 'self' assign). stamps user info.
type MacroService struct {
db DBProvider
db DBProvider
worker *worker.WorkerPool
}
// NewMacroService creates a new MacroService.
@@ -22,6 +24,17 @@ func NewMacroService(db DBProvider) *MacroService {
return &MacroService{db: db}
}
func NewMacroServiceWithWorker(db DBProvider, wp *worker.WorkerPool) *MacroService {
s := NewMacroService(db)
s.SetWorkerPool(wp)
return s
}
func (s *MacroService) SetWorkerPool(wp *worker.WorkerPool) {
s.worker = wp
RegisterActionDeliveryJobs(wp, s.db)
}
// GetByID retrieves a macro by ID.
func (s *MacroService) GetByID(ctx context.Context, id uint) (*Macro, error) {
var macro Macro
@@ -162,7 +175,7 @@ func (s *MacroService) Execute(ctx context.Context, accountID uint, conversation
applogger.L().Infof("executing macro %d (%s) on conversation %d by user %d", macro.ID, macro.Name, conversationID, userID)
actionSvc := NewActionService(s.db)
actionSvc := NewActionServiceWithWorker(s.db, s.worker)
for _, action := range macro.Actions {
// Inject _source_user_id for "self" assignment support
@@ -211,7 +224,7 @@ func (s *MacroService) ExecuteForDisplayIDs(ctx context.Context, accountID uint,
return nil
}
actionSvc := NewActionService(s.db)
actionSvc := NewActionServiceWithWorker(s.db, s.worker)
for _, conversation := range conversations {
applogger.L().Infof("executing macro %d (%s) on conversation %d by user %d", macro.ID, macro.Name, conversation.ID, userID)
for _, action := range macro.Actions {
+15 -2
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"time"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/gorm"
)
@@ -17,7 +18,8 @@ type DBProvider interface {
// AutomationRuleService provides CRUD + condition matching + action execution for automation rules.
// Reference: Chatwoot AutomationRules::ActionService + ConditionsFilterService pattern
type AutomationRuleService struct {
db DBProvider
db DBProvider
worker *worker.WorkerPool
}
// NewAutomationRuleService creates a new AutomationRuleService.
@@ -25,6 +27,17 @@ func NewAutomationRuleService(db DBProvider) *AutomationRuleService {
return &AutomationRuleService{db: db}
}
func NewAutomationRuleServiceWithWorker(db DBProvider, wp *worker.WorkerPool) *AutomationRuleService {
s := NewAutomationRuleService(db)
s.SetWorkerPool(wp)
return s
}
func (s *AutomationRuleService) SetWorkerPool(wp *worker.WorkerPool) {
s.worker = wp
RegisterActionDeliveryJobs(wp, s.db)
}
// GetByID retrieves an automation rule by ID.
func (s *AutomationRuleService) GetByID(ctx context.Context, id uint) (*AutomationRule, error) {
var rule AutomationRule
@@ -233,7 +246,7 @@ func (s *AutomationRuleService) MatchAndExecute(ctx context.Context, accountID u
return fmt.Errorf("failed to load conversation %d: %w", conversationID, err)
}
actionSvc := NewActionService(s.db)
actionSvc := NewActionServiceWithWorker(s.db, s.worker)
logSvc := NewExecutionLogService(s.db)
for _, rule := range rules {