feat(automation): close delayed action parity

This commit is contained in:
2026-06-05 21:42:10 +08:00
parent 43aa2f81eb
commit ddf2694ed3
7 changed files with 464 additions and 39 deletions
+1
View File
@@ -21,6 +21,7 @@ const (
defaultActionDeliveryTimeout = 10 * time.Second
TaskTypeAutomationWebhookDelivery = "automation:webhook_delivery"
TaskTypeAutomationTranscriptDelivery = "automation:transcript_delivery"
TaskTypeAutomationTeamEmailDelivery = "automation:team_email_delivery"
)
// ActionDeliveryResult is copied into AutomationExecution.action_results so
@@ -24,6 +24,13 @@ type automationTranscriptDeliveryJob struct {
Recipient string `json:"recipient"`
}
type automationTeamEmailDeliveryJob struct {
AccountID uint `json:"account_id"`
ConversationID uint `json:"conversation_id"`
TeamID uint `json:"team_id"`
Message string `json:"message"`
}
var actionDeliveryRegistrations sync.Map
// RegisterActionDeliveryJobs wires Chatwoot-style automation external side
@@ -43,6 +50,7 @@ func RegisterActionDeliveryJobs(wp *worker.WorkerPool, db DBProvider) {
}
wp.Register(TaskTypeAutomationWebhookDelivery, runner.performWebhookDelivery)
wp.Register(TaskTypeAutomationTranscriptDelivery, runner.performTranscriptDelivery)
wp.Register(TaskTypeAutomationTeamEmailDelivery, runner.performTeamEmailDelivery)
}
type actionDeliveryJobRunner struct {
@@ -84,3 +92,20 @@ func (r *actionDeliveryJobRunner) performTranscriptDelivery(ctx context.Context,
})
return err
}
func (r *actionDeliveryJobRunner) performTeamEmailDelivery(ctx context.Context, job *model.BackgroundJob) error {
var payload automationTeamEmailDeliveryJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal automation team email job: %w", err)
}
requests, err := NewActionService(r.db).buildTeamEmailRequests(ctx, payload.AccountID, payload.ConversationID, []uint{payload.TeamID}, payload.Message)
if err != nil {
return err
}
for _, req := range requests {
if _, err := r.transcriptDeliverer.DeliverTranscript(ctx, req); err != nil {
return err
}
}
return nil
}
+243
View File
@@ -4,12 +4,14 @@ import (
"context"
"encoding/json"
"fmt"
"strconv"
"strings"
"time"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/gorm"
)
// ActionSource identifies who triggered the action execution.
@@ -101,6 +103,8 @@ func (s *ActionService) ExecuteWithResult(ctx context.Context, accountID uint, c
switch resolvedAction.ActionName {
case "send_message":
err = s.handleSendMessage(ctx, accountID, conversationID, resolvedAction, source, sourceID)
case "send_email_to_team":
deliveryResult, err = s.handleSendEmailToTeam(ctx, accountID, conversationID, resolvedAction)
case "add_label":
err = s.handleAddLabel(ctx, accountID, conversationID, resolvedAction)
case "remove_label":
@@ -141,6 +145,8 @@ func (s *ActionService) ExecuteWithResult(ctx context.Context, accountID uint, c
})
case "change_priority":
err = s.handleChangePriority(ctx, accountID, conversationID, resolvedAction)
case "add_sla":
err = s.handleAddSla(ctx, accountID, conversationID, resolvedAction)
case "send_email_transcript":
deliveryResult, err = s.handleSendEmailTranscript(ctx, accountID, conversationID, resolvedAction)
case "send_attachment":
@@ -210,6 +216,49 @@ func (s *ActionService) handleSendMessage(ctx context.Context, accountID, conver
return s.db.DB().WithContext(ctx).Create(msg).Error
}
// handleSendEmailToTeam sends Chatwoot-style automation team notifications.
// Reference: AutomationRules::ActionService#send_email_to_team receives
// { team_ids, message } and sends one email to each member of each team.
func (s *ActionService) handleSendEmailToTeam(ctx context.Context, accountID, conversationID uint, action Action) (ActionDeliveryResult, error) {
teamIDs := extractUintSlice(action.ActionParams, "team_ids", "team_id")
message := firstStringParam(action.ActionParams, "message", "content")
if len(teamIDs) == 0 {
return ActionDeliveryResult{DeliveryType: "team_email"}, fmt.Errorf("send_email_to_team action requires 'team_ids' param")
}
if s.worker != nil {
for _, teamID := range teamIDs {
_, err := s.worker.Enqueue(ctx, TaskTypeAutomationTeamEmailDelivery, automationTeamEmailDeliveryJob{
AccountID: accountID,
ConversationID: conversationID,
TeamID: teamID,
Message: message,
}, worker.WithQueue("automation"), worker.WithMaxAttempts(defaultActionDeliveryAttempts))
if err != nil {
return ActionDeliveryResult{DeliveryType: "team_email", Target: joinUintTargets(teamIDs)}, err
}
}
return ActionDeliveryResult{DeliveryType: "team_email", Target: joinUintTargets(teamIDs), ResponseBody: "queued", Queued: true}, nil
}
requests, err := s.buildTeamEmailRequests(ctx, accountID, conversationID, teamIDs, message)
if err != nil {
return ActionDeliveryResult{DeliveryType: "team_email", Target: joinUintTargets(teamIDs)}, err
}
aggregate := ActionDeliveryResult{DeliveryType: "team_email", Target: joinUintTargets(teamIDs)}
for _, req := range requests {
result, err := s.transcriptDeliverer.DeliverTranscript(ctx, req)
aggregate.Attempts += result.Attempts
aggregate.ResponseCode = result.ResponseCode
aggregate.ResponseBody = result.ResponseBody
aggregate.Retryable = result.Retryable
if err != nil {
return aggregate, err
}
}
return aggregate, nil
}
// handleAddLabel adds a label to the conversation.
// Reference: Chatwoot add_label action — adds tag/label to conversation
func (s *ActionService) handleAddLabel(ctx context.Context, accountID, conversationID uint, action Action) error {
@@ -409,6 +458,52 @@ func (s *ActionService) handleChangePriority(ctx context.Context, accountID, con
Update("priority", priority).Error
}
// handleAddSla attaches an account-scoped SLA policy when the conversation
// does not already have one. Reference: Enterprise::ActionService#add_sla.
func (s *ActionService) handleAddSla(ctx context.Context, accountID, conversationID uint, action Action) error {
slaPolicyID := extractUintParam(action.ActionParams, "sla_policy_id")
if slaPolicyID == 0 {
return fmt.Errorf("add_sla action requires 'sla_policy_id' param")
}
var policy model.SlaPolicy
if err := s.db.DB().WithContext(ctx).Where("id = ? AND account_id = ?", slaPolicyID, accountID).First(&policy).Error; err != nil {
return fmt.Errorf("sla policy not found: %w", err)
}
var conversation model.Conversation
if err := s.db.DB().WithContext(ctx).Where("id = ? AND account_id = ?", conversationID, accountID).First(&conversation).Error; err != nil {
return fmt.Errorf("conversation not found: %w", err)
}
if conversation.SlaPolicyID != nil && *conversation.SlaPolicyID != 0 {
return nil
}
return s.db.DB().WithContext(ctx).Transaction(func(tx *gorm.DB) error {
result := tx.Model(&model.Conversation{}).
Where("id = ? AND account_id = ? AND (sla_policy_id IS NULL OR sla_policy_id = 0)", conversationID, accountID).
Update("sla_policy_id", slaPolicyID)
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return nil
}
if err := tx.Where("id = ? AND account_id = ?", conversationID, accountID).First(&conversation).Error; err != nil {
return err
}
var count int64
if err := tx.Model(&model.AppliedSLA{}).Where("account_id = ? AND conversation_id = ?", accountID, conversationID).Count(&count).Error; err != nil {
return err
}
if count > 0 {
return nil
}
applied := appliedSlaFromAutomationPolicy(accountID, conversationID, conversation, policy)
return tx.Create(applied).Error
})
}
// handleSendEmailTranscript sends an email transcript of the conversation.
// Reference: Chatwoot send_email_transcript action splits comma-separated emails.
func (s *ActionService) handleSendEmailTranscript(ctx context.Context, accountID, conversationID uint, action Action) (ActionDeliveryResult, error) {
@@ -576,6 +671,73 @@ func extractUintParam(params map[string]interface{}, key string) uint {
}
}
func extractUintSlice(params map[string]interface{}, keys ...string) []uint {
for _, key := range keys {
raw, ok := params[key]
if !ok || raw == nil {
continue
}
switch v := raw.(type) {
case []uint:
return v
case []int:
out := make([]uint, 0, len(v))
for _, item := range v {
if item > 0 {
out = append(out, uint(item))
}
}
return out
case []interface{}:
out := make([]uint, 0, len(v))
for _, item := range v {
if parsed := uintFromAny(item); parsed > 0 {
out = append(out, parsed)
}
}
return out
default:
if parsed := uintFromAny(v); parsed > 0 {
return []uint{parsed}
}
}
}
return nil
}
func uintFromAny(raw interface{}) uint {
switch v := raw.(type) {
case uint:
return v
case int:
if v > 0 {
return uint(v)
}
case int64:
if v > 0 {
return uint(v)
}
case float64:
if v > 0 {
return uint(v)
}
case string:
parsed, err := strconv.ParseUint(strings.TrimSpace(v), 10, 64)
if err == nil && parsed > 0 {
return uint(parsed)
}
}
return 0
}
func joinUintTargets(values []uint) string {
parts := make([]string, 0, len(values))
for _, value := range values {
parts = append(parts, fmt.Sprintf("%d", value))
}
return strings.Join(parts, ",")
}
func applyDeliveryResult(result *ActionExecutionResult, delivery ActionDeliveryResult) {
if delivery.DeliveryType == "" && delivery.Target == "" && delivery.Attempts == 0 && delivery.ResponseCode == 0 && delivery.ResponseBody == "" && !delivery.Retryable {
return
@@ -748,6 +910,87 @@ func (s *ActionService) buildTranscriptEmail(ctx context.Context, accountID, con
return subject, body.String(), nil
}
func (s *ActionService) buildTeamEmailRequests(ctx context.Context, accountID, conversationID uint, teamIDs []uint, message string) ([]AutomationTranscriptRequest, error) {
var conversation model.Conversation
if err := s.db.DB().WithContext(ctx).
Where("id = ? AND account_id = ?", conversationID, accountID).
First(&conversation).Error; err != nil {
return nil, err
}
displayID := conversation.ID
if conversation.DisplayID != nil && *conversation.DisplayID > 0 {
displayID = *conversation.DisplayID
}
subject := fmt.Sprintf("Conversation (#%d) automation notification", displayID)
body := strings.TrimSpace(message)
if body == "" {
body = fmt.Sprintf("Conversation #%d matched an automation rule.", displayID)
}
requests := make([]AutomationTranscriptRequest, 0)
seen := map[string]bool{}
for _, teamID := range teamIDs {
var team model.Team
if err := s.db.DB().WithContext(ctx).Where("id = ? AND account_id = ?", teamID, accountID).First(&team).Error; err != nil {
return nil, err
}
var users []model.User
if err := s.db.DB().WithContext(ctx).
Joins("INNER JOIN team_members ON team_members.user_id = users.id").
Joins("INNER JOIN account_users ON account_users.user_id = users.id AND account_users.account_id = ?", accountID).
Where("team_members.team_id = ? AND users.email <> ''", teamID).
Order("users.id ASC").
Find(&users).Error; err != nil {
return nil, err
}
for _, user := range users {
recipient := strings.TrimSpace(user.Email)
if recipient == "" || seen[recipient] {
continue
}
seen[recipient] = true
requests = append(requests, AutomationTranscriptRequest{
AccountID: accountID,
ConversationID: conversationID,
Recipient: recipient,
Subject: subject,
Body: body,
})
}
}
return requests, nil
}
func appliedSlaFromAutomationPolicy(accountID, conversationID uint, conversation model.Conversation, policy model.SlaPolicy) *model.AppliedSLA {
applied := &model.AppliedSLA{AccountID: accountID, ConversationID: conversationID, SlaPolicyID: policy.ID, SLAStatus: model.SLAStatusActive}
if policy.FirstResponseTimeThreshold > 0 {
frt := conversation.CreatedAt.Add(time.Duration(policy.FirstResponseTimeThreshold) * time.Second)
applied.FRTTargetAt = &frt
}
if policy.NextResponseTimeThreshold > 0 {
if applied.FRTTargetAt != nil {
nrt := applied.FRTTargetAt.Add(time.Duration(policy.NextResponseTimeThreshold) * time.Second)
applied.NRTTargetAt = &nrt
} else {
nrt := conversation.CreatedAt.Add(time.Duration(policy.NextResponseTimeThreshold) * time.Second)
applied.NRTTargetAt = &nrt
}
}
if policy.ResolutionTimeThreshold > 0 {
rt := conversation.CreatedAt.Add(time.Duration(policy.ResolutionTimeThreshold) * time.Second)
applied.RTTargetAt = &rt
}
if conversation.FirstReplyCreatedAt != nil {
frtActual := time.Unix(*conversation.FirstReplyCreatedAt, 0)
applied.FRTActualAt = &frtActual
}
if conversation.ResolvedAt != nil {
applied.RTActualAt = conversation.ResolvedAt
}
return applied
}
func jsonObject(raw []byte) map[string]interface{} {
if len(raw) == 0 {
return map[string]interface{}{}
+106
View File
@@ -4,6 +4,7 @@ import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"strings"
@@ -237,6 +238,111 @@ func TestActionService_SendEmailTranscript_QueuesDurableDeliveries(t *testing.T)
}
}
func TestActionService_SendEmailToTeam_QueuesDurableTeamNotifications(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)
team := &model.Team{AccountID: accountID, Name: "Escalations"}
if err := db.Create(team).Error; err != nil {
t.Fatalf("seed team: %v", err)
}
user := &model.User{AccountID: accountID, Name: "Team Agent", Email: "team-agent@example.com"}
if err := db.Create(user).Error; err != nil {
t.Fatalf("seed user: %v", err)
}
if err := db.Create(&model.AccountUser{AccountID: accountID, UserID: user.ID, Role: "agent"}).Error; err != nil {
t.Fatalf("seed account user: %v", err)
}
if err := db.Create(&model.TeamMember{TeamID: team.ID, UserID: user.ID}).Error; err != nil {
t.Fatalf("seed team member: %v", err)
}
transcript := &recordingTranscriptDeliverer{result: ActionDeliveryResult{DeliveryType: "team_email", Attempts: 1}}
restore := setAutomationActionDeliverersForTest(&recordingWebhookDeliverer{}, transcript)
defer restore()
wp := worker.NewWorkerPool(db)
result, err := NewActionServiceWithWorker(dbProvider, wp).ExecuteWithResult(context.Background(), accountID, conversationID, Action{
ActionName: "send_email_to_team",
ActionParams: map[string]interface{}{
"team_ids": []interface{}{float64(team.ID)},
"message": "Please check this conversation",
},
}, ActionSourceAutomation, 99)
if err != nil {
t.Fatalf("queue team email action: %v", err)
}
if !result.Queued || result.DeliveryType != "team_email" || result.Target != fmt.Sprintf("%d", team.ID) {
t.Fatalf("unexpected queued team email result: %#v", result)
}
if len(transcript.requests) != 0 {
t.Fatalf("team email should not deliver synchronously, got %d", len(transcript.requests))
}
processed, err := wp.ProcessOne(context.Background())
if err != nil || !processed {
t.Fatalf("process team email job: processed=%v err=%v", processed, err)
}
if len(transcript.requests) != 1 || transcript.requests[0].Recipient != "team-agent@example.com" {
t.Fatalf("unexpected durable team email requests: %#v", transcript.requests)
}
if !strings.Contains(transcript.requests[0].Body, "Please check this conversation") {
t.Fatalf("expected team email message body, got %#v", transcript.requests[0])
}
}
func TestActionService_AddSla_AttachesPolicyAndAppliedSLA(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)
policy := &model.SlaPolicy{AccountID: accountID, Name: "Gold", FirstResponseTimeThreshold: 60, ResolutionTimeThreshold: 3600}
if err := db.Create(policy).Error; err != nil {
t.Fatalf("seed sla policy: %v", err)
}
_, err := NewActionService(dbProvider).ExecuteWithResult(context.Background(), accountID, conversationID, Action{
ActionName: "add_sla",
ActionParams: map[string]interface{}{"sla_policy_id": float64(policy.ID)},
}, ActionSourceAutomation, 99)
if err != nil {
t.Fatalf("add sla action: %v", err)
}
var conversation model.Conversation
if err := db.First(&conversation, conversationID).Error; err != nil {
t.Fatalf("reload conversation: %v", err)
}
if conversation.SlaPolicyID == nil || *conversation.SlaPolicyID != policy.ID {
t.Fatalf("expected conversation sla_policy_id %d, got %#v", policy.ID, conversation.SlaPolicyID)
}
var applied model.AppliedSLA
if err := db.Where("conversation_id = ?", conversationID).First(&applied).Error; err != nil {
t.Fatalf("expected applied sla: %v", err)
}
if applied.SlaPolicyID != policy.ID || applied.FRTTargetAt == nil || applied.RTTargetAt == nil {
t.Fatalf("unexpected applied sla: %#v", applied)
}
_, err = NewActionService(dbProvider).ExecuteWithResult(context.Background(), accountID, conversationID, Action{
ActionName: "add_sla",
ActionParams: map[string]interface{}{"sla_policy_id": float64(policy.ID)},
}, ActionSourceAutomation, 99)
if err != nil {
t.Fatalf("repeat add sla action should be idempotent: %v", err)
}
var count int64
if err := db.Model(&model.AppliedSLA{}).Where("conversation_id = ?", conversationID).Count(&count).Error; err != nil {
t.Fatalf("count applied slas: %v", err)
}
if count != 1 {
t.Fatalf("expected one applied sla after repeat, got %d", count)
}
}
func TestAutomationRuleService_MatchAndExecute_RecordsEmailTranscriptFailureMetadata(t *testing.T) {
dbProvider := setupAutomationTestDBProvider(t)
db := dbProvider.DB()
@@ -41,6 +41,8 @@ func setupAutomationTestDB(t *testing.T) *gorm.DB {
&model.Message{},
&model.Team{},
&model.TeamMember{},
&model.SlaPolicy{},
&model.AppliedSLA{},
&AutomationRule{},
&Macro{},
&MacroExecution{},
@@ -426,8 +426,10 @@ func actionArrayToMap(actionName string, values []interface{}) map[string]interf
switch actionName {
case "assign_agent":
params["assignee_id"] = first
case "assign_team", "send_email_to_team":
case "assign_team":
params["team_id"] = first
case "send_email_to_team":
params["team_ids"] = values
case "add_label", "remove_label":
params["labels"] = valuesToStrings(values)
case "change_status":
@@ -450,7 +452,7 @@ func actionArrayToMap(actionName string, values []interface{}) map[string]interf
return params
}
func chatwootActionParams(action automation.Action) []interface{} {
func chatwootActionParams(action automation.Action) interface{} {
params := action.ActionParams
if params == nil {
return []interface{}{}
@@ -458,8 +460,13 @@ func chatwootActionParams(action automation.Action) []interface{} {
switch action.ActionName {
case "assign_agent":
return compactValues(params["assignee_id"], params["agent_id"])
case "assign_team", "send_email_to_team":
case "assign_team":
return compactValues(params["team_id"])
case "send_email_to_team":
return gin.H{
"team_ids": valuesToInterfaces(firstNonNil(params["team_ids"], params["team_id"])),
"message": firstNonNil(params["message"], params["content"]),
}
case "add_label", "remove_label":
if labels, ok := params["labels"]; ok {
return valuesToInterfaces(labels)
@@ -487,6 +494,15 @@ func chatwootActionParams(action automation.Action) []interface{} {
}
}
func firstNonNil(values ...interface{}) interface{} {
for _, value := range values {
if value != nil {
return value
}
}
return nil
}
func firstValue(values []interface{}) interface{} {
if len(values) == 0 {
return nil
@@ -513,6 +529,18 @@ func valuesToInterfaces(value interface{}) []interface{} {
result = append(result, item)
}
return result
case []uint:
result := make([]interface{}, 0, len(typed))
for _, item := range typed {
result = append(result, item)
}
return result
case []int:
result := make([]interface{}, 0, len(typed))
for _, item := range typed {
result = append(result, item)
}
return result
case nil:
return []interface{}{}
default: