From 387330b3a01fb64e5417e4fc46ee2cce4661f4e0 Mon Sep 17 00:00:00 2001 From: Hoshina Date: Tue, 17 Feb 2026 00:21:25 +0800 Subject: [PATCH] refactor(onebot): simplify channel implementation and add emoji reaction - onebot.go from 1179 to 980 lines (~17%) --- pkg/channels/onebot.go | 484 ++++++++++++++--------------------------- 1 file changed, 161 insertions(+), 323 deletions(-) diff --git a/pkg/channels/onebot.go b/pkg/channels/onebot.go index ef7b84c4a..53e82b44d 100644 --- a/pkg/channels/onebot.go +++ b/pkg/channels/onebot.go @@ -22,21 +22,22 @@ import ( type OneBotChannel struct { *BaseChannel - config config.OneBotConfig - conn *websocket.Conn - ctx context.Context - cancel context.CancelFunc - dedup map[string]struct{} - dedupRing []string - dedupIdx int - mu sync.Mutex - writeMu sync.Mutex - echoCounter int64 - selfID int64 - pending map[string]chan json.RawMessage - pendingMu sync.Mutex - transcriber *voice.GroqTranscriber - lastMessageID sync.Map + config config.OneBotConfig + conn *websocket.Conn + ctx context.Context + cancel context.CancelFunc + dedup map[string]struct{} + dedupRing []string + dedupIdx int + mu sync.Mutex + writeMu sync.Mutex + echoCounter int64 + selfID int64 + pending map[string]chan json.RawMessage + pendingMu sync.Mutex + transcriber *voice.GroqTranscriber + lastMessageID sync.Map + pendingEmojiMsg sync.Map } type oneBotRawEvent struct { @@ -85,25 +86,6 @@ type oneBotSender struct { 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 { Action string `json:"action"` Params interface{} `json:"params"` @@ -115,16 +97,6 @@ type oneBotMessageSegment struct { 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) { base := NewBaseChannel("onebot", cfg, messageBus, cfg.AllowFrom) @@ -143,6 +115,22 @@ func (c *OneBotChannel) SetTranscriber(transcriber *voice.GroqTranscriber) { 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 { if c.config.WSUrl == "" { return fmt.Errorf("OneBot ws_url not configured") @@ -159,8 +147,8 @@ func (c *OneBotChannel) Start(ctx context.Context) error { "error": err.Error(), }) } else { - c.fetchSelfID() go c.listen() + c.fetchSelfID() } if c.config.ReconnectInterval > 0 { @@ -238,33 +226,33 @@ func (c *OneBotChannel) fetchSelfID() { return } - var result 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 { + type loginInfo struct { UserID json.RawMessage `json:"user_id"` Nickname string `json:"nickname"` } - if err := json.Unmarshal(resp, &data); err == nil && len(data.UserID) > 0 { - if uid, err := parseJSONInt64(data.UserID); err == nil && uid > 0 { + for _, extract := range []func() (*loginInfo, error){ + 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) logger.InfoCF("onebot", "Bot self ID retrieved", map[string]interface{}{ "self_id": uid, - "nickname": data.Nickname, + "nickname": info.Nickname, }) return } @@ -348,8 +336,8 @@ func (c *OneBotChannel) reconnectLoop() { "error": err.Error(), }) } else { - c.fetchSelfID() go c.listen() + c.fetchSelfID() } } } @@ -423,6 +411,12 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error 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 } @@ -448,41 +442,23 @@ func (c *OneBotChannel) buildMessageSegments(chatID, content string) []oneBotMes func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, interface{}, error) { 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) - return "send_group_msg", oneBotSendGroupMsgParams{ - GroupID: groupID, - Message: segments, - }, nil - } - - if strings.HasPrefix(chatID, "private:") { - userID, err := strconv.ParseInt(chatID[8:], 10, 64) - if err != nil { - return "", nil, fmt.Errorf("invalid user ID in chatID: %s", chatID) - } - segments := c.buildMessageSegments(chatID, msg.Content) - 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 + + var action, idKey string + var rawID string + 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 + } + + id, err := strconv.ParseInt(rawID, 10, 64) + if err != nil { + return "", nil, fmt.Errorf("invalid %s in chatID: %s", idKey, chatID) + } + return action, map[string]interface{}{idKey: id, "message": segments}, nil } func (c *OneBotChannel) listen() { @@ -516,11 +492,6 @@ func (c *OneBotChannel) listen() { _ = 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 if err := json.Unmarshal(message, &raw); err != nil { logger.WarnCF("onebot", "Failed to unmarshal raw event", map[string]interface{}{ @@ -530,6 +501,12 @@ func (c *OneBotChannel) listen() { continue } + logger.DebugCF("onebot", "WebSocket event", map[string]interface{}{ + "length": len(message), + "post_type": raw.PostType, + "sub_type": raw.SubType, + }) + if raw.Echo != "" { c.pendingMu.Lock() ch, ok := c.pending[raw.Echo] @@ -556,13 +533,6 @@ func (c *OneBotChannel) listen() { 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) } } @@ -656,13 +626,16 @@ func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64) } } - case "image": + case "image", "video", "file": if data != nil { url, _ := data["url"].(string) 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 != "" { filename = f + } else if n, ok := data["name"].(string); ok && n != "" { + filename = n } localPath := utils.DownloadFile(url, filename, utils.DownloadOptions{ LoggerPrefix: "onebot", @@ -670,7 +643,7 @@ func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64) if localPath != "" { media = append(media, 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": if data != nil { if id, ok := data["id"]; ok { @@ -784,15 +719,7 @@ func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) { return } } - - 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) + c.handleMessage(raw) case "message_sent": 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) 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) selfID, _ := parseJSONInt64(raw.SelfID) - ts, _ := parseJSONInt64(raw.Time) messageID := parseJSONString(raw.MessageID) if selfID == 0 { @@ -868,117 +824,10 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent } } - logger.DebugCF("onebot", "Normalized message event", map[string]interface{}{ - "message_type": raw.MessageType, - "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 { + // Clean up temp files when done + if len(parsed.LocalFiles) > 0 { defer func() { - for _, f := range evt.LocalFiles { + for _, f := range parsed.LocalFiles { if err := os.Remove(f); err != nil { logger.DebugCF("onebot", "Failed to remove temp file", map[string]interface{}{ "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{}{ - "message_id": evt.MessageID, + "message_id": messageID, }) return } - content := evt.Content if content == "" { logger.DebugCF("onebot", "Received empty message, ignoring", map[string]interface{}{ - "message_id": evt.MessageID, + "message_id": messageID, }) return } - senderID := strconv.FormatInt(evt.UserID, 10) + senderID := strconv.FormatInt(userID, 10) var chatID string metadata := map[string]string{ - "message_id": evt.MessageID, + "message_id": messageID, } - if evt.ReplyTo != "" { - metadata["reply_to_message_id"] = evt.ReplyTo + if parsed.ReplyTo != "" { + metadata["reply_to_message_id"] = parsed.ReplyTo } - switch evt.MessageType { + switch raw.MessageType { case "private": 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": - groupIDStr := strconv.FormatInt(evt.GroupID, 10) + groupIDStr := strconv.FormatInt(groupID, 10) chatID = "group:" + groupIDStr metadata["group_id"] = groupIDStr - senderUserID, _ := parseJSONInt64(evt.Sender.UserID) + senderUserID, _ := parseJSONInt64(sender.UserID) if senderUserID > 0 { metadata["sender_user_id"] = strconv.FormatInt(senderUserID, 10) } - if evt.Sender.Card != "" { - metadata["sender_name"] = evt.Sender.Card - } else if evt.Sender.Nickname != "" { - metadata["sender_name"] = evt.Sender.Nickname + if sender.Card != "" { + metadata["sender_name"] = sender.Card + } else if sender.Nickname != "" { + metadata["sender_name"] = sender.Nickname } - triggered, strippedContent := c.checkGroupTrigger(content, evt.IsBotMentioned) + triggered, strippedContent := c.checkGroupTrigger(content, isBotMentioned) if !triggered { logger.DebugCF("onebot", "Group message ignored (no trigger)", map[string]interface{}{ "sender": senderID, "group": groupIDStr, - "is_mentioned": evt.IsBotMentioned, + "is_mentioned": isBotMentioned, "content": truncate(content, 100), }) return } 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: logger.WarnCF("onebot", "Unknown message type, cannot route", map[string]interface{}{ - "type": evt.MessageType, - "message_id": evt.MessageID, - "user_id": evt.UserID, + "type": raw.MessageType, + "message_id": messageID, + "user_id": userID, }) return } - if evt.Sender.Nickname != "" { - metadata["nickname"] = evt.Sender.Nickname - } - - c.lastMessageID.Store(chatID, evt.MessageID) - - logger.DebugCF("onebot", "Forwarding message to bus", map[string]interface{}{ - "sender_id": senderID, + logger.InfoCF("onebot", "Received "+raw.MessageType+" message", map[string]interface{}{ + "sender": senderID, "chat_id": chatID, + "message_id": messageID, + "length": len(content), "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 {