From 0b7fcd6c48ccb82e3b5d07d8bc4e3d9bcd7c0737 Mon Sep 17 00:00:00 2001 From: NanamiAdmin Date: Wed, 2 Sep 2026 20:47:25 +0800 Subject: [PATCH] chore(telegram): rebuild telegram module with `go-telegram/bot` fix(telegram): fix bot cannot send message in groups --- go.mod | 1 + go.sum | 2 + internal/controller/pipes/telegram/client.go | 291 ------------------ internal/controller/pipes/telegram/send.go | 85 +++++ .../controller/pipes/telegram/telegram.go | 128 ++++++-- 5 files changed, 185 insertions(+), 322 deletions(-) delete mode 100644 internal/controller/pipes/telegram/client.go create mode 100644 internal/controller/pipes/telegram/send.go diff --git a/go.mod b/go.mod index d15e8b6..71a00f1 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module nukumizu-backend go 1.25.0 require ( + github.com/go-telegram/bot v1.25.0 github.com/gorilla/websocket v1.5.3 gopkg.in/mail.v2 v2.3.1 modernc.org/sqlite v1.55.0 diff --git a/go.sum b/go.sum index 0b921ef..c253f56 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,7 @@ github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/go-telegram/bot v1.25.0 h1:C072msZtr+GjJiQsX/UJLirz+jzSNZVWkckxWufR7Ww= +github.com/go-telegram/bot v1.25.0/go.mod h1:i2TRs7fXWIeaceF3z7KzsMt/he0TwkVC680mvdTFYeM= github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e h1:ijClszYn+mADRFY17kjQEVQ1XRhq2/JR1M3sGqeJoxs= github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e/go.mod h1:boTsfXsheKC2y+lKOCMpSfarhxDeIzfZG1jqGcPl3cA= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= diff --git a/internal/controller/pipes/telegram/client.go b/internal/controller/pipes/telegram/client.go deleted file mode 100644 index dd1ceb5..0000000 --- a/internal/controller/pipes/telegram/client.go +++ /dev/null @@ -1,291 +0,0 @@ -package telegram - -import ( - "encoding/json" - "fmt" - "io" - "net/http" - "net/url" - "strconv" - "strings" - "sync" - "time" - - "nukumizu-backend/internal/netproxy" - "nukumizu-backend/postLog" -) - -const apiBase = "https://api.telegram.org/bot" - -// User mirrors a Telegram user. -type User struct { - ID int64 `json:"id"` - IsBot bool `json:"is_bot"` - FirstName string `json:"first_name"` - Username string `json:"username"` -} - -// Chat mirrors a Telegram chat. -type Chat struct { - ID int64 `json:"id"` - Type string `json:"type"` // "private", "group", "supergroup", "channel" - Title string `json:"title"` -} - -// Message mirrors a Telegram message. Only the fields used for command -// handling are modeled. -type Message struct { - MessageID int64 `json:"message_id"` - From *User `json:"from"` - Chat Chat `json:"chat"` - Date int64 `json:"date"` - Text string `json:"text"` - Entities []MessageEntity `json:"entities"` -} - -// MessageEntity mirrors a Telegram message entity (bot_command, mention, ...). -type MessageEntity struct { - Type string `json:"type"` - Offset int `json:"offset"` - Length int `json:"length"` -} - -// Update mirrors a Telegram update. -type Update struct { - UpdateID int64 `json:"update_id"` - Message *Message `json:"message"` -} - -// apiResponse is the Telegram Bot API response envelope. -type apiResponse struct { - OK bool `json:"ok"` - Result json.RawMessage `json:"result"` - ErrorCode int `json:"error_code"` - Description string `json:"description"` -} - -// maxMessageLen is the safe chunk size for outbound messages. Telegram's hard -// limit is 4096 characters; staying under it leaves headroom for encoding. -const maxMessageLen = 4000 - -// pollTimeout and pollLimit are the long-polling parameters passed to -// getUpdates. The server holds the request open for pollTimeout seconds, so the -// HTTP client timeout below must exceed it. -const ( - pollTimeout = 30 - pollLimit = 100 -) - -// retryDelay is the pause between getUpdates attempts after an error. -const retryDelay = 5 * time.Second - -// Client is a thin Telegram Bot API client. It both polls for incoming updates -// (long polling) and issues outbound Bot API calls (sendMessage). -type Client struct { - token string - useProxy bool - httpClient *http.Client - - stopCh chan struct{} - stopOnce sync.Once -} - -// NewClient creates a Telegram Bot API client. When useProxy is set, all API -// calls are routed through the system-wide network proxy. -func NewClient(token string, useProxy bool) *Client { - return &Client{ - token: token, - useProxy: useProxy, - httpClient: netproxy.HTTPClient(useProxy, 5*time.Minute), - stopCh: make(chan struct{}), - } -} - -func (c *Client) baseURL() string { - return apiBase + c.token -} - -// call performs an API call and unmarshals the result field into result (when -// non-nil). Non-OK responses are wrapped as errors, so the 409 Conflict body -// raised when another poller uses the same token surfaces to the caller. -func (c *Client) call(method string, params url.Values, result interface{}) error { - endpoint := c.baseURL() + "/" + method - - var req *http.Request - var err error - if len(params) > 0 { - req, err = http.NewRequest(http.MethodPost, endpoint, strings.NewReader(params.Encode())) - if err != nil { - return fmt.Errorf("failed to create telegram request: %w", err) - } - req.Header.Set("Content-Type", "application/x-www-form-urlencoded") - } else { - req, err = http.NewRequest(http.MethodGet, endpoint, nil) - if err != nil { - return fmt.Errorf("failed to create telegram request: %w", err) - } - } - - resp, err := c.httpClient.Do(req) - if err != nil { - return fmt.Errorf("telegram API %s failed: %w", method, err) - } - defer resp.Body.Close() - - body, err := io.ReadAll(resp.Body) - if err != nil { - return fmt.Errorf("failed to read telegram response: %w", err) - } - - var apiResp apiResponse - if err := json.Unmarshal(body, &apiResp); err != nil { - return fmt.Errorf("failed to unmarshal telegram response (%s): %w\n%s", method, err, string(body)) - } - if !apiResp.OK { - return fmt.Errorf("telegram API %s error: %d %s", method, apiResp.ErrorCode, apiResp.Description) - } - - if result != nil { - if err := json.Unmarshal(apiResp.Result, result); err != nil { - return fmt.Errorf("failed to unmarshal telegram %s result: %w", method, err) - } - } - - return nil -} - -// GetMe verifies the bot token and returns the bot's own user object. -func (c *Client) GetMe() (*User, error) { - var u User - if err := c.call("getMe", nil, &u); err != nil { - return nil, err - } - return &u, nil -} - -// getUpdates fetches incoming updates. offset is the first update to return; -// pass the previous last update_id + 1 to acknowledge the received updates. -func (c *Client) getUpdates(offset int64, limit, timeout int) ([]Update, error) { - params := url.Values{ - "limit": {strconv.Itoa(limit)}, - "timeout": {strconv.Itoa(timeout)}, - } - if offset > 0 { - params.Set("offset", strconv.FormatInt(offset, 10)) - } - - var updates []Update - if err := c.call("getUpdates", params, &updates); err != nil { - return nil, err - } - return updates, nil -} - -// SendMessage sends a text message to a chat, splitting it into chunks that -// fit Telegram's 4096-character limit. All messages are sent with -// parse_mode=Markdown so fenced code blocks and inline formatting render as -// rich text. Templates must stay valid under Telegram's legacy Markdown: -// unpaired '*' or '_' characters (e.g. a lone '*Event: ...' label) make the -// API reject the whole message. -func (c *Client) SendMessage(chatID int64, text string) error { - if strings.TrimSpace(text) == "" { - return nil - } - for _, chunk := range splitMessage(text, maxMessageLen) { - if err := c.sendMessageChunk(chatID, chunk); err != nil { - return err - } - } - return nil -} - -func (c *Client) sendMessageChunk(chatID int64, text string) error { - params := url.Values{ - "chat_id": {strconv.FormatInt(chatID, 10)}, - "text": {text}, - "parse_mode": {"Markdown"}, - } - return c.call("sendMessage", params, nil) -} - -// Listen runs the long-polling loop, calling onUpdate for every incoming -// update. It blocks until Stop is called, retrying on errors (e.g. the 409 -// Conflict raised when another poller uses the same token). Run it in a -// goroutine. -func (c *Client) Listen(onUpdate func(Update)) { - var offset int64 // Next getUpdates offset: previous last update_id + 1. - - for { - select { - case <-c.stopCh: - return - default: - } - - updates, err := c.getUpdates(offset, pollLimit, pollTimeout) - if err != nil { - if strings.Contains(err.Error(), "409") || strings.Contains(err.Error(), "Conflict") { - postLog.Error("Telegram getUpdates conflict (409): another poller is using the same bot token. Ensure only one instance is polling.") - } else { - postLog.Warning("Telegram getUpdates failed: " + err.Error()) - } - - select { - case <-c.stopCh: - return - case <-time.After(retryDelay): - } - continue - } - - for _, u := range updates { - if u.UpdateID >= offset { - offset = u.UpdateID + 1 - } - onUpdate(u) - } - } -} - -// Stop closes the stop channel and unblocks the Listen loop. -func (c *Client) Stop() { - c.stopOnce.Do(func() { close(c.stopCh) }) -} - -// splitMessage splits text into chunks of at most maxLen runes, keeping -// complete lines when possible. -func splitMessage(text string, maxLen int) []string { - if maxLen <= 0 { - return []string{text} - } - - remaining := []rune(text) - if len(remaining) <= maxLen { - return []string{text} - } - - var chunks []string - for len(remaining) > maxLen { - // Prefer the last newline within the window so messages aren't cut - // mid-line; fall back to a hard rune cut for overlong lines. - cut := maxLen - if nl := lastIndexRune(remaining[:maxLen], '\n'); nl > 0 { - cut = nl + 1 - } - chunks = append(chunks, string(remaining[:cut])) - remaining = remaining[cut:] - } - if len(remaining) > 0 { - chunks = append(chunks, string(remaining)) - } - return chunks -} - -func lastIndexRune(s []rune, r rune) int { - for i := len(s) - 1; i >= 0; i-- { - if s[i] == r { - return i - } - } - return -1 -} diff --git a/internal/controller/pipes/telegram/send.go b/internal/controller/pipes/telegram/send.go new file mode 100644 index 0000000..829db93 --- /dev/null +++ b/internal/controller/pipes/telegram/send.go @@ -0,0 +1,85 @@ +package telegram + +import ( + "context" + "strings" + + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" +) + +// maxMessageLen is the safe chunk size for outbound messages. Telegram's hard +// limit is 4096 characters; staying under it leaves headroom for encoding. +const maxMessageLen = 4000 + +// sendMessage sends a text message to a chat, splitting it into chunks that fit +// Telegram's 4096-character limit. All messages are sent with +// parse_mode=Markdown so fenced code blocks and inline formatting render as +// rich text. Templates must stay valid under Telegram's legacy Markdown: +// unpaired '*' or '_' characters (e.g. a lone '*Event: ...' label) make the +// API reject the whole message. +func (t *TelegramController) sendMessage(chatID int64, text string) error { + if t.client == nil { + return nil + } + if strings.TrimSpace(text) == "" { + return nil + } + for _, chunk := range splitMessage(text, maxMessageLen) { + if err := t.sendMessageChunk(chatID, chunk); err != nil { + return err + } + } + return nil +} + +// sendMessageChunk issues a single sendMessage API call for one chunk. +func (t *TelegramController) sendMessageChunk(chatID int64, text string) error { + ctx, cancel := context.WithTimeout(context.Background(), apiTimeout) + defer cancel() + + _, err := t.client.SendMessage(ctx, &bot.SendMessageParams{ + ChatID: chatID, + Text: text, + ParseMode: models.ParseModeMarkdownV1, // Telegram legacy Markdown + }) + return err +} + +// splitMessage splits text into chunks of at most maxLen runes, keeping +// complete lines when possible. +func splitMessage(text string, maxLen int) []string { + if maxLen <= 0 { + return []string{text} + } + + remaining := []rune(text) + if len(remaining) <= maxLen { + return []string{text} + } + + var chunks []string + for len(remaining) > maxLen { + // Prefer the last newline within the window so messages aren't cut + // mid-line; fall back to a hard rune cut for overlong lines. + cut := maxLen + if nl := lastIndexRune(remaining[:maxLen], '\n'); nl > 0 { + cut = nl + 1 + } + chunks = append(chunks, string(remaining[:cut])) + remaining = remaining[cut:] + } + if len(remaining) > 0 { + chunks = append(chunks, string(remaining)) + } + return chunks +} + +func lastIndexRune(s []rune, r rune) int { + for i := len(s) - 1; i >= 0; i-- { + if s[i] == r { + return i + } + } + return -1 +} diff --git a/internal/controller/pipes/telegram/telegram.go b/internal/controller/pipes/telegram/telegram.go index 5fd6954..5353290 100644 --- a/internal/controller/pipes/telegram/telegram.go +++ b/internal/controller/pipes/telegram/telegram.go @@ -1,37 +1,71 @@ package telegram import ( + "context" "encoding/json" + "errors" "fmt" "strconv" "strings" "sync" + "time" + + "github.com/go-telegram/bot" + "github.com/go-telegram/bot/models" "nukumizu-backend/config" "nukumizu-backend/internal/controller" + "nukumizu-backend/internal/netproxy" "nukumizu-backend/internal/node" "nukumizu-backend/internal/template" "nukumizu-backend/postLog" ) -// TelegramController handles Telegram Bot interactions via long polling. +const ( + // pollTimeout is the getUpdates long-polling window requested from the API. + // The framework sends (pollTimeout - 1s) to getUpdates; apiTimeout, which + // also bounds the HTTP client, is kept well above it so a poll request is + // never cut short by the client. + pollTimeout = 30 * time.Second + + // apiTimeout bounds individual Bot API calls (getMe, sendMessage, ...) and + // the initial token check in Start. + apiTimeout = 2 * time.Minute +) + +// TelegramController handles Telegram Bot interactions via long polling. The +// transport and update delivery are provided by the go-telegram/bot framework; +// this type wires it into the unified controller/command pipeline. type TelegramController struct { cfg config.TelegramConfig - client *Client - bot *User // The bot's own user object from getMe. + client *bot.Bot // underlying go-telegram/bot client + self *models.User // the bot's own user object from getMe mu sync.Mutex usernameToID map[string]int64 // resolved @username -> numeric user ID + + pollCtx context.Context // cancelled by Stop to end long polling + pollCancel context.CancelFunc } -// NewTelegramController creates a new Telegram controller. +// NewTelegramController creates a new Telegram controller. When enabled, the +// underlying bot client is built here (with the getMe handshake skipped) so +// outbound messages can be sent as soon as the controller is registered; token +// verification and polling are deferred to Start. func NewTelegramController(cfg config.TelegramConfig) *TelegramController { t := &TelegramController{ cfg: cfg, usernameToID: make(map[string]int64), } - if cfg.Enabled { - t.client = NewClient(cfg.BotToken, cfg.NetworkUseProxy) + t.pollCtx, t.pollCancel = context.WithCancel(context.Background()) + + if cfg.Enabled && strings.TrimSpace(cfg.BotToken) != "" { + b, err := t.newBot() + if err != nil { + postLog.Error("Failed to build Telegram client: " + err.Error()) + return t + } + t.client = b } return t } @@ -47,22 +81,22 @@ func (t *TelegramController) Start() error { return nil } if t.client == nil { - postLog.Warning("Telegram controller enabled but client is nil") - return nil - } - if t.cfg.BotToken == "" { - postLog.Warning("Telegram controller enabled but no bot token configured") + if strings.TrimSpace(t.cfg.BotToken) == "" { + postLog.Warning("Telegram controller enabled but no bot token configured") + } else { + postLog.Warning("Telegram controller enabled but client is nil") + } return nil } - // The tutorial's first step: verify the token with getMe. The returned bot - // identity is used for @mention mode and echo prevention. - bot, err := t.client.GetMe() + ctx, cancel := context.WithTimeout(context.Background(), apiTimeout) + self, err := t.client.GetMe(ctx) + cancel() if err != nil { return fmt.Errorf("failed to verify Telegram bot token (getMe): %w", err) } - t.bot = bot - postLog.Info(fmt.Sprintf("Telegram bot authenticated: @%s (id %d)", bot.Username, bot.ID)) + t.self = self + postLog.Info(fmt.Sprintf("Telegram bot authenticated: @%s (id %d)", self.Username, self.ID)) postLog.Info("Telegram controller started (long polling)") go func() { @@ -71,17 +105,17 @@ func (t *TelegramController) Start() error { postLog.Error(fmt.Sprintf("Telegram poll loop panic recovered: %v", r)) } }() - t.client.Listen(t.handleUpdate) + // Start blocks until pollCtx is cancelled by Stop. Updates are delivered + // to the default handler registered in newBot. + t.client.Start(t.pollCtx) }() return nil } -// Stop shuts down the Telegram controller and its poll loop. +// Stop cancels the long-polling context and shuts down the controller. func (t *TelegramController) Stop() { - if t.client != nil { - t.client.Stop() - } + t.pollCancel() postLog.Info("Telegram controller stopped") } @@ -90,9 +124,20 @@ func (t *TelegramController) IsEnabled() bool { return t.cfg.Enabled } -// handleUpdate processes a single Telegram update received via long polling. -func (t *TelegramController) handleUpdate(update Update) { - if update.Message == nil { +// handleUpdate processes a single Telegram update received via long polling. It +// is installed as the framework's default handler (every update with a Message +// reaches it). Updates are processed sequentially because the bot is created +// with WithNotAsyncHandlers, mirroring the original single-threaded poll loop. +func (t *TelegramController) handleUpdate(_ context.Context, _ *bot.Bot, update *models.Update) { + // Protect the poll loop: a panic while handling one update must not take + // down long polling. + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Telegram update handler panic recovered: %v", r)) + } + }() + + if update == nil || update.Message == nil { return } msg := update.Message @@ -106,11 +151,11 @@ func (t *TelegramController) handleUpdate(update Update) { if msg.From == nil || msg.From.IsBot { return } - if t.bot != nil && msg.From.ID == t.bot.ID { + if t.self != nil && msg.From.ID == t.self.ID { return } - chatType := telegramChatType(msg.Chat.Type) + chatType := telegramChatType(string(msg.Chat.Type)) if chatType == "" { return // channel or other unsupported chat type. } @@ -135,7 +180,7 @@ func (t *TelegramController) handleUpdate(update Update) { return } - if err := t.client.SendMessage(msg.Chat.ID, response); err != nil { + if err := t.sendMessage(msg.Chat.ID, response); err != nil { postLog.Warning("Failed to send Telegram reply: " + err.Error()) } } @@ -148,8 +193,8 @@ func (t *TelegramController) processCommand(cmd controller.Command) string { // In "at" listen mode, require a mention of the bot and strip it before // parsing, mirroring the QQ CQ-at behavior. - if t.cfg.ListenMethod == "at" && t.bot != nil && t.bot.Username != "" { - mention := "@" + t.bot.Username + if t.cfg.ListenMethod == "at" && t.self != nil && t.self.Username != "" { + mention := "@" + t.self.Username if !strings.Contains(text, mention) { return "" // Not mentioned, ignore. } @@ -271,7 +316,7 @@ func stripCommandBotSuffix(text string) string { // recordUser stores the numeric ID for a sender's username so admins configured // by username (rather than numeric ID) can be resolved later. -func (t *TelegramController) recordUser(u *User) { +func (t *TelegramController) recordUser(u *models.User) { if u == nil || u.Username == "" { return } @@ -349,7 +394,28 @@ func (t *TelegramController) sendAdminMessage(admin string, message string) { // sendToChat sends a message to a chat ID, logging failures. func (t *TelegramController) sendToChat(chatID int64, message string) { - if err := t.client.SendMessage(chatID, message); err != nil { + if err := t.sendMessage(chatID, message); err != nil { postLog.Warning(fmt.Sprintf("Failed to send Telegram message to %d: %v", chatID, err)) } } + +// newBot builds the underlying go-telegram/bot client. getMe is skipped here so +// the constructor stays non-blocking (outbound sends must work before Start +// runs); the token is actually verified in Start. Requests are routed through +// the network proxy when enabled, and poll errors are logged via postLog. +func (t *TelegramController) newBot() (*bot.Bot, error) { + opts := []bot.Option{ + bot.WithSkipGetMe(), + bot.WithHTTPClient(pollTimeout, netproxy.HTTPClient(t.cfg.NetworkUseProxy, apiTimeout)), + bot.WithDefaultHandler(t.handleUpdate), + bot.WithNotAsyncHandlers(), + bot.WithErrorsHandler(func(err error) { + if errors.Is(err, bot.ErrorConflict) { + postLog.Error("Telegram getUpdates conflict (409): another poller is using the same bot token. Ensure only one instance is polling.") + } else { + postLog.Warning("Telegram Bot API error: " + err.Error()) + } + }), + } + return bot.New(t.cfg.BotToken, opts...) +}