Add Telegram Group Topics (Forum) support

- Add threadIDs sync.Map to TelegramChannel struct to store thread ID mappings
- Implement parseTelegramChatID() to parse chatID:threadID format (similar to Slack)
- Update handleMessage() to extract and store MessageThreadID from incoming messages
- Encode thread ID into chatID string when messages come from topics
- Update Send() method to parse and set MessageThreadID on outgoing messages
- Set MessageThreadID on "Thinking..." placeholder messages
- Set MessageThreadID on SendChatAction (typing indicator)
- Add metadata fields: message_thread_id and is_topic_message
- Add comprehensive unit tests for parseTelegramChatID()

Co-authored-by: zhaopengme <1415418+zhaopengme@users.noreply.github.com>
This commit is contained in:
copilot-swe-agent[bot] 2026-02-15 09:11:44 +00:00
parent 6d55011835
commit fb2179d1b5
2 changed files with 189 additions and 6 deletions

View file

@ -29,6 +29,7 @@ type TelegramChannel struct {
transcriber *voice.GroqTranscriber transcriber *voice.GroqTranscriber
placeholders sync.Map // chatID -> messageID placeholders sync.Map // chatID -> messageID
stopThinking sync.Map // chatID -> thinkingCancel stopThinking sync.Map // chatID -> thinkingCancel
threadIDs sync.Map // chatIDStr -> MessageThreadID
} }
type thinkingCancel struct { type thinkingCancel struct {
@ -124,7 +125,7 @@ func (c *TelegramChannel) Send(ctx context.Context, msg bus.OutboundMessage) err
return fmt.Errorf("telegram bot not running") return fmt.Errorf("telegram bot not running")
} }
chatID, err := parseChatID(msg.ChatID) chatID, threadID, err := parseTelegramChatID(msg.ChatID)
if err != nil { if err != nil {
return fmt.Errorf("invalid chat ID: %w", err) return fmt.Errorf("invalid chat ID: %w", err)
} }
@ -144,6 +145,7 @@ func (c *TelegramChannel) Send(ctx context.Context, msg bus.OutboundMessage) err
c.placeholders.Delete(msg.ChatID) c.placeholders.Delete(msg.ChatID)
editMsg := tu.EditMessageText(tu.ID(chatID), pID.(int), htmlContent) editMsg := tu.EditMessageText(tu.ID(chatID), pID.(int), htmlContent)
editMsg.ParseMode = telego.ModeHTML editMsg.ParseMode = telego.ModeHTML
// Note: EditMessageText doesn't require MessageThreadID as it edits existing message
if _, err = c.bot.EditMessageText(ctx, editMsg); err == nil { if _, err = c.bot.EditMessageText(ctx, editMsg); err == nil {
return nil return nil
@ -154,6 +156,11 @@ func (c *TelegramChannel) Send(ctx context.Context, msg bus.OutboundMessage) err
tgMsg := tu.Message(tu.ID(chatID), htmlContent) tgMsg := tu.Message(tu.ID(chatID), htmlContent)
tgMsg.ParseMode = telego.ModeHTML tgMsg.ParseMode = telego.ModeHTML
// Set thread ID if present
if threadID != 0 {
tgMsg.MessageThreadID = threadID
}
if _, err = c.bot.SendMessage(ctx, tgMsg); err != nil { if _, err = c.bot.SendMessage(ctx, tgMsg); err != nil {
logger.ErrorCF("telegram", "HTML parse failed, falling back to plain text", map[string]interface{}{ logger.ErrorCF("telegram", "HTML parse failed, falling back to plain text", map[string]interface{}{
"error": err.Error(), "error": err.Error(),
@ -195,6 +202,16 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, update telego.Updat
chatID := message.Chat.ID chatID := message.Chat.ID
c.chatIDs[senderID] = chatID c.chatIDs[senderID] = chatID
// Check for message thread (forum topic)
messageThreadID := message.MessageThreadID
chatIDStr := fmt.Sprintf("%d", chatID)
if messageThreadID != 0 {
// Store thread ID for later use
c.threadIDs.Store(chatIDStr, messageThreadID)
// Encode thread ID into chatID string (similar to Slack pattern)
chatIDStr = fmt.Sprintf("%d:%d", chatID, messageThreadID)
}
content := "" content := ""
mediaPaths := []string{} mediaPaths := []string{}
localFiles := []string{} // 跟踪需要清理的本地文件 localFiles := []string{} // 跟踪需要清理的本地文件
@ -301,11 +318,16 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, update telego.Updat
logger.DebugCF("telegram", "Received message", map[string]interface{}{ logger.DebugCF("telegram", "Received message", map[string]interface{}{
"sender_id": senderID, "sender_id": senderID,
"chat_id": fmt.Sprintf("%d", chatID), "chat_id": fmt.Sprintf("%d", chatID),
"thread_id": messageThreadID,
"preview": utils.Truncate(content, 50), "preview": utils.Truncate(content, 50),
}) })
// Thinking indicator // Thinking indicator - include thread ID if present
err := c.bot.SendChatAction(ctx, tu.ChatAction(tu.ID(chatID), telego.ChatActionTyping)) chatAction := tu.ChatAction(tu.ID(chatID), telego.ChatActionTyping)
if messageThreadID != 0 {
chatAction.MessageThreadID = messageThreadID
}
err := c.bot.SendChatAction(ctx, chatAction)
if err != nil { if err != nil {
logger.ErrorCF("telegram", "Failed to send chat action", map[string]interface{}{ logger.ErrorCF("telegram", "Failed to send chat action", map[string]interface{}{
"error": err.Error(), "error": err.Error(),
@ -313,7 +335,6 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, update telego.Updat
} }
// Stop any previous thinking animation // Stop any previous thinking animation
chatIDStr := fmt.Sprintf("%d", chatID)
if prevStop, ok := c.stopThinking.Load(chatIDStr); ok { if prevStop, ok := c.stopThinking.Load(chatIDStr); ok {
if cf, ok := prevStop.(*thinkingCancel); ok && cf != nil { if cf, ok := prevStop.(*thinkingCancel); ok && cf != nil {
cf.Cancel() cf.Cancel()
@ -324,7 +345,12 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, update telego.Updat
_, thinkCancel := context.WithTimeout(ctx, 5*time.Minute) _, thinkCancel := context.WithTimeout(ctx, 5*time.Minute)
c.stopThinking.Store(chatIDStr, &thinkingCancel{fn: thinkCancel}) c.stopThinking.Store(chatIDStr, &thinkingCancel{fn: thinkCancel})
pMsg, err := c.bot.SendMessage(ctx, tu.Message(tu.ID(chatID), "Thinking... 💭")) // Send "Thinking..." message - include thread ID if present
thinkingMsg := tu.Message(tu.ID(chatID), "Thinking... 💭")
if messageThreadID != 0 {
thinkingMsg.MessageThreadID = messageThreadID
}
pMsg, err := c.bot.SendMessage(ctx, thinkingMsg)
if err == nil { if err == nil {
pID := pMsg.MessageID pID := pMsg.MessageID
c.placeholders.Store(chatIDStr, pID) c.placeholders.Store(chatIDStr, pID)
@ -338,7 +364,12 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, update telego.Updat
"is_group": fmt.Sprintf("%t", message.Chat.Type != "private"), "is_group": fmt.Sprintf("%t", message.Chat.Type != "private"),
} }
c.HandleMessage(senderID, fmt.Sprintf("%d", chatID), content, mediaPaths, metadata) if messageThreadID != 0 {
metadata["message_thread_id"] = fmt.Sprintf("%d", messageThreadID)
metadata["is_topic_message"] = "true"
}
c.HandleMessage(senderID, chatIDStr, content, mediaPaths, metadata)
} }
func (c *TelegramChannel) downloadPhoto(ctx context.Context, fileID string) string { func (c *TelegramChannel) downloadPhoto(ctx context.Context, fileID string) string {
@ -386,6 +417,30 @@ func parseChatID(chatIDStr string) (int64, error) {
return id, err return id, err
} }
// parseTelegramChatID extracts chatID and threadID from a combined chatID string
// Format: "chatID" or "chatID:threadID"
func parseTelegramChatID(chatIDStr string) (chatID int64, threadID int, err error) {
parts := strings.SplitN(chatIDStr, ":", 2)
var id int64
_, err = fmt.Sscanf(parts[0], "%d", &id)
if err != nil {
return 0, 0, err
}
chatID = id
if len(parts) > 1 {
var tid int
_, err = fmt.Sscanf(parts[1], "%d", &tid)
if err != nil {
return chatID, 0, fmt.Errorf("invalid thread ID: %w", err)
}
threadID = tid
}
return chatID, threadID, nil
}
func markdownToTelegramHTML(text string) string { func markdownToTelegramHTML(text string) string {
if text == "" { if text == "" {
return "" return ""

View file

@ -0,0 +1,128 @@
package channels
import (
"testing"
)
func TestParseTelegramChatID(t *testing.T) {
tests := []struct {
name string
chatIDStr string
wantChatID int64
wantThreadID int
wantErr bool
}{
{
name: "chat only",
chatIDStr: "123456789",
wantChatID: 123456789,
wantThreadID: 0,
wantErr: false,
},
{
name: "negative chat ID (private chat)",
chatIDStr: "-987654321",
wantChatID: -987654321,
wantThreadID: 0,
wantErr: false,
},
{
name: "chat with thread",
chatIDStr: "123456789:42",
wantChatID: 123456789,
wantThreadID: 42,
wantErr: false,
},
{
name: "negative chat with thread",
chatIDStr: "-987654321:100",
wantChatID: -987654321,
wantThreadID: 100,
wantErr: false,
},
{
name: "invalid chat ID",
chatIDStr: "invalid",
wantChatID: 0,
wantThreadID: 0,
wantErr: true,
},
{
name: "invalid thread ID",
chatIDStr: "123456789:invalid",
wantChatID: 123456789,
wantThreadID: 0,
wantErr: true,
},
{
name: "empty string",
chatIDStr: "",
wantChatID: 0,
wantThreadID: 0,
wantErr: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
chatID, threadID, err := parseTelegramChatID(tt.chatIDStr)
if (err != nil) != tt.wantErr {
t.Errorf("parseTelegramChatID(%q) error = %v, wantErr %v", tt.chatIDStr, err, tt.wantErr)
return
}
if !tt.wantErr {
if chatID != tt.wantChatID {
t.Errorf("parseTelegramChatID(%q) chatID = %d, want %d", tt.chatIDStr, chatID, tt.wantChatID)
}
if threadID != tt.wantThreadID {
t.Errorf("parseTelegramChatID(%q) threadID = %d, want %d", tt.chatIDStr, threadID, tt.wantThreadID)
}
}
})
}
}
func TestParseChatID(t *testing.T) {
tests := []struct {
name string
chatIDStr string
want int64
wantErr bool
}{
{
name: "positive chat ID",
chatIDStr: "123456789",
want: 123456789,
wantErr: false,
},
{
name: "negative chat ID",
chatIDStr: "-987654321",
want: -987654321,
wantErr: false,
},
{
name: "invalid chat ID",
chatIDStr: "invalid",
want: 0,
wantErr: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := parseChatID(tt.chatIDStr)
if (err != nil) != tt.wantErr {
t.Errorf("parseChatID(%q) error = %v, wantErr %v", tt.chatIDStr, err, tt.wantErr)
return
}
if !tt.wantErr && got != tt.want {
t.Errorf("parseChatID(%q) = %d, want %d", tt.chatIDStr, got, tt.want)
}
})
}
}