优化实时聊天:稳定 WS、按 seq 增量同步,并修复访客评价推送

- 客服/访客 WebSocket 不再因切换会话反复重连,断线自动恢复
- 新增 after_seq 增量拉取,重连与消息空洞时 catch-up
- 结束会话同时推送给访客,预览页可弹出评价并支持兜底状态同步
This commit is contained in:
yml2213
2026-07-15 14:38:08 +08:00
parent ea3901d8b0
commit 1359bfe597
7 changed files with 758 additions and 107 deletions
+79 -3
View File
@@ -2,6 +2,7 @@ package handler
import (
"net/http"
"strconv"
"strings"
"time"
"unicode/utf8"
@@ -145,9 +146,12 @@ func broadcastSessionMessage(session *model.Session, message model.Message) {
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)
if err != nil {
return
}
// 坐席列表需要刷新;访客也必须收到 ended/active,才能弹出评价或更新接待状态
ws.DefaultHub.BroadcastToTenantStaff(session.TenantID, payload)
ws.DefaultHub.BroadcastToVisitor(session.TenantID, session.ID, payload)
}
func loadAssignableAgent(tenantID, agentID uint) error {
@@ -351,7 +355,79 @@ func (h *SessionHandler) Get(c *gin.Context) {
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})
maxSeq := 0
for _, m := range messages {
if m.Seq > maxSeq {
maxSeq = m.Seq
}
}
middleware.JSON(c, gin.H{
"session": session,
"messages": messages,
"events": events,
"pending_count": pendingCount,
"max_seq": maxSeq,
})
}
// ListMessages 按 seq 增量拉取消息:after_seq 之后的新消息(重连 catch-up)。
// GET /sessions/:id/messages?after_seq=12
func (h *SessionHandler) ListMessages(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
}
afterSeq := 0
if raw := strings.TrimSpace(c.Query("after_seq")); raw != "" {
parsed, err := strconv.Atoi(raw)
if err != nil || parsed < 0 {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "after_seq 无效"})
return
}
afterSeq = parsed
}
limit := 200
if raw := strings.TrimSpace(c.Query("limit")); raw != "" {
if parsed, err := strconv.Atoi(raw); err == nil && parsed > 0 {
if parsed > 500 {
parsed = 500
}
limit = parsed
}
}
query := model.DB.Where("session_id = ?", session.ID)
if afterSeq > 0 {
query = query.Where("seq > ?", afterSeq)
}
var messages []model.Message
if err := query.Order("seq asc").Limit(limit).Find(&messages).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "查询消息失败"})
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
}
middleware.JSON(c, gin.H{
"messages": messages,
"after_seq": afterSeq,
"max_seq": maxSeq,
"has_more": len(messages) >= limit && (len(messages) == 0 || messages[len(messages)-1].Seq < maxSeq),
})
}
func (h *SessionHandler) MarkRead(c *gin.Context) {