feat(eventbus): implement event handling for ticket and conversation assignments
This commit is contained in:
@@ -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()
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -1,12 +0,0 @@
|
||||
package eventbus
|
||||
|
||||
type UserCreated struct {
|
||||
UserID int64
|
||||
Name string
|
||||
}
|
||||
|
||||
type OrderCreated struct {
|
||||
OrderID int64
|
||||
UserID int64
|
||||
Amount int64
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
package event_handlers
|
||||
|
||||
import "sync"
|
||||
|
||||
var registerOnce sync.Once
|
||||
|
||||
func Register() {
|
||||
registerOnce.Do(func() {
|
||||
registerWxWorkNotifyEventHandlers()
|
||||
})
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"}
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user