mirror of
https://github.com/mmahdium/TGSS.git
synced 2026-08-12 08:32:48 +03:30
feat: Add authorization handlers and improve message retrieval limits
- Introduced NewAuthHandler and integrated it into the server routes. - Added limit query parameter with default value handling in GetLatestMessages. - Improved error handling for invalid limit values. - Refactored Config struct to remove unnecessary fields and clean up loading logic. - Updated Telegram client connection handling to enhance stability.
This commit is contained in:
@@ -18,6 +18,7 @@ func main() {
|
||||
server.NewGin,
|
||||
|
||||
handlers.NewChannelHandler,
|
||||
handlers.NewAuthHandler,
|
||||
),
|
||||
telegram.Module,
|
||||
fx.Invoke(server.RegisterRoutes),
|
||||
|
||||
@@ -29,6 +29,7 @@ require (
|
||||
github.com/goccy/go-json v0.10.5 // indirect
|
||||
github.com/goccy/go-yaml v1.19.2 // indirect
|
||||
github.com/google/uuid v1.6.0 // indirect
|
||||
github.com/gotd/contrib v0.21.1 // indirect
|
||||
github.com/gotd/ige v0.2.2 // indirect
|
||||
github.com/gotd/neo v0.1.5 // indirect
|
||||
github.com/joho/godotenv v1.5.1 // indirect
|
||||
|
||||
@@ -62,6 +62,8 @@ github.com/goccy/go-yaml v1.19.2/go.mod h1:XBurs7gK8ATbW4ZPGKgcbrY1Br56PdM69F7Lk
|
||||
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
|
||||
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
|
||||
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
|
||||
github.com/gotd/contrib v0.21.1 h1:NSF+0YEnosQ34QEo2o4s6MA5YFDAor1LVvLhN1L3H1M=
|
||||
github.com/gotd/contrib v0.21.1/go.mod h1:trVJBP9Q/TJbjmJbVnLc0cnX/8T4N0RpQBULVa3BNnE=
|
||||
github.com/gotd/ige v0.2.2 h1:XQ9dJZwBfDnOGSTxKXBGP4gMud3Qku2ekScRjDWWfEk=
|
||||
github.com/gotd/ige v0.2.2/go.mod h1:tuCRb+Y5Y3eNTo3ypIfNpQ4MFjrnONiL2jN2AKZXmb0=
|
||||
github.com/gotd/neo v0.1.5 h1:oj0iQfMbGClP8xI59x7fE/uHoTJD7NZH9oV1WNuPukQ=
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
package api
|
||||
@@ -12,7 +12,6 @@ import (
|
||||
|
||||
type Config struct {
|
||||
TgPhoneNumber string
|
||||
TgBotToken string // Only for getting authenticated without phone number
|
||||
TgAppId int
|
||||
TgAppHash string
|
||||
|
||||
@@ -21,21 +20,12 @@ type Config struct {
|
||||
}
|
||||
|
||||
func Load(logger *zap.Logger) *Config {
|
||||
// TODO: get app port and host
|
||||
err := godotenv.Load()
|
||||
if err != nil {
|
||||
logger.Warn("Error loading .env file - Using environment variables instead")
|
||||
}
|
||||
|
||||
phoneNumber := os.Getenv("TG_PHONE_NUMBER")
|
||||
if phoneNumber == "" {
|
||||
logger.Warn("No phone number provided")
|
||||
}
|
||||
|
||||
botToken := os.Getenv("TG_BOT_TOKEN")
|
||||
if botToken == "" {
|
||||
logger.Fatal("No bot token provided")
|
||||
}
|
||||
|
||||
appId := os.Getenv("TG_APP_ID")
|
||||
if appId == "" {
|
||||
logger.Fatal("No APP ID provided")
|
||||
@@ -64,8 +54,6 @@ func Load(logger *zap.Logger) *Config {
|
||||
proxyURL := os.Getenv("TG_PROXY_URL")
|
||||
|
||||
return &Config{
|
||||
TgPhoneNumber: strings.TrimSpace(phoneNumber),
|
||||
TgBotToken: strings.TrimSpace(botToken),
|
||||
TgAppId: func() int { i, _ := strconv.Atoi(appId); return i }(),
|
||||
TgAppHash: strings.TrimSpace(appHash),
|
||||
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"tgss/internal/telegram"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
type AuthHandler struct {
|
||||
logger *zap.Logger
|
||||
tg *telegram.Service
|
||||
}
|
||||
|
||||
func NewAuthHandler(tg *telegram.Service, logger *zap.Logger) *AuthHandler {
|
||||
return &AuthHandler{tg: tg, logger: logger}
|
||||
}
|
||||
|
||||
type phoneReq struct {
|
||||
Phone string `json:"phone"`
|
||||
}
|
||||
|
||||
type codeReq struct {
|
||||
Code string `json:"code"`
|
||||
}
|
||||
|
||||
type passwordReq struct {
|
||||
Password string `json:"password"`
|
||||
}
|
||||
|
||||
func (h *AuthHandler) SendCode(c *gin.Context) {
|
||||
var req phoneReq
|
||||
if err := c.BindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := h.tg.SendCode(ctx, req.Phone); err != nil {
|
||||
c.JSON(500, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(200, gin.H{"status": "code sent"})
|
||||
}
|
||||
|
||||
func (h *AuthHandler) Verify(c *gin.Context) {
|
||||
var req codeReq
|
||||
if err := c.BindJSON(&req); err != nil {
|
||||
c.JSON(400, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := h.tg.VerifyCode(ctx, req.Code); err != nil {
|
||||
c.JSON(500, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(200, gin.H{"status": "authenticated"})
|
||||
}
|
||||
|
||||
func (h *AuthHandler) Password(c *gin.Context) {
|
||||
var req passwordReq
|
||||
if err := c.BindJSON(&req); err != nil {
|
||||
c.JSON(400, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := h.tg.Password(ctx, req.Password); err != nil {
|
||||
c.JSON(500, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(200, gin.H{"status": "authenticated"})
|
||||
}
|
||||
@@ -2,6 +2,7 @@ package handlers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strconv"
|
||||
"tgss/internal/telegram"
|
||||
"time"
|
||||
|
||||
@@ -25,12 +26,28 @@ func NewChannelHandler(
|
||||
}
|
||||
|
||||
func (ch *ChannelHandler) GetLatestMessages(c *gin.Context) {
|
||||
chennelId := c.Param("id")
|
||||
channelId := c.Param("id")
|
||||
limit := 5
|
||||
|
||||
if limitStr := c.Query("limit"); limitStr != "" {
|
||||
var err error
|
||||
limit, err = strconv.Atoi(limitStr)
|
||||
if err != nil || limit > 50 {
|
||||
ch.logger.Warn("Invalid limit value provided", zap.String("limit", limitStr))
|
||||
c.JSON(400, gin.H{"error": "Limit value provided is invalid or bigger than 50"})
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
m, err := ch.tgService.LastMessages(ctx, chennelId, 10)
|
||||
|
||||
messages, err := ch.tgService.LastMessages(ctx, channelId, limit)
|
||||
if err != nil {
|
||||
ch.logger.Error("Error getting channel messages:", zap.Error(err))
|
||||
c.JSON(500, gin.H{"error": "Unable to get channel messages"})
|
||||
return
|
||||
}
|
||||
c.JSON(200, m)
|
||||
|
||||
c.JSON(200, messages)
|
||||
}
|
||||
|
||||
@@ -18,6 +18,7 @@ type RouterParams struct {
|
||||
Logger *zap.Logger
|
||||
|
||||
ChannelHandler *handlers.ChannelHandler
|
||||
AuthHandler *handlers.AuthHandler
|
||||
}
|
||||
|
||||
func NewGin() *gin.Engine {
|
||||
@@ -42,6 +43,10 @@ func RegisterRoutes(params RouterParams) {
|
||||
|
||||
params.GinEngine.GET("/channel/:id", params.ChannelHandler.GetLatestMessages)
|
||||
|
||||
params.GinEngine.POST("/auth/send-code", params.AuthHandler.SendCode)
|
||||
params.GinEngine.POST("/auth/verify", params.AuthHandler.Verify)
|
||||
params.GinEngine.POST("/auth/password", params.AuthHandler.Password)
|
||||
|
||||
params.Logger.Info("Starting server")
|
||||
go params.GinEngine.Run(":3000")
|
||||
params.Logger.Info("Server is running on port 8080")
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"tgss/internal/config"
|
||||
"time"
|
||||
|
||||
"github.com/gotd/contrib/bg"
|
||||
"github.com/gotd/td/session"
|
||||
"github.com/gotd/td/telegram"
|
||||
"github.com/gotd/td/telegram/dcs"
|
||||
@@ -17,15 +18,18 @@ import (
|
||||
|
||||
func NewTelegramClient(cfg *config.Config, logger *zap.Logger) *telegram.Client {
|
||||
opts := telegram.Options{
|
||||
DC: 2,
|
||||
DCList: dcs.Prod(),
|
||||
Logger: logger,
|
||||
SessionStorage: &session.FileStorage{Path: cfg.SessionPath},
|
||||
DialTimeout: 20*time.Second,
|
||||
ExchangeTimeout: 20*time.Second,
|
||||
MigrationTimeout: 20*time.Second,
|
||||
DC: 2,
|
||||
DCList: dcs.Prod(),
|
||||
Logger: logger,
|
||||
SessionStorage: &session.FileStorage{Path: cfg.SessionPath},
|
||||
DialTimeout: 20 * time.Second,
|
||||
ExchangeTimeout: 20 * time.Second,
|
||||
MigrationTimeout: 20 * time.Second,
|
||||
}
|
||||
|
||||
// TODO: Add pebble and bbolt https://github.com/gotd/td/blob/6f8e63c553210a2901f5c3586f6b88d6524fe9b3/examples/userbot/main.go#L106
|
||||
// TODO: Add ratelimit https://github.com/gotd/td/blob/6f8e63c553210a2901f5c3586f6b88d6524fe9b3/examples/userbot/main.go#L149
|
||||
// TODO: Cleanup contexts over the codebase and organize timeout values
|
||||
if cfg.ProxyURL != "" {
|
||||
proxyURL, err := url.Parse(cfg.ProxyURL)
|
||||
if err != nil {
|
||||
@@ -48,42 +52,25 @@ func NewTelegramClient(cfg *config.Config, logger *zap.Logger) *telegram.Client
|
||||
return telegram.NewClient(cfg.TgAppId, cfg.TgAppHash, opts)
|
||||
}
|
||||
|
||||
func NewService(client *telegram.Client, logger *zap.Logger) *Service {
|
||||
return &Service{client: client, log: logger}
|
||||
}
|
||||
func RunClient(lc fx.Lifecycle, client *telegram.Client, logger *zap.Logger) {
|
||||
var stop func() error
|
||||
|
||||
func RunClient(lc fx.Lifecycle, client *telegram.Client, cfg *config.Config, logger *zap.Logger) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(context.Context) error {
|
||||
errCh := make(chan error, 1)
|
||||
readyCh := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
if err := client.Run(ctx, func(ctx context.Context) error {
|
||||
if _, err := client.Auth().Bot(ctx, cfg.TgBotToken); err != nil {
|
||||
return err
|
||||
}
|
||||
logger.Info("Telegram bot authenticated and running")
|
||||
close(readyCh)
|
||||
<-ctx.Done()
|
||||
return nil
|
||||
}); err != nil {
|
||||
errCh <- err
|
||||
}
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-readyCh:
|
||||
return nil
|
||||
case err := <-errCh:
|
||||
OnStart: func(ctx context.Context) error {
|
||||
s, err := bg.Connect(client)
|
||||
if err != nil {
|
||||
return err
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
stop = s
|
||||
logger.Info("telegram client connected")
|
||||
return nil
|
||||
},
|
||||
OnStop: func(context.Context) error {
|
||||
cancel()
|
||||
OnStop: func(ctx context.Context) error {
|
||||
if stop != nil {
|
||||
logger.Info("telegram client stopping")
|
||||
return stop()
|
||||
}
|
||||
return nil
|
||||
},
|
||||
})
|
||||
|
||||
@@ -3,8 +3,10 @@ package telegram
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
|
||||
"github.com/gotd/td/telegram"
|
||||
"github.com/gotd/td/telegram/auth"
|
||||
"github.com/gotd/td/tg"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
@@ -12,37 +14,91 @@ import (
|
||||
type Service struct {
|
||||
client *telegram.Client
|
||||
log *zap.Logger
|
||||
|
||||
mu sync.Mutex
|
||||
phone string
|
||||
phoneCodeHash string
|
||||
}
|
||||
|
||||
func NewService(client *telegram.Client, logger *zap.Logger) *Service {
|
||||
return &Service{
|
||||
client: client,
|
||||
log: logger,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Service) AuthStatus(ctx context.Context) (bool, error) {
|
||||
_, err := s.client.Auth().Status(ctx)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (s *Service) SendCode(ctx context.Context, phone string) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
resp, err := s.client.Auth().SendCode(ctx, phone, auth.SendCodeOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s.phone = phone
|
||||
s.phoneCodeHash = resp.(*tg.AuthSentCode).PhoneCodeHash
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) VerifyCode(ctx context.Context, code string) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
if s.phone == "" {
|
||||
return errors.New("phone not initialized")
|
||||
}
|
||||
|
||||
_, err := s.client.Auth().SignIn(ctx, s.phone, code, s.phoneCodeHash)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Service) Password(ctx context.Context, password string) error {
|
||||
_, err := s.client.Auth().Password(ctx, password)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Service) LastMessages(ctx context.Context, username string, limit int) ([]tg.MessageClass, error) {
|
||||
api := s.client.API()
|
||||
|
||||
resolved, err := api.ContactsResolveUsername(ctx, &tg.ContactsResolveUsernameRequest{
|
||||
Username: username,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if len(resolved.Chats) == 0 {
|
||||
return nil, errors.New("channel not found")
|
||||
}
|
||||
|
||||
channel, ok := resolved.Chats[0].(*tg.Channel)
|
||||
if !ok {
|
||||
return nil, errors.New("resolved peer is not a channel")
|
||||
}
|
||||
|
||||
history, err := s.client.API().MessagesGetHistory(ctx, &tg.MessagesGetHistoryRequest{
|
||||
history, err := api.MessagesGetHistory(ctx, &tg.MessagesGetHistoryRequest{
|
||||
Peer: &tg.InputPeerChannel{
|
||||
ChannelID: channel.ID,
|
||||
AccessHash: channel.AccessHash,
|
||||
},
|
||||
Limit: 100,
|
||||
Limit: limit,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Extract messages
|
||||
var msgs []tg.MessageClass
|
||||
|
||||
switch h := history.(type) {
|
||||
case *tg.MessagesMessages:
|
||||
msgs = h.Messages
|
||||
@@ -51,5 +107,6 @@ func (s *Service) LastMessages(ctx context.Context, username string, limit int)
|
||||
case *tg.MessagesChannelMessages:
|
||||
msgs = h.Messages
|
||||
}
|
||||
|
||||
return msgs, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user