refactor(bus,channels): promote peer and messageID from metadata to structured fields

Add bus.Peer struct and explicit Peer/MessageID fields to InboundMessage,
replacing the implicit peer_kind/peer_id/message_id metadata convention.

- Add Peer{Kind, ID} type to pkg/bus/types.go
- Extend InboundMessage with Peer and MessageID fields
- Change BaseChannel.HandleMessage signature to accept peer and messageID
- Adapt all 12 channel implementations to pass structured peer/messageID
- Simplify agent extractPeer() to read msg.Peer directly
- extractParentPeer unchanged (parent_peer still via metadata)
This commit is contained in:
Hoshina 2026-02-22 21:57:12 +08:00
parent 00fd70e1aa
commit 153198e0f3
16 changed files with 108 additions and 84 deletions

View file

@ -1122,21 +1122,20 @@ func (al *AgentLoop) handleCommand(ctx context.Context, msg bus.InboundMessage)
return "", false return "", false
} }
// extractPeer extracts the routing peer from inbound message metadata. // extractPeer extracts the routing peer from the inbound message's structured Peer field.
func extractPeer(msg bus.InboundMessage) *routing.RoutePeer { func extractPeer(msg bus.InboundMessage) *routing.RoutePeer {
peerKind := msg.Metadata["peer_kind"] if msg.Peer.Kind == "" {
if peerKind == "" {
return nil return nil
} }
peerID := msg.Metadata["peer_id"] peerID := msg.Peer.ID
if peerID == "" { if peerID == "" {
if peerKind == "direct" { if msg.Peer.Kind == "direct" {
peerID = msg.SenderID peerID = msg.SenderID
} else { } else {
peerID = msg.ChatID peerID = msg.ChatID
} }
} }
return &routing.RoutePeer{Kind: peerKind, ID: peerID} return &routing.RoutePeer{Kind: msg.Peer.Kind, ID: peerID}
} }
// extractParentPeer extracts the parent peer (reply-to) from inbound message metadata. // extractParentPeer extracts the parent peer (reply-to) from inbound message metadata.

View file

@ -1,11 +1,19 @@
package bus package bus
// Peer identifies the routing peer for a message (direct, group, channel, etc.)
type Peer struct {
Kind string `json:"kind"` // "direct" | "group" | "channel" | ""
ID string `json:"id"`
}
type InboundMessage struct { type InboundMessage struct {
Channel string `json:"channel"` Channel string `json:"channel"`
SenderID string `json:"sender_id"` SenderID string `json:"sender_id"`
ChatID string `json:"chat_id"` ChatID string `json:"chat_id"`
Content string `json:"content"` Content string `json:"content"`
Media []string `json:"media,omitempty"` Media []string `json:"media,omitempty"`
Peer Peer `json:"peer"` // routing peer
MessageID string `json:"message_id,omitempty"` // platform message ID
SessionKey string `json:"session_key"` SessionKey string `json:"session_key"`
Metadata map[string]string `json:"metadata,omitempty"` Metadata map[string]string `json:"metadata,omitempty"`
} }

View file

@ -81,18 +81,25 @@ func (c *BaseChannel) IsAllowed(senderID string) bool {
return false return false
} }
func (c *BaseChannel) HandleMessage(senderID, chatID, content string, media []string, metadata map[string]string) { func (c *BaseChannel) HandleMessage(
peer bus.Peer,
messageID, senderID, chatID, content string,
media []string,
metadata map[string]string,
) {
if !c.IsAllowed(senderID) { if !c.IsAllowed(senderID) {
return return
} }
msg := bus.InboundMessage{ msg := bus.InboundMessage{
Channel: c.name, Channel: c.name,
SenderID: senderID, SenderID: senderID,
ChatID: chatID, ChatID: chatID,
Content: content, Content: content,
Media: media, Media: media,
Metadata: metadata, Peer: peer,
MessageID: messageID,
Metadata: metadata,
} }
c.bus.PublishInbound(msg) c.bus.PublishInbound(msg)

View file

@ -160,12 +160,11 @@ func (c *DingTalkChannel) onChatBotMessageReceived(
"session_webhook": data.SessionWebhook, "session_webhook": data.SessionWebhook,
} }
var peer bus.Peer
if data.ConversationType == "1" { if data.ConversationType == "1" {
metadata["peer_kind"] = "direct" peer = bus.Peer{Kind: "direct", ID: senderID}
metadata["peer_id"] = senderID
} else { } else {
metadata["peer_kind"] = "group" peer = bus.Peer{Kind: "group", ID: data.ConversationId}
metadata["peer_id"] = data.ConversationId
} }
logger.DebugCF("dingtalk", "Received message", map[string]any{ logger.DebugCF("dingtalk", "Received message", map[string]any{
@ -175,7 +174,7 @@ func (c *DingTalkChannel) onChatBotMessageReceived(
}) })
// Handle the message through the base channel // Handle the message through the base channel
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(peer, "", senderID, chatID, content, nil, metadata)
// Return nil to indicate we've handled the message asynchronously // Return nil to indicate we've handled the message asynchronously
// The response will be sent through the message bus // The response will be sent through the message bus

View file

@ -294,19 +294,18 @@ func (c *DiscordChannel) handleMessage(s *discordgo.Session, m *discordgo.Messag
peerID = senderID peerID = senderID
} }
peer := bus.Peer{Kind: peerKind, ID: peerID}
metadata := map[string]string{ metadata := map[string]string{
"message_id": m.ID,
"user_id": senderID, "user_id": senderID,
"username": m.Author.Username, "username": m.Author.Username,
"display_name": senderName, "display_name": senderName,
"guild_id": m.GuildID, "guild_id": m.GuildID,
"channel_id": m.ChannelID, "channel_id": m.ChannelID,
"is_dm": fmt.Sprintf("%t", m.GuildID == ""), "is_dm": fmt.Sprintf("%t", m.GuildID == ""),
"peer_kind": peerKind,
"peer_id": peerID,
} }
c.HandleMessage(senderID, m.ChannelID, content, mediaPaths, metadata) c.HandleMessage(peer, m.ID, senderID, m.ChannelID, content, mediaPaths, metadata)
} }
// startTyping starts a continuous typing indicator loop for the given chatID. // startTyping starts a continuous typing indicator loop for the given chatID.

View file

@ -18,7 +18,9 @@ type FeishuChannel struct {
// NewFeishuChannel returns an error on 32-bit architectures where the Feishu SDK is not supported // NewFeishuChannel returns an error on 32-bit architectures where the Feishu SDK is not supported
func NewFeishuChannel(cfg config.FeishuConfig, bus *bus.MessageBus) (*FeishuChannel, error) { func NewFeishuChannel(cfg config.FeishuConfig, bus *bus.MessageBus) (*FeishuChannel, error) {
return nil, errors.New("feishu channel is not supported on 32-bit architectures (armv7l, 386, etc.). Please use a 64-bit system or disable feishu in your config") return nil, errors.New(
"feishu channel is not supported on 32-bit architectures (armv7l, 386, etc.). Please use a 64-bit system or disable feishu in your config",
)
} }
// Start is a stub method to satisfy the Channel interface // Start is a stub method to satisfy the Channel interface

View file

@ -153,8 +153,9 @@ func (c *FeishuChannel) handleMessageReceive(_ context.Context, event *larkim.P2
} }
metadata := map[string]string{} metadata := map[string]string{}
if messageID := stringValue(message.MessageId); messageID != "" { messageID := ""
metadata["message_id"] = messageID if mid := stringValue(message.MessageId); mid != "" {
messageID = mid
} }
if messageType := stringValue(message.MessageType); messageType != "" { if messageType := stringValue(message.MessageType); messageType != "" {
metadata["message_type"] = messageType metadata["message_type"] = messageType
@ -167,12 +168,11 @@ func (c *FeishuChannel) handleMessageReceive(_ context.Context, event *larkim.P2
} }
chatType := stringValue(message.ChatType) chatType := stringValue(message.ChatType)
var peer bus.Peer
if chatType == "p2p" { if chatType == "p2p" {
metadata["peer_kind"] = "direct" peer = bus.Peer{Kind: "direct", ID: senderID}
metadata["peer_id"] = senderID
} else { } else {
metadata["peer_kind"] = "group" peer = bus.Peer{Kind: "group", ID: chatID}
metadata["peer_id"] = chatID
} }
logger.InfoCF("feishu", "Feishu message received", map[string]any{ logger.InfoCF("feishu", "Feishu message received", map[string]any{
@ -181,7 +181,7 @@ func (c *FeishuChannel) handleMessageReceive(_ context.Context, event *larkim.P2
"preview": utils.Truncate(content, 80), "preview": utils.Truncate(content, 80),
}) })
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(peer, messageID, senderID, chatID, content, nil, metadata)
return nil return nil
} }

View file

@ -364,15 +364,13 @@ func (c *LINEChannel) processEvent(event lineEvent) {
metadata := map[string]string{ metadata := map[string]string{
"platform": "line", "platform": "line",
"source_type": event.Source.Type, "source_type": event.Source.Type,
"message_id": msg.ID,
} }
var peer bus.Peer
if isGroup { if isGroup {
metadata["peer_kind"] = "group" peer = bus.Peer{Kind: "group", ID: chatID}
metadata["peer_id"] = chatID
} else { } else {
metadata["peer_kind"] = "direct" peer = bus.Peer{Kind: "direct", ID: senderID}
metadata["peer_id"] = senderID
} }
logger.DebugCF("line", "Received message", map[string]any{ logger.DebugCF("line", "Received message", map[string]any{
@ -386,7 +384,7 @@ func (c *LINEChannel) processEvent(event lineEvent) {
// Show typing/loading indicator (requires user ID, not group ID) // Show typing/loading indicator (requires user ID, not group ID)
c.sendLoading(senderID) c.sendLoading(senderID)
c.HandleMessage(senderID, chatID, content, mediaPaths, metadata) c.HandleMessage(peer, msg.ID, senderID, chatID, content, mediaPaths, metadata)
} }
// isBotMentioned checks if the bot is mentioned in the message. // isBotMentioned checks if the bot is mentioned in the message.

View file

@ -171,11 +171,9 @@ func (c *MaixCamChannel) handlePersonDetection(msg MaixCamMessage) {
"y": fmt.Sprintf("%.0f", y), "y": fmt.Sprintf("%.0f", y),
"w": fmt.Sprintf("%.0f", w), "w": fmt.Sprintf("%.0f", w),
"h": fmt.Sprintf("%.0f", h), "h": fmt.Sprintf("%.0f", h),
"peer_kind": "channel",
"peer_id": "default",
} }
c.HandleMessage(senderID, chatID, content, []string{}, metadata) c.HandleMessage(bus.Peer{Kind: "channel", ID: "default"}, "", senderID, chatID, content, []string{}, metadata)
} }
func (c *MaixCamChannel) handleStatusUpdate(msg MaixCamMessage) { func (c *MaixCamChannel) handleStatusUpdate(msg MaixCamMessage) {

View file

@ -856,9 +856,9 @@ func (c *OneBotChannel) handleMessage(raw *oneBotRawEvent) {
senderID := strconv.FormatInt(userID, 10) senderID := strconv.FormatInt(userID, 10)
var chatID string var chatID string
metadata := map[string]string{ var peer bus.Peer
"message_id": messageID,
} metadata := map[string]string{}
if parsed.ReplyTo != "" { if parsed.ReplyTo != "" {
metadata["reply_to_message_id"] = parsed.ReplyTo metadata["reply_to_message_id"] = parsed.ReplyTo
@ -867,14 +867,12 @@ func (c *OneBotChannel) handleMessage(raw *oneBotRawEvent) {
switch raw.MessageType { switch raw.MessageType {
case "private": case "private":
chatID = "private:" + senderID chatID = "private:" + senderID
metadata["peer_kind"] = "direct" peer = bus.Peer{Kind: "direct", ID: senderID}
metadata["peer_id"] = senderID
case "group": case "group":
groupIDStr := strconv.FormatInt(groupID, 10) groupIDStr := strconv.FormatInt(groupID, 10)
chatID = "group:" + groupIDStr chatID = "group:" + groupIDStr
metadata["peer_kind"] = "group" peer = bus.Peer{Kind: "group", ID: groupIDStr}
metadata["peer_id"] = groupIDStr
metadata["group_id"] = groupIDStr metadata["group_id"] = groupIDStr
senderUserID, _ := parseJSONInt64(sender.UserID) senderUserID, _ := parseJSONInt64(sender.UserID)
@ -929,7 +927,7 @@ func (c *OneBotChannel) handleMessage(raw *oneBotRawEvent) {
c.pendingEmojiMsg.Store(chatID, messageID) c.pendingEmojiMsg.Store(chatID, messageID)
} }
c.HandleMessage(senderID, chatID, content, parsed.Media, metadata) c.HandleMessage(peer, messageID, senderID, chatID, content, parsed.Media, metadata)
} }
func (c *OneBotChannel) isDuplicate(messageID string) bool { func (c *OneBotChannel) isDuplicate(messageID string) bool {

View file

@ -164,13 +164,17 @@ func (c *QQChannel) handleC2CMessage() event.C2CMessageEventHandler {
}) })
// 转发到消息总线 // 转发到消息总线
metadata := map[string]string{ metadata := map[string]string{}
"message_id": data.ID,
"peer_kind": "direct",
"peer_id": senderID,
}
c.HandleMessage(senderID, senderID, content, []string{}, metadata) c.HandleMessage(
bus.Peer{Kind: "direct", ID: senderID},
data.ID,
senderID,
senderID,
content,
[]string{},
metadata,
)
return nil return nil
} }
@ -208,13 +212,18 @@ func (c *QQChannel) handleGroupATMessage() event.GroupATMessageEventHandler {
// 转发到消息总线(使用 GroupID 作为 ChatID // 转发到消息总线(使用 GroupID 作为 ChatID
metadata := map[string]string{ metadata := map[string]string{
"message_id": data.ID, "group_id": data.GroupID,
"group_id": data.GroupID,
"peer_kind": "group",
"peer_id": data.GroupID,
} }
c.HandleMessage(senderID, data.GroupID, content, []string{}, metadata) c.HandleMessage(
bus.Peer{Kind: "group", ID: data.GroupID},
data.ID,
senderID,
data.GroupID,
content,
[]string{},
metadata,
)
return nil return nil
} }

View file

@ -284,13 +284,13 @@ func (c *SlackChannel) handleMessageEvent(ev *slackevents.MessageEvent) {
peerID = senderID peerID = senderID
} }
peer := bus.Peer{Kind: peerKind, ID: peerID}
metadata := map[string]string{ metadata := map[string]string{
"message_ts": messageTS, "message_ts": messageTS,
"channel_id": channelID, "channel_id": channelID,
"thread_ts": threadTS, "thread_ts": threadTS,
"platform": "slack", "platform": "slack",
"peer_kind": peerKind,
"peer_id": peerID,
"team_id": c.teamID, "team_id": c.teamID,
} }
@ -301,7 +301,7 @@ func (c *SlackChannel) handleMessageEvent(ev *slackevents.MessageEvent) {
"has_thread": threadTS != "", "has_thread": threadTS != "",
}) })
c.HandleMessage(senderID, chatID, content, mediaPaths, metadata) c.HandleMessage(peer, messageTS, senderID, chatID, content, mediaPaths, metadata)
} }
func (c *SlackChannel) handleAppMention(ev *slackevents.AppMentionEvent) { func (c *SlackChannel) handleAppMention(ev *slackevents.AppMentionEvent) {
@ -351,18 +351,18 @@ func (c *SlackChannel) handleAppMention(ev *slackevents.AppMentionEvent) {
mentionPeerID = senderID mentionPeerID = senderID
} }
mentionPeer := bus.Peer{Kind: mentionPeerKind, ID: mentionPeerID}
metadata := map[string]string{ metadata := map[string]string{
"message_ts": messageTS, "message_ts": messageTS,
"channel_id": channelID, "channel_id": channelID,
"thread_ts": threadTS, "thread_ts": threadTS,
"platform": "slack", "platform": "slack",
"is_mention": "true", "is_mention": "true",
"peer_kind": mentionPeerKind,
"peer_id": mentionPeerID,
"team_id": c.teamID, "team_id": c.teamID,
} }
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(mentionPeer, messageTS, senderID, chatID, content, nil, metadata)
} }
func (c *SlackChannel) handleSlashCommand(event socketmode.Event) { func (c *SlackChannel) handleSlashCommand(event socketmode.Event) {
@ -396,8 +396,6 @@ func (c *SlackChannel) handleSlashCommand(event socketmode.Event) {
"platform": "slack", "platform": "slack",
"is_command": "true", "is_command": "true",
"trigger_id": cmd.TriggerID, "trigger_id": cmd.TriggerID,
"peer_kind": "channel",
"peer_id": channelID,
"team_id": c.teamID, "team_id": c.teamID,
} }
@ -407,7 +405,7 @@ func (c *SlackChannel) handleSlashCommand(event socketmode.Event) {
"text": utils.Truncate(content, 50), "text": utils.Truncate(content, 50),
}) })
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(bus.Peer{Kind: "channel", ID: channelID}, "", senderID, chatID, content, nil, metadata)
} }
func (c *SlackChannel) downloadSlackFile(file slack.File) string { func (c *SlackChannel) downloadSlackFile(file slack.File) string {

View file

@ -362,17 +362,25 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, message *telego.Mes
peerID = fmt.Sprintf("%d", chatID) peerID = fmt.Sprintf("%d", chatID)
} }
peer := bus.Peer{Kind: peerKind, ID: peerID}
messageID := fmt.Sprintf("%d", message.MessageID)
metadata := map[string]string{ metadata := map[string]string{
"message_id": fmt.Sprintf("%d", message.MessageID),
"user_id": fmt.Sprintf("%d", user.ID), "user_id": fmt.Sprintf("%d", user.ID),
"username": user.Username, "username": user.Username,
"first_name": user.FirstName, "first_name": user.FirstName,
"is_group": fmt.Sprintf("%t", message.Chat.Type != "private"), "is_group": fmt.Sprintf("%t", message.Chat.Type != "private"),
"peer_kind": peerKind,
"peer_id": peerID,
} }
c.HandleMessage(fmt.Sprintf("%d", user.ID), fmt.Sprintf("%d", chatID), content, mediaPaths, metadata) c.HandleMessage(
peer,
messageID,
fmt.Sprintf("%d", user.ID),
fmt.Sprintf("%d", chatID),
content,
mediaPaths,
metadata,
)
return nil return nil
} }

View file

@ -425,6 +425,9 @@ func (c *WeComAppChannel) processMessage(ctx context.Context, msg WeComXMLMessag
// Build metadata // Build metadata
// WeCom App only supports direct messages (private chat) // WeCom App only supports direct messages (private chat)
peer := bus.Peer{Kind: "direct", ID: senderID}
messageID := fmt.Sprintf("%d", msg.MsgId)
metadata := map[string]string{ metadata := map[string]string{
"msg_type": msg.MsgType, "msg_type": msg.MsgType,
"msg_id": fmt.Sprintf("%d", msg.MsgId), "msg_id": fmt.Sprintf("%d", msg.MsgId),
@ -432,8 +435,6 @@ func (c *WeComAppChannel) processMessage(ctx context.Context, msg WeComXMLMessag
"platform": "wecom_app", "platform": "wecom_app",
"media_id": msg.MediaId, "media_id": msg.MediaId,
"create_time": fmt.Sprintf("%d", msg.CreateTime), "create_time": fmt.Sprintf("%d", msg.CreateTime),
"peer_kind": "direct",
"peer_id": senderID,
} }
content := msg.Content content := msg.Content
@ -445,7 +446,7 @@ func (c *WeComAppChannel) processMessage(ctx context.Context, msg WeComXMLMessag
}) })
// Handle the message through the base channel // Handle the message through the base channel
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(peer, messageID, senderID, chatID, content, nil, metadata)
} }
// tokenRefreshLoop periodically refreshes the access token // tokenRefreshLoop periodically refreshes the access token

View file

@ -378,12 +378,12 @@ func (c *WeComBotChannel) processMessage(ctx context.Context, msg WeComBotMessag
} }
// Build metadata // Build metadata
peer := bus.Peer{Kind: peerKind, ID: peerID}
metadata := map[string]string{ metadata := map[string]string{
"msg_type": msg.MsgType, "msg_type": msg.MsgType,
"msg_id": msg.MsgID, "msg_id": msg.MsgID,
"platform": "wecom", "platform": "wecom",
"peer_kind": peerKind,
"peer_id": peerID,
"response_url": msg.ResponseURL, "response_url": msg.ResponseURL,
} }
if isGroupChat { if isGroupChat {
@ -400,7 +400,7 @@ func (c *WeComBotChannel) processMessage(ctx context.Context, msg WeComBotMessag
}) })
// Handle the message through the base channel // Handle the message through the base channel
c.HandleMessage(senderID, chatID, content, nil, metadata) c.HandleMessage(peer, msg.MsgID, senderID, chatID, content, nil, metadata)
} }
// sendWebhookReply sends a reply using the webhook URL // sendWebhookReply sends a reply using the webhook URL

View file

@ -172,22 +172,22 @@ func (c *WhatsAppChannel) handleIncomingMessage(msg map[string]any) {
} }
metadata := make(map[string]string) metadata := make(map[string]string)
if messageID, ok := msg["id"].(string); ok { var messageID string
metadata["message_id"] = messageID if mid, ok := msg["id"].(string); ok {
messageID = mid
} }
if userName, ok := msg["from_name"].(string); ok { if userName, ok := msg["from_name"].(string); ok {
metadata["user_name"] = userName metadata["user_name"] = userName
} }
var peer bus.Peer
if chatID == senderID { if chatID == senderID {
metadata["peer_kind"] = "direct" peer = bus.Peer{Kind: "direct", ID: senderID}
metadata["peer_id"] = senderID
} else { } else {
metadata["peer_kind"] = "group" peer = bus.Peer{Kind: "group", ID: chatID}
metadata["peer_id"] = chatID
} }
log.Printf("WhatsApp message from %s: %s...", senderID, utils.Truncate(content, 50)) log.Printf("WhatsApp message from %s: %s...", senderID, utils.Truncate(content, 50))
c.HandleMessage(senderID, chatID, content, mediaPaths, metadata) c.HandleMessage(peer, messageID, senderID, chatID, content, mediaPaths, metadata)
} }