1
0

tgbot.go 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594
  1. package tgbot
  2. import (
  3. "context"
  4. "crypto/rand"
  5. "embed"
  6. "math/big"
  7. "net/url"
  8. "os"
  9. "regexp"
  10. "slices"
  11. "strconv"
  12. "strings"
  13. "sync"
  14. "time"
  15. "github.com/mhsanaei/3x-ui/v3/internal/eventbus"
  16. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  17. "github.com/mhsanaei/3x-ui/v3/internal/util/common"
  18. "github.com/mhsanaei/3x-ui/v3/internal/web/global"
  19. "github.com/mhsanaei/3x-ui/v3/internal/web/locale"
  20. "github.com/mhsanaei/3x-ui/v3/internal/web/service"
  21. "github.com/mymmrac/telego"
  22. th "github.com/mymmrac/telego/telegohandler"
  23. "github.com/valyala/fasthttp"
  24. "github.com/valyala/fasthttp/fasthttpproxy"
  25. )
  26. var (
  27. bot *telego.Bot
  28. // botCancel stores the function to cancel the context, stopping Long Polling gracefully.
  29. botCancel context.CancelFunc
  30. // tgBotMutex protects concurrent access to botCancel variable
  31. tgBotMutex sync.Mutex
  32. // botWG waits for the OnReceive Long Polling goroutine to finish.
  33. botWG sync.WaitGroup
  34. botHandler *th.BotHandler
  35. adminIds []int64
  36. isRunning bool
  37. hostname string
  38. hashStorage *global.HashStorage
  39. // EventBus is set from web layer to publish login/security events.
  40. EventBus *eventbus.Bus
  41. // Performance improvements
  42. messageWorkerPool chan struct{} // Semaphore for limiting concurrent message processing
  43. // Simple cache for frequently accessed data
  44. statusCache struct {
  45. data *service.Status
  46. timestamp time.Time
  47. mutex sync.RWMutex
  48. }
  49. serverStatsCache struct {
  50. data string
  51. timestamp time.Time
  52. mutex sync.RWMutex
  53. }
  54. )
  55. // clientDraft is one chat's add-client wizard state. Per-protocol secrets are
  56. // filled per-inbound on submit, so only the universal fields live here.
  57. type clientDraft struct {
  58. sync.Mutex
  59. receiverInboundID int
  60. receiverInboundIDs []int
  61. email string
  62. limitIP int
  63. totalGB int64
  64. expiryTime int64
  65. enable bool
  66. tgID string
  67. subID string
  68. comment string
  69. reset int
  70. }
  71. // chatUser names the admin a wizard belongs to. A private chat's ids are equal;
  72. // in a group they are not, and each admin at its keyboard fills in their own.
  73. type chatUser struct {
  74. chatID int64
  75. userID int64
  76. }
  77. // messageActor reads the sender off a message. A post without one (a channel)
  78. // keys to user 0, an id no admin can hold.
  79. func messageActor(message telego.Message) chatUser {
  80. if message.From == nil {
  81. return chatUser{chatID: message.Chat.ID}
  82. }
  83. return chatUser{chatID: message.Chat.ID, userID: message.From.ID}
  84. }
  85. // callbackActor reads the admin who tapped the button, not the chat the keyboard
  86. // sits in: every admin in a group sees the same keyboard.
  87. func callbackActor(callbackQuery *telego.CallbackQuery) chatUser {
  88. return chatUser{chatID: callbackQuery.Message.GetChat().ID, userID: callbackQuery.From.ID}
  89. }
  90. // clientDrafts keys a draft by the admin filling it in: the steps arrive on the
  91. // worker pool, so one draft let two admins fill in one client between them.
  92. type clientDrafts struct {
  93. mu sync.Mutex
  94. drafts map[chatUser]*clientDraft
  95. }
  96. var addClientDrafts = &clientDrafts{drafts: make(map[chatUser]*clientDraft)}
  97. func (s *clientDrafts) forActor(actor chatUser) *clientDraft {
  98. s.mu.Lock()
  99. defer s.mu.Unlock()
  100. draft, ok := s.drafts[actor]
  101. if !ok {
  102. draft = &clientDraft{}
  103. s.drafts[actor] = draft
  104. }
  105. return draft
  106. }
  107. func (s *clientDrafts) reset(actor chatUser) {
  108. s.mu.Lock()
  109. defer s.mu.Unlock()
  110. delete(s.drafts, actor)
  111. }
  112. // isAddClientStep reports whether callback data belongs to the add-client
  113. // wizard, the only flow that reads or writes a draft.
  114. func isAddClientStep(data string) bool {
  115. return strings.HasPrefix(data, "add_client")
  116. }
  117. func (s *clientDrafts) resetAll() {
  118. s.mu.Lock()
  119. defer s.mu.Unlock()
  120. s.drafts = make(map[chatUser]*clientDraft)
  121. }
  122. // userStateStore guards the per-admin conversation states. The Telegram command
  123. // and callback handlers run on a worker-pool goroutine while the message handler
  124. // runs on the dispatch goroutine, so a bare map would be a concurrent-map-write
  125. // crash. It also expires abandoned conversations so a user who starts a flow and
  126. // goes silent doesn't leave an entry forever.
  127. type userStateStore struct {
  128. mu sync.Mutex
  129. states map[chatUser]userStateEntry
  130. lastPrune time.Time
  131. }
  132. type userStateEntry struct {
  133. state string
  134. at time.Time
  135. }
  136. var userStateMgr = &userStateStore{states: make(map[chatUser]userStateEntry)}
  137. func (s *userStateStore) set(actor chatUser, state string) {
  138. s.mu.Lock()
  139. s.states[actor] = userStateEntry{state: state, at: time.Now()}
  140. s.mu.Unlock()
  141. }
  142. func (s *userStateStore) get(actor chatUser) (string, bool) {
  143. s.mu.Lock()
  144. defer s.mu.Unlock()
  145. e, ok := s.states[actor]
  146. return e.state, ok
  147. }
  148. func (s *userStateStore) clear(actor chatUser) {
  149. s.mu.Lock()
  150. delete(s.states, actor)
  151. s.mu.Unlock()
  152. }
  153. func (s *userStateStore) reset() {
  154. s.mu.Lock()
  155. s.states = make(map[chatUser]userStateEntry)
  156. s.mu.Unlock()
  157. }
  158. // maybePrune drops conversations older than maxAge, at most once per maxAge so a
  159. // busy bot doesn't sweep the whole map on every message.
  160. func (s *userStateStore) maybePrune(maxAge time.Duration) {
  161. s.mu.Lock()
  162. defer s.mu.Unlock()
  163. now := time.Now()
  164. if now.Sub(s.lastPrune) < maxAge {
  165. return
  166. }
  167. s.lastPrune = now
  168. for id, e := range s.states {
  169. if now.Sub(e.at) > maxAge {
  170. delete(s.states, id)
  171. }
  172. }
  173. }
  174. // LoginStatus represents the result of a login attempt.
  175. type LoginStatus byte
  176. // Login status constants
  177. const (
  178. LoginSuccess LoginStatus = 1 // Login was successful
  179. LoginFail LoginStatus = 0 // Login failed
  180. EmptyTelegramUserID = int64(0) // Default value for empty Telegram user ID
  181. )
  182. // LoginAttempt contains safe metadata for panel login notifications.
  183. // It intentionally does not include attempted passwords.
  184. type LoginAttempt struct {
  185. Username string
  186. IP string
  187. Time string
  188. Status LoginStatus
  189. Reason string
  190. }
  191. // Tgbot provides business logic for Telegram bot integration.
  192. // It handles bot commands, user interactions, and status reporting via Telegram.
  193. type Tgbot struct {
  194. inboundService service.InboundService
  195. clientService service.ClientService
  196. settingService service.SettingService
  197. serverService service.ServerService
  198. xrayService service.XrayService
  199. lastStatus *service.Status
  200. }
  201. // NewTgbot creates a new Tgbot instance.
  202. func (t *Tgbot) NewTgbot() *Tgbot {
  203. return new(Tgbot)
  204. }
  205. // I18nBot retrieves a localized message for the bot interface.
  206. func (t *Tgbot) I18nBot(name string, params ...string) string {
  207. return locale.I18n(locale.Bot, name, params...)
  208. }
  209. // GetHashStorage returns the hash storage instance for callback queries.
  210. func (t *Tgbot) GetHashStorage() *global.HashStorage {
  211. return hashStorage
  212. }
  213. // getCachedStatus returns cached server status if it's fresh enough (less than 5 seconds old)
  214. func (t *Tgbot) getCachedStatus() (*service.Status, bool) {
  215. statusCache.mutex.RLock()
  216. defer statusCache.mutex.RUnlock()
  217. if statusCache.data != nil && time.Since(statusCache.timestamp) < 5*time.Second {
  218. return statusCache.data, true
  219. }
  220. return nil, false
  221. }
  222. // setCachedStatus updates the status cache
  223. func (t *Tgbot) setCachedStatus(status *service.Status) {
  224. statusCache.mutex.Lock()
  225. defer statusCache.mutex.Unlock()
  226. statusCache.data = status
  227. statusCache.timestamp = time.Now()
  228. }
  229. // getCachedServerStats returns cached server stats if it's fresh enough (less than 10 seconds old)
  230. func (t *Tgbot) getCachedServerStats() (string, bool) {
  231. serverStatsCache.mutex.RLock()
  232. defer serverStatsCache.mutex.RUnlock()
  233. if serverStatsCache.data != "" && time.Since(serverStatsCache.timestamp) < 10*time.Second {
  234. return serverStatsCache.data, true
  235. }
  236. return "", false
  237. }
  238. // setCachedServerStats updates the server stats cache
  239. func (t *Tgbot) setCachedServerStats(stats string) {
  240. serverStatsCache.mutex.Lock()
  241. defer serverStatsCache.mutex.Unlock()
  242. serverStatsCache.data = stats
  243. serverStatsCache.timestamp = time.Now()
  244. }
  245. // Start initializes and starts the Telegram bot with the provided translation files.
  246. func (t *Tgbot) Start(i18nFS embed.FS) error {
  247. // Initialize localizer
  248. err := locale.InitLocalizer(i18nFS, &t.settingService)
  249. if err != nil {
  250. return err
  251. }
  252. // If Start is called again (e.g. during reload), ensure any previous long-polling
  253. // loop is stopped before creating a new bot / receiver.
  254. StopBot()
  255. // Initialize hash storage to store callback queries
  256. hashStorage = global.NewHashStorage(20 * time.Minute)
  257. // Initialize worker pool for concurrent message processing (max 10 concurrent handlers)
  258. messageWorkerPool = make(chan struct{}, 10)
  259. t.SetHostname()
  260. // Get Telegram bot token
  261. tgBotToken, err := t.settingService.GetTgBotToken()
  262. if err != nil || tgBotToken == "" {
  263. logger.Warning("Failed to get Telegram bot token:", err)
  264. return err
  265. }
  266. // Get Telegram bot chat ID(s)
  267. tgBotID, err := t.settingService.GetTgBotChatId()
  268. if err != nil {
  269. logger.Warning("Failed to get Telegram bot chat ID:", err)
  270. return err
  271. }
  272. parsedAdminIds := make([]int64, 0)
  273. // Parse admin IDs from comma-separated string
  274. if tgBotID != "" {
  275. for adminID := range strings.SplitSeq(tgBotID, ",") {
  276. id, err := strconv.ParseInt(adminID, 10, 64)
  277. if err != nil {
  278. logger.Warning("Failed to parse admin ID from Telegram bot chat ID:", err)
  279. return err
  280. }
  281. parsedAdminIds = append(parsedAdminIds, id)
  282. }
  283. }
  284. tgBotMutex.Lock()
  285. adminIds = parsedAdminIds
  286. tgBotMutex.Unlock()
  287. // Get Telegram bot proxy URL
  288. tgBotProxy, err := t.settingService.GetTgBotProxy()
  289. if err != nil {
  290. logger.Warning("Failed to get Telegram bot proxy URL:", err)
  291. }
  292. // Fall back to the panel-wide egress bridge when no dedicated bot proxy is
  293. // set. Resolved once at bot start: if Xray comes up later, the bot keeps
  294. // its direct connection until it is restarted.
  295. if tgBotProxy == "" {
  296. if egress := t.settingService.PanelEgressProxyURL(); egress != "" && isSupportedBotProxyScheme(egress) {
  297. tgBotProxy = egress
  298. }
  299. }
  300. // Get Telegram bot API server URL
  301. tgBotAPIServer, err := t.settingService.GetTgBotAPIServer()
  302. if err != nil {
  303. logger.Warning("Failed to get Telegram bot API server URL:", err)
  304. }
  305. // Create new Telegram bot instance
  306. bot, err = t.NewBot(tgBotToken, tgBotProxy, tgBotAPIServer)
  307. if err != nil {
  308. logger.Error("Failed to initialize Telegram bot API:", err)
  309. return err
  310. }
  311. t.trySetBotCommands(bot)
  312. // Start receiving Telegram bot messages
  313. tgBotMutex.Lock()
  314. alreadyRunning := isRunning || botCancel != nil
  315. tgBotMutex.Unlock()
  316. if !alreadyRunning {
  317. logger.Info("Telegram bot receiver started")
  318. go t.OnReceive()
  319. }
  320. return nil
  321. }
  322. func (t *Tgbot) trySetBotCommands(bot *telego.Bot) {
  323. defer func() {
  324. if r := recover(); r != nil {
  325. logger.Warning("Failed to register bot commands (Telegram may be rate-limiting); bot will continue without them:", r)
  326. }
  327. }()
  328. err := bot.SetMyCommands(context.Background(), &telego.SetMyCommandsParams{
  329. Commands: []telego.BotCommand{
  330. {Command: "start", Description: t.I18nBot("tgbot.commands.startDesc")},
  331. {Command: "help", Description: t.I18nBot("tgbot.commands.helpDesc")},
  332. {Command: "status", Description: t.I18nBot("tgbot.commands.statusDesc")},
  333. {Command: "id", Description: t.I18nBot("tgbot.commands.idDesc")},
  334. {Command: "usage", Description: t.I18nBot("tgbot.commands.usageDesc")},
  335. {Command: "inbound", Description: t.I18nBot("tgbot.commands.inboundDesc")},
  336. {Command: "restart", Description: t.I18nBot("tgbot.commands.restartDesc")},
  337. {Command: "clearall", Description: t.I18nBot("tgbot.commands.clearallDesc")},
  338. {Command: "broadcast", Description: t.I18nBot("tgbot.commands.broadcastDesc")},
  339. },
  340. })
  341. if err != nil {
  342. logger.Warning("Failed to set bot commands:", err)
  343. }
  344. }
  345. func isSupportedBotProxyScheme(proxyUrl string) bool {
  346. return strings.HasPrefix(proxyUrl, "socks5://") ||
  347. strings.HasPrefix(proxyUrl, "http://") ||
  348. strings.HasPrefix(proxyUrl, "https://")
  349. }
  350. // createRobustFastHTTPClient creates a fasthttp.Client with proper connection handling
  351. func (t *Tgbot) createRobustFastHTTPClient(proxyUrl string) *fasthttp.Client {
  352. client := &fasthttp.Client{
  353. // Connection timeouts
  354. ReadTimeout: 30 * time.Second,
  355. WriteTimeout: 30 * time.Second,
  356. MaxIdleConnDuration: 60 * time.Second,
  357. MaxConnDuration: 0, // unlimited, but controlled by MaxIdleConnDuration
  358. MaxIdemponentCallAttempts: 3,
  359. ReadBufferSize: 4096,
  360. WriteBufferSize: 4096,
  361. MaxConnsPerHost: 100,
  362. MaxConnWaitTimeout: 10 * time.Second,
  363. DisableHeaderNamesNormalizing: false,
  364. DisablePathNormalizing: false,
  365. // resetTimeout stays false to keep the pre-RetryIfErr retry timing.
  366. RetryIfErr: func(request *fasthttp.Request, _ int, _ error) (bool, bool) {
  367. method := string(request.Header.Method())
  368. return false, method == "GET" || method == "POST"
  369. },
  370. }
  371. if proxyUrl != "" {
  372. if strings.HasPrefix(proxyUrl, "socks5://") {
  373. client.Dial = fasthttpproxy.FasthttpSocksDialer(proxyUrl)
  374. } else {
  375. client.Dial = fasthttpproxy.FasthttpHTTPDialer(proxyUrl)
  376. }
  377. }
  378. return client
  379. }
  380. // NewBot creates a new Telegram bot instance with optional proxy and API server settings.
  381. func (t *Tgbot) NewBot(token string, proxyUrl string, apiServerUrl string) (*telego.Bot, error) {
  382. // Validate proxy URL if provided
  383. if proxyUrl != "" {
  384. if !isSupportedBotProxyScheme(proxyUrl) {
  385. logger.Warning("Unsupported proxy scheme (want socks5:// or http(s)://), ignoring proxy")
  386. proxyUrl = "" // Clear invalid proxy
  387. } else if _, err := url.Parse(proxyUrl); err != nil {
  388. logger.Warningf("Can't parse proxy URL, ignoring proxy: %v", err)
  389. proxyUrl = ""
  390. }
  391. }
  392. // Validate API server URL if provided
  393. if apiServerUrl != "" {
  394. safeURL, err := service.SanitizePublicHTTPURL(apiServerUrl, false)
  395. if err != nil {
  396. logger.Warningf("Invalid or blocked API server URL, using default: %v", err)
  397. apiServerUrl = ""
  398. } else {
  399. apiServerUrl = safeURL
  400. }
  401. }
  402. // Create robust fasthttp client
  403. client := t.createRobustFastHTTPClient(proxyUrl)
  404. // Build bot options
  405. var options []telego.BotOption
  406. options = append(options, telego.WithFastHTTPClient(client))
  407. if apiServerUrl != "" {
  408. options = append(options, telego.WithAPIServer(apiServerUrl))
  409. }
  410. return telego.NewBot(token, options...)
  411. }
  412. // IsRunning checks if the Telegram bot is currently running.
  413. func (t *Tgbot) IsRunning() bool {
  414. tgBotMutex.Lock()
  415. defer tgBotMutex.Unlock()
  416. return isRunning
  417. }
  418. // adminSnapshot returns the admin chat list under the mutex Start and Stop
  419. // replace it under: a torn slice header is not a harmless race.
  420. func adminSnapshot() []int64 {
  421. tgBotMutex.Lock()
  422. defer tgBotMutex.Unlock()
  423. return slices.Clone(adminIds)
  424. }
  425. // SetHostname sets the hostname for the bot.
  426. func (t *Tgbot) SetHostname() {
  427. host, err := os.Hostname()
  428. if err != nil {
  429. logger.Error("get hostname error:", err)
  430. hostname = ""
  431. return
  432. }
  433. hostname = host
  434. }
  435. // Stop safely stops the Telegram bot's Long Polling operation.
  436. // This method now calls the global StopBot function and cleans up other resources.
  437. func (t *Tgbot) Stop() {
  438. StopBot()
  439. logger.Info("Stop Telegram receiver ...")
  440. tgBotMutex.Lock()
  441. adminIds = nil
  442. tgBotMutex.Unlock()
  443. }
  444. // StopBot safely stops the Telegram bot's Long Polling operation by cancelling its context.
  445. // This is the global function called from main.go's signal handler and t.Stop().
  446. func StopBot() {
  447. // Don't hold the mutex while cancelling/waiting.
  448. tgBotMutex.Lock()
  449. cancel := botCancel
  450. botCancel = nil
  451. handler := botHandler
  452. botHandler = nil
  453. isRunning = false
  454. tgBotMutex.Unlock()
  455. userStateMgr.reset()
  456. addClientDrafts.resetAll()
  457. broadcastResetAll()
  458. if handler != nil {
  459. _ = handler.Stop()
  460. }
  461. if cancel != nil {
  462. logger.Info("Sending cancellation signal to Telegram bot...")
  463. // Cancels the context passed to UpdatesViaLongPolling; this closes updates channel
  464. // and lets botHandler.Start() exit cleanly.
  465. cancel()
  466. botWG.Wait()
  467. logger.Info("Telegram bot successfully stopped.")
  468. }
  469. }
  470. // encodeQuery encodes the query string if it's longer than 64 characters.
  471. func (t *Tgbot) encodeQuery(query string) string {
  472. // NOTE: we only need to hash for more than 64 chars
  473. if len(query) <= 64 {
  474. return query
  475. }
  476. return hashStorage.SaveHash(query)
  477. }
  478. // decodeQuery decodes a hashed query string back to its original form.
  479. func (t *Tgbot) decodeQuery(query string) (string, error) {
  480. if !hashStorage.IsMD5(query) {
  481. return query, nil
  482. }
  483. decoded, exists := hashStorage.GetValue(query)
  484. if !exists {
  485. return "", common.NewError("hash not found in storage!")
  486. }
  487. return decoded, nil
  488. }
  489. // randomLowerAndNum generates a random string of lowercase letters and numbers.
  490. func (t *Tgbot) randomLowerAndNum(length int) string {
  491. charset := "abcdefghijklmnopqrstuvwxyz0123456789"
  492. bytes := make([]byte, length)
  493. for i := range bytes {
  494. randomIndex, _ := rand.Int(rand.Reader, big.NewInt(int64(len(charset))))
  495. bytes[i] = charset[randomIndex.Int64()]
  496. }
  497. return string(bytes)
  498. }
  499. // int64Contains checks if an int64 slice contains a specific item.
  500. func int64Contains(slice []int64, item int64) bool {
  501. return slices.Contains(slice, item)
  502. }
  503. // isSingleWord checks if the text contains only a single word.
  504. func (t *Tgbot) isSingleWord(text string) bool {
  505. text = strings.TrimSpace(text)
  506. re := regexp.MustCompile(`\s+`)
  507. return re.MatchString(text)
  508. }