mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-07-23 04:56:07 +00:00
41645255f1
* refactor(service): split client.go into focused files
client.go had grown to 4455 lines mixing ~10 responsibilities. Split it
verbatim into cohesive same-package files (no behavior change):
client.go foundation: ClientService, ClientWithAttachments,
ClientCreatePayload, ErrClientNotInInbound, sqlInChunk
client_locks.go inbound mutation locks, delete tombstones, compactOrphans
client_lookup.go read-only lookups (GetByID, List, EffectiveFlow, ...)
client_link.go inbound association sync (SyncInbound, DetachInbound, ...)
client_crud.go single-client CRUD + validation + protocol defaults
client_inbound_apply.go low-level inbound-settings mutators + by-email setters
client_bulk.go bulk attach/detach/adjust/delete/create + DelDepleted
client_traffic.go traffic-reset paths
client_groups.go client group management
client_paging.go paged listing, filtering, sorting, summary
Every declaration moved unchanged (verified: identical func/type/const/var
signature set before vs after). Imports redistributed per file via goimports.
go build ./..., go vet, and go test ./web/service/... all pass.
* refactor(service): split inbound.go into focused files
inbound.go was 4100 lines. Split it verbatim into cohesive same-package
files (no behavior change):
inbound.go core inbound CRUD + InboundService (keeps pkg doc)
inbound_protocol.go protocol / stream capability helpers
inbound_node.go node/runtime/remote coordination + online tracking
inbound_traffic.go traffic accounting, reset, client stats
inbound_client_ips.go per-client IP tracking
inbound_clients.go client lookups within inbounds + copy-clients
inbound_disable.go auto-disable invalid inbounds/clients
inbound_migration.go DB migrations
inbound_sublink.go subscription link providers
inbound_util.go generic slice/string helpers
Identical func/type/const/var signature set before vs after; package doc
comment preserved on inbound.go. Imports redistributed via goimports.
Build, vet, and go test ./web/service/... all pass.
* refactor(service): split tgbot.go into focused files
tgbot.go was 3738 lines dominated by a 1246-line answerCallback. Split it
verbatim into cohesive same-package files (no behavior change):
tgbot.go lifecycle, bot setup, caches, small utils
tgbot_router.go incoming update / command / callback dispatch
tgbot_send.go outbound messaging primitives
tgbot_client.go client views, actions, subscription links
tgbot_inbound.go inbound listing / pickers
tgbot_report.go server usage, exhausted, online, backups, notifications
Identical func/type/const/var signature set before vs after. Imports
redistributed via goimports. Build, vet, and go test ./web/service/... pass.
* refactor(client): dedupe single-field by-email setters
ResetClientIpLimitByEmail, ResetClientExpiryTimeByEmail, and
ResetClientTrafficLimitByEmail shared an identical ~50-line body that
resolves the inbound by email, confirms the client exists, rewrites a
single-client settings payload, and delegates to UpdateInboundClient.
Extract that into applyClientFieldByEmail(inboundSvc, email, mutate) and
reduce each setter to a 3-line wrapper. Behavior is unchanged: same checks
and error strings, same single-client payload contract, same totalGB guard.
SetClientTelegramUserID (resolves by traffic id, different error text) and
ToggleClientEnableByEmail/SetClientEnableByEmail (different return shape and
a pre-read of the old state) intentionally keep their own bodies.
* refactor(service): extract panel/ subpackage
Move the panel-administration leaf services out of the flat service
package into web/service/panel/ (package panel):
user.go UserService (auth / 2FA / LDAP)
panel.go PanelService (restart / self-update) + version helpers
panel_other.go non-unix RestartPanel
panel_unix.go unix RestartPanel
api_token.go ApiTokenService
websocket.go WebSocketService
panel_test.go version/shellQuote unit tests
These are leaves: they depend on core (SettingService, Release) but no
core file references them, so the extraction creates no import cycle.
Core references are now qualified (service.SettingService, service.Release);
callers in main.go, web/web.go, and web/controller/* updated to panel.*.
Build, vet, and go test ./web/... pass.
* refactor(service): extract integration/ subpackage
Move the external-provider integration leaves into web/service/integration/
(package integration):
warp.go WarpService (Cloudflare WARP)
nord.go NordService (NordVPN)
custom_geo.go CustomGeoService (custom geo asset management)
*_test.go custom_geo / panel-proxy tests
These depend on core (SettingService, ServerService, XraySettingService) but
no core file references them. xray_setting.go stays in core because it calls
the unexported SettingService.saveSetting. The shared isBlockedIP SSRF helper
(used by core url_safety.go and by custom_geo) now has a small copy in each
package rather than being exported. Core references qualified; callers in
web/web.go, web/job/*, and web/controller/* updated to integration.*.
Build, vet, and go test ./web/... pass.
* refactor(service): extract tgbot/ subpackage
Move the Telegram bot (6 files + test) into web/service/tgbot/ (package
tgbot). It is a leaf: it embeds five core services (Inbound/Client/Setting/
Server/Xray) and the core never references it, so no import cycle.
To support the package boundary without changing behavior:
- core exposes XrayProcess() *xray.Process so tgbot keeps calling the
exact same running-process methods it used via the package-level `p`;
- three core methods tgbot calls are exported: ClientService.checkIs-
EnabledByEmail -> CheckIsEnabledByEmail, InboundService.getAllEmails ->
GetAllEmails (callers updated in-package);
- tgbot's embedded-field types and the few core type refs (Status,
ClientCreatePayload, SanitizePublicHTTPURL) are now service-qualified.
Callers in main.go, web/web.go, web/job/*, and web/controller/* updated to
tgbot.*. Build, vet, and go test ./web/... pass.
* refactor(service): extract outbound/ subpackage
OutboundService (outbound.go) imports only neutral packages (config,
database, model, xray) and its production code is referenced by no core or
sibling service file — only by web/controller/xray_setting.go and
web/job/xray_traffic_job.go. Move it to web/service/outbound/ (package
outbound); no core qualification needed inside. Callers updated to outbound.*.
The one coupling was a tiny pure test helper, outboundsContainTag, used by
both outbound.go and the core outbound_subscription_test.go; it now has a
small copy in that test file rather than being shared across the boundary.
Build, vet, and go test ./web/... pass.
* refactor(util): move wireguard into its own subpackage
util/wireguard.go was the lone file of the root `util` package (24 lines,
one exported func GenerateWireguardKeypair), while every other util concern
lives in a focused subpackage (util/common, util/crypto, util/netsafe, ...).
Move it to util/wireguard/ (package wireguard) for consistency; its only
importer, web/service/integration/warp.go, is updated. The root `util`
package no longer exists.
* refactor(sub): drop redundant sub prefix from filenames
Inside package sub the subXxx.go prefix just repeats the package name
(like client_*.go did inside service). Rename for consistency; content and
type names are unchanged:
subController.go -> controller.go
subService.go -> service.go
subClashService.go -> clash_service.go
subJsonService.go -> json_service.go
(+ matching _test.go files)
* refactor(controller): rename xui.go -> spa.go
XUIController serves the panel's single-page-app shell; spa.go names that
role plainly (the other controller files are domain-named). File rename only
— the type stays XUIController. api_docs_test.go keys route base paths by
filename, so its "xui.go" case is updated to "spa.go".
* refactor: move backend packages under internal/
Adopt the idiomatic Go application layout: the backend packages now live
under internal/ (a boundary the toolchain enforces), signalling private
implementation instead of a library-style flat root. No runtime behavior
changes — only import paths and a few build/config paths move.
Moved: config, database, logger, mtproto, sub, util, web, xray -> internal/.
main.go stays at the repo root and tools/openapigen stays under tools/ (both
still import internal/* because the internal rule keys off the module root).
The module path github.com/mhsanaei/3x-ui/v3 is unchanged; 149 .go files had
their import prefix rewritten to .../internal/<pkg>.
Couplings the Go compiler can't see, updated to the new layout:
- frontend i18n imports of web/translation (react.ts, setup.components.ts)
- vite outDir + eslint/tsconfig ignore globs -> internal/web/dist
- Dockerfile COPY paths for web/dist and web/translation
- locale.go os.DirFS("web") disk fallback -> "internal/web"
- .gitignore and ci.yml go:embed stub for internal/web/dist
- api_docs_test.go repo-root relative walk (one level deeper)
- tools/openapigen filesystem package paths; ApiTokenView repointed to the
web/service/panel subpackage and codegen regenerated (clears a stale
type the ci.yml codegen check was failing on)
Verified: go build/vet/test (all packages), and frontend typecheck, lint,
vitest (478 tests), and production build into internal/web/dist.
* fix(config): keep test runs from writing logs into the source tree
GetLogFolder() returns a CWD-relative "./log" on Windows. Under `go test`
the working directory is each package's own folder, so InitLogger (called by
tests in web/job, web/service, xray, web/websocket) created stray log/
directories scattered through the source tree (e.g. internal/web/job/log/).
Redirect to a shared temp folder when testing.Testing() reports a test run.
Production behavior is unchanged: Windows still uses ./log next to the binary
and Linux /var/log/x-ui. The log files were always gitignored (*.log) and
never committed; this just stops the noise at the source.
* docs: move subscription-template guide out of root into docs/
sub_templates/ was a top-level folder holding only a README and no actual
templates (3x-ui ships none by design), referenced nowhere and unlinked from
any doc — it read like an empty placeholder cluttering the repo root.
Move the guide to docs/custom-subscription-templates.md (a proper docs home),
reword its intro to read as documentation rather than a folder note, link it
from the Features list in README.md, and drop the empty sub_templates/ folder.
* fix: update stale web/ path references after the internal/ move
The internal/ migration rewrote Go import paths but left some references to
the old top-level layout in docs, comments, and a few runtime disk paths.
Functional (dev-mode only): the disk-serving fallbacks that read the Vite
build from disk when running from source still pointed at web/dist/, which
moved to internal/web/dist/ — so `os.DirFS`/`os.Stat`/`os.ReadFile` in
internal/web/web.go and internal/sub/{sub,controller}.go are corrected.
Production was unaffected (it serves the embedded FS; verified by the Docker
build), but `go run` with a live frontend build silently fell back to embed.
Docs/comments: frontend/README.md, CONTRIBUTING.md, the claude-issue-bot and
release workflows, the openapigen -root help text, and assorted Go comments
now reference internal/web, internal/database, internal/sub, internal/xray,
etc. Package-name mentions (the "web" package), root paths (main.go,
frontend/, install scripts, /etc/x-ui), routes (/panel/api/xray), and the
historical "web/assets no longer exists" note were intentionally left as-is.
* refactor(web): remove the legacy /xui -> /panel redirect middleware
RedirectMiddleware existed only for backward compatibility with the old
`/xui` URL scheme (301-redirecting /xui and /xui/API to /panel and
/panel/api). That cutover was long ago, so drop the middleware, its
registration in initRouter, and the now-inaccurate "URL redirection"
mention in the middleware package doc. Old /xui URLs now 404 like any other
unknown path. HTTPS auto-redirect and auth redirects are unrelated and stay.
* build: fix .dockerignore for internal/ layout and exclude runtime dir
- web/dist -> internal/web/dist: the embedded frontend moved under internal/,
so the stale exclude no longer matched and the locally-built dist could be
sent to the build context (the frontend stage rebuilds it fresh anyway).
- exclude x-ui/: the local runtime directory (SQLite db, geo .dat files, xray
binaries, certs — ~150MB) was being shipped into the build context for no
reason. Verified the pattern excludes only the directory and still keeps
x-ui.sh, which the Dockerfile copies to /usr/bin/x-ui.
471 lines
14 KiB
Go
471 lines
14 KiB
Go
package tgbot
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"embed"
|
|
"math/big"
|
|
"net/http"
|
|
"net/url"
|
|
"os"
|
|
"regexp"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/mhsanaei/3x-ui/v3/internal/logger"
|
|
"github.com/mhsanaei/3x-ui/v3/internal/util/common"
|
|
"github.com/mhsanaei/3x-ui/v3/internal/web/global"
|
|
"github.com/mhsanaei/3x-ui/v3/internal/web/locale"
|
|
"github.com/mhsanaei/3x-ui/v3/internal/web/service"
|
|
|
|
"github.com/mymmrac/telego"
|
|
th "github.com/mymmrac/telego/telegohandler"
|
|
"github.com/valyala/fasthttp"
|
|
"github.com/valyala/fasthttp/fasthttpproxy"
|
|
)
|
|
|
|
var (
|
|
bot *telego.Bot
|
|
|
|
// botCancel stores the function to cancel the context, stopping Long Polling gracefully.
|
|
botCancel context.CancelFunc
|
|
// tgBotMutex protects concurrent access to botCancel variable
|
|
tgBotMutex sync.Mutex
|
|
// botWG waits for the OnReceive Long Polling goroutine to finish.
|
|
botWG sync.WaitGroup
|
|
|
|
botHandler *th.BotHandler
|
|
adminIds []int64
|
|
isRunning bool
|
|
hostname string
|
|
hashStorage *global.HashStorage
|
|
|
|
// Performance improvements
|
|
messageWorkerPool chan struct{} // Semaphore for limiting concurrent message processing
|
|
optimizedHTTPClient *http.Client // HTTP client with connection pooling and timeouts
|
|
|
|
// Simple cache for frequently accessed data
|
|
statusCache struct {
|
|
data *service.Status
|
|
timestamp time.Time
|
|
mutex sync.RWMutex
|
|
}
|
|
|
|
serverStatsCache struct {
|
|
data string
|
|
timestamp time.Time
|
|
mutex sync.RWMutex
|
|
}
|
|
|
|
// clients data to adding new client. receiver_inbound_IDs is the set of
|
|
// inbounds the new client will be attached to; receiver_inbound_ID mirrors
|
|
// the primary pick for the legacy attach-picker entry point. Per-protocol
|
|
// secrets (UUID, password, flow, method) are filled per-inbound on submit
|
|
// by ClientService.fillProtocolDefaults, so the bot only tracks universal
|
|
// client fields here.
|
|
receiver_inbound_ID int
|
|
receiver_inbound_IDs []int
|
|
client_Email string
|
|
client_LimitIP int
|
|
client_TotalGB int64
|
|
client_ExpiryTime int64
|
|
client_Enable bool
|
|
client_TgID string
|
|
client_SubID string
|
|
client_Comment string
|
|
client_Reset int
|
|
)
|
|
|
|
var userStates = make(map[int64]string)
|
|
|
|
// LoginStatus represents the result of a login attempt.
|
|
type LoginStatus byte
|
|
|
|
// Login status constants
|
|
const (
|
|
LoginSuccess LoginStatus = 1 // Login was successful
|
|
LoginFail LoginStatus = 0 // Login failed
|
|
EmptyTelegramUserID = int64(0) // Default value for empty Telegram user ID
|
|
)
|
|
|
|
// LoginAttempt contains safe metadata for panel login notifications.
|
|
// It intentionally does not include attempted passwords.
|
|
type LoginAttempt struct {
|
|
Username string
|
|
IP string
|
|
Time string
|
|
Status LoginStatus
|
|
Reason string
|
|
}
|
|
|
|
// Tgbot provides business logic for Telegram bot integration.
|
|
// It handles bot commands, user interactions, and status reporting via Telegram.
|
|
type Tgbot struct {
|
|
inboundService service.InboundService
|
|
clientService service.ClientService
|
|
settingService service.SettingService
|
|
serverService service.ServerService
|
|
xrayService service.XrayService
|
|
lastStatus *service.Status
|
|
}
|
|
|
|
// NewTgbot creates a new Tgbot instance.
|
|
func (t *Tgbot) NewTgbot() *Tgbot {
|
|
return new(Tgbot)
|
|
}
|
|
|
|
// I18nBot retrieves a localized message for the bot interface.
|
|
func (t *Tgbot) I18nBot(name string, params ...string) string {
|
|
return locale.I18n(locale.Bot, name, params...)
|
|
}
|
|
|
|
// GetHashStorage returns the hash storage instance for callback queries.
|
|
func (t *Tgbot) GetHashStorage() *global.HashStorage {
|
|
return hashStorage
|
|
}
|
|
|
|
// getCachedStatus returns cached server status if it's fresh enough (less than 5 seconds old)
|
|
func (t *Tgbot) getCachedStatus() (*service.Status, bool) {
|
|
statusCache.mutex.RLock()
|
|
defer statusCache.mutex.RUnlock()
|
|
|
|
if statusCache.data != nil && time.Since(statusCache.timestamp) < 5*time.Second {
|
|
return statusCache.data, true
|
|
}
|
|
return nil, false
|
|
}
|
|
|
|
// setCachedStatus updates the status cache
|
|
func (t *Tgbot) setCachedStatus(status *service.Status) {
|
|
statusCache.mutex.Lock()
|
|
defer statusCache.mutex.Unlock()
|
|
|
|
statusCache.data = status
|
|
statusCache.timestamp = time.Now()
|
|
}
|
|
|
|
// getCachedServerStats returns cached server stats if it's fresh enough (less than 10 seconds old)
|
|
func (t *Tgbot) getCachedServerStats() (string, bool) {
|
|
serverStatsCache.mutex.RLock()
|
|
defer serverStatsCache.mutex.RUnlock()
|
|
|
|
if serverStatsCache.data != "" && time.Since(serverStatsCache.timestamp) < 10*time.Second {
|
|
return serverStatsCache.data, true
|
|
}
|
|
return "", false
|
|
}
|
|
|
|
// setCachedServerStats updates the server stats cache
|
|
func (t *Tgbot) setCachedServerStats(stats string) {
|
|
serverStatsCache.mutex.Lock()
|
|
defer serverStatsCache.mutex.Unlock()
|
|
|
|
serverStatsCache.data = stats
|
|
serverStatsCache.timestamp = time.Now()
|
|
}
|
|
|
|
// Start initializes and starts the Telegram bot with the provided translation files.
|
|
func (t *Tgbot) Start(i18nFS embed.FS) error {
|
|
// Initialize localizer
|
|
err := locale.InitLocalizer(i18nFS, &t.settingService)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// If Start is called again (e.g. during reload), ensure any previous long-polling
|
|
// loop is stopped before creating a new bot / receiver.
|
|
StopBot()
|
|
|
|
// Initialize hash storage to store callback queries
|
|
hashStorage = global.NewHashStorage(20 * time.Minute)
|
|
|
|
// Initialize worker pool for concurrent message processing (max 10 concurrent handlers)
|
|
messageWorkerPool = make(chan struct{}, 10)
|
|
|
|
// Initialize optimized HTTP client with connection pooling
|
|
optimizedHTTPClient = &http.Client{
|
|
Timeout: 15 * time.Second,
|
|
Transport: &http.Transport{
|
|
MaxIdleConns: 100,
|
|
MaxIdleConnsPerHost: 10,
|
|
IdleConnTimeout: 30 * time.Second,
|
|
DisableKeepAlives: false,
|
|
},
|
|
}
|
|
|
|
t.SetHostname()
|
|
|
|
// Get Telegram bot token
|
|
tgBotToken, err := t.settingService.GetTgBotToken()
|
|
if err != nil || tgBotToken == "" {
|
|
logger.Warning("Failed to get Telegram bot token:", err)
|
|
return err
|
|
}
|
|
|
|
// Get Telegram bot chat ID(s)
|
|
tgBotID, err := t.settingService.GetTgBotChatId()
|
|
if err != nil {
|
|
logger.Warning("Failed to get Telegram bot chat ID:", err)
|
|
return err
|
|
}
|
|
|
|
parsedAdminIds := make([]int64, 0)
|
|
// Parse admin IDs from comma-separated string
|
|
if tgBotID != "" {
|
|
for adminID := range strings.SplitSeq(tgBotID, ",") {
|
|
id, err := strconv.ParseInt(adminID, 10, 64)
|
|
if err != nil {
|
|
logger.Warning("Failed to parse admin ID from Telegram bot chat ID:", err)
|
|
return err
|
|
}
|
|
parsedAdminIds = append(parsedAdminIds, int64(id))
|
|
}
|
|
}
|
|
tgBotMutex.Lock()
|
|
adminIds = parsedAdminIds
|
|
tgBotMutex.Unlock()
|
|
|
|
// Get Telegram bot proxy URL
|
|
tgBotProxy, err := t.settingService.GetTgBotProxy()
|
|
if err != nil {
|
|
logger.Warning("Failed to get Telegram bot proxy URL:", err)
|
|
}
|
|
|
|
// Fall back to the panel-wide proxy when no dedicated bot proxy is set.
|
|
if tgBotProxy == "" {
|
|
panelProxy, perr := t.settingService.GetPanelProxy()
|
|
if perr != nil {
|
|
logger.Warning("Failed to get panel proxy URL:", perr)
|
|
} else if isSupportedBotProxyScheme(panelProxy) {
|
|
tgBotProxy = panelProxy
|
|
}
|
|
}
|
|
|
|
// Get Telegram bot API server URL
|
|
tgBotAPIServer, err := t.settingService.GetTgBotAPIServer()
|
|
if err != nil {
|
|
logger.Warning("Failed to get Telegram bot API server URL:", err)
|
|
}
|
|
|
|
// Create new Telegram bot instance
|
|
bot, err = t.NewBot(tgBotToken, tgBotProxy, tgBotAPIServer)
|
|
if err != nil {
|
|
logger.Error("Failed to initialize Telegram bot API:", err)
|
|
return err
|
|
}
|
|
|
|
t.trySetBotCommands(bot)
|
|
|
|
// Start receiving Telegram bot messages
|
|
tgBotMutex.Lock()
|
|
alreadyRunning := isRunning || botCancel != nil
|
|
tgBotMutex.Unlock()
|
|
if !alreadyRunning {
|
|
logger.Info("Telegram bot receiver started")
|
|
go t.OnReceive()
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (t *Tgbot) trySetBotCommands(bot *telego.Bot) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
logger.Warning("Failed to register bot commands (Telegram may be rate-limiting); bot will continue without them:", r)
|
|
}
|
|
}()
|
|
|
|
err := bot.SetMyCommands(context.Background(), &telego.SetMyCommandsParams{
|
|
Commands: []telego.BotCommand{
|
|
{Command: "start", Description: t.I18nBot("tgbot.commands.startDesc")},
|
|
{Command: "help", Description: t.I18nBot("tgbot.commands.helpDesc")},
|
|
{Command: "status", Description: t.I18nBot("tgbot.commands.statusDesc")},
|
|
{Command: "id", Description: t.I18nBot("tgbot.commands.idDesc")},
|
|
},
|
|
})
|
|
if err != nil {
|
|
logger.Warning("Failed to set bot commands:", err)
|
|
}
|
|
}
|
|
|
|
func isSupportedBotProxyScheme(proxyUrl string) bool {
|
|
return strings.HasPrefix(proxyUrl, "socks5://") ||
|
|
strings.HasPrefix(proxyUrl, "http://") ||
|
|
strings.HasPrefix(proxyUrl, "https://")
|
|
}
|
|
|
|
// createRobustFastHTTPClient creates a fasthttp.Client with proper connection handling
|
|
func (t *Tgbot) createRobustFastHTTPClient(proxyUrl string) *fasthttp.Client {
|
|
client := &fasthttp.Client{
|
|
// Connection timeouts
|
|
ReadTimeout: 30 * time.Second,
|
|
WriteTimeout: 30 * time.Second,
|
|
MaxIdleConnDuration: 60 * time.Second,
|
|
MaxConnDuration: 0, // unlimited, but controlled by MaxIdleConnDuration
|
|
MaxIdemponentCallAttempts: 3,
|
|
ReadBufferSize: 4096,
|
|
WriteBufferSize: 4096,
|
|
MaxConnsPerHost: 100,
|
|
MaxConnWaitTimeout: 10 * time.Second,
|
|
DisableHeaderNamesNormalizing: false,
|
|
DisablePathNormalizing: false,
|
|
// Retry on connection errors
|
|
RetryIf: func(request *fasthttp.Request) bool {
|
|
// Retry on connection errors for GET requests
|
|
return string(request.Header.Method()) == "GET" || string(request.Header.Method()) == "POST"
|
|
},
|
|
}
|
|
|
|
if proxyUrl != "" {
|
|
if strings.HasPrefix(proxyUrl, "socks5://") {
|
|
client.Dial = fasthttpproxy.FasthttpSocksDialer(proxyUrl)
|
|
} else {
|
|
client.Dial = fasthttpproxy.FasthttpHTTPDialer(proxyUrl)
|
|
}
|
|
}
|
|
|
|
return client
|
|
}
|
|
|
|
// NewBot creates a new Telegram bot instance with optional proxy and API server settings.
|
|
func (t *Tgbot) NewBot(token string, proxyUrl string, apiServerUrl string) (*telego.Bot, error) {
|
|
// Validate proxy URL if provided
|
|
if proxyUrl != "" {
|
|
if !isSupportedBotProxyScheme(proxyUrl) {
|
|
logger.Warning("Unsupported proxy scheme (want socks5:// or http(s)://), ignoring proxy")
|
|
proxyUrl = "" // Clear invalid proxy
|
|
} else if _, err := url.Parse(proxyUrl); err != nil {
|
|
logger.Warningf("Can't parse proxy URL, ignoring proxy: %v", err)
|
|
proxyUrl = ""
|
|
}
|
|
}
|
|
|
|
// Validate API server URL if provided
|
|
if apiServerUrl != "" {
|
|
safeURL, err := service.SanitizePublicHTTPURL(apiServerUrl, false)
|
|
if err != nil {
|
|
logger.Warningf("Invalid or blocked API server URL, using default: %v", err)
|
|
apiServerUrl = ""
|
|
} else {
|
|
apiServerUrl = safeURL
|
|
}
|
|
}
|
|
|
|
// Create robust fasthttp client
|
|
client := t.createRobustFastHTTPClient(proxyUrl)
|
|
|
|
// Build bot options
|
|
var options []telego.BotOption
|
|
options = append(options, telego.WithFastHTTPClient(client))
|
|
|
|
if apiServerUrl != "" {
|
|
options = append(options, telego.WithAPIServer(apiServerUrl))
|
|
}
|
|
|
|
return telego.NewBot(token, options...)
|
|
}
|
|
|
|
// IsRunning checks if the Telegram bot is currently running.
|
|
func (t *Tgbot) IsRunning() bool {
|
|
tgBotMutex.Lock()
|
|
defer tgBotMutex.Unlock()
|
|
return isRunning
|
|
}
|
|
|
|
// SetHostname sets the hostname for the bot.
|
|
func (t *Tgbot) SetHostname() {
|
|
host, err := os.Hostname()
|
|
if err != nil {
|
|
logger.Error("get hostname error:", err)
|
|
hostname = ""
|
|
return
|
|
}
|
|
hostname = host
|
|
}
|
|
|
|
// Stop safely stops the Telegram bot's Long Polling operation.
|
|
// This method now calls the global StopBot function and cleans up other resources.
|
|
func (t *Tgbot) Stop() {
|
|
StopBot()
|
|
logger.Info("Stop Telegram receiver ...")
|
|
tgBotMutex.Lock()
|
|
adminIds = nil
|
|
tgBotMutex.Unlock()
|
|
}
|
|
|
|
// StopBot safely stops the Telegram bot's Long Polling operation by cancelling its context.
|
|
// This is the global function called from main.go's signal handler and t.Stop().
|
|
func StopBot() {
|
|
// Don't hold the mutex while cancelling/waiting.
|
|
tgBotMutex.Lock()
|
|
cancel := botCancel
|
|
botCancel = nil
|
|
handler := botHandler
|
|
botHandler = nil
|
|
isRunning = false
|
|
tgBotMutex.Unlock()
|
|
|
|
if handler != nil {
|
|
handler.Stop()
|
|
}
|
|
|
|
if cancel != nil {
|
|
logger.Info("Sending cancellation signal to Telegram bot...")
|
|
// Cancels the context passed to UpdatesViaLongPolling; this closes updates channel
|
|
// and lets botHandler.Start() exit cleanly.
|
|
cancel()
|
|
botWG.Wait()
|
|
logger.Info("Telegram bot successfully stopped.")
|
|
}
|
|
}
|
|
|
|
// encodeQuery encodes the query string if it's longer than 64 characters.
|
|
func (t *Tgbot) encodeQuery(query string) string {
|
|
// NOTE: we only need to hash for more than 64 chars
|
|
if len(query) <= 64 {
|
|
return query
|
|
}
|
|
|
|
return hashStorage.SaveHash(query)
|
|
}
|
|
|
|
// decodeQuery decodes a hashed query string back to its original form.
|
|
func (t *Tgbot) decodeQuery(query string) (string, error) {
|
|
if !hashStorage.IsMD5(query) {
|
|
return query, nil
|
|
}
|
|
|
|
decoded, exists := hashStorage.GetValue(query)
|
|
if !exists {
|
|
return "", common.NewError("hash not found in storage!")
|
|
}
|
|
|
|
return decoded, nil
|
|
}
|
|
|
|
// randomLowerAndNum generates a random string of lowercase letters and numbers.
|
|
func (t *Tgbot) randomLowerAndNum(length int) string {
|
|
charset := "abcdefghijklmnopqrstuvwxyz0123456789"
|
|
bytes := make([]byte, length)
|
|
for i := range bytes {
|
|
randomIndex, _ := rand.Int(rand.Reader, big.NewInt(int64(len(charset))))
|
|
bytes[i] = charset[randomIndex.Int64()]
|
|
}
|
|
return string(bytes)
|
|
}
|
|
|
|
// int64Contains checks if an int64 slice contains a specific item.
|
|
func int64Contains(slice []int64, item int64) bool {
|
|
return slices.Contains(slice, item)
|
|
}
|
|
|
|
// isSingleWord checks if the text contains only a single word.
|
|
func (t *Tgbot) isSingleWord(text string) bool {
|
|
text = strings.TrimSpace(text)
|
|
re := regexp.MustCompile(`\s+`)
|
|
return re.MatchString(text)
|
|
}
|