package handler import ( "net/http" "strings" "time" "unicode/utf8" "github.com/gin-gonic/gin" "gorm.io/gorm" "kefu-sys/server/internal/middleware" "kefu-sys/server/internal/model" "kefu-sys/server/internal/ws" ) type SessionHandler struct{} func NewSessionHandler() *SessionHandler { return &SessionHandler{} } type SendMessageReq struct { Content string `json:"content" binding:"required"` Type string `json:"type"` } type SessionListItem struct { model.Session UnreadCount int `json:"unread_count"` LastMessage string `json:"last_message"` LastMessageAt *time.Time `json:"last_message_at,omitempty"` MessageCount int `json:"message_count"` CustomerName string `json:"customer_name"` AgentName string `json:"agent_name"` ChannelName string `json:"channel_name"` ChannelType string `json:"channel_type"` } type CreateNoteReq struct { Content string `json:"content" binding:"required"` } func (h *SessionHandler) SendMessage(c *gin.Context) { userID := middleware.GetUserID(c) id := c.Param("id") var req SendMessageReq if err := c.ShouldBindJSON(&req); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "参数错误"}) return } if req.Type == "" { req.Type = "text" } content, err := validateMessageContent(req.Type, req.Content) if err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": err.Error()}) return } session, ok := loadTenantSession(c, id) if !ok { return } if !canOperateSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权向该会话发送消息"}) return } if session.Status == "ended" || session.Status == "archived" { c.JSON(http.StatusConflict, gin.H{"code": 409, "message": "会话已结束"}) return } msg := model.Message{ SessionID: session.ID, SenderType: "agent", SenderID: &userID, Content: content, Type: req.Type, SentAt: time.Now(), } if err := model.CreateMessage(&msg); err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "发送失败"}) return } if middleware.GetRole(c) == "agent" { model.DB.Model(&model.Session{}).Where("id = ?", session.ID).Update("last_read_seq", msg.Seq) session.LastReadSeq = msg.Seq } broadcastSessionMessage(session, msg) middleware.JSON(c, msg) } type CreateSessionReq struct { ChannelID uint `json:"channel_id"` CustomerID uint `json:"customer_id"` Priority string `json:"priority"` } type AssignSessionReq struct { AgentID uint `json:"agent_id"` } func isTenantManager(c *gin.Context) bool { return middleware.HasAnyRole(c, "admin", "supervisor") } func loadTenantSession(c *gin.Context, id string) (*model.Session, bool) { var session model.Session if err := model.DB.First(&session, id).Error; err != nil { c.JSON(http.StatusNotFound, gin.H{"code": 404, "message": "会话不存在"}) return nil, false } if session.TenantID != middleware.GetTenantID(c) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权访问其他租户会话"}) return nil, false } return &session, true } func canReadSession(c *gin.Context, session *model.Session) bool { if isTenantManager(c) { return true } if middleware.GetRole(c) == "agent" { userID := middleware.GetUserID(c) return session.Status == "waiting" || (session.AgentID != nil && *session.AgentID == userID) } return false } func canOperateSession(c *gin.Context, session *model.Session) bool { if isTenantManager(c) { return true } return middleware.GetRole(c) == "agent" && session.AgentID != nil && *session.AgentID == middleware.GetUserID(c) } func broadcastSessionMessage(session *model.Session, message model.Message) { payload, err := ws.NewEvent("message", session.ID, message) if err == nil { ws.DefaultHub.BroadcastToSession(session.TenantID, session.ID, session.AgentID, payload) } } func broadcastSessionUpdate(session *model.Session) { payload, err := ws.NewEvent("session_updated", session.ID, gin.H{"status": session.Status}) if err == nil { ws.DefaultHub.BroadcastToTenantStaff(session.TenantID, payload) } } func loadAssignableAgent(tenantID, agentID uint) error { var agent model.User if err := model.DB.First(&agent, agentID).Error; err != nil { return err } if agent.TenantID != tenantID || agent.Role != "agent" || agent.Status == "disabled" { return gorm.ErrRecordNotFound } return nil } func unreadCount(session model.Session) int { var count int64 model.DB.Model(&model.Message{}). Where("session_id = ? AND sender_type = ? AND seq > ?", session.ID, "visitor", session.LastReadSeq). Count(&count) return int(count) } func (h *SessionHandler) List(c *gin.Context) { tenantID := middleware.GetTenantID(c) page, pageSize := middleware.GetPageParams(c) status := c.Query("status") priority := c.Query("priority") agentID := c.Query("agent_id") channelID := c.Query("channel_id") search := strings.TrimSpace(c.Query("search")) from := c.Query("from") to := c.Query("to") var sessions []model.Session var total int64 query := model.DB.Model(&model.Session{}).Where("sessions.tenant_id = ?", tenantID) if middleware.GetRole(c) == "agent" { query = query.Where("sessions.agent_id = ? OR sessions.status = ?", middleware.GetUserID(c), "waiting") } if status != "" { query = query.Where("sessions.status = ?", status) } if priority != "" { query = query.Where("sessions.priority = ?", priority) } if agentID != "" { query = query.Where("sessions.agent_id = ?", agentID) } if channelID != "" { query = query.Where("sessions.channel_id = ?", channelID) } if from != "" { if t, err := time.ParseInLocation("2006-01-02", from, time.Local); err == nil { query = query.Where("sessions.created_at >= ?", t) } } if to != "" { if t, err := time.ParseInLocation("2006-01-02", to, time.Local); err == nil { query = query.Where("sessions.created_at < ?", t.Add(24*time.Hour)) } } if search != "" { like := "%" + search + "%" query = query.Joins("LEFT JOIN customers ON customers.id = sessions.customer_id"). Where("customers.name LIKE ? OR CAST(sessions.id AS TEXT) LIKE ? OR COALESCE(sessions.end_reason, '') LIKE ?", like, like, like) } query.Count(&total) if err := query.Order("sessions.created_at desc").Offset((page - 1) * pageSize).Limit(pageSize).Find(&sessions).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "查询会话失败"}) return } // 批量补全客户 / 客服 / 渠道 / 消息数 customerIDs := make([]uint, 0) agentIDs := make([]uint, 0) channelIDs := make([]uint, 0) sessionIDs := make([]uint, 0, len(sessions)) for _, s := range sessions { sessionIDs = append(sessionIDs, s.ID) customerIDs = append(customerIDs, s.CustomerID) channelIDs = append(channelIDs, s.ChannelID) if s.AgentID != nil { agentIDs = append(agentIDs, *s.AgentID) } } customerMap := map[uint]model.Customer{} if len(customerIDs) > 0 { var customers []model.Customer model.DB.Where("id IN ?", uniqueUint(customerIDs)).Find(&customers) for _, cu := range customers { customerMap[cu.ID] = cu } } agentMap := map[uint]model.User{} if len(agentIDs) > 0 { var users []model.User model.DB.Where("id IN ?", uniqueUint(agentIDs)).Find(&users) for _, u := range users { agentMap[u.ID] = u } } channelMap := map[uint]model.Channel{} if len(channelIDs) > 0 { var channels []model.Channel model.DB.Where("id IN ?", uniqueUint(channelIDs)).Find(&channels) for _, ch := range channels { channelMap[ch.ID] = ch } } msgCountMap := map[uint]int{} if len(sessionIDs) > 0 { type row struct { SessionID uint Cnt int } var rows []row model.DB.Model(&model.Message{}). Select("session_id, COUNT(*) as cnt"). Where("session_id IN ?", sessionIDs). Group("session_id"). Scan(&rows) for _, r := range rows { msgCountMap[r.SessionID] = r.Cnt } } items := make([]SessionListItem, 0, len(sessions)) for _, session := range sessions { item := SessionListItem{Session: session, MessageCount: msgCountMap[session.ID]} if middleware.GetRole(c) == "agent" && session.AgentID != nil && *session.AgentID == middleware.GetUserID(c) { item.UnreadCount = unreadCount(session) } if cu, ok := customerMap[session.CustomerID]; ok { item.CustomerName = cu.Name } if session.AgentID != nil { if u, ok := agentMap[*session.AgentID]; ok { item.AgentName = u.Nickname } } if ch, ok := channelMap[session.ChannelID]; ok { item.ChannelName = ch.Name item.ChannelType = ch.Type } var lastMsg model.Message if err := model.DB.Where("session_id = ?", session.ID).Order("seq desc").First(&lastMsg).Error; err == nil { if lastMsg.Type == "image" { item.LastMessage = "[图片]" } else { item.LastMessage = lastMsg.Content if utf8.RuneCountInString(item.LastMessage) > 40 { runes := []rune(item.LastMessage) item.LastMessage = string(runes[:40]) + "…" } } t := lastMsg.SentAt item.LastMessageAt = &t } items = append(items, item) } middleware.JSONList(c, items, total, page, pageSize) } func uniqueUint(ids []uint) []uint { seen := make(map[uint]struct{}, len(ids)) out := make([]uint, 0, len(ids)) for _, id := range ids { if id == 0 { continue } if _, ok := seen[id]; ok { continue } seen[id] = struct{}{} out = append(out, id) } return out } func (h *SessionHandler) Get(c *gin.Context) { id := c.Param("id") session, ok := loadTenantSession(c, id) if !ok { return } if !canReadSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权查看该会话"}) return } var messages []model.Message if err := model.DB.Where("session_id = ?", session.ID).Order("seq asc").Find(&messages).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "查询消息失败"}) return } var events []model.SessionEvent model.DB.Where("session_id = ?", session.ID).Order("created_at asc").Find(&events) var pendingCount int64 model.DB.Model(&model.Session{}).Where("customer_id = ? AND tenant_id = ? AND status = ?", session.CustomerID, session.TenantID, "waiting").Count(&pendingCount) middleware.JSON(c, gin.H{"session": session, "messages": messages, "events": events, "pending_count": pendingCount}) } func (h *SessionHandler) MarkRead(c *gin.Context) { session, ok := loadTenantSession(c, c.Param("id")) if !ok { return } if !canReadSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权标记该会话已读"}) return } if session.AgentID == nil || middleware.GetRole(c) != "agent" || *session.AgentID != middleware.GetUserID(c) { middleware.JSON(c, gin.H{"last_read_seq": session.LastReadSeq}) return } var maxSeq int if err := model.DB.Model(&model.Message{}).Where("session_id = ?", session.ID).Select("COALESCE(MAX(seq), 0)").Scan(&maxSeq).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "标记已读失败"}) return } if err := model.DB.Model(&model.Session{}).Where("id = ?", session.ID).Update("last_read_seq", maxSeq).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "标记已读失败"}) return } middleware.JSON(c, gin.H{"last_read_seq": maxSeq}) } func (h *SessionHandler) AddNote(c *gin.Context) { session, ok := loadTenantSession(c, c.Param("id")) if !ok { return } if !canOperateSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权添加内部备注"}) return } var req CreateNoteReq if err := c.ShouldBindJSON(&req); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "参数错误"}) return } content := strings.TrimSpace(req.Content) if content == "" || utf8.RuneCountInString(content) > 500 { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "内部备注需为 1 至 500 字"}) return } event := model.SessionEvent{SessionID: session.ID, OperatorID: middleware.GetUserID(c), Action: "note", Detail: content} if err := model.DB.Create(&event).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "保存内部备注失败"}) return } middleware.JSON(c, event) } func (h *SessionHandler) ListAvailableAgents(c *gin.Context) { type agentItem struct { ID uint `json:"id"` Nickname string `json:"nickname"` Status string `json:"status"` } // all=1 返回租户全部坐席(含离线),用于对话记录筛选;默认仅在线(工作台转接) q := model.DB.Where("tenant_id = ? AND role IN ?", middleware.GetTenantID(c), []string{"agent", "supervisor", "admin"}) if c.Query("all") != "1" { q = q.Where("role = ? AND status = ?", "agent", "online") } var users []model.User if err := q.Order("nickname asc").Find(&users).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "查询客服失败"}) return } items := make([]agentItem, 0, len(users)) for _, user := range users { items = append(items, agentItem{ID: user.ID, Nickname: user.Nickname, Status: user.Status}) } middleware.JSON(c, items) } func (h *SessionHandler) Create(c *gin.Context) { if !isTenantManager(c) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "仅主管或管理员可创建会话"}) return } var req CreateSessionReq if err := c.ShouldBindJSON(&req); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "参数错误"}) return } tenantID := middleware.GetTenantID(c) var channel model.Channel if err := model.DB.Where("id = ? AND tenant_id = ? AND status = ?", req.ChannelID, tenantID, "enabled").First(&channel).Error; err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "渠道不存在或未启用"}) return } var customer model.Customer if err := model.DB.Where("id = ? AND tenant_id = ?", req.CustomerID, tenantID).First(&customer).Error; err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "客户不存在"}) return } session := model.Session{ TenantID: tenantID, ChannelID: req.ChannelID, CustomerID: req.CustomerID, Priority: req.Priority, Status: "waiting", } if session.Priority == "" { session.Priority = "normal" } if err := model.DB.Create(&session).Error; err != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "创建会话失败"}) return } model.DB.Model(&model.Customer{}).Where("id = ? AND tenant_id = ?", customer.ID, tenantID). Updates(map[string]interface{}{"conversation_count": gorm.Expr("conversation_count + 1"), "last_contact_at": time.Now()}) middleware.JSON(c, session) } func (h *SessionHandler) Assign(c *gin.Context) { id := c.Param("id") session, ok := loadTenantSession(c, id) if !ok { return } var req AssignSessionReq if err := c.ShouldBindJSON(&req); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "参数错误"}) return } if session.Status != "waiting" { c.JSON(http.StatusConflict, gin.H{"code": 409, "message": "会话已被分配"}) return } if middleware.GetRole(c) == "agent" { req.AgentID = middleware.GetUserID(c) } else if !isTenantManager(c) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权分配会话"}) return } if req.AgentID == 0 { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "请选择目标客服"}) return } if err := loadAssignableAgent(session.TenantID, req.AgentID); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "目标客服不存在或不可用"}) return } result := model.DB.Model(&model.Session{}). Where("id = ? AND tenant_id = ? AND status = ?", session.ID, session.TenantID, "waiting"). Updates(map[string]interface{}{"agent_id": req.AgentID, "status": "active", "last_read_seq": 0}) if result.Error != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "分配失败"}) return } if result.RowsAffected == 0 { c.JSON(http.StatusConflict, gin.H{"code": 409, "message": "会话已被分配"}) return } model.DB.Create(&model.SessionEvent{SessionID: session.ID, OperatorID: middleware.GetUserID(c), Action: "assign", Detail: "会话分配"}) session.AgentID = &req.AgentID session.Status = "active" broadcastSessionUpdate(session) middleware.JSON(c, gin.H{"message": "分配成功"}) } func (h *SessionHandler) Transfer(c *gin.Context) { id := c.Param("id") session, ok := loadTenantSession(c, id) if !ok { return } if !canOperateSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权转接该会话"}) return } var req AssignSessionReq if err := c.ShouldBindJSON(&req); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "参数错误"}) return } if err := loadAssignableAgent(session.TenantID, req.AgentID); err != nil { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "目标客服不存在或不可用"}) return } result := model.DB.Model(&model.Session{}). Where("id = ? AND tenant_id = ? AND status = ?", session.ID, session.TenantID, "active"). Updates(map[string]interface{}{"agent_id": req.AgentID, "last_read_seq": 0}) if result.Error != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "转接失败"}) return } if result.RowsAffected == 0 { c.JSON(http.StatusConflict, gin.H{"code": 409, "message": "会话不可转接"}) return } model.DB.Create(&model.SessionEvent{ SessionID: session.ID, OperatorID: middleware.GetUserID(c), Action: "transfer", Detail: "会话转接", }) session.AgentID = &req.AgentID broadcastSessionUpdate(session) middleware.JSON(c, gin.H{"message": "转接成功"}) } func (h *SessionHandler) End(c *gin.Context) { id := c.Param("id") reason := c.Query("reason") session, ok := loadTenantSession(c, id) if !ok { return } if !canOperateSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权结束该会话"}) return } if reason == "" { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "请填写结束原因"}) return } if reason != "resolved" && reason != "no_response" && reason != "visitor_left" && reason != "transferred" && reason != "other" { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "结束原因无效"}) return } now := time.Now() result := model.DB.Model(&model.Session{}). Where("id = ? AND tenant_id = ? AND status <> ?", session.ID, session.TenantID, "ended"). Updates(map[string]interface{}{"status": "ended", "end_reason": reason, "ended_at": now}) if result.RowsAffected == 0 { c.JSON(http.StatusNotFound, gin.H{"code": 404, "message": "会话不存在"}) return } model.DB.Create(&model.SessionEvent{ SessionID: session.ID, OperatorID: middleware.GetUserID(c), Action: "end", Detail: "结束会话: " + reason, }) session.Status = "ended" session.EndedAt = &now broadcastSessionUpdate(session) middleware.JSON(c, gin.H{"message": "已结束"}) } func (h *SessionHandler) UpdatePriority(c *gin.Context) { id := c.Param("id") priority := c.Query("priority") session, ok := loadTenantSession(c, id) if !ok { return } if !canOperateSession(c, session) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "无权更新该会话"}) return } if priority != "urgent" && priority != "normal" { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "优先级无效"}) return } result := model.DB.Model(&model.Session{}). Where("id = ? AND tenant_id = ?", session.ID, session.TenantID). Update("priority", priority) if result.Error != nil || result.RowsAffected == 0 { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "更新失败"}) return } session.Priority = priority broadcastSessionUpdate(session) middleware.JSON(c, gin.H{"message": "已更新"}) } // Archive 将已结束会话归档(主管/管理员) func (h *SessionHandler) Archive(c *gin.Context) { if !isTenantManager(c) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "仅主管或管理员可归档"}) return } session, ok := loadTenantSession(c, c.Param("id")) if !ok { return } if session.Status != "ended" { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "仅已结束的会话可归档"}) return } result := model.DB.Model(&model.Session{}). Where("id = ? AND tenant_id = ? AND status = ?", session.ID, session.TenantID, "ended"). Update("status", "archived") if result.Error != nil || result.RowsAffected == 0 { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "归档失败"}) return } model.DB.Create(&model.SessionEvent{ SessionID: session.ID, OperatorID: middleware.GetUserID(c), Action: "archive", Detail: "会话已归档", }) session.Status = "archived" broadcastSessionUpdate(session) middleware.JSON(c, gin.H{"message": "已归档", "session": session}) } // BatchArchive 批量归档已结束会话 func (h *SessionHandler) BatchArchive(c *gin.Context) { if !isTenantManager(c) { c.JSON(http.StatusForbidden, gin.H{"code": 403, "message": "仅主管或管理员可归档"}) return } var req struct { IDs []uint `json:"ids" binding:"required"` } if err := c.ShouldBindJSON(&req); err != nil || len(req.IDs) == 0 { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "请选择要归档的会话"}) return } if len(req.IDs) > 100 { c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "单次最多归档 100 条"}) return } tenantID := middleware.GetTenantID(c) result := model.DB.Model(&model.Session{}). Where("tenant_id = ? AND status = ? AND id IN ?", tenantID, "ended", req.IDs). Update("status", "archived") if result.Error != nil { c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "批量归档失败"}) return } for _, id := range req.IDs { model.DB.Create(&model.SessionEvent{ SessionID: id, OperatorID: middleware.GetUserID(c), Action: "archive", Detail: "会话已归档", }) } middleware.JSON(c, gin.H{"message": "已归档", "count": result.RowsAffected}) }