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
This commit is contained in:
Hoshina 2026-02-15 14:09:06 +08:00
parent 50628b386a
commit 50e4edc810
2 changed files with 509 additions and 76 deletions

View file

@ -617,6 +617,12 @@ func gatewayCmd() {
logger.InfoC("voice", "Groq transcription attached to Slack channel") 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() enabledChannels := channelManager.GetEnabledChannels()

View file

@ -4,9 +4,11 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
"os"
"strconv" "strconv"
"strings" "strings"
"sync" "sync"
"sync/atomic"
"time" "time"
"github.com/gorilla/websocket" "github.com/gorilla/websocket"
@ -14,20 +16,27 @@ import (
"github.com/sipeed/picoclaw/pkg/bus" "github.com/sipeed/picoclaw/pkg/bus"
"github.com/sipeed/picoclaw/pkg/config" "github.com/sipeed/picoclaw/pkg/config"
"github.com/sipeed/picoclaw/pkg/logger" "github.com/sipeed/picoclaw/pkg/logger"
"github.com/sipeed/picoclaw/pkg/utils"
"github.com/sipeed/picoclaw/pkg/voice"
) )
type OneBotChannel struct { type OneBotChannel struct {
*BaseChannel *BaseChannel
config config.OneBotConfig config config.OneBotConfig
conn *websocket.Conn conn *websocket.Conn
ctx context.Context ctx context.Context
cancel context.CancelFunc cancel context.CancelFunc
dedup map[string]struct{} dedup map[string]struct{}
dedupRing []string dedupRing []string
dedupIdx int dedupIdx int
mu sync.Mutex mu sync.Mutex
writeMu sync.Mutex writeMu sync.Mutex
echoCounter int64 echoCounter int64
selfID int64
pending map[string]chan json.RawMessage
pendingMu sync.Mutex
transcriber *voice.GroqTranscriber
lastMessageID sync.Map
} }
type oneBotRawEvent struct { type oneBotRawEvent struct {
@ -43,9 +52,11 @@ type oneBotRawEvent struct {
SelfID json.RawMessage `json:"self_id"` SelfID json.RawMessage `json:"self_id"`
Time json.RawMessage `json:"time"` Time json.RawMessage `json:"time"`
MetaEventType string `json:"meta_event_type"` MetaEventType string `json:"meta_event_type"`
NoticeType string `json:"notice_type"`
Echo string `json:"echo"` Echo string `json:"echo"`
RetCode json.RawMessage `json:"retcode"` RetCode json.RawMessage `json:"retcode"`
Status json.RawMessage `json:"status"` Status json.RawMessage `json:"status"`
Data json.RawMessage `json:"data"`
} }
type BotStatus struct { type BotStatus struct {
@ -88,6 +99,9 @@ type oneBotEvent struct {
SelfID int64 SelfID int64
Time int64 Time int64
MetaEventType string MetaEventType string
Media []string
LocalFiles []string
ReplyTo string
} }
type oneBotAPIRequest struct { type oneBotAPIRequest struct {
@ -96,14 +110,19 @@ type oneBotAPIRequest struct {
Echo string `json:"echo,omitempty"` Echo string `json:"echo,omitempty"`
} }
type oneBotMessageSegment struct {
Type string `json:"type"`
Data map[string]interface{} `json:"data"`
}
type oneBotSendPrivateMsgParams struct { type oneBotSendPrivateMsgParams struct {
UserID int64 `json:"user_id"` UserID int64 `json:"user_id"`
Message string `json:"message"` Message interface{} `json:"message"`
} }
type oneBotSendGroupMsgParams struct { type oneBotSendGroupMsgParams struct {
GroupID int64 `json:"group_id"` GroupID int64 `json:"group_id"`
Message string `json:"message"` Message interface{} `json:"message"`
} }
func NewOneBotChannel(cfg config.OneBotConfig, messageBus *bus.MessageBus) (*OneBotChannel, error) { 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), dedup: make(map[string]struct{}, dedupSize),
dedupRing: make([]string, dedupSize), dedupRing: make([]string, dedupSize),
dedupIdx: 0, dedupIdx: 0,
pending: make(map[string]chan json.RawMessage),
}, nil }, nil
} }
func (c *OneBotChannel) SetTranscriber(transcriber *voice.GroqTranscriber) {
c.transcriber = transcriber
}
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")
@ -135,13 +159,13 @@ func (c *OneBotChannel) Start(ctx context.Context) error {
"error": err.Error(), "error": err.Error(),
}) })
} else { } else {
c.fetchSelfID()
go c.listen() go c.listen()
} }
if c.config.ReconnectInterval > 0 { if c.config.ReconnectInterval > 0 {
go c.reconnectLoop() go c.reconnectLoop()
} else { } else {
// If reconnect is disabled but initial connection failed, we cannot recover
if c.conn == nil { if c.conn == nil {
return fmt.Errorf("failed to connect to OneBot and reconnect is disabled") return fmt.Errorf("failed to connect to OneBot and reconnect is disabled")
} }
@ -167,14 +191,141 @@ func (c *OneBotChannel) connect() error {
return err 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.mu.Lock()
c.conn = conn c.conn = conn
c.mu.Unlock() c.mu.Unlock()
go c.pinger(conn)
logger.InfoC("onebot", "WebSocket connected") logger.InfoC("onebot", "WebSocket connected")
return nil 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() { func (c *OneBotChannel) reconnectLoop() {
interval := time.Duration(c.config.ReconnectInterval) * time.Second interval := time.Duration(c.config.ReconnectInterval) * time.Second
if interval < 5*time.Second { if interval < 5*time.Second {
@ -197,6 +348,7 @@ func (c *OneBotChannel) reconnectLoop() {
"error": err.Error(), "error": err.Error(),
}) })
} else { } else {
c.fetchSelfID()
go c.listen() go c.listen()
} }
} }
@ -212,6 +364,13 @@ func (c *OneBotChannel) Stop(ctx context.Context) error {
c.cancel() c.cancel()
} }
c.pendingMu.Lock()
for echo, ch := range c.pending {
close(ch)
delete(c.pending, echo)
}
c.pendingMu.Unlock()
c.mu.Lock() c.mu.Lock()
if c.conn != nil { if c.conn != nil {
c.conn.Close() c.conn.Close()
@ -240,10 +399,7 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error
return err return err
} }
c.writeMu.Lock() echo := fmt.Sprintf("send_%d", atomic.AddInt64(&c.echoCounter, 1))
c.echoCounter++
echo := fmt.Sprintf("send_%d", c.echoCounter)
c.writeMu.Unlock()
req := oneBotAPIRequest{ req := oneBotAPIRequest{
Action: action, Action: action,
@ -270,28 +426,50 @@ func (c *OneBotChannel) Send(ctx context.Context, msg bus.OutboundMessage) error
return nil 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) { func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, interface{}, error) {
chatID := msg.ChatID chatID := msg.ChatID
if len(chatID) > 6 && chatID[:6] == "group:" { if strings.HasPrefix(chatID, "group:") {
groupID, err := strconv.ParseInt(chatID[6:], 10, 64) groupID, err := strconv.ParseInt(chatID[6:], 10, 64)
if err != nil { if err != nil {
return "", nil, fmt.Errorf("invalid group ID in chatID: %s", chatID) return "", nil, fmt.Errorf("invalid group ID in chatID: %s", chatID)
} }
segments := c.buildMessageSegments(chatID, msg.Content)
return "send_group_msg", oneBotSendGroupMsgParams{ return "send_group_msg", oneBotSendGroupMsgParams{
GroupID: groupID, GroupID: groupID,
Message: msg.Content, Message: segments,
}, nil }, nil
} }
if len(chatID) > 8 && chatID[:8] == "private:" { if strings.HasPrefix(chatID, "private:") {
userID, err := strconv.ParseInt(chatID[8:], 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 user ID in chatID: %s", chatID)
} }
segments := c.buildMessageSegments(chatID, msg.Content)
return "send_private_msg", oneBotSendPrivateMsgParams{ return "send_private_msg", oneBotSendPrivateMsgParams{
UserID: userID, UserID: userID,
Message: msg.Content, Message: segments,
}, nil }, nil
} }
@ -300,34 +478,35 @@ func (c *OneBotChannel) buildSendRequest(msg bus.OutboundMessage) (string, inter
return "", nil, fmt.Errorf("invalid chatID for OneBot: %s", chatID) return "", nil, fmt.Errorf("invalid chatID for OneBot: %s", chatID)
} }
segments := c.buildMessageSegments(chatID, msg.Content)
return "send_private_msg", oneBotSendPrivateMsgParams{ return "send_private_msg", oneBotSendPrivateMsgParams{
UserID: userID, UserID: userID,
Message: msg.Content, Message: segments,
}, nil }, nil
} }
func (c *OneBotChannel) listen() { 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 { for {
select { select {
case <-c.ctx.Done(): case <-c.ctx.Done():
return return
default: 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() _, message, err := conn.ReadMessage()
if err != nil { if err != nil {
logger.ErrorCF("onebot", "WebSocket read error", map[string]interface{}{ logger.ErrorCF("onebot", "WebSocket read error", map[string]interface{}{
"error": err.Error(), "error": err.Error(),
}) })
c.mu.Lock() c.mu.Lock()
if c.conn != nil { if c.conn == conn {
c.conn.Close() c.conn.Close()
c.conn = nil c.conn = nil
} }
@ -335,6 +514,8 @@ func (c *OneBotChannel) listen() {
return return
} }
_ = conn.SetReadDeadline(time.Now().Add(60 * time.Second))
logger.DebugCF("onebot", "Raw WebSocket message received", map[string]interface{}{ logger.DebugCF("onebot", "Raw WebSocket message received", map[string]interface{}{
"length": len(message), "length": len(message),
"payload": string(message), "payload": string(message),
@ -349,9 +530,27 @@ func (c *OneBotChannel) listen() {
continue continue
} }
if raw.Echo != "" || isAPIResponse(raw.Status) { if raw.Echo != "" {
logger.DebugCF("onebot", "Received API response, skipping", map[string]interface{}{ c.pendingMu.Lock()
"echo": raw.Echo, 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), "status": string(raw.Status),
}) })
continue continue
@ -401,9 +600,12 @@ func parseJSONString(raw json.RawMessage) string {
type parseMessageResult struct { type parseMessageResult struct {
Text string Text string
IsBotMentioned bool 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 { if len(raw) == 0 {
return parseMessageResult{} return parseMessageResult{}
} }
@ -423,32 +625,152 @@ func parseMessageContentEx(raw json.RawMessage, selfID int64) parseMessageResult
} }
var segments []map[string]interface{} var segments []map[string]interface{}
if err := json.Unmarshal(raw, &segments); err == nil { if err := json.Unmarshal(raw, &segments); err != nil {
var text string return parseMessageResult{}
mentioned := false }
selfIDStr := strconv.FormatInt(selfID, 10)
for _, seg := range segments { var textParts []string
segType, _ := seg["type"].(string) mentioned := false
data, _ := seg["data"].(map[string]interface{}) selfIDStr := strconv.FormatInt(selfID, 10)
switch segType { var media []string
case "text": var localFiles []string
if data != nil { var replyTo string
if t, ok := data["text"].(string); ok {
text += t 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"]) case "at":
if qqVal == selfIDStr || qqVal == "all" { if data != nil && selfID > 0 {
mentioned = true 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) { func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) {
@ -462,21 +784,30 @@ func (c *OneBotChannel) handleRawEvent(raw *oneBotRawEvent) {
return return
} }
c.handleMessage(evt) 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": case "meta_event":
c.handleMetaEvent(raw) c.handleMetaEvent(raw)
case "notice": case "notice":
logger.DebugCF("onebot", "Notice event received", map[string]interface{}{ c.handleNoticeEvent(raw)
"sub_type": raw.SubType,
})
case "request": case "request":
logger.DebugCF("onebot", "Request event received", map[string]interface{}{ logger.DebugCF("onebot", "Request event received", map[string]interface{}{
"sub_type": raw.SubType, "sub_type": raw.SubType,
}) })
case "": case "":
logger.DebugCF("onebot", "Event with empty post_type (possibly API response)", map[string]interface{}{ logger.DebugCF("onebot", "Event with empty post_type (possibly API response)", map[string]interface{}{
"echo": raw.Echo, "echo": raw.Echo,
"status": raw.Status, "status": raw.Status,
}) })
default: default:
logger.DebugCF("onebot", "Unknown post_type", map[string]interface{}{ logger.DebugCF("onebot", "Unknown post_type", map[string]interface{}{
"post_type": raw.PostType, "post_type": raw.PostType,
@ -495,7 +826,11 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent
ts, _ := parseJSONInt64(raw.Time) ts, _ := parseJSONInt64(raw.Time)
messageID := parseJSONString(raw.MessageID) 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 isBotMentioned := parsed.IsBotMentioned
content := raw.RawMessage 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 var sender oneBotSender
if len(raw.Sender) > 0 { if len(raw.Sender) > 0 {
if err := json.Unmarshal(raw.Sender, &sender); err != nil { 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{}{ logger.DebugCF("onebot", "Normalized message event", map[string]interface{}{
"message_type": raw.MessageType, "message_type": raw.MessageType,
"user_id": userID, "user_id": userID,
"group_id": groupID, "group_id": groupID,
"message_id": messageID, "message_id": messageID,
"content_len": len(content), "content_len": len(content),
"nickname": sender.Nickname, "nickname": sender.Nickname,
"media_count": len(parsed.Media),
"has_reply": parsed.ReplyTo != "",
"has_localfiles": len(parsed.LocalFiles) > 0,
}) })
return &oneBotEvent{ return &oneBotEvent{
@ -543,6 +885,9 @@ func (c *OneBotChannel) normalizeMessageEvent(raw *oneBotRawEvent) (*oneBotEvent
SelfID: selfID, SelfID: selfID,
Time: ts, Time: ts,
MetaEventType: raw.MetaEventType, MetaEventType: raw.MetaEventType,
Media: parsed.Media,
LocalFiles: parsed.LocalFiles,
ReplyTo: parsed.ReplyTo,
}, nil }, 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) { 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) { if c.isDuplicate(evt.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": evt.MessageID,
@ -584,14 +1002,19 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) {
"message_id": evt.MessageID, "message_id": evt.MessageID,
} }
if evt.ReplyTo != "" {
metadata["reply_to_message_id"] = evt.ReplyTo
}
switch evt.MessageType { switch evt.MessageType {
case "private": case "private":
chatID = "private:" + senderID chatID = "private:" + senderID
logger.InfoCF("onebot", "Received private message", map[string]interface{}{ logger.InfoCF("onebot", "Received private message", map[string]interface{}{
"sender": senderID, "sender": senderID,
"message_id": evt.MessageID, "message_id": evt.MessageID,
"length": len(content), "length": len(content),
"content": truncate(content, 100), "content": truncate(content, 100),
"media_count": len(evt.Media),
}) })
case "group": case "group":
@ -629,6 +1052,7 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) {
"is_mentioned": evt.IsBotMentioned, "is_mentioned": evt.IsBotMentioned,
"length": len(content), "length": len(content),
"content": truncate(content, 100), "content": truncate(content, 100),
"media_count": len(evt.Media),
}) })
default: default:
@ -644,13 +1068,16 @@ func (c *OneBotChannel) handleMessage(evt *oneBotEvent) {
metadata["nickname"] = evt.Sender.Nickname metadata["nickname"] = evt.Sender.Nickname
} }
c.lastMessageID.Store(chatID, evt.MessageID)
logger.DebugCF("onebot", "Forwarding message to bus", map[string]interface{}{ logger.DebugCF("onebot", "Forwarding message to bus", map[string]interface{}{
"sender_id": senderID, "sender_id": senderID,
"chat_id": chatID, "chat_id": chatID,
"content": truncate(content, 100), "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 { func (c *OneBotChannel) isDuplicate(messageID string) bool {