From 50e4edc8104e81054eeb7c9259bc8429d0c7803f Mon Sep 17 00:00:00 2001 From: Hoshina Date: Sun, 15 Feb 2026 14:09:06 +0800 Subject: [PATCH] feat(onebot): add rich media, API callback, keepalive and voice transcription Comprehensive improvements to the OneBot channel for better NapCatQQ compatibility: - Add echo-based API callback mechanism (sendAPIRequest) for request/response correlation via pending map - Add WebSocket ping/pong keepalive (30s ping, 60s read deadline) - Fetch bot self ID via get_login_info on connect/reconnect - Refactor parseMessageContentEx into parseMessageSegments supporting image, record, video, file, reply, face, forward segments - Add voice transcription via Groq transcriber (SetTranscriber) - Switch to message segment array format for sending with auto reply quote via lastMessageID tracking - Add message_sent event handling and detailed notice event processing (recall, poke, group increase/decrease, friend add, etc.) - Use sync/atomic for echoCounter, optimize listen() lock pattern - Clean up pending callbacks on Stop(), defer temp file cleanup - Mount Groq transcriber on OneBot channel in main.go gateway --- cmd/picoclaw/main.go | 6 + pkg/channels/onebot.go | 579 +++++++++++++++++++++++++++++++++++------ 2 files changed, 509 insertions(+), 76 deletions(-) diff --git a/cmd/picoclaw/main.go b/cmd/picoclaw/main.go index 2129662d7..c21443919 100644 --- a/cmd/picoclaw/main.go +++ b/cmd/picoclaw/main.go @@ -617,6 +617,12 @@ func gatewayCmd() { logger.InfoC("voice", "Groq transcription attached to Slack channel") } } + if onebotChannel, ok := channelManager.GetChannel("onebot"); ok { + if oc, ok := onebotChannel.(*channels.OneBotChannel); ok { + oc.SetTranscriber(transcriber) + logger.InfoC("voice", "Groq transcription attached to OneBot channel") + } + } } enabledChannels := channelManager.GetEnabledChannels() diff --git a/pkg/channels/onebot.go b/pkg/channels/onebot.go index f301646d8..f3932d796 100644 --- a/pkg/channels/onebot.go +++ b/pkg/channels/onebot.go @@ -4,9 +4,11 @@ import ( "context" "encoding/json" "fmt" + "os" "strconv" "strings" "sync" + "sync/atomic" "time" "github.com/gorilla/websocket" @@ -14,20 +16,27 @@ import ( "github.com/sipeed/picoclaw/pkg/bus" "github.com/sipeed/picoclaw/pkg/config" "github.com/sipeed/picoclaw/pkg/logger" + "github.com/sipeed/picoclaw/pkg/utils" + "github.com/sipeed/picoclaw/pkg/voice" ) 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 + 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 } type oneBotRawEvent struct { @@ -43,9 +52,11 @@ type oneBotRawEvent struct { SelfID json.RawMessage `json:"self_id"` Time json.RawMessage `json:"time"` MetaEventType string `json:"meta_event_type"` + NoticeType string `json:"notice_type"` Echo string `json:"echo"` RetCode json.RawMessage `json:"retcode"` Status json.RawMessage `json:"status"` + Data json.RawMessage `json:"data"` } type BotStatus struct { @@ -88,6 +99,9 @@ type oneBotEvent struct { SelfID int64 Time int64 MetaEventType string + Media []string + LocalFiles []string + ReplyTo string } type oneBotAPIRequest struct { @@ -96,14 +110,19 @@ type oneBotAPIRequest struct { Echo string `json:"echo,omitempty"` } +type oneBotMessageSegment struct { + Type string `json:"type"` + Data map[string]interface{} `json:"data"` +} + type oneBotSendPrivateMsgParams struct { - UserID int64 `json:"user_id"` - Message string `json:"message"` + UserID int64 `json:"user_id"` + Message interface{} `json:"message"` } type oneBotSendGroupMsgParams struct { - GroupID int64 `json:"group_id"` - Message string `json:"message"` + GroupID int64 `json:"group_id"` + Message interface{} `json:"message"` } func NewOneBotChannel(cfg config.OneBotConfig, messageBus *bus.MessageBus) (*OneBotChannel, error) { @@ -116,9 +135,14 @@ func NewOneBotChannel(cfg config.OneBotConfig, messageBus *bus.MessageBus) (*One dedup: make(map[string]struct{}, dedupSize), dedupRing: make([]string, dedupSize), dedupIdx: 0, + pending: make(map[string]chan json.RawMessage), }, nil } +func (c *OneBotChannel) SetTranscriber(transcriber *voice.GroqTranscriber) { + c.transcriber = transcriber +} + func (c *OneBotChannel) Start(ctx context.Context) error { if c.config.WSUrl == "" { return fmt.Errorf("OneBot ws_url not configured") @@ -135,13 +159,13 @@ func (c *OneBotChannel) Start(ctx context.Context) error { "error": err.Error(), }) } else { + c.fetchSelfID() go c.listen() } if c.config.ReconnectInterval > 0 { go c.reconnectLoop() } else { - // If reconnect is disabled but initial connection failed, we cannot recover if c.conn == nil { return fmt.Errorf("failed to connect to OneBot and reconnect is disabled") } @@ -167,14 +191,141 @@ func (c *OneBotChannel) connect() error { return err } + conn.SetPongHandler(func(appData string) error { + _ = conn.SetReadDeadline(time.Now().Add(60 * time.Second)) + return nil + }) + _ = conn.SetReadDeadline(time.Now().Add(60 * time.Second)) + c.mu.Lock() c.conn = conn c.mu.Unlock() + go c.pinger(conn) + logger.InfoC("onebot", "WebSocket connected") return nil } +func (c *OneBotChannel) pinger(conn *websocket.Conn) { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + for { + select { + case <-c.ctx.Done(): + return + case <-ticker.C: + c.writeMu.Lock() + err := conn.WriteMessage(websocket.PingMessage, nil) + c.writeMu.Unlock() + if err != nil { + logger.DebugCF("onebot", "Ping write failed, stopping pinger", map[string]interface{}{ + "error": err.Error(), + }) + return + } + } + } +} + +func (c *OneBotChannel) fetchSelfID() { + resp, err := c.sendAPIRequest("get_login_info", nil, 5*time.Second) + if err != nil { + logger.WarnCF("onebot", "Failed to get_login_info", map[string]interface{}{ + "error": err.Error(), + }) + 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 { + 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 { + atomic.StoreInt64(&c.selfID, uid) + logger.InfoCF("onebot", "Bot self ID retrieved", map[string]interface{}{ + "self_id": uid, + "nickname": data.Nickname, + }) + return + } + } + + logger.WarnCF("onebot", "Could not parse self ID from get_login_info response", map[string]interface{}{ + "response": string(resp), + }) +} + +func (c *OneBotChannel) sendAPIRequest(action string, params interface{}, timeout time.Duration) (json.RawMessage, error) { + c.mu.Lock() + conn := c.conn + c.mu.Unlock() + + if conn == nil { + return nil, fmt.Errorf("WebSocket not connected") + } + + echo := fmt.Sprintf("api_%d_%d", time.Now().UnixNano(), atomic.AddInt64(&c.echoCounter, 1)) + + ch := make(chan json.RawMessage, 1) + c.pendingMu.Lock() + c.pending[echo] = ch + c.pendingMu.Unlock() + + defer func() { + c.pendingMu.Lock() + delete(c.pending, echo) + c.pendingMu.Unlock() + }() + + req := oneBotAPIRequest{ + Action: action, + Params: params, + Echo: echo, + } + + data, err := json.Marshal(req) + if err != nil { + return nil, fmt.Errorf("failed to marshal API request: %w", err) + } + + c.writeMu.Lock() + err = conn.WriteMessage(websocket.TextMessage, data) + c.writeMu.Unlock() + + if err != nil { + return nil, fmt.Errorf("failed to write API request: %w", err) + } + + select { + case resp := <-ch: + return resp, nil + case <-time.After(timeout): + return nil, fmt.Errorf("API request %s timed out after %v", action, timeout) + case <-c.ctx.Done(): + return nil, fmt.Errorf("context cancelled") + } +} + func (c *OneBotChannel) reconnectLoop() { interval := time.Duration(c.config.ReconnectInterval) * time.Second if interval < 5*time.Second { @@ -197,6 +348,7 @@ func (c *OneBotChannel) reconnectLoop() { "error": err.Error(), }) } else { + c.fetchSelfID() go c.listen() } } @@ -212,6 +364,13 @@ func (c *OneBotChannel) Stop(ctx context.Context) error { c.cancel() } + c.pendingMu.Lock() + for echo, ch := range c.pending { + close(ch) + delete(c.pending, echo) + } + c.pendingMu.Unlock() + c.mu.Lock() if c.conn != nil { c.conn.Close() @@ -240,10 +399,7 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error return err } - c.writeMu.Lock() - c.echoCounter++ - echo := fmt.Sprintf("send_%d", c.echoCounter) - c.writeMu.Unlock() + echo := fmt.Sprintf("send_%d", atomic.AddInt64(&c.echoCounter, 1)) req := oneBotAPIRequest{ Action: action, @@ -270,28 +426,50 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error return nil } +func (c *OneBotChannel) buildMessageSegments(chatID, content string) []oneBotMessageSegment { + var segments []oneBotMessageSegment + + if lastMsgID, ok := c.lastMessageID.Load(chatID); ok { + if msgID, ok := lastMsgID.(string); ok && msgID != "" { + segments = append(segments, oneBotMessageSegment{ + Type: "reply", + Data: map[string]interface{}{"id": msgID}, + }) + } + } + + segments = append(segments, oneBotMessageSegment{ + Type: "text", + Data: map[string]interface{}{"text": content}, + }) + + return segments +} + func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, interface{}, error) { chatID := msg.ChatID - if len(chatID) > 6 && chatID[:6] == "group:" { + 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: msg.Content, + Message: segments, }, nil } - if len(chatID) > 8 && chatID[:8] == "private:" { + 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: msg.Content, + Message: segments, }, nil } @@ -300,34 +478,35 @@ func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, inter return "", nil, fmt.Errorf("invalid chatID for OneBot: %s", chatID) } + segments := c.buildMessageSegments(chatID, msg.Content) return "send_private_msg", oneBotSendPrivateMsgParams{ UserID: userID, - Message: msg.Content, + Message: segments, }, nil } func (c *OneBotChannel) listen() { + c.mu.Lock() + conn := c.conn + c.mu.Unlock() + + if conn == nil { + logger.WarnC("onebot", "WebSocket connection is nil, listener exiting") + return + } + for { select { case <-c.ctx.Done(): return default: - c.mu.Lock() - conn := c.conn - c.mu.Unlock() - - if conn == nil { - logger.WarnC("onebot", "WebSocket connection is nil, listener exiting") - return - } - _, message, err := conn.ReadMessage() if err != nil { logger.ErrorCF("onebot", "WebSocket read error", map[string]interface{}{ "error": err.Error(), }) c.mu.Lock() - if c.conn != nil { + if c.conn == conn { c.conn.Close() c.conn = nil } @@ -335,6 +514,8 @@ func (c *OneBotChannel) listen() { return } + _ = conn.SetReadDeadline(time.Now().Add(60 * time.Second)) + logger.DebugCF("onebot", "Raw WebSocket message received", map[string]interface{}{ "length": len(message), "payload": string(message), @@ -349,9 +530,27 @@ func (c *OneBotChannel) listen() { continue } - if raw.Echo != "" || isAPIResponse(raw.Status) { - logger.DebugCF("onebot", "Received API response, skipping", map[string]interface{}{ - "echo": raw.Echo, + if raw.Echo != "" { + c.pendingMu.Lock() + ch, ok := c.pending[raw.Echo] + c.pendingMu.Unlock() + + if ok { + select { + case ch <- message: + default: + } + } else { + logger.DebugCF("onebot", "Received API response (no waiter)", map[string]interface{}{ + "echo": raw.Echo, + "status": string(raw.Status), + }) + } + continue + } + + if isAPIResponse(raw.Status) { + logger.DebugCF("onebot", "Received API response without echo, skipping", map[string]interface{}{ "status": string(raw.Status), }) continue @@ -401,9 +600,12 @@ func parseJSONString(raw json.RawMessage) string { type parseMessageResult struct { Text string IsBotMentioned bool + Media []string + LocalFiles []string + ReplyTo string } -func parseMessageContentEx(raw json.RawMessage, selfID int64) parseMessageResult { +func (c *OneBotChannel) parseMessageSegments(raw json.RawMessage, selfID int64) parseMessageResult { if len(raw) == 0 { return parseMessageResult{} } @@ -423,32 +625,152 @@ func parseMessageContentEx(raw json.RawMessage, selfID int64) parseMessageResult } var segments []map[string]interface{} - if err := json.Unmarshal(raw, &segments); err == nil { - var text string - mentioned := false - selfIDStr := strconv.FormatInt(selfID, 10) - for _, seg := range segments { - segType, _ := seg["type"].(string) - data, _ := seg["data"].(map[string]interface{}) - switch segType { - case "text": - if data != nil { - if t, ok := data["text"].(string); ok { - text += t - } + if err := json.Unmarshal(raw, &segments); err != nil { + return parseMessageResult{} + } + + var textParts []string + mentioned := false + selfIDStr := strconv.FormatInt(selfID, 10) + var media []string + var localFiles []string + var replyTo string + + for _, seg := range segments { + segType, _ := seg["type"].(string) + data, _ := seg["data"].(map[string]interface{}) + + switch segType { + case "text": + if data != nil { + if t, ok := data["text"].(string); ok { + textParts = append(textParts, t) } - case "at": - if data != nil && selfID > 0 { - qqVal := fmt.Sprintf("%v", data["qq"]) - if qqVal == selfIDStr || qqVal == "all" { - mentioned = true + } + + case "at": + if data != nil && selfID > 0 { + qqVal := fmt.Sprintf("%v", data["qq"]) + if qqVal == selfIDStr || qqVal == "all" { + mentioned = true + } + } + + case "image": + if data != nil { + url, _ := data["url"].(string) + if url != "" { + filename := "image.jpg" + 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, "[image]") } } } + + case "record": + if data != nil { + url, _ := data["url"].(string) + if url != "" { + localPath := utils.DownloadFile(url, "voice.amr", utils.DownloadOptions{ + LoggerPrefix: "onebot", + }) + if localPath != "" { + localFiles = append(localFiles, localPath) + if c.transcriber != nil && c.transcriber.IsAvailable() { + tctx, tcancel := context.WithTimeout(c.ctx, 30*time.Second) + result, err := c.transcriber.Transcribe(tctx, localPath) + tcancel() + if err != nil { + logger.WarnCF("onebot", "Voice transcription failed", map[string]interface{}{ + "error": err.Error(), + }) + textParts = append(textParts, "[voice (transcription failed)]") + media = append(media, localPath) + } else { + textParts = append(textParts, fmt.Sprintf("[voice transcription: %s]", result.Text)) + } + } else { + textParts = append(textParts, "[voice]") + media = append(media, localPath) + } + } + } + } + + 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 { + replyTo = fmt.Sprintf("%v", id) + } + } + + case "face": + if data != nil { + faceID, _ := data["id"] + textParts = append(textParts, fmt.Sprintf("[face:%v]", faceID)) + } + + case "forward": + textParts = append(textParts, "[forward message]") + + default: + } - return parseMessageResult{Text: strings.TrimSpace(text), IsBotMentioned: mentioned} } - return parseMessageResult{} + + return parseMessageResult{ + Text: strings.TrimSpace(strings.Join(textParts, "")), + IsBotMentioned: mentioned, + Media: media, + LocalFiles: localFiles, + ReplyTo: replyTo, + } } func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) { @@ -462,21 +784,30 @@ func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) { return } c.handleMessage(evt) + + case "message_sent": + logger.DebugCF("onebot", "Bot sent message event", map[string]interface{}{ + "message_type": raw.MessageType, + "message_id": parseJSONString(raw.MessageID), + }) + case "meta_event": c.handleMetaEvent(raw) + case "notice": - logger.DebugCF("onebot", "Notice event received", map[string]interface{}{ - "sub_type": raw.SubType, - }) + c.handleNoticeEvent(raw) + case "request": logger.DebugCF("onebot", "Request event received", map[string]interface{}{ "sub_type": raw.SubType, }) + case "": logger.DebugCF("onebot", "Event with empty post_type (possibly API response)", map[string]interface{}{ "echo": raw.Echo, "status": raw.Status, }) + default: logger.DebugCF("onebot", "Unknown post_type", map[string]interface{}{ "post_type": raw.PostType, @@ -495,7 +826,11 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent ts, _ := parseJSONInt64(raw.Time) messageID := parseJSONString(raw.MessageID) - parsed := parseMessageContentEx(raw.Message, selfID) + if selfID == 0 { + selfID = atomic.LoadInt64(&c.selfID) + } + + parsed := c.parseMessageSegments(raw.Message, selfID) isBotMentioned := parsed.IsBotMentioned content := raw.RawMessage @@ -510,6 +845,10 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent } } + if parsed.Text != "" && content != parsed.Text && (len(parsed.Media) > 0 || parsed.ReplyTo != "") { + content = parsed.Text + } + var sender oneBotSender if len(raw.Sender) > 0 { if err := json.Unmarshal(raw.Sender, &sender); err != nil { @@ -521,12 +860,15 @@ 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, + "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{ @@ -543,6 +885,9 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent SelfID: selfID, Time: ts, MetaEventType: raw.MetaEventType, + Media: parsed.Media, + LocalFiles: parsed.LocalFiles, + ReplyTo: parsed.ReplyTo, }, nil } @@ -561,7 +906,80 @@ func (c *OneBotChannel) handleMetaEvent(raw *oneBotRawEvent) { } } +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() { + for _, f := range evt.LocalFiles { + if err := os.Remove(f); err != nil { + logger.DebugCF("onebot", "Failed to remove temp file", map[string]interface{}{ + "path": f, + "error": err.Error(), + }) + } + } + }() + } + if c.isDuplicate(evt.MessageID) { logger.DebugCF("onebot", "Duplicate message, skipping", map[string]interface{}{ "message_id": evt.MessageID, @@ -584,14 +1002,19 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) { "message_id": evt.MessageID, } + if evt.ReplyTo != "" { + metadata["reply_to_message_id"] = evt.ReplyTo + } + switch evt.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), + "sender": senderID, + "message_id": evt.MessageID, + "length": len(content), + "content": truncate(content, 100), + "media_count": len(evt.Media), }) case "group": @@ -629,6 +1052,7 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) { "is_mentioned": evt.IsBotMentioned, "length": len(content), "content": truncate(content, 100), + "media_count": len(evt.Media), }) default: @@ -644,13 +1068,16 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) { metadata["nickname"] = evt.Sender.Nickname } + c.lastMessageID.Store(chatID, evt.MessageID) + logger.DebugCF("onebot", "Forwarding message to bus", map[string]interface{}{ - "sender_id": senderID, - "chat_id": chatID, - "content": truncate(content, 100), + "sender_id": senderID, + "chat_id": chatID, + "content": truncate(content, 100), + "media_count": len(evt.Media), }) - c.HandleMessage(senderID, chatID, content, []string{}, metadata) + c.HandleMessage(senderID, chatID, content, evt.Media, metadata) } func (c *OneBotChannel) isDuplicate(messageID string) bool {