refactor(onebot): simplify channel implementation and add emoji reaction

- onebot.go from 1179 to 980 lines (~17%)
This commit is contained in:
Hoshina 2026-02-17 00:21:25 +08:00
parent 43a443e64b
commit 387330b3a0

View file

@ -37,6 +37,7 @@ type OneBotChannel struct {
pendingMu sync.Mutex pendingMu sync.Mutex
transcriber *voice.GroqTranscriber transcriber *voice.GroqTranscriber
lastMessageID sync.Map lastMessageID sync.Map
pendingEmojiMsg sync.Map
} }
type oneBotRawEvent struct { type oneBotRawEvent struct {
@ -85,25 +86,6 @@ type oneBotSender struct {
Card string `json:"card"` Card string `json:"card"`
} }
type oneBotEvent struct {
PostType string
MessageType string
SubType string
MessageID string
UserID int64
GroupID int64
Content string
RawContent string
IsBotMentioned bool
Sender oneBotSender
SelfID int64
Time int64
MetaEventType string
Media []string
LocalFiles []string
ReplyTo string
}
type oneBotAPIRequest struct { type oneBotAPIRequest struct {
Action string `json:"action"` Action string `json:"action"`
Params interface{} `json:"params"` Params interface{} `json:"params"`
@ -115,16 +97,6 @@ type oneBotMessageSegment struct {
Data map[string]interface{} `json:"data"` Data map[string]interface{} `json:"data"`
} }
type oneBotSendPrivateMsgParams struct {
UserID int64 `json:"user_id"`
Message interface{} `json:"message"`
}
type oneBotSendGroupMsgParams struct {
GroupID int64 `json:"group_id"`
Message interface{} `json:"message"`
}
func NewOneBotChannel(cfg config.OneBotConfig, messageBus *bus.MessageBus) (*OneBotChannel, error) { func NewOneBotChannel(cfg config.OneBotConfig, messageBus *bus.MessageBus) (*OneBotChannel, error) {
base := NewBaseChannel("onebot", cfg, messageBus, cfg.AllowFrom) base := NewBaseChannel("onebot", cfg, messageBus, cfg.AllowFrom)
@ -143,6 +115,22 @@ func (c *OneBotChannel) SetTranscriber(transcriber *voice.GroqTranscriber) {
c.transcriber = transcriber c.transcriber = transcriber
} }
func (c *OneBotChannel) setMsgEmojiLike(messageID string, emojiID int, set bool) {
go func() {
_, err := c.sendAPIRequest("set_msg_emoji_like", map[string]interface{}{
"message_id": messageID,
"emoji_id": emojiID,
"set": set,
}, 5*time.Second)
if err != nil {
logger.DebugCF("onebot", "Failed to set emoji like", map[string]interface{}{
"message_id": messageID,
"error": err.Error(),
})
}
}()
}
func (c *OneBotChannel) Start(ctx context.Context) error { func (c *OneBotChannel) Start(ctx context.Context) error {
if c.config.WSUrl == "" { if c.config.WSUrl == "" {
return fmt.Errorf("OneBot ws_url not configured") return fmt.Errorf("OneBot ws_url not configured")
@ -159,8 +147,8 @@ func (c *OneBotChannel) Start(ctx context.Context) error {
"error": err.Error(), "error": err.Error(),
}) })
} else { } else {
c.fetchSelfID()
go c.listen() go c.listen()
c.fetchSelfID()
} }
if c.config.ReconnectInterval > 0 { if c.config.ReconnectInterval > 0 {
@ -238,33 +226,33 @@ func (c *OneBotChannel) fetchSelfID() {
return return
} }
var result struct { type loginInfo struct {
Data struct {
UserID json.RawMessage `json:"user_id"`
Nickname string `json:"nickname"`
} `json:"data"`
}
if err := json.Unmarshal(resp, &result); err == nil && len(result.Data.UserID) > 0 {
if uid, err := parseJSONInt64(result.Data.UserID); err == nil && uid > 0 {
atomic.StoreInt64(&c.selfID, uid)
logger.InfoCF("onebot", "Bot self ID retrieved", map[string]interface{}{
"self_id": uid,
"nickname": result.Data.Nickname,
})
return
}
}
var data struct {
UserID json.RawMessage `json:"user_id"` UserID json.RawMessage `json:"user_id"`
Nickname string `json:"nickname"` Nickname string `json:"nickname"`
} }
if err := json.Unmarshal(resp, &data); err == nil && len(data.UserID) > 0 { for _, extract := range []func() (*loginInfo, error){
if uid, err := parseJSONInt64(data.UserID); err == nil && uid > 0 { func() (*loginInfo, error) {
var w struct {
Data loginInfo `json:"data"`
}
err := json.Unmarshal(resp, &w)
return &w.Data, err
},
func() (*loginInfo, error) {
var f loginInfo
err := json.Unmarshal(resp, &f)
return &f, err
},
} {
info, err := extract()
if err != nil || len(info.UserID) == 0 {
continue
}
if uid, err := parseJSONInt64(info.UserID); err == nil && uid > 0 {
atomic.StoreInt64(&c.selfID, uid) atomic.StoreInt64(&c.selfID, uid)
logger.InfoCF("onebot", "Bot self ID retrieved", map[string]interface{}{ logger.InfoCF("onebot", "Bot self ID retrieved", map[string]interface{}{
"self_id": uid, "self_id": uid,
"nickname": data.Nickname, "nickname": info.Nickname,
}) })
return return
} }
@ -348,8 +336,8 @@ func (c *OneBotChannel) reconnectLoop() {
"error": err.Error(), "error": err.Error(),
}) })
} else { } else {
c.fetchSelfID()
go c.listen() go c.listen()
c.fetchSelfID()
} }
} }
} }
@ -423,6 +411,12 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error
return err return err
} }
if msgID, ok := c.pendingEmojiMsg.LoadAndDelete(msg.ChatID); ok {
if mid, ok := msgID.(string); ok && mid != "" {
c.setMsgEmojiLike(mid, 289, false)
}
}
return nil return nil
} }
@ -448,41 +442,23 @@ func (c *OneBotChannel) buildMessageSegments(chatID, content string) []oneBotMes
func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, interface{}, error) { func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, interface{}, error) {
chatID := msg.ChatID chatID := msg.ChatID
if strings.HasPrefix(chatID, "group:") {
groupID, err := strconv.ParseInt(chatID[6:], 10, 64)
if err != nil {
return "", nil, fmt.Errorf("invalid group ID in chatID: %s", chatID)
}
segments := c.buildMessageSegments(chatID, msg.Content) segments := c.buildMessageSegments(chatID, msg.Content)
return "send_group_msg", oneBotSendGroupMsgParams{
GroupID: groupID, var action, idKey string
Message: segments, var rawID string
}, nil if rest, ok := strings.CutPrefix(chatID, "group:"); ok {
action, idKey, rawID = "send_group_msg", "group_id", rest
} else if rest, ok := strings.CutPrefix(chatID, "private:"); ok {
action, idKey, rawID = "send_private_msg", "user_id", rest
} else {
action, idKey, rawID = "send_private_msg", "user_id", chatID
} }
if strings.HasPrefix(chatID, "private:") { id, err := strconv.ParseInt(rawID, 10, 64)
userID, err := strconv.ParseInt(chatID[8:], 10, 64)
if err != nil { if err != nil {
return "", nil, fmt.Errorf("invalid user ID in chatID: %s", chatID) return "", nil, fmt.Errorf("invalid %s in chatID: %s", idKey, chatID)
} }
segments := c.buildMessageSegments(chatID, msg.Content) return action, map[string]interface{}{idKey: id, "message": segments}, nil
return "send_private_msg", oneBotSendPrivateMsgParams{
UserID: userID,
Message: segments,
}, nil
}
userID, err := strconv.ParseInt(chatID, 10, 64)
if err != nil {
return "", nil, fmt.Errorf("invalid chatID for OneBot: %s", chatID)
}
segments := c.buildMessageSegments(chatID, msg.Content)
return "send_private_msg", oneBotSendPrivateMsgParams{
UserID: userID,
Message: segments,
}, nil
} }
func (c *OneBotChannel) listen() { func (c *OneBotChannel) listen() {
@ -516,11 +492,6 @@ func (c *OneBotChannel) listen() {
_ = conn.SetReadDeadline(time.Now().Add(60 * time.Second)) _ = conn.SetReadDeadline(time.Now().Add(60 * time.Second))
logger.DebugCF("onebot", "Raw WebSocket message received", map[string]interface{}{
"length": len(message),
"payload": string(message),
})
var raw oneBotRawEvent var raw oneBotRawEvent
if err := json.Unmarshal(message, &raw); err != nil { if err := json.Unmarshal(message, &raw); err != nil {
logger.WarnCF("onebot", "Failed to unmarshal raw event", map[string]interface{}{ logger.WarnCF("onebot", "Failed to unmarshal raw event", map[string]interface{}{
@ -530,6 +501,12 @@ func (c *OneBotChannel) listen() {
continue continue
} }
logger.DebugCF("onebot", "WebSocket event", map[string]interface{}{
"length": len(message),
"post_type": raw.PostType,
"sub_type": raw.SubType,
})
if raw.Echo != "" { if raw.Echo != "" {
c.pendingMu.Lock() c.pendingMu.Lock()
ch, ok := c.pending[raw.Echo] ch, ok := c.pending[raw.Echo]
@ -556,13 +533,6 @@ func (c *OneBotChannel) listen() {
continue continue
} }
logger.DebugCF("onebot", "Parsed raw event", map[string]interface{}{
"post_type": raw.PostType,
"message_type": raw.MessageType,
"sub_type": raw.SubType,
"meta_event_type": raw.MetaEventType,
})
c.handleRawEvent(&raw) c.handleRawEvent(&raw)
} }
} }
@ -656,13 +626,16 @@ func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64)
} }
} }
case "image": case "image", "video", "file":
if data != nil { if data != nil {
url, _ := data["url"].(string) url, _ := data["url"].(string)
if url != "" { if url != "" {
filename := "image.jpg" defaults := map[string]string{"image": "image.jpg", "video": "video.mp4", "file": "file"}
filename := defaults[segType]
if f, ok := data["file"].(string); ok && f != "" { if f, ok := data["file"].(string); ok && f != "" {
filename = f filename = f
} else if n, ok := data["name"].(string); ok && n != "" {
filename = n
} }
localPath := utils.DownloadFile(url, filename, utils.DownloadOptions{ localPath := utils.DownloadFile(url, filename, utils.DownloadOptions{
LoggerPrefix: "onebot", LoggerPrefix: "onebot",
@ -670,7 +643,7 @@ func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64)
if localPath != "" { if localPath != "" {
media = append(media, localPath) media = append(media, localPath)
localFiles = append(localFiles, localPath) localFiles = append(localFiles, localPath)
textParts = append(textParts, "[image]") textParts = append(textParts, fmt.Sprintf("[%s]", segType))
} }
} }
} }
@ -705,44 +678,6 @@ func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64)
} }
} }
case "video":
if data != nil {
url, _ := data["url"].(string)
if url != "" {
filename := "video.mp4"
if f, ok := data["file"].(string); ok && f != "" {
filename = f
}
localPath := utils.DownloadFile(url, filename, utils.DownloadOptions{
LoggerPrefix: "onebot",
})
if localPath != "" {
media = append(media, localPath)
localFiles = append(localFiles, localPath)
textParts = append(textParts, "[video]")
}
}
}
case "file":
if data != nil {
url, _ := data["url"].(string)
name, _ := data["name"].(string)
if url != "" {
if name == "" {
name = "file"
}
localPath := utils.DownloadFile(url, name, utils.DownloadOptions{
LoggerPrefix: "onebot",
})
if localPath != "" {
media = append(media, localPath)
localFiles = append(localFiles, localPath)
textParts = append(textParts, fmt.Sprintf("[file: %s]", name))
}
}
}
case "reply": case "reply":
if data != nil { if data != nil {
if id, ok := data["id"]; ok { if id, ok := data["id"]; ok {
@ -784,15 +719,7 @@ func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) {
return return
} }
} }
c.handleMessage(raw)
evt, err := c.normalizeMessageEvent(raw)
if err != nil {
logger.WarnCF("onebot", "Failed to normalize message event", map[string]interface{}{
"error": err.Error(),
})
return
}
c.handleMessage(evt)
case "message_sent": case "message_sent":
logger.DebugCF("onebot", "Bot sent message event", map[string]interface{}{ logger.DebugCF("onebot", "Bot sent message event", map[string]interface{}{
@ -824,15 +751,44 @@ func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) {
} }
} }
func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent, error) { func (c *OneBotChannel) handleMetaEvent(raw *oneBotRawEvent) {
if raw.MetaEventType == "lifecycle" {
logger.InfoCF("onebot", "Lifecycle event", map[string]interface{}{"sub_type": raw.SubType})
} else if raw.MetaEventType != "heartbeat" {
logger.DebugCF("onebot", "Meta event: "+raw.MetaEventType, nil)
}
}
func (c *OneBotChannel) handleNoticeEvent(raw *oneBotRawEvent) {
fields := map[string]interface{}{
"notice_type": raw.NoticeType,
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
"message_id": parseJSONString(raw.MessageID),
}
switch raw.NoticeType {
case "group_recall", "group_increase", "group_decrease",
"friend_add", "group_admin", "group_ban":
logger.InfoCF("onebot", "Notice: "+raw.NoticeType, fields)
default:
logger.DebugCF("onebot", "Notice: "+raw.NoticeType, fields)
}
}
func (c *OneBotChannel) handleMessage(raw *oneBotRawEvent) {
// Parse fields from raw event
userID, err := parseJSONInt64(raw.UserID) userID, err := parseJSONInt64(raw.UserID)
if err != nil { if err != nil {
return nil, fmt.Errorf("parse user_id: %w (raw: %s)", err, string(raw.UserID)) logger.WarnCF("onebot", "Failed to parse user_id", map[string]interface{}{
"error": err.Error(),
"raw": string(raw.UserID),
})
return
} }
groupID, _ := parseJSONInt64(raw.GroupID) groupID, _ := parseJSONInt64(raw.GroupID)
selfID, _ := parseJSONInt64(raw.SelfID) selfID, _ := parseJSONInt64(raw.SelfID)
ts, _ := parseJSONInt64(raw.Time)
messageID := parseJSONString(raw.MessageID) messageID := parseJSONString(raw.MessageID)
if selfID == 0 { if selfID == 0 {
@ -868,117 +824,10 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent
} }
} }
logger.DebugCF("onebot", "Normalized message event", map[string]interface{}{ // Clean up temp files when done
"message_type": raw.MessageType, if len(parsed.LocalFiles) > 0 {
"user_id": userID,
"group_id": groupID,
"message_id": messageID,
"content_len": len(content),
"nickname": sender.Nickname,
"media_count": len(parsed.Media),
"has_reply": parsed.ReplyTo != "",
"has_localfiles": len(parsed.LocalFiles) > 0,
})
return &oneBotEvent{
PostType: raw.PostType,
MessageType: raw.MessageType,
SubType: raw.SubType,
MessageID: messageID,
UserID: userID,
GroupID: groupID,
Content: content,
RawContent: raw.RawMessage,
IsBotMentioned: isBotMentioned,
Sender: sender,
SelfID: selfID,
Time: ts,
MetaEventType: raw.MetaEventType,
Media: parsed.Media,
LocalFiles: parsed.LocalFiles,
ReplyTo: parsed.ReplyTo,
}, nil
}
func (c *OneBotChannel) handleMetaEvent(raw *oneBotRawEvent) {
switch raw.MetaEventType {
case "lifecycle":
logger.InfoCF("onebot", "Lifecycle event", map[string]interface{}{
"sub_type": raw.SubType,
})
case "heartbeat":
logger.DebugC("onebot", "Heartbeat received")
default:
logger.DebugCF("onebot", "Unknown meta_event_type", map[string]interface{}{
"meta_event_type": raw.MetaEventType,
})
}
}
func (c *OneBotChannel) handleNoticeEvent(raw *oneBotRawEvent) {
switch raw.NoticeType {
case "group_recall":
logger.InfoCF("onebot", "Group message recalled", map[string]interface{}{
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
"message_id": parseJSONString(raw.MessageID),
})
case "notify":
logger.DebugCF("onebot", "Notify event", map[string]interface{}{
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
case "group_increase":
logger.InfoCF("onebot", "Group member increased", map[string]interface{}{
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
case "group_decrease":
logger.InfoCF("onebot", "Group member decreased", map[string]interface{}{
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
case "friend_add":
logger.InfoCF("onebot", "New friend added", map[string]interface{}{
"user_id": parseJSONString(raw.UserID),
})
case "friend_recall":
logger.DebugCF("onebot", "Friend message recalled", map[string]interface{}{
"user_id": parseJSONString(raw.UserID),
"message_id": parseJSONString(raw.MessageID),
})
case "group_upload":
logger.DebugCF("onebot", "Group file uploaded", map[string]interface{}{
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
case "group_admin":
logger.InfoCF("onebot", "Group admin changed", map[string]interface{}{
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
case "group_ban":
logger.InfoCF("onebot", "Group ban event", map[string]interface{}{
"sub_type": raw.SubType,
"group_id": parseJSONString(raw.GroupID),
"user_id": parseJSONString(raw.UserID),
})
default:
logger.DebugCF("onebot", "Notice event received", map[string]interface{}{
"notice_type": raw.NoticeType,
"sub_type": raw.SubType,
})
}
}
func (c *OneBotChannel) handleMessage(evt *oneBotEvent) {
if len(evt.LocalFiles) > 0 {
defer func() { defer func() {
for _, f := range evt.LocalFiles { for _, f := range parsed.LocalFiles {
if err := os.Remove(f); err != nil { if err := os.Remove(f); err != nil {
logger.DebugCF("onebot", "Failed to remove temp file", map[string]interface{}{ logger.DebugCF("onebot", "Failed to remove temp file", map[string]interface{}{
"path": f, "path": f,
@ -989,104 +838,93 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) {
}() }()
} }
if c.isDuplicate(evt.MessageID) { if c.isDuplicate(messageID) {
logger.DebugCF("onebot", "Duplicate message, skipping", map[string]interface{}{ logger.DebugCF("onebot", "Duplicate message, skipping", map[string]interface{}{
"message_id": evt.MessageID, "message_id": messageID,
}) })
return return
} }
content := evt.Content
if content == "" { if content == "" {
logger.DebugCF("onebot", "Received empty message, ignoring", map[string]interface{}{ logger.DebugCF("onebot", "Received empty message, ignoring", map[string]interface{}{
"message_id": evt.MessageID, "message_id": messageID,
}) })
return return
} }
senderID := strconv.FormatInt(evt.UserID, 10) senderID := strconv.FormatInt(userID, 10)
var chatID string var chatID string
metadata := map[string]string{ metadata := map[string]string{
"message_id": evt.MessageID, "message_id": messageID,
} }
if evt.ReplyTo != "" { if parsed.ReplyTo != "" {
metadata["reply_to_message_id"] = evt.ReplyTo metadata["reply_to_message_id"] = parsed.ReplyTo
} }
switch evt.MessageType { switch raw.MessageType {
case "private": case "private":
chatID = "private:" + senderID chatID = "private:" + senderID
logger.InfoCF("onebot", "Received private message", map[string]interface{}{
"sender": senderID,
"message_id": evt.MessageID,
"length": len(content),
"content": truncate(content, 100),
"media_count": len(evt.Media),
})
case "group": case "group":
groupIDStr := strconv.FormatInt(evt.GroupID, 10) groupIDStr := strconv.FormatInt(groupID, 10)
chatID = "group:" + groupIDStr chatID = "group:" + groupIDStr
metadata["group_id"] = groupIDStr metadata["group_id"] = groupIDStr
senderUserID, _ := parseJSONInt64(evt.Sender.UserID) senderUserID, _ := parseJSONInt64(sender.UserID)
if senderUserID > 0 { if senderUserID > 0 {
metadata["sender_user_id"] = strconv.FormatInt(senderUserID, 10) metadata["sender_user_id"] = strconv.FormatInt(senderUserID, 10)
} }
if evt.Sender.Card != "" { if sender.Card != "" {
metadata["sender_name"] = evt.Sender.Card metadata["sender_name"] = sender.Card
} else if evt.Sender.Nickname != "" { } else if sender.Nickname != "" {
metadata["sender_name"] = evt.Sender.Nickname metadata["sender_name"] = sender.Nickname
} }
triggered, strippedContent := c.checkGroupTrigger(content, evt.IsBotMentioned) triggered, strippedContent := c.checkGroupTrigger(content, isBotMentioned)
if !triggered { if !triggered {
logger.DebugCF("onebot", "Group message ignored (no trigger)", map[string]interface{}{ logger.DebugCF("onebot", "Group message ignored (no trigger)", map[string]interface{}{
"sender": senderID, "sender": senderID,
"group": groupIDStr, "group": groupIDStr,
"is_mentioned": evt.IsBotMentioned, "is_mentioned": isBotMentioned,
"content": truncate(content, 100), "content": truncate(content, 100),
}) })
return return
} }
content = strippedContent content = strippedContent
logger.InfoCF("onebot", "Received group message", map[string]interface{}{
"sender": senderID,
"group": groupIDStr,
"message_id": evt.MessageID,
"is_mentioned": evt.IsBotMentioned,
"length": len(content),
"content": truncate(content, 100),
"media_count": len(evt.Media),
})
default: default:
logger.WarnCF("onebot", "Unknown message type, cannot route", map[string]interface{}{ logger.WarnCF("onebot", "Unknown message type, cannot route", map[string]interface{}{
"type": evt.MessageType, "type": raw.MessageType,
"message_id": evt.MessageID, "message_id": messageID,
"user_id": evt.UserID, "user_id": userID,
}) })
return return
} }
if evt.Sender.Nickname != "" { logger.InfoCF("onebot", "Received "+raw.MessageType+" message", map[string]interface{}{
metadata["nickname"] = evt.Sender.Nickname "sender": senderID,
}
c.lastMessageID.Store(chatID, evt.MessageID)
logger.DebugCF("onebot", "Forwarding message to bus", map[string]interface{}{
"sender_id": senderID,
"chat_id": chatID, "chat_id": chatID,
"message_id": messageID,
"length": len(content),
"content": truncate(content, 100), "content": truncate(content, 100),
"media_count": len(evt.Media), "media_count": len(parsed.Media),
}) })
c.HandleMessage(senderID, chatID, content, evt.Media, metadata) if sender.Nickname != "" {
metadata["nickname"] = sender.Nickname
}
c.lastMessageID.Store(chatID, messageID)
if raw.MessageType == "group" && messageID != "" && messageID != "0" {
c.setMsgEmojiLike(messageID, 289, true)
c.pendingEmojiMsg.Store(chatID, messageID)
}
c.HandleMessage(senderID, chatID, content, parsed.Media, metadata)
} }
func (c *OneBotChannel) isDuplicate(messageID string) bool { func (c *OneBotChannel) isDuplicate(messageID string) bool {