Files
2026-08-07 13:36:56 +08:00

144 lines
4.3 KiB
Go

package modules
import (
"encoding/base64"
"fmt"
"net/http"
"time"
"github.com/gin-gonic/gin"
"github.com/lineup/app-server/internal/wkapi"
)
// MessageModule handles message sending and syncing.
type MessageModule struct {
wk *wkapi.Client
}
// NewMessageModule creates a new MessageModule.
func NewMessageModule(wk *wkapi.Client) *MessageModule {
return &MessageModule{wk: wk}
}
// RegisterRoutes registers message routes.
func (m *MessageModule) RegisterRoutes(r *gin.Engine) {
r.POST("/messages/send", m.handleSend)
r.GET("/messages/sync", m.handleSync)
r.POST("/messages/sync", m.handleSyncPost)
}
type sendMessageRequest struct {
FromUID string `json:"from_uid" binding:"required"`
ChannelID string `json:"channel_id" binding:"required"`
ChannelType uint8 `json:"channel_type" binding:"required"`
Payload string `json:"payload" binding:"required"` // raw text, will be base64-encoded
}
func (m *MessageModule) handleSend(c *gin.Context) {
var req sendMessageRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "参数错误: from_uid, channel_id, channel_type, payload必填"})
return
}
// Base64-encode the payload for WuKongIM
encoded := base64.StdEncoding.EncodeToString([]byte(req.Payload))
resp, err := m.wk.SendMessage(wkapi.SendMessageRequest{
FromUID: req.FromUID,
ChannelID: req.ChannelID,
ChannelType: req.ChannelType,
ClientMsgNo: fmt.Sprintf("msg_%d", time.Now().UnixNano()),
Payload: encoded,
})
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "发送消息失败: " + err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{
"message_id": resp.MessageID,
"message_seq": resp.MessageSeq,
})
}
type syncMessagesRequest struct {
ChannelID string `form:"channel_id" binding:"required"`
ChannelType uint8 `form:"channel_type" binding:"required"`
StartMessageSeq uint64 `form:"start_message_seq"`
Limit int `form:"limit"`
}
func (m *MessageModule) handleSync(c *gin.Context) {
var req syncMessagesRequest
if err := c.ShouldBindQuery(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "参数错误: channel_id, channel_type必填"})
return
}
m.doSync(c, req.ChannelID, req.ChannelType, req.StartMessageSeq, req.Limit, "")
}
type syncMessagesPostRequest struct {
ChannelID string `json:"channel_id" binding:"required"`
ChannelType uint8 `json:"channel_type" binding:"required"`
LoginUID string `json:"login_uid"`
StartMessageSeq uint64 `json:"start_message_seq"`
Limit int `json:"limit"`
}
func (m *MessageModule) handleSyncPost(c *gin.Context) {
var req syncMessagesPostRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "参数错误: channel_id, channel_type必填"})
return
}
m.doSync(c, req.ChannelID, req.ChannelType, req.StartMessageSeq, req.Limit, req.LoginUID)
}
func (m *MessageModule) doSync(c *gin.Context, channelID string, channelType uint8, startSeq uint64, limit int, loginUID string) {
if limit <= 0 {
limit = 20
}
resp, err := m.wk.SyncChannelMessages(wkapi.SyncChannelMessagesRequest{
LoginUID: loginUID,
ChannelID: channelID,
ChannelType: channelType,
StartMessageSeq: startSeq,
Limit: limit,
PullMode: 1,
})
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "同步消息失败: " + err.Error()})
return
}
// Decode base64 payloads for the response
type msgItem struct {
MessageID uint64 `json:"message_id"`
MessageSeq uint64 `json:"message_seq"`
FromUID string `json:"from_uid"`
ChannelID string `json:"channel_id"`
Payload string `json:"payload"`
Timestamp int32 `json:"timestamp"`
}
messages := make([]msgItem, 0, len(resp.Messages))
for _, m := range resp.Messages {
decoded, _ := base64.StdEncoding.DecodeString(m.Payload)
messages = append(messages, msgItem{
MessageID: m.MessageID,
MessageSeq: m.MessageSeq,
FromUID: m.FromUID,
ChannelID: m.ChannelID,
Payload: string(decoded),
Timestamp: m.Timestamp,
})
}
c.JSON(http.StatusOK, gin.H{
"start_message_seq": resp.StartMessageSeq,
"end_message_seq": resp.EndMessageSeq,
"more": resp.More,
"messages": messages,
})
}