144 lines
4.3 KiB
Go
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,
|
|
})
|
|
}
|