package services import ( "context" "encoding/json" "log/slog" "code.tczkiot.com/wlw/ai-agent/identity" "code.tczkiot.com/wlw/ai-agent/internal/events" "code.tczkiot.com/wlw/ai-agent/internal/models" "code.tczkiot.com/wlw/ai-agent/internal/pkg/dto" "code.tczkiot.com/wlw/ai-agent/internal/pkg/dto/request" "code.tczkiot.com/wlw/ai-agent/internal/pkg/enums" "code.tczkiot.com/wlw/ai-agent/internal/pkg/errorsx" "code.tczkiot.com/wlw/ai-agent/internal/pkg/eventbus" "code.tczkiot.com/wlw/ai-agent/internal/pkg/openidentity" "code.tczkiot.com/wlw/ai-agent/internal/pkg/utils" "code.tczkiot.com/wlw/ai-agent/internal/repositories" "strings" "time" "github.com/mlogclub/simple/common/strs" "github.com/mlogclub/simple/sqls" "gorm.io/gorm" ) var ConversationService = newConversationService() func newConversationService() *conversationService { return &conversationService{} } type conversationService struct { } func (s *conversationService) Get(id int64) *models.Conversation { if id <= 0 { return nil } return repositories.ConversationRepository.Get(sqls.DB(), id) } func (s *conversationService) Find(cnd *sqls.Cnd) []models.Conversation { return repositories.ConversationRepository.Find(sqls.DB(), cnd) } func (s *conversationService) FindOne(cnd *sqls.Cnd) *models.Conversation { return repositories.ConversationRepository.FindOne(sqls.DB(), cnd) } func (s *conversationService) FindPageByCnd(cnd *sqls.Cnd) (list []models.Conversation, paging *sqls.Paging) { return repositories.ConversationRepository.FindPageByCnd(sqls.DB(), cnd) } func (s *conversationService) ListConversations(userID int64, filter request.AgentConversationFilter, keyword string, paging *sqls.Paging) ([]models.Conversation, *sqls.Paging, error) { cnd := sqls.NewCnd().Page(paging.Page, paging.Limit) if strs.IsNotBlank(keyword) { keyword = strings.TrimSpace(keyword) keywordLike := "%" + keyword + "%" cnd.Where("customer_name LIKE ? OR last_message_summary LIKE ?", keywordLike, keywordLike) } switch filter { case request.AgentConversationFilterAIServing: cnd.Eq("current_assignee_id", 0).Eq("status", enums.IMConversationStatusAIServing).Desc("last_active_at").Desc("id") case request.AgentConversationFilterMine: cnd.Eq("current_assignee_id", userID).Desc("last_active_at").Desc("id") case request.AgentConversationFilterActive: cnd.Eq("current_assignee_id", userID).Eq("status", enums.IMConversationStatusActive).Desc("last_active_at").Desc("id") case request.AgentConversationFilterPending: cnd.Eq("current_assignee_id", 0).Eq("status", enums.IMConversationStatusPending).Asc("last_active_at").Desc("id") case request.AgentConversationFilterClosed: cnd.Eq("current_assignee_id", userID).Eq("status", enums.IMConversationStatusClosed).Desc("last_active_at").Desc("id") default: return nil, nil, errorsx.InvalidParamI18n("error.e0121") } list, paging := repositories.ConversationRepository.FindPageByCnd(sqls.DB(), cnd) return list, paging, nil } func (s *conversationService) Updates(id int64, columns map[string]interface{}) error { return repositories.ConversationRepository.Updates(sqls.DB(), id, columns) } func (s *conversationService) getLatestNotFinishedByExternalUser(db *gorm.DB, externalUser openidentity.ExternalUser, channelID int64) *models.Conversation { externalID := strings.TrimSpace(externalUser.ExternalID) if externalID == "" || channelID <= 0 { return nil } cnd := sqls.NewCnd() cnd.Eq("channel_id", channelID) cnd.Eq("customer_type", externalCustomerType(externalUser)) cnd.Eq("customer_external_id", externalID) cnd.In("status", []enums.IMConversationStatus{ enums.IMConversationStatusAIServing, enums.IMConversationStatusPending, enums.IMConversationStatusActive, }) cnd.Desc("id") return repositories.ConversationRepository.FindOne(db, cnd) } func (s *conversationService) Create(externalUser openidentity.ExternalUser, channelID, aiAgentID int64) (*models.Conversation, error) { serviceMode := enums.IMConversationServiceModeHumanOnly var aiAgent *models.AIAgent if aiAgentID > 0 { aiAgent = AIAgentService.Get(aiAgentID) if aiAgent == nil || aiAgent.Status != enums.StatusOk { return nil, errorsx.InvalidParamI18n("error.e0002") } serviceMode = aiAgent.ServiceMode } var conversation *models.Conversation var welcomeMessage *models.Message created := false reconfigured := false if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { customerType := externalCustomerType(externalUser) customerID := externalUser.SubjectID customerName := strings.TrimSpace(externalUser.ExternalName) existing := s.getLatestNotFinishedByExternalUser(ctx.Tx, externalUser, channelID) // A conversation already being handled by a human keeps its original // service contract. If the channel was switched to another Agent/mode, // start a new conversation with the latest config instead of silently // reusing the stale human conversation. if existing != nil && (existing.CurrentAssigneeID > 0 || existing.HandoffAt != nil) && (existing.AIAgentID != aiAgentID || existing.ServiceMode != serviceMode) { existing = nil } if existing != nil { conversation = existing updates := make(map[string]any) if customerName != "" && existing.CustomerName != customerName { updates["customer_name"] = customerName conversation.CustomerName = customerName } // A channel binding may change after a conversation was created. Keep an // unassigned conversation aligned with the latest channel/Agent config, // while never taking a conversation away from a human or a handoff flow. if existing.CurrentAssigneeID == 0 && existing.HandoffAt == nil && (existing.ChannelID != channelID || existing.AIAgentID != aiAgentID || existing.ServiceMode != serviceMode) { updates["channel_id"] = channelID updates["ai_agent_id"] = aiAgentID updates["service_mode"] = serviceMode updates["status"] = s.resolveInitialStatus(serviceMode) conversation.ChannelID = channelID conversation.AIAgentID = aiAgentID conversation.ServiceMode = serviceMode conversation.Status = s.resolveInitialStatus(serviceMode) reconfigured = true } if len(updates) > 0 { updates["updated_at"] = time.Now() if err := repositories.ConversationRepository.Updates(ctx.Tx, existing.ID, updates); err != nil { return err } } return nil } created = true now := time.Now() conversation = &models.Conversation{ AIAgentID: aiAgentID, ChannelID: channelID, CustomerType: customerType, CustomerID: customerID, CustomerExternalID: strings.TrimSpace(externalUser.ExternalID), CustomerName: customerName, Status: s.resolveInitialStatus(serviceMode), ServiceMode: serviceMode, Priority: 0, CurrentAssigneeID: 0, CurrentTeamID: 0, LastMessageAt: now, LastActiveAt: now, AuditFields: utils.BuildAuditFields(nil), } if err := ctx.Tx.Create(conversation).Error; err != nil { return err } if err := ConversationParticipantService.CreateCustomerParticipant(ctx, conversation.ID, externalUser); err != nil { return err } if err := ConversationEventLogService.CreateEvent(ctx, conversation.ID, enums.IMEventTypeCreate, enums.IMSenderTypeCustomer, 0, "用户创建会话", ""); err != nil { return err } if aiAgent != nil { var welcomeErr error welcomeMessage, welcomeErr = MessageService.createAIWelcomeMessage(ctx, conversation, aiAgent, now) return welcomeErr } return nil }); err != nil { return nil, err } if conversation == nil { return nil, errorsx.BusinessErrorI18n(1, "error.conversation.createFailed") } if !created { if reconfigured { WsService.PublishConversationChanged(conversation, enums.IMRealtimeEventConversationUpdated) } return s.Get(conversation.ID), nil } // 推送会话创建事件 WsService.PublishConversationChanged(conversation, enums.IMRealtimeEventConversationCreated) if welcomeMessage != nil { if updatedConversation := s.Get(conversation.ID); updatedConversation != nil { WsService.PublishMessageCreated(updatedConversation, welcomeMessage) WsService.PublishConversationChanged(updatedConversation, enums.IMRealtimeEventConversationUpdated) } } if serviceMode == enums.IMConversationServiceModeHumanOnly { var err error if aiAgent == nil { _, err = ConversationHumanDispatchService.ApplyHumanChannelCreate(conversation.ID) } else { _, err = ConversationHumanDispatchService.ApplyHumanOnlyCreate(conversation.ID, *aiAgent) } if err != nil { return nil, err } } return s.Get(conversation.ID), nil } func (s *conversationService) AssignConversation(req request.AssignConversationRequest, operator *dto.AuthPrincipal) error { if operator == nil { return errorsx.UnauthorizedI18n("error.auth.expired") } targetProfile := AgentProfileService.GetByUserID(req.AssigneeID) if targetProfile == nil || targetProfile.Status != enums.StatusOk { return errorsx.InvalidParamI18n("error.e0276") } var assignedEvent events.ConversationAssignedEvent var previousTeamID int64 if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { conversation := repositories.ConversationRepository.Get(ctx.Tx, req.ConversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if conversation.Status != enums.IMConversationStatusPending { return errorsx.InvalidParamI18n("error.e0135") } previousTeamID = conversation.CurrentTeamID now := time.Now() if err := ConversationAssignmentService.FinishActiveAssignments(ctx, req.ConversationID, now); err != nil { return err } if err := ConversationAssignmentService.CreateAssignment(ctx, req.ConversationID, conversation.CurrentAssigneeID, req.AssigneeID, enums.IMAssignmentTypeAssign, req.Reason, operator, now); err != nil { return err } if err := repositories.ConversationRepository.Updates(ctx.Tx, req.ConversationID, map[string]any{ "current_assignee_id": req.AssigneeID, "current_team_id": targetProfile.TeamID, "status": enums.IMConversationStatusActive, "update_user_id": operator.UserID, "update_user_name": operator.Username, "updated_at": now, }); err != nil { return err } if err := ConversationEventLogService.CreateEvent(ctx, req.ConversationID, enums.IMEventTypeAssign, enums.IMSenderTypeAgent, operator.UserID, "会话已分配", s.buildEventPayload(map[string]any{ "from_status": conversation.Status, "to_status": enums.IMConversationStatusActive, "from_assignee_id": conversation.CurrentAssigneeID, "to_assignee_id": req.AssigneeID, "to_team_id": targetProfile.TeamID, "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) } ConversationQueueService.PublishPoolUpdates(previousTeamID) eventbus.PublishAsync(context.Background(), assignedEvent) return nil } func (s *conversationService) AutoAssignConversation(conversationID int64, operator *dto.AuthPrincipal) error { if operator == nil { return errorsx.UnauthorizedI18n("error.auth.expired") } conversation := s.Get(conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if conversation.Status != enums.IMConversationStatusPending { return errorsx.InvalidParamI18n("error.e0136") } if conversation.CurrentAssigneeID > 0 { return errorsx.InvalidParamI18n("error.e0190") } if conversation.AIAgentID > 0 { aiAgent := AIAgentService.Get(conversation.AIAgentID) if aiAgent == nil || aiAgent.Status != enums.StatusOk { return errorsx.InvalidParamI18n("error.e0003") } result, err := ConversationHumanDispatchService.DispatchPendingConversation(conversationID, *aiAgent) if err != nil { return err } if result == nil || result.Decision == HandoffDecisionOffHours { return errorsx.InvalidParamI18n("error.e0194") } return nil } result, err := ConversationDispatchService.DispatchConversation(conversationID) if err != nil { return err } if result == nil { return errorsx.InvalidParamI18n("error.e0194") } return nil } func (s *conversationService) TransferConversation(conversationID, toUserID int64, reason string, operator *dto.AuthPrincipal) error { if operator == nil { return errorsx.UnauthorizedI18n("error.auth.expired") } if toUserID <= 0 { return errorsx.InvalidParamI18n("error.e0278") } targetProfile := AgentProfileService.GetByUserID(toUserID) if targetProfile == nil || targetProfile.Status != enums.StatusOk { return errorsx.InvalidParamI18n("error.e0276") } var assignedEvent events.ConversationAssignedEvent if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { conversation := repositories.ConversationRepository.Get(ctx.Tx, conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if !s.canTransferConversation(conversation, operator) { return errorsx.ForbiddenI18n("error.e0223") } if conversation.Status != enums.IMConversationStatusActive { return errorsx.InvalidParamI18n("error.e0134") } if conversation.CurrentAssigneeID <= 0 { return errorsx.InvalidParamI18n("error.e0193") } if conversation.CurrentAssigneeID == toUserID { return errorsx.InvalidParamI18n("error.e0277") } now := time.Now() if err := ConversationAssignmentService.FinishActiveAssignments(ctx, conversationID, now); err != nil { return err } if err := ConversationAssignmentService.CreateAssignment(ctx, conversationID, conversation.CurrentAssigneeID, toUserID, enums.IMAssignmentTypeTransfer, reason, operator, now); err != nil { return err } if err := repositories.ConversationRepository.Updates(ctx.Tx, conversationID, map[string]any{ "current_assignee_id": toUserID, "status": enums.IMConversationStatusActive, "update_user_id": operator.UserID, "update_user_name": operator.Username, "updated_at": now, }); err != nil { return err } if err := ConversationEventLogService.CreateEvent(ctx, conversationID, enums.IMEventTypeTransfer, enums.IMSenderTypeAgent, operator.UserID, "会话已转接", s.buildEventPayload(map[string]any{ "from_status": conversation.Status, "to_status": enums.IMConversationStatusActive, "from_assignee_id": conversation.CurrentAssigneeID, "to_assignee_id": 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) } eventbus.PublishAsync(context.Background(), assignedEvent) return nil } func (s *conversationService) HandoffByAI(conversationID int64, aiAgent models.AIAgent, reason string) error { return s.HandoffByAIWithRequestID(conversationID, aiAgent, reason, "") } func (s *conversationService) HandoffByAIWithRequestID(conversationID int64, aiAgent models.AIAgent, reason string, requestID string) error { if conversationID <= 0 { return errorsx.InvalidParamI18n("error.e0116") } _, err := ConversationHumanDispatchService.HandoffByAIWithRequestID(conversationID, aiAgent, reason, requestID) if err != nil { slog.Warn("schedule-aware ai handoff failed", "requestId", requestID, "conversation_id", conversationID, "ai_agent_id", aiAgent.ID, "error", err) } return err } func (s *conversationService) TryOffHoursHandoffByAI(conversationID int64, aiAgent models.AIAgent, reason string) (bool, error) { return s.TryOffHoursHandoffByAIWithRequestID(conversationID, aiAgent, reason, "") } func (s *conversationService) TryOffHoursHandoffByAIWithRequestID(conversationID int64, aiAgent models.AIAgent, reason string, requestID string) (bool, error) { if conversationID <= 0 { return false, errorsx.InvalidParamI18n("error.e0116") } handled, err := ConversationHumanDispatchService.TryOffHoursHandoffByAIWithRequestID(conversationID, aiAgent, reason, requestID) if err != nil { slog.Warn("off-hours ai handoff failed", "requestId", requestID, "conversation_id", conversationID, "ai_agent_id", aiAgent.ID, "error", err) } return handled, err } func (s *conversationService) CloseConversation(conversationID int64, closeReason string, operator *dto.AuthPrincipal) error { if operator == nil { return errorsx.UnauthorizedI18n("error.auth.expired") } return s.closeConversation(conversationID, enums.IMSenderTypeAgent, closeReason, operator) } func (s *conversationService) CloseCustomerConversation(conversationID int64, externalUser openidentity.ExternalUser) error { conversation := s.Get(conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if !s.IsCustomerConversationOwner(conversation, externalUser) { return errorsx.ForbiddenI18n("error.e0222") } return s.closeConversation(conversationID, enums.IMSenderTypeCustomer, "", nil) } func (s *conversationService) closeConversation(conversationID int64, senderType enums.IMSenderType, closeReason string, operator *dto.AuthPrincipal) error { var queuedTeamID int64 var wasQueued bool if err := sqls.WithTransaction(func(ctx *sqls.TxContext) error { conversation := repositories.ConversationRepository.Get(ctx.Tx, conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if conversation.Status == enums.IMConversationStatusClosed { return nil } if conversation.Status != enums.IMConversationStatusAIServing && conversation.Status != enums.IMConversationStatusPending && conversation.Status != enums.IMConversationStatusActive { return errorsx.InvalidParamI18n("error.e0197") } wasQueued = conversation.Status == enums.IMConversationStatusPending && conversation.CurrentAssigneeID == 0 queuedTeamID = conversation.CurrentTeamID var ( now = time.Now() eventDesc = "会话已关闭" operatorID int64 operatorName string ) closeReason = strings.TrimSpace(closeReason) if senderType == enums.IMSenderTypeCustomer { eventDesc = "客户关闭会话" } else { if operator == nil { return errorsx.InvalidParamI18n("error.e0226") } if closeReason == "" { return errorsx.InvalidParamI18n("error.e0128") } if !s.canCloseConversation(conversation, operator) { return errorsx.ForbiddenI18n("error.e0221") } operatorID = operator.UserID operatorName = operator.Nickname } if err := ConversationAssignmentService.FinishActiveAssignments(ctx, conversationID, now); err != nil { return err } if err := repositories.ConversationRepository.Updates(ctx.Tx, conversationID, map[string]any{ "status": enums.IMConversationStatusClosed, "closed_at": now, "closed_by": operatorID, "close_reason": closeReason, "update_user_id": operatorID, "update_user_name": operatorName, "updated_at": now, }); err != nil { return err } return ConversationEventLogService.CreateEvent(ctx, conversationID, enums.IMEventTypeClose, senderType, operatorID, eventDesc, s.buildEventPayload(map[string]any{ "from_status": conversation.Status, "to_status": enums.IMConversationStatusClosed, "from_assignee_id": conversation.CurrentAssigneeID, "to_assignee_id": conversation.CurrentAssigneeID, "close_reason": closeReason, })) }); err != nil { return err } if conversation := s.Get(conversationID); conversation != nil { WsService.PublishConversationChanged(conversation, enums.IMRealtimeEventConversationClosed) } if wasQueued { ConversationQueueService.PublishPoolUpdates(queuedTeamID) } return nil } // MarkAgentConversationReadToMessage 控制台客服将会话已读推进到指定消息。 func (s *conversationService) MarkAgentConversationReadToMessage(conversationID, messageID int64, operator *dto.AuthPrincipal) error { if operator == nil { return errorsx.UnauthorizedI18n("error.auth.expired") } conversation := s.Get(conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } changed, err := s.markConversationReadWithActor(conversation, messageID, agentConversationReadActor{operator: operator}) if err != nil { return err } if changed { if updated := s.Get(conversationID); updated != nil { WsService.PublishConversationChanged(updated, enums.IMRealtimeEventConversationRead) } } return nil } // MarkCustomerConversationReadToMessage IM 客户将会话已读推进到指定消息(需为会话归属外部身份)。 func (s *conversationService) MarkCustomerConversationReadToMessage(conversationID, messageID int64, external *openidentity.ExternalUser) error { if external == nil || strings.TrimSpace(external.ExternalID) == "" { return errorsx.UnauthorizedI18n("error.e0149") } conversation := s.Get(conversationID) if conversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if !s.IsCustomerConversationOwner(conversation, *external) { return errorsx.ForbiddenI18n("error.e0222") } changed, err := s.markConversationReadWithActor(conversation, messageID, customerConversationReadActor{external: external}) if err != nil { return err } if changed { if updated := s.Get(conversationID); updated != nil { WsService.PublishConversationChanged(updated, enums.IMRealtimeEventConversationRead) } } return nil } func displayExternalName(ext *openidentity.ExternalUser) string { if ext == nil { return "" } if n := strings.TrimSpace(ext.ExternalName); n != "" { return n } return strings.TrimSpace(ext.ExternalID) } // conversationReadActor 抽象「读者身份」,供 markConversationReadWithActor 共用(包内私有)。 type conversationReadActor interface { isAgentSide() bool getReadState(conversationID int64) *models.ConversationReadState markRead(ctx *sqls.TxContext, conversation *models.Conversation, targetMessage *models.Message) error conversationUpdateAudit() (userID int64, userName string) } type agentConversationReadActor struct { operator *dto.AuthPrincipal } func (a agentConversationReadActor) isAgentSide() bool { return true } func (a agentConversationReadActor) getReadState(conversationID int64) *models.ConversationReadState { return ConversationReadStateService.GetByAgentReader(conversationID, a.operator) } func (a agentConversationReadActor) markRead(ctx *sqls.TxContext, conversation *models.Conversation, targetMessage *models.Message) error { _, err := ConversationReadStateService.MarkAgentRead(ctx, conversation, a.operator, targetMessage) return err } func (a agentConversationReadActor) conversationUpdateAudit() (int64, string) { if a.operator == nil { return 0, "" } return a.operator.UserID, a.operator.Username } type customerConversationReadActor struct { external *openidentity.ExternalUser } func (a customerConversationReadActor) isAgentSide() bool { return false } func (a customerConversationReadActor) getReadState(conversationID int64) *models.ConversationReadState { return ConversationReadStateService.GetByCustomerReader(conversationID, a.external) } func (a customerConversationReadActor) markRead(ctx *sqls.TxContext, conversation *models.Conversation, targetMessage *models.Message) error { _, err := ConversationReadStateService.MarkCustomerRead(ctx, conversation, a.external, targetMessage) return err } func (a customerConversationReadActor) conversationUpdateAudit() (int64, string) { return 0, displayExternalName(a.external) } func (s *conversationService) markConversationReadWithActor(conversation *models.Conversation, messageID int64, actor conversationReadActor) (bool, error) { if conversation == nil { return false, errorsx.InvalidParamI18n("error.e0116") } targetMessage, err := MessageService.GetConversationReadTarget(conversation.ID, messageID) if err != nil { return false, err } if targetMessage == nil { if actor.isAgentSide() && conversation.AgentUnreadCount == 0 { return false, nil } if !actor.isAgentSide() && conversation.CustomerUnreadCount == 0 { return false, nil } now := time.Now() updateUserID, updateUserName := actor.conversationUpdateAudit() updates := map[string]any{ "update_user_id": updateUserID, "update_user_name": updateUserName, "updated_at": now, } if actor.isAgentSide() { updates["agent_unread_count"] = 0 } else { updates["customer_unread_count"] = 0 } return true, s.Updates(conversation.ID, updates) } currentReadState := actor.getReadState(conversation.ID) if currentReadState != nil && currentReadState.LastReadMessageID >= targetMessage.ID { if actor.isAgentSide() && conversation.AgentUnreadCount == 0 { return false, nil } if !actor.isAgentSide() && conversation.CustomerUnreadCount == 0 { return false, nil } } err = sqls.WithTransaction(func(ctx *sqls.TxContext) error { currentConversation := repositories.ConversationRepository.Get(ctx.Tx, conversation.ID) if currentConversation == nil { return errorsx.InvalidParamI18n("error.e0116") } if err := actor.markRead(ctx, currentConversation, targetMessage); err != nil { return err } agentReadState, customerReadState := ConversationReadStateService.getConversationReadStates(ctx.Tx, currentConversation.ID) agentUnreadCount, err := s.countUnreadByState(ctx, currentConversation.ID, agentReadState, enums.IMSenderTypeCustomer) if err != nil { return err } customerUnreadCount, err := s.countUnreadByState(ctx, currentConversation.ID, customerReadState, enums.IMSenderTypeAgent, enums.IMSenderTypeAI) if err != nil { return err } if actor.isAgentSide() && currentConversation.AgentUnreadCount == agentUnreadCount && currentReadState != nil && currentReadState.LastReadMessageID >= targetMessage.ID { return nil } if !actor.isAgentSide() && currentConversation.CustomerUnreadCount == customerUnreadCount && currentReadState != nil && currentReadState.LastReadMessageID >= targetMessage.ID { return nil } updateUserID, updateUserName := actor.conversationUpdateAudit() return repositories.ConversationRepository.Updates(ctx.Tx, currentConversation.ID, map[string]any{ "agent_unread_count": agentUnreadCount, "customer_unread_count": customerUnreadCount, "update_user_id": updateUserID, "update_user_name": updateUserName, "updated_at": time.Now(), }) }) if err != nil { return false, err } return true, nil } func (s *conversationService) countUnreadByState(ctx *sqls.TxContext, conversationID int64, state *models.ConversationReadState, senderTypes ...enums.IMSenderType) (int, error) { lastReadMessageID := int64(0) if state != nil { lastReadMessageID = state.LastReadMessageID } normalizedSenderTypes := make([]enums.IMSenderType, 0, len(senderTypes)) for _, senderType := range senderTypes { normalizedSenderTypes = append(normalizedSenderTypes, senderType) } count, err := ConversationReadStateService.CountUnreadMessages(ctx, conversationID, lastReadMessageID, normalizedSenderTypes...) return int(count), err } func (s *conversationService) IsCustomerConversationOwner(conversation *models.Conversation, externalUser openidentity.ExternalUser) bool { if conversation == nil { return false } extID := strings.TrimSpace(externalUser.ExternalID) if extID == "" || strings.TrimSpace(string(externalUser.ExternalSource)) == "" { return false } return conversation.CustomerType == externalCustomerType(externalUser) && strings.TrimSpace(conversation.CustomerExternalID) == extID } func (s *conversationService) BuildConversationSummary(conversation *models.Conversation) string { if conversation == nil { return "" } if strings.TrimSpace(conversation.LastMessageSummary) != "" { return conversation.LastMessageSummary } return strings.TrimSpace(conversation.CustomerName) } func externalCustomerType(externalUser openidentity.ExternalUser) string { if externalUser.SubjectType != "" { return string(externalUser.SubjectType) } return string(externalUser.ExternalSource) } func (s *conversationService) canCloseConversation(conversation *models.Conversation, operator *dto.AuthPrincipal) bool { if conversation == nil || operator == nil { return false } if s.isAdmin(operator) { return true } return conversation.Status == enums.IMConversationStatusActive && conversation.CurrentAssigneeID > 0 && conversation.CurrentAssigneeID == operator.UserID } func (s *conversationService) canTransferConversation(conversation *models.Conversation, operator *dto.AuthPrincipal) bool { if conversation == nil || operator == nil { return false } if s.isAdmin(operator) { return true } return conversation.Status == enums.IMConversationStatusActive && conversation.CurrentAssigneeID > 0 && conversation.CurrentAssigneeID == operator.UserID } func (s *conversationService) isAdmin(operator *dto.AuthPrincipal) bool { if operator == nil { return false } return operator.SubjectType == identity.SubjectAdmin } func (s *conversationService) buildEventPayload(payload map[string]any) string { if len(payload) == 0 { return "" } data, err := json.Marshal(payload) if err != nil { return "" } return string(data) } func (s *conversationService) GetConversationExternalIdentity(conversation *models.Conversation) *openidentity.ExternalUser { if conversation == nil || strings.TrimSpace(conversation.CustomerExternalID) == "" { return nil } external := &openidentity.ExternalUser{ ExternalID: strings.TrimSpace(conversation.CustomerExternalID), ExternalName: strings.TrimSpace(conversation.CustomerName), SubjectID: conversation.CustomerID, } switch conversation.CustomerType { case string(identity.SubjectCard), string(identity.SubjectDevice), string(identity.SubjectMallUser): external.ExternalSource = enums.ExternalSourceUser external.SubjectType = identity.SubjectType(conversation.CustomerType) default: external.ExternalSource = enums.ExternalSource(conversation.CustomerType) } return external } func (s *conversationService) resolveInitialStatus(serviceMode enums.IMConversationServiceMode) enums.IMConversationStatus { switch serviceMode { case enums.IMConversationServiceModeHumanOnly: return enums.IMConversationStatusPending case enums.IMConversationServiceModeAIOnly, enums.IMConversationServiceModeAIFirst: return enums.IMConversationStatusAIServing default: return enums.IMConversationStatusAIServing } }