diff --git a/internal/bootstrap/init.go b/internal/bootstrap/init.go index 5d6ed1f..7a44d91 100644 --- a/internal/bootstrap/init.go +++ b/internal/bootstrap/init.go @@ -5,6 +5,7 @@ import ( "cs-agent/internal/pkg/config" "cs-agent/internal/pkg/logx" "cs-agent/internal/services/cronx" + "cs-agent/internal/services/event_handlers" "cs-agent/internal/wxwork" "log/slog" ) @@ -36,6 +37,8 @@ func Init(configPath string) error { return err } + event_handlers.Register() + // 启动任务调度器 cronx.Init() diff --git a/internal/events/notification_events.go b/internal/events/notification_events.go new file mode 100644 index 0000000..b804c51 --- /dev/null +++ b/internal/events/notification_events.go @@ -0,0 +1,29 @@ +package events + +const ( + ConversationAssignTypeAssign = "assign" + ConversationAssignTypeTransfer = "transfer" + ConversationAssignTypeAutoAssign = "auto_assign" +) + +type TicketCreatedEvent struct { + TicketID int64 + OperatorID int64 +} + +type TicketAssignedEvent struct { + TicketID int64 + FromUserID int64 + ToUserID int64 + OperatorID int64 + Reason string +} + +type ConversationAssignedEvent struct { + ConversationID int64 + FromUserID int64 + ToUserID int64 + OperatorID int64 + Reason string + AssignType string +} diff --git a/internal/pkg/eventbus/events.go b/internal/pkg/eventbus/events.go deleted file mode 100644 index 3155a74..0000000 --- a/internal/pkg/eventbus/events.go +++ /dev/null @@ -1,12 +0,0 @@ -package eventbus - -type UserCreated struct { - UserID int64 - Name string -} - -type OrderCreated struct { - OrderID int64 - UserID int64 - Amount int64 -} diff --git a/internal/pkg/eventbus/manager.go b/internal/pkg/eventbus/manager.go index 23caf7f..87e73b9 100644 --- a/internal/pkg/eventbus/manager.go +++ b/internal/pkg/eventbus/manager.go @@ -23,7 +23,7 @@ func Register[T any](opts ...Option[T]) *Bus[T] { return bus.(*Bus[T]) } - created := New[T](opts...) + created := New(opts...) buss[key] = created return created } diff --git a/internal/pkg/eventbus/manager_test.go b/internal/pkg/eventbus/manager_test.go index 3b24e75..6e95e32 100644 --- a/internal/pkg/eventbus/manager_test.go +++ b/internal/pkg/eventbus/manager_test.go @@ -7,11 +7,19 @@ import ( "testing" ) +type userCreatedEvent struct { + UserID int64 +} + +type orderCreatedEvent struct { + OrderID int64 +} + func TestGetReturnsSameBusForSameType(t *testing.T) { resetManagerForTest(t) - first := Get[UserCreated]() - second := Get[UserCreated]() + first := Get[userCreatedEvent]() + second := Get[userCreatedEvent]() if first != second { t.Fatalf("expected same bus instance for same event type") @@ -21,8 +29,8 @@ func TestGetReturnsSameBusForSameType(t *testing.T) { func TestGetReturnsDifferentBusForDifferentTypes(t *testing.T) { resetManagerForTest(t) - userBus := Get[UserCreated]() - orderBus := Get[OrderCreated]() + userBus := Get[userCreatedEvent]() + orderBus := Get[orderCreatedEvent]() if userBus == any(orderBus) { t.Fatalf("expected different bus instances for different event types") @@ -34,21 +42,21 @@ func TestGetCreatesOnlyOneBusUnderConcurrency(t *testing.T) { const workers = 32 - results := make(chan *Bus[UserCreated], workers) + results := make(chan *Bus[userCreatedEvent], workers) var wg sync.WaitGroup for range workers { wg.Add(1) go func() { defer wg.Done() - results <- Get[UserCreated]() + results <- Get[userCreatedEvent]() }() } wg.Wait() close(results) - var first *Bus[UserCreated] + var first *Bus[userCreatedEvent] for bus := range results { if first == nil { first = bus diff --git a/internal/services/conversation_dispatch_service.go b/internal/services/conversation_dispatch_service.go index 44225a0..c4383c5 100644 --- a/internal/services/conversation_dispatch_service.go +++ b/internal/services/conversation_dispatch_service.go @@ -1,11 +1,7 @@ package services import ( - "cs-agent/internal/models" - "cs-agent/internal/pkg/dto" - "cs-agent/internal/pkg/enums" - "cs-agent/internal/pkg/utils" - "cs-agent/internal/repositories" + "context" "errors" "log/slog" "math" @@ -14,6 +10,14 @@ import ( "sync/atomic" "time" + "cs-agent/internal/events" + "cs-agent/internal/models" + "cs-agent/internal/pkg/dto" + "cs-agent/internal/pkg/enums" + "cs-agent/internal/pkg/eventbus" + "cs-agent/internal/pkg/utils" + "cs-agent/internal/repositories" + "github.com/mlogclub/simple/sqls" ) @@ -121,7 +125,13 @@ func (s *conversationDispatchService) DispatchPendingConversation(conversation * "requested_team_ids", report.RequestedTeamIDs, ) WsService.PublishConversationChanged(dispatched, enums.IMRealtimeEventConversationAssigned) - WxWorkNotifyService.NotifyConversationAssigned(dispatched.ID, dispatched.CurrentAssigneeID, "自动分配") + eventbus.PublishAsync(context.Background(), events.ConversationAssignedEvent{ + ConversationID: dispatched.ID, + ToUserID: dispatched.CurrentAssigneeID, + OperatorID: systemDispatchPrincipal().UserID, + Reason: "自动分配", + AssignType: events.ConversationAssignTypeAutoAssign, + }) return dispatched, nil } } diff --git a/internal/services/conversation_service.go b/internal/services/conversation_service.go index b57fecc..013b693 100644 --- a/internal/services/conversation_service.go +++ b/internal/services/conversation_service.go @@ -1,18 +1,21 @@ package services import ( + "context" "crypto/md5" "encoding/hex" "encoding/json" "fmt" "log/slog" + "cs-agent/internal/events" "cs-agent/internal/models" "cs-agent/internal/pkg/constants" "cs-agent/internal/pkg/dto" "cs-agent/internal/pkg/dto/request" "cs-agent/internal/pkg/enums" "cs-agent/internal/pkg/errorsx" + "cs-agent/internal/pkg/eventbus" "cs-agent/internal/pkg/openidentity" "cs-agent/internal/pkg/utils" "cs-agent/internal/repositories" @@ -162,6 +165,7 @@ func (s *conversationService) AssignConversation(req request.AssignConversationR if targetProfile == nil || targetProfile.Status != enums.StatusOk { return errorsx.InvalidParam("目标客服不存在") } + var assignedEvent events.ConversationAssignedEvent if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { conversation := repositories.ConversationRepository.Get(ctx.Tx, req.ConversationID) if conversation == nil { @@ -186,20 +190,31 @@ func (s *conversationService) AssignConversation(req request.AssignConversationR }); err != nil { return err } - return ConversationEventLogService.CreateEvent(ctx, req.ConversationID, enums.IMEventTypeAssign, enums.IMSenderTypeAgent, operator.UserID, "会话已分配", s.buildEventPayload(map[string]any{ + if err := ConversationEventLogService.CreateEvent(ctx, req.ConversationID, enums.IMEventTypeAssign, enums.IMSenderTypeAgent, operator.UserID, "会话已分配", s.buildEventPayload(map[string]any{ "fromStatus": conversation.Status, "toStatus": enums.IMConversationStatusActive, "fromAssigneeId": conversation.CurrentAssigneeID, "toAssigneeId": req.AssigneeID, "reason": strings.TrimSpace(req.Reason), - })) + })); err != nil { + return err + } + assignedEvent = events.ConversationAssignedEvent{ + ConversationID: req.ConversationID, + FromUserID: conversation.CurrentAssigneeID, + ToUserID: req.AssigneeID, + OperatorID: operator.UserID, + Reason: strings.TrimSpace(req.Reason), + AssignType: events.ConversationAssignTypeAssign, + } + return nil }); err != nil { return err } if conversation := s.Get(req.ConversationID); conversation != nil { WsService.PublishConversationChanged(conversation, enums.IMRealtimeEventConversationAssigned) } - WxWorkNotifyService.NotifyConversationAssigned(req.ConversationID, req.AssigneeID, req.Reason) + eventbus.PublishAsync(context.Background(), assignedEvent) return nil } @@ -240,6 +255,7 @@ func (s *conversationService) TransferConversation(conversationID, toUserID int6 if targetProfile == nil || targetProfile.Status != enums.StatusOk { return errorsx.InvalidParam("目标客服不存在") } + var assignedEvent events.ConversationAssignedEvent if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { conversation := repositories.ConversationRepository.Get(ctx.Tx, conversationID) if conversation == nil { @@ -273,20 +289,31 @@ func (s *conversationService) TransferConversation(conversationID, toUserID int6 }); err != nil { return err } - return ConversationEventLogService.CreateEvent(ctx, conversationID, enums.IMEventTypeTransfer, enums.IMSenderTypeAgent, operator.UserID, "会话已转接", s.buildEventPayload(map[string]any{ + if err := ConversationEventLogService.CreateEvent(ctx, conversationID, enums.IMEventTypeTransfer, enums.IMSenderTypeAgent, operator.UserID, "会话已转接", s.buildEventPayload(map[string]any{ "fromStatus": conversation.Status, "toStatus": enums.IMConversationStatusActive, "fromAssigneeId": conversation.CurrentAssigneeID, "toAssigneeId": toUserID, "reason": strings.TrimSpace(reason), - })) + })); err != nil { + return err + } + assignedEvent = events.ConversationAssignedEvent{ + ConversationID: conversationID, + FromUserID: conversation.CurrentAssigneeID, + ToUserID: toUserID, + OperatorID: operator.UserID, + Reason: strings.TrimSpace(reason), + AssignType: events.ConversationAssignTypeTransfer, + } + return nil }); err != nil { return err } if conversation := s.Get(conversationID); conversation != nil { WsService.PublishConversationChanged(conversation, enums.IMRealtimeEventConversationTransferred) } - WxWorkNotifyService.NotifyConversationAssigned(conversationID, toUserID, reason) + eventbus.PublishAsync(context.Background(), assignedEvent) return nil } diff --git a/internal/services/event_handlers/registry.go b/internal/services/event_handlers/registry.go new file mode 100644 index 0000000..384bdac --- /dev/null +++ b/internal/services/event_handlers/registry.go @@ -0,0 +1,11 @@ +package event_handlers + +import "sync" + +var registerOnce sync.Once + +func Register() { + registerOnce.Do(func() { + registerWxWorkNotifyEventHandlers() + }) +} diff --git a/internal/services/event_handlers/wxwork_notify_event_handler.go b/internal/services/event_handlers/wxwork_notify_event_handler.go new file mode 100644 index 0000000..660e0d6 --- /dev/null +++ b/internal/services/event_handlers/wxwork_notify_event_handler.go @@ -0,0 +1,156 @@ +package event_handlers + +import ( + "context" + "fmt" + "log/slog" + "strings" + "time" + + "cs-agent/internal/events" + "cs-agent/internal/models" + "cs-agent/internal/pkg/enums" + "cs-agent/internal/pkg/eventbus" + "cs-agent/internal/services" +) + +func registerWxWorkNotifyEventHandlers() { + eventbus. + Register(eventbus.WithErrorHandler[events.TicketCreatedEvent](handleWxWorkNotifyEventError)). + Subscribe(handleTicketCreatedNotify) + eventbus. + Register(eventbus.WithErrorHandler[events.TicketAssignedEvent](handleWxWorkNotifyEventError)). + Subscribe(handleTicketAssignedNotify) + eventbus. + Register(eventbus.WithErrorHandler[events.ConversationAssignedEvent](handleWxWorkNotifyEventError)). + Subscribe(handleConversationAssignedNotify) +} + +func handleWxWorkNotifyEventError(ctx context.Context, err error) { + slog.Warn("handle wxwork notify event failed", "error", err) +} + +func handleTicketCreatedNotify(ctx context.Context, event events.TicketCreatedEvent) error { + if event.TicketID <= 0 { + return nil + } + ticket := services.TicketService.Get(event.TicketID) + if ticket == nil { + return nil + } + return services.WxWorkNotifyService.SendTextToAssigneeOrDefault(ticket.CurrentAssigneeID, "工单创建提醒", buildTicketCreatedNotifyBody(ticket)) +} + +func handleTicketAssignedNotify(ctx context.Context, event events.TicketAssignedEvent) error { + if event.TicketID <= 0 || event.ToUserID <= 0 { + return nil + } + ticket := services.TicketService.Get(event.TicketID) + if ticket == nil { + return nil + } + return services.WxWorkNotifyService.SendTextToAssigneeOrDefault(event.ToUserID, "工单指派提醒", buildTicketAssignedNotifyBody(ticket, event.ToUserID, event.Reason)) +} + +func handleConversationAssignedNotify(ctx context.Context, event events.ConversationAssignedEvent) error { + if event.ConversationID <= 0 || event.ToUserID <= 0 { + return nil + } + conversation := services.ConversationService.Get(event.ConversationID) + if conversation == nil { + return nil + } + return services.WxWorkNotifyService.SendTextToAssigneeOrDefault(event.ToUserID, conversationAssignedNotifyTitle(event.AssignType), buildConversationAssignedNotifyBody(conversation, event.ToUserID, event.Reason, event.AssignType)) +} + +func conversationAssignedNotifyTitle(assignType string) string { + switch strings.TrimSpace(assignType) { + case events.ConversationAssignTypeTransfer: + return "会话转接提醒" + case events.ConversationAssignTypeAutoAssign: + return "会话自动分配提醒" + default: + return "会话分配提醒" + } +} + +func buildConversationAssignedNotifyBody(conversation *models.Conversation, assigneeID int64, reason string, assignType string) string { + if conversation == nil { + return "" + } + reasonLabel := "分配原因" + if strings.TrimSpace(assignType) == events.ConversationAssignTypeTransfer { + reasonLabel = "转接原因" + } + lines := []string{ + fmt.Sprintf("会话ID: #%d", conversation.ID), + fmt.Sprintf("会话主题: %s", defaultIfBlank(conversation.Subject, "-")), + fmt.Sprintf("接入渠道: %s", enums.GetExternalSourceLabel(conversation.ExternalSource)), + fmt.Sprintf("当前状态: %s", enums.GetIMConversationStatusLabel(conversation.Status)), + fmt.Sprintf("处理人: %s", resolveNotifyUserLabel(assigneeID)), + } + if strings.TrimSpace(reason) != "" { + lines = append(lines, fmt.Sprintf("%s: %s", reasonLabel, strings.TrimSpace(reason))) + } + lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) + return strings.Join(lines, "\n") +} + +func buildTicketCreatedNotifyBody(ticket *models.Ticket) string { + if ticket == nil { + return "" + } + lines := []string{ + fmt.Sprintf("工单号: %s", defaultIfBlank(ticket.TicketNo, fmt.Sprintf("#%d", ticket.ID))), + fmt.Sprintf("工单标题: %s", defaultIfBlank(ticket.Title, "-")), + fmt.Sprintf("工单来源: %s", defaultIfBlank(string(ticket.Source), "-")), + fmt.Sprintf("当前状态: %s", enums.GetTicketStatusLabel(ticket.Status)), + } + if ticket.CurrentAssigneeID > 0 { + lines = append(lines, fmt.Sprintf("处理人: %s", resolveNotifyUserLabel(ticket.CurrentAssigneeID))) + } + lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) + return strings.Join(lines, "\n") +} + +func buildTicketAssignedNotifyBody(ticket *models.Ticket, assigneeID int64, reason string) string { + if ticket == nil { + return "" + } + lines := []string{ + fmt.Sprintf("工单号: %s", defaultIfBlank(ticket.TicketNo, fmt.Sprintf("#%d", ticket.ID))), + fmt.Sprintf("工单标题: %s", defaultIfBlank(ticket.Title, "-")), + fmt.Sprintf("当前状态: %s", enums.GetTicketStatusLabel(ticket.Status)), + fmt.Sprintf("处理人: %s", resolveNotifyUserLabel(assigneeID)), + } + if strings.TrimSpace(reason) != "" { + lines = append(lines, fmt.Sprintf("指派原因: %s", strings.TrimSpace(reason))) + } + lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) + return strings.Join(lines, "\n") +} + +func resolveNotifyUserLabel(userID int64) string { + if userID <= 0 { + return "-" + } + user := services.UserService.Get(userID) + if user == nil { + return fmt.Sprintf("用户#%d", userID) + } + if nickname := strings.TrimSpace(user.Nickname); nickname != "" { + return nickname + } + if username := strings.TrimSpace(user.Username); username != "" { + return username + } + return fmt.Sprintf("用户#%d", userID) +} + +func defaultIfBlank(value, fallback string) string { + value = strings.TrimSpace(value) + if value != "" { + return value + } + return strings.TrimSpace(fallback) +} diff --git a/internal/services/event_handlers/wxwork_notify_event_handler_test.go b/internal/services/event_handlers/wxwork_notify_event_handler_test.go new file mode 100644 index 0000000..8238e0a --- /dev/null +++ b/internal/services/event_handlers/wxwork_notify_event_handler_test.go @@ -0,0 +1,40 @@ +package event_handlers + +import ( + "strings" + "testing" + + "cs-agent/internal/events" + "cs-agent/internal/models" + "cs-agent/internal/pkg/enums" +) + +func TestWxWorkNotifyBuildTicketCreatedNotifyBody(t *testing.T) { + body := buildTicketCreatedNotifyBody(&models.Ticket{ + ID: 12, + TicketNo: "T-12", + Title: "登录失败", + Source: enums.TicketSourceManual, + Status: enums.TicketStatusNew, + }) + + for _, want := range []string{"工单号: T-12", "工单标题: 登录失败", "当前状态:"} { + if !strings.Contains(body, want) { + t.Fatalf("expected body to contain %q, got %q", want, body) + } + } +} + +func TestWxWorkNotifyBuildConversationAssignedNotifyBody(t *testing.T) { + body := buildConversationAssignedNotifyBody(&models.Conversation{ + ID: 7, + Subject: "售后咨询", + Status: enums.IMConversationStatusActive, + }, 0, "客户等待中", events.ConversationAssignTypeTransfer) + + for _, want := range []string{"会话ID: #7", "会话主题: 售后咨询", "转接原因: 客户等待中"} { + if !strings.Contains(body, want) { + t.Fatalf("expected body to contain %q, got %q", want, body) + } + } +} diff --git a/internal/services/ticket_service.go b/internal/services/ticket_service.go index 5610dee..df6fa40 100644 --- a/internal/services/ticket_service.go +++ b/internal/services/ticket_service.go @@ -1,16 +1,19 @@ package services import ( + "context" "encoding/json" "fmt" "strings" "time" + "cs-agent/internal/events" "cs-agent/internal/models" "cs-agent/internal/pkg/dto" "cs-agent/internal/pkg/dto/request" "cs-agent/internal/pkg/enums" "cs-agent/internal/pkg/errorsx" + "cs-agent/internal/pkg/eventbus" "cs-agent/internal/pkg/utils" "cs-agent/internal/repositories" @@ -832,7 +835,10 @@ func (s *ticketService) CreateTicket(req request.CreateTicketRequest, operator * return nil, err } current := s.Get(ticket.ID) - WxWorkNotifyService.NotifyTicketCreated(ticket.ID) + eventbus.PublishAsync(context.Background(), events.TicketCreatedEvent{ + TicketID: ticket.ID, + OperatorID: operator.UserID, + }) return current, nil } @@ -983,29 +989,37 @@ func (s *ticketService) LinkTicketCustomer(ticketID, customerID int64, operator } func (s *ticketService) AssignTicket(req request.AssignTicketRequest, operator *dto.AuthPrincipal) error { + var assignedEvent *events.TicketAssignedEvent if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { - return s.assignTicketTx(ctx.Tx, req, operator) + event, err := s.assignTicketTx(ctx.Tx, req, operator) + if err != nil { + return err + } + assignedEvent = event + return nil }); err != nil { return err } - WxWorkNotifyService.NotifyTicketAssigned(req.TicketID, req.ToUserID, req.Reason) + if assignedEvent != nil { + eventbus.PublishAsync(context.Background(), *assignedEvent) + } return nil } -func (s *ticketService) assignTicketTx(tx *gorm.DB, req request.AssignTicketRequest, operator *dto.AuthPrincipal) error { +func (s *ticketService) assignTicketTx(tx *gorm.DB, req request.AssignTicketRequest, operator *dto.AuthPrincipal) (*events.TicketAssignedEvent, error) { if operator == nil { - return errorsx.Unauthorized("未登录或登录已过期") + return nil, errorsx.Unauthorized("未登录或登录已过期") } ticket := repositories.TicketRepository.Get(tx, req.TicketID) if ticket == nil { - return errorsx.InvalidParam("工单不存在") + return nil, errorsx.InvalidParam("工单不存在") } teamID, assigneeID, err := s.normalizeAssignmentTx(tx, req.ToTeamID, req.ToUserID) if err != nil { - return err + return nil, err } if assigneeID <= 0 { - return errorsx.InvalidParam("目标处理人不能为空") + return nil, errorsx.InvalidParam("目标处理人不能为空") } now := time.Now() eventType := enums.TicketEventTypeAssigned @@ -1026,9 +1040,18 @@ func (s *ticketService) assignTicketTx(tx *gorm.DB, req request.AssignTicketRequ "update_user_name": operator.Username, "updated_at": now, }); err != nil { - return err + return nil, err } - return s.logEvent(tx, ticket.ID, eventType, operator, fmt.Sprintf("%d", ticket.CurrentAssigneeID), fmt.Sprintf("%d", assigneeID), strings.TrimSpace(content), strings.TrimSpace(req.Reason)) + if err := s.logEvent(tx, ticket.ID, eventType, operator, fmt.Sprintf("%d", ticket.CurrentAssigneeID), fmt.Sprintf("%d", assigneeID), strings.TrimSpace(content), strings.TrimSpace(req.Reason)); err != nil { + return nil, err + } + return &events.TicketAssignedEvent{ + TicketID: ticket.ID, + FromUserID: ticket.CurrentAssigneeID, + ToUserID: assigneeID, + OperatorID: operator.UserID, + Reason: strings.TrimSpace(req.Reason), + }, nil } func (s *ticketService) ChangeStatus(req request.ChangeTicketStatusRequest, operator *dto.AuthPrincipal) error { @@ -1457,19 +1480,30 @@ func (s *ticketService) BatchAssignTickets(req request.BatchAssignTicketRequest, if len(ticketIDs) == 0 { return errorsx.InvalidParam("请选择工单") } - return sqls.WithTransaction(func(ctx *sqls.TxContext) error { + assignedEvents := make([]events.TicketAssignedEvent, 0, len(ticketIDs)) + if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { for _, ticketID := range ticketIDs { - if err := s.assignTicketTx(ctx.Tx, request.AssignTicketRequest{ + event, err := s.assignTicketTx(ctx.Tx, request.AssignTicketRequest{ TicketID: ticketID, ToUserID: req.ToUserID, ToTeamID: req.ToTeamID, Reason: req.Reason, - }, operator); err != nil { + }, operator) + if err != nil { return err } + if event != nil { + assignedEvents = append(assignedEvents, *event) + } } return nil - }) + }); err != nil { + return err + } + for _, event := range assignedEvents { + eventbus.PublishAsync(context.Background(), event) + } + return nil } func (s *ticketService) BatchChangeStatus(req request.BatchChangeTicketStatusRequest, operator *dto.AuthPrincipal) error { diff --git a/internal/services/ticket_service_test.go b/internal/services/ticket_service_test.go index b0d537c..a6a44d7 100644 --- a/internal/services/ticket_service_test.go +++ b/internal/services/ticket_service_test.go @@ -1,6 +1,7 @@ package services_test import ( + "context" "fmt" "path/filepath" "strings" @@ -10,11 +11,13 @@ import ( "cs-agent/internal/bootstrap" "cs-agent/internal/builders" + "cs-agent/internal/events" "cs-agent/internal/models" "cs-agent/internal/pkg/config" "cs-agent/internal/pkg/dto" "cs-agent/internal/pkg/dto/request" "cs-agent/internal/pkg/enums" + "cs-agent/internal/pkg/eventbus" "cs-agent/internal/repositories" "cs-agent/internal/services" @@ -61,6 +64,34 @@ func TestCreateTicketSetsTicketNoAndDeadlines(t *testing.T) { } } +func TestCreateTicketPublishesTicketCreatedEvent(t *testing.T) { + setupTicketTestDB(t) + operator := &dto.AuthPrincipal{UserID: 1, Username: "admin"} + eventsCh := make(chan events.TicketCreatedEvent, 1) + _, unsubscribe := eventbus.Subscribe(func(ctx context.Context, event events.TicketCreatedEvent) error { + eventsCh <- event + return nil + }) + defer unsubscribe() + + created, err := services.TicketService.CreateTicket(createTestTicketRequest("event-ticket"), operator) + if err != nil { + t.Fatalf("CreateTicket() error = %v", err) + } + + select { + case event := <-eventsCh: + if event.TicketID != created.ID { + t.Fatalf("expected ticket id %d, got %d", created.ID, event.TicketID) + } + if event.OperatorID != operator.UserID { + t.Fatalf("expected operator id %d, got %d", operator.UserID, event.OperatorID) + } + case <-time.After(time.Second): + t.Fatalf("expected ticket created event") + } +} + func TestAddInternalNoteAllowsMentionSameUserAcrossTickets(t *testing.T) { setupTicketTestDB(t) operator := &dto.AuthPrincipal{UserID: 1, Username: "admin"} @@ -89,6 +120,55 @@ func TestAddInternalNoteAllowsMentionSameUserAcrossTickets(t *testing.T) { } } +func TestBatchAssignTicketsPublishesTicketAssignedEvents(t *testing.T) { + setupTicketTestDB(t) + operator := &dto.AuthPrincipal{UserID: 1, Username: "admin"} + teamID, assigneeID := createTestAgentProfile(t, "batch-event-assignee") + first, err := services.TicketService.CreateTicket(createTestTicketRequest("batch-event-1"), operator) + if err != nil { + t.Fatalf("CreateTicket() first error = %v", err) + } + second, err := services.TicketService.CreateTicket(createTestTicketRequest("batch-event-2"), operator) + if err != nil { + t.Fatalf("CreateTicket() second error = %v", err) + } + + eventsCh := make(chan events.TicketAssignedEvent, 2) + _, unsubscribe := eventbus.Subscribe(func(ctx context.Context, event events.TicketAssignedEvent) error { + eventsCh <- event + return nil + }) + defer unsubscribe() + + if err := services.TicketService.BatchAssignTickets(request.BatchAssignTicketRequest{ + TicketIDs: []int64{first.ID, second.ID}, + ToUserID: assigneeID, + ToTeamID: teamID, + Reason: "batch assign event", + }, operator); err != nil { + t.Fatalf("BatchAssignTickets() error = %v", err) + } + + got := map[int64]events.TicketAssignedEvent{} + for len(got) < 2 { + select { + case event := <-eventsCh: + got[event.TicketID] = event + case <-time.After(time.Second): + t.Fatalf("expected 2 ticket assigned events, got %d", len(got)) + } + } + for _, ticketID := range []int64{first.ID, second.ID} { + event, ok := got[ticketID] + if !ok { + t.Fatalf("missing event for ticket %d", ticketID) + } + if event.ToUserID != assigneeID { + t.Fatalf("expected assignee id %d, got %d", assigneeID, event.ToUserID) + } + } +} + func TestBatchChangeStatusRollsBackOnFailure(t *testing.T) { setupTicketTestDB(t) operator := &dto.AuthPrincipal{UserID: 1, Username: "admin"} diff --git a/internal/services/wxwork_notify_service.go b/internal/services/wxwork_notify_service.go index e6d449e..769e2ff 100644 --- a/internal/services/wxwork_notify_service.go +++ b/internal/services/wxwork_notify_service.go @@ -2,11 +2,8 @@ package services import ( "fmt" - "log/slog" "strings" - "time" - "cs-agent/internal/models" "cs-agent/internal/pkg/config" "cs-agent/internal/pkg/enums" "cs-agent/internal/repositories" @@ -43,57 +40,6 @@ func newWxWorkNotifyService() *wxWorkNotifyService { } } -func (s *wxWorkNotifyService) NotifyConversationAssigned(conversationID, assigneeID int64, reason string) { - if conversationID <= 0 { - return - } - conversation := ConversationService.Get(conversationID) - if conversation == nil { - return - } - if err := s.sendToAssigneeOrDefault(assigneeID, "会话分配提醒", s.buildConversationAssignedBody(conversation, assigneeID, reason)); err != nil { - slog.Warn("send wxwork conversation assignment notify failed", - "conversation_id", conversationID, - "assignee_id", assigneeID, - "error", err, - ) - } -} - -func (s *wxWorkNotifyService) NotifyTicketCreated(ticketID int64) { - if ticketID <= 0 { - return - } - ticket := TicketService.Get(ticketID) - if ticket == nil { - return - } - if err := s.sendToAssigneeOrDefault(ticket.CurrentAssigneeID, "工单创建提醒", s.buildTicketCreatedBody(ticket)); err != nil { - slog.Warn("send wxwork ticket created notify failed", - "ticket_id", ticketID, - "assignee_id", ticket.CurrentAssigneeID, - "error", err, - ) - } -} - -func (s *wxWorkNotifyService) NotifyTicketAssigned(ticketID, assigneeID int64, reason string) { - if ticketID <= 0 || assigneeID <= 0 { - return - } - ticket := TicketService.Get(ticketID) - if ticket == nil { - return - } - if err := s.sendToAssigneeOrDefault(assigneeID, "工单指派提醒", s.buildTicketAssignedBody(ticket, assigneeID, reason)); err != nil { - slog.Warn("send wxwork ticket assigned notify failed", - "ticket_id", ticketID, - "assignee_id", assigneeID, - "error", err, - ) - } -} - func (s *wxWorkNotifyService) Enabled() bool { if !wxwork.Enabled() { return false @@ -101,7 +47,7 @@ func (s *wxWorkNotifyService) Enabled() bool { return config.Current().WxWork.Notify.Enabled } -func (s *wxWorkNotifyService) sendToAssigneeOrDefault(assigneeID int64, title, body string) error { +func (s *wxWorkNotifyService) SendTextToAssigneeOrDefault(assigneeID int64, title, body string) error { if !s.Enabled() { return nil } @@ -175,75 +121,6 @@ func (s *wxWorkNotifyService) defaultRecipients() wxWorkNotifyRecipients { } } -func (s *wxWorkNotifyService) buildConversationAssignedBody(conversation *models.Conversation, assigneeID int64, reason string) string { - if conversation == nil { - return "" - } - lines := []string{ - fmt.Sprintf("会话ID: #%d", conversation.ID), - fmt.Sprintf("会话主题: %s", defaultIfBlank(conversation.Subject, "-")), - fmt.Sprintf("接入渠道: %s", enums.GetExternalSourceLabel(conversation.ExternalSource)), - fmt.Sprintf("当前状态: %s", enums.GetIMConversationStatusLabel(conversation.Status)), - fmt.Sprintf("处理人: %s", s.resolveUserLabel(assigneeID)), - } - if strings.TrimSpace(reason) != "" { - lines = append(lines, fmt.Sprintf("分配原因: %s", strings.TrimSpace(reason))) - } - lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) - return strings.Join(lines, "\n") -} - -func (s *wxWorkNotifyService) buildTicketCreatedBody(ticket *models.Ticket) string { - if ticket == nil { - return "" - } - lines := []string{ - fmt.Sprintf("工单号: %s", defaultIfBlank(ticket.TicketNo, fmt.Sprintf("#%d", ticket.ID))), - fmt.Sprintf("工单标题: %s", defaultIfBlank(ticket.Title, "-")), - fmt.Sprintf("工单来源: %s", defaultIfBlank(string(ticket.Source), "-")), - fmt.Sprintf("当前状态: %s", enums.GetTicketStatusLabel(ticket.Status)), - } - if ticket.CurrentAssigneeID > 0 { - lines = append(lines, fmt.Sprintf("处理人: %s", s.resolveUserLabel(ticket.CurrentAssigneeID))) - } - lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) - return strings.Join(lines, "\n") -} - -func (s *wxWorkNotifyService) buildTicketAssignedBody(ticket *models.Ticket, assigneeID int64, reason string) string { - if ticket == nil { - return "" - } - lines := []string{ - fmt.Sprintf("工单号: %s", defaultIfBlank(ticket.TicketNo, fmt.Sprintf("#%d", ticket.ID))), - fmt.Sprintf("工单标题: %s", defaultIfBlank(ticket.Title, "-")), - fmt.Sprintf("当前状态: %s", enums.GetTicketStatusLabel(ticket.Status)), - fmt.Sprintf("处理人: %s", s.resolveUserLabel(assigneeID)), - } - if strings.TrimSpace(reason) != "" { - lines = append(lines, fmt.Sprintf("指派原因: %s", strings.TrimSpace(reason))) - } - lines = append(lines, fmt.Sprintf("时间: %s", time.Now().Format("2006-01-02 15:04:05"))) - return strings.Join(lines, "\n") -} - -func (s *wxWorkNotifyService) resolveUserLabel(userID int64) string { - if userID <= 0 { - return "-" - } - user := UserService.Get(userID) - if user == nil { - return fmt.Sprintf("用户#%d", userID) - } - if nickname := strings.TrimSpace(user.Nickname); nickname != "" { - return nickname - } - if username := strings.TrimSpace(user.Username); username != "" { - return username - } - return fmt.Sprintf("用户#%d", userID) -} - func (s *wxWorkNotifyService) buildTextContent(title, body string) string { title = strings.TrimSpace(title) body = strings.TrimSpace(body)