feat(qq): 支持文件处理。
This commit is contained in:
parent
0adb4fc7e9
commit
2f91c7c72f
1 changed files with 146 additions and 87 deletions
|
|
@ -6,15 +6,13 @@ import (
|
|||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"github.com/sipeed/picoclaw/pkg/identity"
|
||||
"github.com/sipeed/picoclaw/pkg/utils"
|
||||
"github.com/tencent-connect/botgo/constant"
|
||||
"github.com/tencent-connect/botgo/openapi/options"
|
||||
"math"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"path"
|
||||
"fmt"
|
||||
"math"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"strings"
|
||||
|
|
@ -23,18 +21,24 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/sipeed/picoclaw/pkg/media"
|
||||
"github.com/sipeed/picoclaw/pkg/utils"
|
||||
"github.com/tidwall/gjson"
|
||||
|
||||
"github.com/tencent-connect/botgo"
|
||||
"github.com/tencent-connect/botgo/constant"
|
||||
"github.com/tencent-connect/botgo/dto"
|
||||
"github.com/tencent-connect/botgo/event"
|
||||
"github.com/tencent-connect/botgo/openapi/options"
|
||||
"github.com/tencent-connect/botgo/token"
|
||||
"golang.org/x/oauth2"
|
||||
|
||||
"github.com/sipeed/picoclaw/pkg/bus"
|
||||
"github.com/sipeed/picoclaw/pkg/channels"
|
||||
"github.com/sipeed/picoclaw/pkg/config"
|
||||
"github.com/sipeed/picoclaw/pkg/identity"
|
||||
"github.com/sipeed/picoclaw/pkg/logger"
|
||||
"github.com/sipeed/picoclaw/pkg/media"
|
||||
"github.com/sipeed/picoclaw/pkg/utils"
|
||||
)
|
||||
|
||||
const (
|
||||
|
|
@ -48,6 +52,10 @@ const (
|
|||
|
||||
type kindType string
|
||||
|
||||
func (k kindType) String() string {
|
||||
return string(k)
|
||||
}
|
||||
|
||||
const (
|
||||
kindDirect kindType = "direct"
|
||||
kindGroup kindType = "group"
|
||||
|
|
@ -56,22 +64,10 @@ const (
|
|||
var emojiRegexp = regexp.MustCompile(`<[^<]*?ext="([^"]+)"[^<]*?faceType=(\d+)[^<]*?>|<[^<]*?faceType=(\d+)[^<]*?ext="([^"]+)"[^<]*?>`)
|
||||
var extRegexp = regexp.MustCompile(`ext="([^"]+)"`)
|
||||
|
||||
|
||||
type qqAPI interface {
|
||||
WS(ctx context.Context, params map[string]string, body string) (*dto.WebsocketAP, error)
|
||||
PostGroupMessage(
|
||||
ctx context.Context, groupID string, msg dto.APIMessage, opt ...options.Option,
|
||||
) (*dto.Message, error)
|
||||
PostC2CMessage(
|
||||
ctx context.Context, userID string, msg dto.APIMessage, opt ...options.Option,
|
||||
) (*dto.Message, error)
|
||||
Transport(ctx context.Context, method, url string, body any) ([]byte, error)
|
||||
}
|
||||
|
||||
type QQChannel struct {
|
||||
*channels.BaseChannel
|
||||
config config.QQConfig
|
||||
api qqAPI
|
||||
api qqAPI
|
||||
tokenSource oauth2.TokenSource
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
|
|
@ -114,7 +110,7 @@ func NewQQChannel(cfg config.QQConfig, messageBus *bus.MessageBus) (*QQChannel,
|
|||
}
|
||||
|
||||
func (c *QQChannel) Start(ctx context.Context) error {
|
||||
if c.config.AppID == "" || c.config.AppSecret.String() == "" {
|
||||
if c.config.AppID == "" || c.config.AppSecret() == "" {
|
||||
return fmt.Errorf("QQ app_id and app_secret not configured")
|
||||
}
|
||||
|
||||
|
|
@ -128,7 +124,7 @@ func (c *QQChannel) Start(ctx context.Context) error {
|
|||
// create token source
|
||||
credentials := &token.QQBotCredentials{
|
||||
AppID: c.config.AppID,
|
||||
AppSecret: c.config.AppSecret.String(),
|
||||
AppSecret: c.config.AppSecret(),
|
||||
}
|
||||
c.tokenSource = token.NewQQBotTokenSource(credentials)
|
||||
|
||||
|
|
@ -178,7 +174,7 @@ func (c *QQChannel) Start(ctx context.Context) error {
|
|||
// Pre-register reasoning_channel_id as group chat if configured,
|
||||
// so outbound-only destinations are routed correctly.
|
||||
if c.config.ReasoningChannelID != "" {
|
||||
c.chatType.Store(c.config.ReasoningChannelID, kindGroup)
|
||||
c.chatType.Store(c.config.ReasoningChannelID, kindGroup.String())
|
||||
}
|
||||
|
||||
c.SetRunning(true)
|
||||
|
|
@ -233,7 +229,21 @@ func (c *QQChannel) Send(ctx context.Context, msg bus.OutboundMessage) error {
|
|||
return channels.ErrNotRunning
|
||||
}
|
||||
|
||||
chatKind := c.getChatKind(msg.ChatID)
|
||||
|
||||
|
||||
c.applyPassiveReplyMetadata(msg.ChatID, msgToCreate)
|
||||
|
||||
// Sanitize URLs in group messages to avoid QQ's URL blacklist rejection.
|
||||
if chatKind == "group" {
|
||||
if msgToCreate.Content != "" {
|
||||
msgToCreate.Content = sanitizeURLs(msgToCreate.Content)
|
||||
}
|
||||
if msgToCreate.Markdown != nil && msgToCreate.Markdown.Content != "" {
|
||||
msgToCreate.Markdown.Content = sanitizeURLs(msgToCreate.Markdown.Content)
|
||||
}
|
||||
}
|
||||
|
||||
// Route to group or C2C.
|
||||
mdMsg, textMsg := c.genReplyMsg(ctx, msg)
|
||||
var err error
|
||||
for _, replyMsg := range []dto.MessageToCreate{mdMsg, textMsg} {
|
||||
|
|
@ -320,14 +330,13 @@ func (c *QQChannel) StartTyping(ctx context.Context, chatID string) (func(), err
|
|||
|
||||
// SendMedia implements the channels.MediaSender interface.
|
||||
// QQ group/C2C media sending is a two-step flow:
|
||||
// 1. Upload media to /files using a remote URL or local bytes.
|
||||
// 1. Upload media to /files using a remote URL or base64-encoded local bytes.
|
||||
// 2. Send a msg_type=7 message using the returned file_info.
|
||||
func (c *QQChannel) SendMedia(ctx context.Context, msg bus.OutboundMediaMessage) error {
|
||||
|
||||
if !c.IsRunning() {
|
||||
return channels.ErrNotRunning
|
||||
}
|
||||
chatKind := c.getChatKind(msg.ChatID)
|
||||
var err error
|
||||
for _, part := range msg.Parts {
|
||||
fileInfo, err := c.uploadMedia(ctx, chatKind, msg.ChatID, part)
|
||||
|
|
@ -343,11 +352,15 @@ func (c *QQChannel) SendMedia(ctx context.Context, msg bus.OutboundMediaMessage)
|
|||
return fmt.Errorf("qq send media: %w", channels.ErrTemporary)
|
||||
}
|
||||
|
||||
if err = c.sendUploadedMedia(ctx, chatKind, msg.ChatID, part, fileInfo); err != nil {
|
||||
if err := c.sendUploadedMedia(ctx, chatKind, msg.ChatID, part, fileInfo); err != nil {
|
||||
logger.ErrorCF("qq", "Failed to send media", map[string]any{
|
||||
"type": part.Type,
|
||||
"chat_id": msg.ChatID,
|
||||
"error": err.Error(),
|
||||
if err = c.sendOneMedia(ctx, msg.ChatID, part); err != nil {
|
||||
logger.ErrorCF("qq", "Failed to send media", map[string]any{
|
||||
"part": part,
|
||||
"error": err.Error(),
|
||||
})
|
||||
continue
|
||||
}
|
||||
|
|
@ -355,9 +368,19 @@ func (c *QQChannel) SendMedia(ctx context.Context, msg bus.OutboundMediaMessage)
|
|||
return err
|
||||
}
|
||||
|
||||
func (c *QQChannel) uploadMedia( ctx context.Context,
|
||||
chatKind kindType, chatID string, part bus.MediaPart) ([]byte, error) {
|
||||
type qqMediaUpload struct {
|
||||
FileType uint64 `json:"file_type"`
|
||||
URL string `json:"url,omitempty"`
|
||||
FileData string `json:"file_data,omitempty"`
|
||||
FileName string `json:"file_name,omitempty"`
|
||||
SrvSendMsg bool `json:"srv_send_msg,omitempty"`
|
||||
}
|
||||
|
||||
func (c *QQChannel) uploadMedia(
|
||||
ctx context.Context,
|
||||
chatKind, chatID string,
|
||||
part bus.MediaPart,
|
||||
) ([]byte, error) {
|
||||
payload, err := c.buildMediaUpload(part)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -379,8 +402,8 @@ func (c *QQChannel) uploadMedia( ctx context.Context,
|
|||
return uploaded.FileInfo, nil
|
||||
}
|
||||
|
||||
func (c *QQChannel) buildMediaUpload(part bus.MediaPart) (*RichMediaMessage, error) {
|
||||
payload := &RichMediaMessage{}
|
||||
func (c *QQChannel) buildMediaUpload(part bus.MediaPart) (*qqMediaUpload, error) {
|
||||
payload := &qqMediaUpload{}
|
||||
|
||||
mediaRef := part.Ref
|
||||
if isHTTPURL(mediaRef) {
|
||||
|
|
@ -405,12 +428,15 @@ func (c *QQChannel) buildMediaUpload(part bus.MediaPart) (*RichMediaMessage, err
|
|||
if part.ContentType == "" {
|
||||
part.ContentType = meta.ContentType
|
||||
}
|
||||
payload.FileType = qqFileType(c.outboundMediaType(part, resolved))
|
||||
payload.FileName = qqUploadFilename(part, resolved, payload.FileType)
|
||||
|
||||
if isHTTPURL(resolved) {
|
||||
payload.FileType = qqFileType(c.outboundMediaType(part, ""))
|
||||
payload.URL = resolved
|
||||
payload.FileName = qqUploadFilename(part, resolved, payload.FileType)
|
||||
return payload, nil
|
||||
}
|
||||
payload.FileType = qqFileType(c.outboundMediaType(part, resolved))
|
||||
payload.FileName = qqUploadFilename(part, resolved, payload.FileType)
|
||||
|
||||
if limitBytes := c.maxBase64FileSizeBytes(); limitBytes > 0 {
|
||||
info, statErr := os.Stat(resolved)
|
||||
|
|
@ -432,7 +458,8 @@ func (c *QQChannel) buildMediaUpload(part bus.MediaPart) (*RichMediaMessage, err
|
|||
if err != nil {
|
||||
return nil, fmt.Errorf("qq read local media %q: %v: %w", resolved, err, channels.ErrSendFailed)
|
||||
}
|
||||
payload.FileData = data
|
||||
|
||||
payload.FileData = base64.StdEncoding.EncodeToString(data)
|
||||
return payload, nil
|
||||
}
|
||||
|
||||
|
|
@ -500,12 +527,12 @@ func (c *QQChannel) outboundMediaType(part bus.MediaPart, localPath string) stri
|
|||
return "audio"
|
||||
}
|
||||
|
||||
|
||||
|
||||
// Fix the sendUploadedMedia function syntax error
|
||||
func (c *QQChannel) sendUploadedMedia(ctx context.Context, chatKind kindType, chatID string, part bus.MediaPart,
|
||||
fileInfo []byte) error {
|
||||
|
||||
func (c *QQChannel) sendUploadedMedia(
|
||||
ctx context.Context,
|
||||
chatKind, chatID string,
|
||||
part bus.MediaPart,
|
||||
fileInfo []byte,
|
||||
) error {
|
||||
msg := &dto.MessageToCreate{
|
||||
Content: part.Caption,
|
||||
MsgType: dto.RichMediaMsg,
|
||||
|
|
@ -513,13 +540,13 @@ func (c *QQChannel) sendUploadedMedia(ctx context.Context, chatKind kindType, ch
|
|||
FileInfo: fileInfo,
|
||||
},
|
||||
}
|
||||
c.applyPassiveReplyMetadata(chatID, msg)
|
||||
|
||||
msg.MsgID, msg.MsgSeq = c.getReplyExtInfo(ctx, chatID)
|
||||
if chatKind == "group" && msg.Content != "" {
|
||||
msg.Content = sanitizeURLs(msg.Content)
|
||||
}
|
||||
|
||||
if chatKind == kindGroup {
|
||||
if msg.Content != "" {
|
||||
msg.Content = sanitizeURLs(msg.Content)
|
||||
}
|
||||
if chatKind == "group" {
|
||||
_, err := c.api.PostGroupMessage(ctx, chatID, msg)
|
||||
return err
|
||||
}
|
||||
|
|
@ -527,9 +554,25 @@ func (c *QQChannel) sendUploadedMedia(ctx context.Context, chatKind kindType, ch
|
|||
return err
|
||||
}
|
||||
|
||||
func (c *QQChannel) mediaUploadURL(chatKind kindType, chatID string) string {
|
||||
func (c *QQChannel) applyPassiveReplyMetadata(chatID string, msg *dto.MessageToCreate) {
|
||||
if v, ok := c.lastMsgID.Load(chatID); ok {
|
||||
if msgID, ok := v.(string); ok && msgID != "" {
|
||||
msg.MsgID = msgID
|
||||
|
||||
// Increment msg_seq atomically for multi-part replies.
|
||||
if counterVal, ok := c.msgSeqCounters.Load(chatID); ok {
|
||||
if counter, ok := counterVal.(*atomic.Uint64); ok {
|
||||
seq := counter.Add(1)
|
||||
msg.MsgSeq = uint32(seq)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (c *QQChannel) mediaUploadURL(chatKind, chatID string) string {
|
||||
base := constant.APIDomain
|
||||
if chatKind == kindGroup {
|
||||
if chatKind == "group" {
|
||||
return fmt.Sprintf("%s/v2/groups/%s/files", base, chatID)
|
||||
}
|
||||
return fmt.Sprintf("%s/v2/users/%s/files", base, chatID)
|
||||
|
|
@ -588,8 +631,8 @@ func (c *QQChannel) handleC2CMessage() event.C2CMessageEventHandler {
|
|||
content = appendContent(content, note)
|
||||
}
|
||||
if content == "" && len(mediaPaths) == 0 {
|
||||
logger.DebugC("qq", "Received empty C2C message with no content/attachments, ignoring")
|
||||
return nil
|
||||
logger.DebugC("qq", "Received empty C2C message with no attachments, ignoring")
|
||||
Username: data.Author.Username,
|
||||
}
|
||||
|
||||
if !c.IsAllowedSender(sender) {
|
||||
|
|
@ -600,38 +643,36 @@ func (c *QQChannel) handleC2CMessage() event.C2CMessageEventHandler {
|
|||
}
|
||||
|
||||
scope := channels.BuildMediaScope("qq", senderID, data.ID)
|
||||
content, mediaPaths = c.decodeMessage(context.Background(), event, (*dto.Message)(data), scope)
|
||||
content, mediaPaths := c.decodeMessage(context.Background(), event, (*dto.Message)(data), scope)
|
||||
if content == "" {
|
||||
logger.DebugC("qq", "Received empty C2C message, ignoring")
|
||||
return nil
|
||||
}
|
||||
|
||||
logger.InfoCF("qq", "Received C2C message", map[string]any{
|
||||
"sender": senderID,
|
||||
"length": len(content),
|
||||
"media_count": len(mediaPaths),
|
||||
"sender": senderID,
|
||||
"length": len(content),
|
||||
logger.InfoCF("qq", "Received C2C message", map[string]any{
|
||||
"sender": senderID,
|
||||
"length": len(content),
|
||||
"media_count": len(mediaPaths),
|
||||
})
|
||||
|
||||
// Store chat routing context.
|
||||
c.saveChatKind(senderID, kindDirect)
|
||||
c.lastMsgID.Store(senderID, data.ID)
|
||||
|
||||
metadata := map[string]string{
|
||||
"account_id": senderID,
|
||||
}
|
||||
metadata := map[string]string{
|
||||
"account_id": senderID,
|
||||
}
|
||||
|
||||
c.HandleMessage(c.ctx,
|
||||
bus.Peer{Kind: string(kindDirect), ID: senderID},
|
||||
data.ID,
|
||||
senderID,
|
||||
senderID,
|
||||
content,
|
||||
mediaPaths,
|
||||
metadata,
|
||||
sender,
|
||||
)
|
||||
c.HandleMessage(c.ctx,
|
||||
bus.Peer{Kind: kindDirect.String(), ID: senderID},
|
||||
data.ID,
|
||||
senderID,
|
||||
senderID,
|
||||
content,
|
||||
mediaPaths,
|
||||
metadata,
|
||||
sender,
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
@ -658,7 +699,6 @@ func (c *QQChannel) handleGroupATMessage() event.GroupATMessageEventHandler {
|
|||
Platform: "qq",
|
||||
PlatformID: data.Author.ID,
|
||||
CanonicalID: identity.BuildCanonicalID("qq", data.Author.ID),
|
||||
Username: data.Author.Username,
|
||||
}
|
||||
|
||||
if !c.IsAllowedSender(sender) {
|
||||
|
|
@ -671,6 +711,9 @@ func (c *QQChannel) handleGroupATMessage() event.GroupATMessageEventHandler {
|
|||
content = appendContent(content, note)
|
||||
}
|
||||
|
||||
// GroupAT event means bot is always mentioned; apply group trigger filtering.
|
||||
Username: data.Author.Username,
|
||||
}
|
||||
|
||||
if !c.IsAllowedSender(sender) {
|
||||
logger.Infof("qq", "Received group message from unauthorized sender", map[string]any{
|
||||
|
|
@ -681,7 +724,7 @@ func (c *QQChannel) handleGroupATMessage() event.GroupATMessageEventHandler {
|
|||
|
||||
scope := channels.BuildMediaScope("qq", data.GroupID, data.ID)
|
||||
|
||||
content, mediaPaths = c.decodeMessage(context.Background(), event, (*dto.Message)(data), scope)
|
||||
content, mediaPaths := c.decodeMessage(context.Background(), event, (*dto.Message)(data), scope)
|
||||
if content == "" {
|
||||
logger.DebugC("qq", "Received empty group message, ignoring")
|
||||
return nil
|
||||
|
|
@ -713,24 +756,25 @@ func (c *QQChannel) handleGroupATMessage() event.GroupATMessageEventHandler {
|
|||
"group_id": data.GroupID,
|
||||
}
|
||||
|
||||
c.HandleMessage(c.ctx,
|
||||
bus.Peer{Kind: string(kindGroup), ID: data.GroupID},
|
||||
data.ID,
|
||||
senderID,
|
||||
data.GroupID,
|
||||
content,
|
||||
mediaPaths,
|
||||
metadata,
|
||||
sender,
|
||||
)
|
||||
c.HandleMessage(c.ctx,
|
||||
bus.Peer{Kind: kindGroup.String(), ID: data.GroupID},
|
||||
data.ID,
|
||||
senderID,
|
||||
data.GroupID,
|
||||
content,
|
||||
mediaPaths,
|
||||
metadata,
|
||||
sender,
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func (c *QQChannel) extractInboundAttachments( chatID, messageID string,
|
||||
attachments []*dto.MessageAttachment ) ([]string, []string) {
|
||||
|
||||
func (c *QQChannel) extractInboundAttachments(
|
||||
chatID, messageID string,
|
||||
attachments []*dto.MessageAttachment,
|
||||
) ([]string, []string) {
|
||||
if len(attachments) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
|
@ -958,7 +1002,7 @@ func (c *QQChannel) genReplyMsg(ctx context.Context, msg bus.OutboundMessage) (d
|
|||
return mdMsg, textMsg
|
||||
}
|
||||
|
||||
func (c *QQChannel) getReplyExtInfo(ctx context.Context, chatID string) (replyID string, seq uint32) {
|
||||
func (c *QQChannel) getReplyExtInfo(_ context.Context, chatID string) (replyID string, seq uint32) {
|
||||
// Attach passive reply msg_id and msg_seq if available.
|
||||
if v, ok := c.lastMsgID.Load(chatID); ok {
|
||||
if msgID, ok := v.(string); ok && msgID != "" {
|
||||
|
|
@ -1150,7 +1194,7 @@ func (c *QQChannel) processAttachments(ctx context.Context, attachments []Messag
|
|||
|
||||
for _, attachment := range attachments {
|
||||
attachmentType := c.getAttachmentType(attachment)
|
||||
localPath := c.downloadAttachment(attachment.URL, attachment.FileName)
|
||||
localPath := c.downloadAttachment(ctx, attachment)
|
||||
if localPath == "" {
|
||||
mediaPaths = append(mediaPaths, attachment.URL)
|
||||
content = appendContent(content, fmt.Sprintf("[%v: %s]", attachment.ContentType, attachment.URL))
|
||||
|
|
@ -1168,10 +1212,19 @@ func (c *QQChannel) processAttachments(ctx context.Context, attachments []Messag
|
|||
return mediaPaths, content
|
||||
}
|
||||
|
||||
|
||||
// downloadAttachment downloads an attachment from QQ server
|
||||
func (c *QQChannel) downloadAttachment(ctx context.Context, attachment MessageAttachment) string {
|
||||
logger.InfoCF("qq", "Downloading attachment", map[string]any{
|
||||
"attachment": attachment,
|
||||
})
|
||||
return utils.DownloadFile(attachment.URL, attachment.FileName, utils.DownloadOptions{
|
||||
LoggerPrefix: "qq",
|
||||
})
|
||||
}
|
||||
|
||||
// getAttachmentType determines the type of attachment (image, audio, video, file)
|
||||
func (c *QQChannel) getAttachmentType(attachment MessageAttachment) string {
|
||||
|
||||
if strings.HasPrefix(attachment.ContentType, "image") {
|
||||
return "image"
|
||||
}
|
||||
|
|
@ -1238,7 +1291,13 @@ func sanitizeURLs(text string) string {
|
|||
})
|
||||
}
|
||||
|
||||
|
||||
// appendContent safely appends content to existing text
|
||||
func appendContent(content, suffix string) string {
|
||||
if content == "" {
|
||||
return suffix
|
||||
}
|
||||
return content + "\n" + suffix
|
||||
}
|
||||
|
||||
func getVoiceInfo(event *dto.WSPayload) (string, string) {
|
||||
_raw, err := json.Marshal(event.Data)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue