REPOSITORY / ScuroNeko/Laniakea

Compare commits

DIFF REPOSITORY

Compare commits

...
12 Commits
Author SHA1 Message Date
ScuroNeko ae7426c36a 1.0.0 beta 4 2026-03-01 23:08:22 +03:00
ScuroNeko 61562e8a3b 1.0.0 beta 3 2026-03-01 23:01:06 +03:00
ScuroNeko a84e24ff25 small fix 2026-02-27 13:53:00 +03:00
ScuroNeko c0a26024f4 v1.0.0 beta 2 2026-02-26 15:15:35 +03:00
ScuroNeko 786da652e6 v1.0.0 beta 1 2026-02-26 15:12:36 +03:00
ScuroNeko 28ec2b7ca9 0.8.0 beta 4 2026-02-26 14:31:03 +03:00
ScuroNeko da122a3be4 0.8.0 beta 3 2026-02-19 13:58:34 +03:00
ScuroNeko 1bf7499496 0.8.0 beta 2 2026-02-19 13:33:27 +03:00
ScuroNeko 7b9292557e 0.8.0 beta 1 2026-02-19 13:27:03 +03:00
ScuroNeko 466093e39b version fix 2026-02-19 12:07:03 +03:00
ScuroNeko 0e0f8a0813 v0.7.0; support for test server and local bot api 2026-02-19 11:49:04 +03:00
ScuroNeko d84b0a1b55 small fixes 2026-02-18 14:05:36 +03:00
21 changed files with 846 additions and 519 deletions
+215 -152
View File
@@ -1,43 +1,65 @@
package laniakea package laniakea
import ( import (
"context"
"fmt" "fmt"
"os" "os"
"sort" "sort"
"strconv"
"strings" "strings"
"time" "sync"
"git.nix13.pw/scuroneko/extypes" "git.nix13.pw/scuroneko/extypes"
"git.nix13.pw/scuroneko/laniakea/tgapi" "git.nix13.pw/scuroneko/laniakea/tgapi"
"git.nix13.pw/scuroneko/slog" "git.nix13.pw/scuroneko/slog"
"github.com/redis/go-redis/v9" "github.com/alitto/pond/v2"
"github.com/vinovest/sqlx" "golang.org/x/time/rate"
"go.mongodb.org/mongo-driver/v2/mongo"
) )
type BotSettings struct { type BotOpts struct {
Token string Token string
UpdateTypes []string
Debug bool Debug bool
ErrorTemplate string ErrorTemplate string
Prefixes []string Prefixes []string
UpdateTypes []string
LoggerBasePath string LoggerBasePath string
UseRequestLogger bool UseRequestLogger bool
WriteToFile bool WriteToFile bool
UseTestServer bool
APIUrl string
RateLimit int
DropRLOverflow bool
} }
func LoadSettingsFromEnv() *BotSettings { func NewOpts() *BotOpts { return new(BotOpts) }
return &BotSettings{ func LoadOptsFromEnv() *BotOpts {
rateLimit := 30
if rl := os.Getenv("RATE_LIMIT"); rl != "" {
rateLimit, _ = strconv.Atoi(rl)
}
return &BotOpts{
Token: os.Getenv("TG_TOKEN"), Token: os.Getenv("TG_TOKEN"),
UpdateTypes: strings.Split(os.Getenv("UPDATE_TYPES"), ";"),
Debug: os.Getenv("DEBUG") == "true", Debug: os.Getenv("DEBUG") == "true",
ErrorTemplate: os.Getenv("ERROR_TEMPLATE"), ErrorTemplate: os.Getenv("ERROR_TEMPLATE"),
Prefixes: LoadPrefixesFromEnv(), Prefixes: LoadPrefixesFromEnv(),
UpdateTypes: strings.Split(os.Getenv("UPDATE_TYPES"), ";"),
UseRequestLogger: os.Getenv("USE_REQ_LOG") == "true", UseRequestLogger: os.Getenv("USE_REQ_LOG") == "true",
WriteToFile: os.Getenv("WRITE_TO_FILE") == "true", WriteToFile: os.Getenv("WRITE_TO_FILE") == "true",
UseTestServer: os.Getenv("USE_TEST_SERVER") == "true",
APIUrl: os.Getenv("API_URL"),
RateLimit: rateLimit,
DropRLOverflow: os.Getenv("DROP_RL_OVERFLOW") == "true",
} }
} }
func LoadPrefixesFromEnv() []string { func LoadPrefixesFromEnv() []string {
prefixesS, exists := os.LookupEnv("PREFIXES") prefixesS, exists := os.LookupEnv("PREFIXES")
if !exists { if !exists {
@@ -46,57 +68,107 @@ func LoadPrefixesFromEnv() []string {
return strings.Split(prefixesS, ";") return strings.Split(prefixesS, ";")
} }
type Bot struct { type DbContext interface{}
type NoDB struct{ DbContext }
type Bot[T DbContext] struct {
token string token string
debug bool debug bool
errorTemplate string errorTemplate string
logger *slog.Logger logger *slog.Logger
RequestLogger *slog.Logger RequestLogger *slog.Logger
extraLoggers extypes.Slice[*slog.Logger]
plugins []Plugin plugins []Plugin[T]
middlewares []Middleware middlewares []Middleware[T]
prefixes []string prefixes []string
runners []Runner runners []Runner[T]
dbContext *DatabaseContext
api *tgapi.API api *tgapi.API
l10n L10n uploader *tgapi.Uploader
dbContext *T
dbWriterRequested extypes.Slice[*slog.Logger] l10n *L10n
draftProvider *DraftProvider
updateOffsetMu sync.Mutex
updateOffset int updateOffset int
updateTypes []tgapi.UpdateType updateTypes []tgapi.UpdateType
updateQueue *extypes.Queue[*tgapi.Update] updateQueue chan *tgapi.Update
} }
func NewBot(settings *BotSettings) *Bot { func NewBot[T any](opts *BotOpts) *Bot[T] {
updateQueue := extypes.CreateQueue[*tgapi.Update](256) updateQueue := make(chan *tgapi.Update, 512)
api := tgapi.NewAPI(settings.Token)
bot := &Bot{
updateOffset: 0, plugins: make([]Plugin, 0), debug: settings.Debug, errorTemplate: "%s",
prefixes: settings.Prefixes, updateTypes: make([]tgapi.UpdateType, 0), runners: make([]Runner, 0),
updateQueue: updateQueue, api: api, dbWriterRequested: make([]*slog.Logger, 0),
token: settings.Token, l10n: L10n{},
}
bot.dbWriterRequested = bot.dbWriterRequested.Push(api.Logger)
if len(settings.ErrorTemplate) > 0 { var limiter *rate.Limiter
bot.errorTemplate = settings.ErrorTemplate if opts.RateLimit > 0 {
} limiter = rate.NewLimiter(rate.Limit(opts.RateLimit), opts.RateLimit)
if len(settings.LoggerBasePath) == 0 {
settings.LoggerBasePath = "./"
} }
apiOpts := tgapi.NewAPIOpts(opts.Token).SetAPIUrl(opts.APIUrl).UseTestServer(opts.UseTestServer).SetLimiter(limiter)
api := tgapi.NewAPI(apiOpts)
uploader := tgapi.NewUploader(api)
bot := &Bot[T]{
updateOffset: 0,
errorTemplate: "%s",
updateQueue: updateQueue,
api: api,
uploader: uploader,
debug: opts.Debug,
prefixes: opts.Prefixes,
token: opts.Token,
plugins: make([]Plugin[T], 0),
updateTypes: make([]tgapi.UpdateType, 0),
runners: make([]Runner[T], 0),
extraLoggers: make([]*slog.Logger, 0),
l10n: &L10n{},
draftProvider: NewRandomDraftProvider(api),
}
bot.extraLoggers = bot.extraLoggers.Push(api.GetLogger()).Push(uploader.GetLogger())
if len(opts.ErrorTemplate) > 0 {
bot.errorTemplate = opts.ErrorTemplate
}
if len(opts.LoggerBasePath) == 0 {
opts.LoggerBasePath = "./"
}
bot.initLoggers(opts)
u, err := api.GetMe()
if err != nil {
_ = bot.Close()
bot.logger.Fatal(err)
}
bot.logger.Infof("Authorized as %s\n", u.FirstName)
return bot
}
func (bot *Bot[T]) Close() error {
if err := bot.uploader.Close(); err != nil {
bot.logger.Errorln(err)
}
if err := bot.api.CloseApi(); err != nil {
bot.logger.Errorln(err)
}
if err := bot.RequestLogger.Close(); err != nil {
bot.logger.Errorln(err)
}
if err := bot.logger.Close(); err != nil {
return err
}
return nil
}
func (bot *Bot[T]) initLoggers(opts *BotOpts) {
level := slog.FATAL level := slog.FATAL
if settings.Debug { if opts.Debug {
level = slog.DEBUG level = slog.DEBUG
} }
bot.logger = slog.CreateLogger().Level(level).Prefix("BOT") bot.logger = slog.CreateLogger().Level(level).Prefix("BOT")
bot.logger.AddWriter(bot.logger.CreateJsonStdoutWriter()) bot.logger.AddWriter(bot.logger.CreateJsonStdoutWriter())
if settings.WriteToFile { if opts.WriteToFile {
path := fmt.Sprintf("%s/main.log", strings.TrimRight(settings.LoggerBasePath, "/")) path := fmt.Sprintf("%s/main.log", strings.TrimRight(opts.LoggerBasePath, "/"))
fileWriter, err := bot.logger.CreateTextFileWriter(path) fileWriter, err := bot.logger.CreateTextFileWriter(path)
if err != nil { if err != nil {
bot.logger.Fatal(err) bot.logger.Fatal(err)
@@ -104,11 +176,11 @@ func NewBot(settings *BotSettings) *Bot {
bot.logger.AddWriter(fileWriter) bot.logger.AddWriter(fileWriter)
} }
if settings.UseRequestLogger { if opts.UseRequestLogger {
bot.RequestLogger = slog.CreateLogger().Level(level).Prefix("REQUESTS") bot.RequestLogger = slog.CreateLogger().Level(level).Prefix("REQUESTS")
bot.RequestLogger.AddWriter(bot.RequestLogger.CreateJsonStdoutWriter()) bot.RequestLogger.AddWriter(bot.RequestLogger.CreateJsonStdoutWriter())
if settings.WriteToFile { if opts.WriteToFile {
path := fmt.Sprintf("%s/requests.log", strings.TrimRight(settings.LoggerBasePath, "/")) path := fmt.Sprintf("%s/requests.log", strings.TrimRight(opts.LoggerBasePath, "/"))
fileWriter, err := bot.RequestLogger.CreateTextFileWriter(path) fileWriter, err := bot.RequestLogger.CreateTextFileWriter(path)
if err != nil { if err != nil {
bot.logger.Fatal(err) bot.logger.Fatal(err)
@@ -116,162 +188,153 @@ func NewBot(settings *BotSettings) *Bot {
bot.RequestLogger.AddWriter(fileWriter) bot.RequestLogger.AddWriter(fileWriter)
} }
} }
}
u, err := api.GetMe() func (bot *Bot[T]) GetUpdateOffset() int {
if err != nil { bot.updateOffsetMu.Lock()
bot.logger.Fatal(err) defer bot.updateOffsetMu.Unlock()
} return bot.updateOffset
bot.logger.Infof("Authorized as %s\n", u.FirstName) }
func (bot *Bot[T]) SetUpdateOffset(offset int) {
bot.updateOffsetMu.Lock()
defer bot.updateOffsetMu.Unlock()
bot.updateOffset = offset
}
func (bot *Bot[T]) GetUpdateTypes() []tgapi.UpdateType { return bot.updateTypes }
func (bot *Bot[T]) GetLogger() *slog.Logger { return bot.logger }
func (bot *Bot[T]) GetDBContext() *T { return bot.dbContext }
func (bot *Bot[T]) L10n(lang, key string) string { return bot.l10n.Translate(lang, key) }
func (bot *Bot[T]) SetDraftProvider(p *DraftProvider) *Bot[T] {
bot.draftProvider = p
return bot return bot
} }
func (b *Bot) Close() error { type DbLogger[T DbContext] func(db *T) slog.LoggerWriter
err := b.logger.Close()
if err != nil { func (bot *Bot[T]) AddDatabaseLoggerWriter(writer DbLogger[T]) *Bot[T] {
return err w := writer(bot.dbContext)
bot.logger.AddWriter(w)
if bot.RequestLogger != nil {
bot.RequestLogger.AddWriter(w)
} }
err = b.RequestLogger.Close() for _, l := range bot.extraLoggers {
return err
}
func (b *Bot) GetUpdateOffset() int { return b.updateOffset }
func (b *Bot) SetUpdateOffset(offset int) { b.updateOffset = offset }
func (b *Bot) GetUpdateTypes() []tgapi.UpdateType { return b.updateTypes }
func (b *Bot) GetQueue() *extypes.Queue[*tgapi.Update] { return b.updateQueue }
type DatabaseContext struct {
PostgresSQL *sqlx.DB
MongoDB *mongo.Client
Redis *redis.Client
}
func (b *Bot) AddDatabaseLogger(writer func(db *DatabaseContext) slog.LoggerWriter) *Bot {
w := writer(b.dbContext)
b.logger.AddWriter(w)
if b.RequestLogger != nil {
b.RequestLogger.AddWriter(w)
}
for _, l := range b.dbWriterRequested {
l.AddWriter(w) l.AddWriter(w)
} }
return b return bot
} }
func (b *Bot) DatabaseContext(ctx *DatabaseContext) *Bot { func (bot *Bot[T]) DatabaseContext(ctx *T) *Bot[T] {
b.dbContext = ctx bot.dbContext = ctx
return b return bot
} }
func (b *Bot) UpdateTypes(t ...tgapi.UpdateType) *Bot { func (bot *Bot[T]) UpdateTypes(t ...tgapi.UpdateType) *Bot[T] {
b.updateTypes = make([]tgapi.UpdateType, 0) bot.updateTypes = make([]tgapi.UpdateType, 0)
b.updateTypes = append(b.updateTypes, t...) bot.updateTypes = append(bot.updateTypes, t...)
return b return bot
} }
func (b *Bot) AddUpdateType(t ...tgapi.UpdateType) *Bot { func (bot *Bot[T]) AddUpdateType(t ...tgapi.UpdateType) *Bot[T] {
b.updateTypes = append(b.updateTypes, t...) bot.updateTypes = append(bot.updateTypes, t...)
return b return bot
} }
func (b *Bot) AddPrefixes(prefixes ...string) *Bot { func (bot *Bot[T]) AddPrefixes(prefixes ...string) *Bot[T] {
b.prefixes = append(b.prefixes, prefixes...) bot.prefixes = append(bot.prefixes, prefixes...)
return b return bot
} }
func (b *Bot) ErrorTemplate(s string) *Bot { func (bot *Bot[T]) ErrorTemplate(s string) *Bot[T] {
b.errorTemplate = s bot.errorTemplate = s
return b return bot
} }
func (b *Bot) Debug(debug bool) *Bot { func (bot *Bot[T]) Debug(debug bool) *Bot[T] {
b.debug = debug bot.debug = debug
return b return bot
} }
func (b *Bot) AddPlugins(plugin ...*Plugin) *Bot { func (bot *Bot[T]) AddPlugins(plugin ...*Plugin[T]) *Bot[T] {
for _, p := range plugin { for _, p := range plugin {
b.plugins = append(b.plugins, *p) bot.plugins = append(bot.plugins, *p)
b.logger.Debugln(fmt.Sprintf("plugins with name \"%s\" registered", p.Name)) bot.logger.Debugln(fmt.Sprintf("plugins with name \"%s\" registered", p.name))
} }
return b return bot
} }
func (b *Bot) AddMiddleware(middleware ...Middleware) *Bot { func (bot *Bot[T]) AddMiddleware(middleware ...Middleware[T]) *Bot[T] {
b.middlewares = append(b.middlewares, middleware...) bot.middlewares = append(bot.middlewares, middleware...)
for _, m := range middleware { for _, m := range middleware {
b.logger.Debugln(fmt.Sprintf("middleware with name \"%s\" registered", m.name)) bot.logger.Debugln(fmt.Sprintf("middleware with name \"%s\" registered", m.name))
} }
sort.Slice(b.middlewares, func(i, j int) bool { sort.Slice(bot.middlewares, func(i, j int) bool {
first := b.middlewares[i] first := bot.middlewares[i]
second := b.middlewares[j] second := bot.middlewares[j]
if first.order == second.order { if first.order == second.order {
return first.name < second.name return first.name < second.name
} }
return first.order < second.order return first.order < second.order
}) })
return b return bot
} }
func (b *Bot) AddRunner(runner Runner) *Bot { func (bot *Bot[T]) AddRunner(runner Runner[T]) *Bot[T] {
b.runners = append(b.runners, runner) bot.runners = append(bot.runners, runner)
b.logger.Debugln(fmt.Sprintf("runner with name \"%s\" registered", runner.Name)) bot.logger.Debugln(fmt.Sprintf("runner with name \"%s\" registered", runner.name))
return b return bot
} }
func (b *Bot) AddL10n(l L10n) *Bot { func (bot *Bot[T]) AddL10n(l *L10n) *Bot[T] {
b.l10n = l bot.l10n = l
return b return bot
}
func (b *Bot) L10n(lang, key string) string {
return b.l10n.Translate(lang, key)
}
func (b *Bot) Logger() *slog.Logger {
return b.logger
}
func (b *Bot) GetDBContext() *DatabaseContext {
return b.dbContext
} }
func (b *Bot) Run() { func (bot *Bot[T]) enqueueUpdate(u *tgapi.Update) error {
if len(b.prefixes) == 0 { select {
b.logger.Fatalln("no prefixes defined") case bot.updateQueue <- u:
return nil
default:
return extypes.QueueFullErr
}
}
func (bot *Bot[T]) RunWithContext(ctx context.Context) {
if len(bot.prefixes) == 0 {
bot.logger.Fatalln("no prefixes defined")
return return
} }
if len(b.plugins) == 0 { if len(bot.plugins) == 0 {
b.logger.Fatalln("no plugins defined") bot.logger.Fatalln("no plugins defined")
return return
} }
b.logger.Infoln("Executing runners...") bot.ExecRunners()
b.ExecRunners()
b.logger.Infoln("Bot running. Press CTRL+C to exit.") bot.logger.Infoln("Bot running. Press CTRL+C to exit.")
go func() { go func() {
for { for {
_, err := b.Updates() select {
case <-ctx.Done():
return
default:
updates, err := bot.Updates()
if err != nil { if err != nil {
b.logger.Errorln(err) bot.logger.Errorln(err)
continue
}
for _, u := range updates {
select {
case bot.updateQueue <- new(u):
case <-ctx.Done():
return
}
}
} }
} }
}() }()
for { pool := pond.NewPool(16)
queue := b.updateQueue for update := range bot.updateQueue {
if queue.IsEmpty() { update := update
time.Sleep(time.Millisecond * 25) pool.Submit(func() {
continue bot.handle(update)
} })
u := queue.Dequeue()
if u == nil {
b.logger.Errorln("update is nil")
continue
}
ctx := &MsgContext{Bot: b, Update: *u, Api: b.api}
for _, middleware := range b.middlewares {
middleware.Execute(ctx, b.dbContext)
}
if u.CallbackQuery != nil {
b.handleCallback(u, ctx)
} else {
b.handleMessage(u, ctx)
}
} }
} }
func (bot *Bot[T]) Run() {
bot.RunWithContext(context.Background())
}
+19 -9
View File
@@ -1,13 +1,14 @@
package laniakea package laniakea
import ( import (
"errors"
"fmt" "fmt"
"strings" "strings"
"git.nix13.pw/scuroneko/laniakea/tgapi" "git.nix13.pw/scuroneko/laniakea/tgapi"
) )
func generateBotCommand(cmd Command) tgapi.BotCommand { func generateBotCommand[T any](cmd Command[T]) tgapi.BotCommand {
desc := cmd.command desc := cmd.command
if len(cmd.description) > 0 { if len(cmd.description) > 0 {
desc = cmd.description desc = cmd.description
@@ -24,9 +25,9 @@ func generateBotCommand(cmd Command) tgapi.BotCommand {
return tgapi.BotCommand{Command: cmd.command, Description: desc} return tgapi.BotCommand{Command: cmd.command, Description: desc}
} }
func generateBotCommandForPlugin(pl Plugin) []tgapi.BotCommand { func generateBotCommandForPlugin[T any](pl Plugin[T]) []tgapi.BotCommand {
commands := make([]tgapi.BotCommand, 0) commands := make([]tgapi.BotCommand, 0)
for _, cmd := range pl.Commands { for _, cmd := range pl.commands {
if cmd.skipAutoCmd { if cmd.skipAutoCmd {
continue continue
} }
@@ -35,28 +36,37 @@ func generateBotCommandForPlugin(pl Plugin) []tgapi.BotCommand {
return commands return commands
} }
func (b *Bot) AutoGenerateCommands() error { var ErrTooManyCommands = errors.New("too many commands. max 100")
_, err := b.api.DeleteMyCommands(tgapi.DeleteMyCommandsP{})
func (bot *Bot[T]) AutoGenerateCommands() error {
_, err := bot.api.DeleteMyCommands(tgapi.DeleteMyCommandsP{})
if err != nil { if err != nil {
return err return err
} }
commands := make([]tgapi.BotCommand, 0) commands := make([]tgapi.BotCommand, 0)
for _, pl := range b.plugins { for _, pl := range bot.plugins {
if pl.skipAutoCmd {
continue
}
commands = append(commands, generateBotCommandForPlugin(pl)...) commands = append(commands, generateBotCommandForPlugin(pl)...)
} }
if len(commands) > 100 {
return ErrTooManyCommands
}
privateChatsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopePrivateType} privateChatsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopePrivateType}
groupChatsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopeGroupType} groupChatsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopeGroupType}
chatAdminsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopeAllChatAdministratorsType} chatAdminsScope := &tgapi.BotCommandScope{Type: tgapi.BotCommandScopeAllChatAdministratorsType}
_, err = b.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: privateChatsScope}) _, err = bot.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: privateChatsScope})
if err != nil { if err != nil {
return err return err
} }
_, err = b.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: groupChatsScope}) _, err = bot.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: groupChatsScope})
if err != nil { if err != nil {
return err return err
} }
_, err = b.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: chatAdminsScope}) _, err = bot.api.SetMyCommands(tgapi.SetMyCommandsP{Commands: commands, Scope: chatAdminsScope})
return err return err
} }
+89
View File
@@ -0,0 +1,89 @@
package laniakea
import (
"math"
"math/rand/v2"
"sync/atomic"
"git.nix13.pw/scuroneko/laniakea/tgapi"
)
type draftIdGenerator interface {
Next() uint64
}
type RandomDraftIdGenerator struct {
draftIdGenerator
}
func (g *RandomDraftIdGenerator) Next() uint64 {
return rand.Uint64N(math.MaxUint64)
}
type LinearDraftIdGenerator struct {
draftIdGenerator
lastId uint64
}
func (g *LinearDraftIdGenerator) Next() uint64 {
return atomic.AddUint64(&g.lastId, 1)
}
type DraftProvider struct {
api *tgapi.API
chatID int
messageThreadID int
parseMode tgapi.ParseMode
entities []tgapi.MessageEntity
drafts map[uint64]*Draft
generator draftIdGenerator
}
type Draft struct {
api *tgapi.API
chatID int
messageThreadID int
parseMode tgapi.ParseMode
entities []tgapi.MessageEntity
ID uint64
Message string
}
func NewRandomDraftProvider(api *tgapi.API) *DraftProvider {
return &DraftProvider{
api: api, generator: &RandomDraftIdGenerator{},
drafts: make(map[uint64]*Draft),
}
}
func NewLinearDraftProvider(api *tgapi.API, startValue uint64) *DraftProvider {
return &DraftProvider{
api: api,
generator: &LinearDraftIdGenerator{lastId: startValue},
drafts: make(map[uint64]*Draft),
}
}
func (d *DraftProvider) NewDraft() *Draft {
id := d.generator.Next()
draft := &Draft{d.api, d.chatID, d.messageThreadID, d.parseMode, d.entities, id, ""}
d.drafts[id] = draft
return draft
}
func (d *Draft) Push(newText string) error {
d.Message += newText
params := tgapi.SendMessageDraftP{
ChatID: d.chatID,
DraftID: d.ID,
Text: d.Message,
ParseMode: d.parseMode,
Entities: d.entities,
}
if d.messageThreadID > 0 {
params.MessageThreadID = d.messageThreadID
}
_, err := d.api.SendMessageDraft(params)
return err
}
+3 -17
View File
@@ -3,29 +3,15 @@ module git.nix13.pw/scuroneko/laniakea
go 1.26 go 1.26
require ( require (
git.nix13.pw/scuroneko/extypes v1.2.0 git.nix13.pw/scuroneko/extypes v1.2.1
git.nix13.pw/scuroneko/slog v1.0.2 git.nix13.pw/scuroneko/slog v1.0.2
github.com/redis/go-redis/v9 v9.18.0 github.com/alitto/pond/v2 v2.6.2
github.com/vinovest/sqlx v1.7.1 golang.org/x/time v0.14.0
go.mongodb.org/mongo-driver/v2 v2.5.0
) )
require ( require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
github.com/fatih/color v1.18.0 // indirect github.com/fatih/color v1.18.0 // indirect
github.com/klauspost/compress v1.18.4 // indirect
github.com/mattn/go-colorable v0.1.14 // indirect github.com/mattn/go-colorable v0.1.14 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect github.com/mattn/go-isatty v0.0.20 // indirect
github.com/muir/list v1.2.1 // indirect
github.com/muir/sqltoken v0.3.0 // indirect
github.com/xdg-go/pbkdf2 v1.0.0 // indirect
github.com/xdg-go/scram v1.2.0 // indirect
github.com/xdg-go/stringprep v1.0.4 // indirect
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect
go.uber.org/atomic v1.11.0 // indirect
golang.org/x/crypto v0.48.0 // indirect
golang.org/x/sync v0.19.0 // indirect
golang.org/x/sys v0.41.0 // indirect golang.org/x/sys v0.41.0 // indirect
golang.org/x/text v0.34.0 // indirect
) )
+6 -84
View File
@@ -1,95 +1,17 @@
filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= git.nix13.pw/scuroneko/extypes v1.2.1 h1:IYrOjnWKL2EAuJYtYNa+luB1vBe6paE8VY/YD+5/RpQ=
filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= git.nix13.pw/scuroneko/extypes v1.2.1/go.mod h1:uZVs8Yo3RrYAG9dMad6qR6lsYY67t+459D9c65QAYAw=
git.nix13.pw/scuroneko/extypes v1.2.0 h1:2n2hD6KsMAted+6MGhAyeWyli2Qzc9G2y+pQNB7C1dM=
git.nix13.pw/scuroneko/extypes v1.2.0/go.mod h1:uZVs8Yo3RrYAG9dMad6qR6lsYY67t+459D9c65QAYAw=
git.nix13.pw/scuroneko/slog v1.0.2 h1:vZyUROygxC2d5FJHUQM/30xFEHY1JT/aweDZXA4rm2g= git.nix13.pw/scuroneko/slog v1.0.2 h1:vZyUROygxC2d5FJHUQM/30xFEHY1JT/aweDZXA4rm2g=
git.nix13.pw/scuroneko/slog v1.0.2/go.mod h1:3Qm2wzkR5KjwOponMfG7TcGSDjmYaFqRAmLvSPTuWJI= git.nix13.pw/scuroneko/slog v1.0.2/go.mod h1:3Qm2wzkR5KjwOponMfG7TcGSDjmYaFqRAmLvSPTuWJI=
github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/alitto/pond/v2 v2.6.2 h1:Sphe40g0ILeM1pA2c2K+Th0DGU+pt0A/Kprr+WB24Pw=
github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= github.com/alitto/pond/v2 v2.6.2/go.mod h1:xkjYEgQ05RSpWdfSd1nM3OVv7TBhLdy7rMp3+2Nq+yE=
github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA=
github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78=
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc=
github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM= github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM=
github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU=
github.com/go-sql-driver/mysql v1.9.0 h1:Y0zIbQXhQKmQgTp44Y1dp3wTXcn804QoTptLZT1vtvo=
github.com/go-sql-driver/mysql v1.9.0/go.mod h1:pDetrLJeA3oMujJuvXc8RJoasr589B6A9fwzD3QMrqw=
github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI=
github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY=
github.com/klauspost/compress v1.18.4 h1:RPhnKRAQ4Fh8zU2FY/6ZFDwTVTxgJ/EMydqSTzE9a2c=
github.com/klauspost/compress v1.18.4/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4=
github.com/klauspost/cpuid/v2 v2.0.9 h1:lgaqFMSdTdQYdZ04uHyN2d/eKdOMyi2YLSvlQIBFYa4=
github.com/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg=
github.com/lib/pq v1.10.9 h1:YXG7RB+JIjhP29X+OtkiDnYaXQwpS4JEWq7dtCCRUEw=
github.com/lib/pq v1.10.9/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE=
github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
github.com/mattn/go-sqlite3 v1.14.16 h1:yOQRA0RpS5PFz/oikGwBEqvAWhWg5ufRz4ETLjwpU1Y=
github.com/mattn/go-sqlite3 v1.14.16/go.mod h1:2eHXhiwb8IkHr+BDWZGa96P6+rkvnG63S2DGjv9HUNg=
github.com/muir/list v1.2.1 h1:lmF8fz2B1WbXkzHr/Eh0oWPJArDBzWqIifOwbA4gWSo=
github.com/muir/list v1.2.1/go.mod h1:v0l2f997MxCohQlD7PTejJqyYKwFVz/i3mTpDl4LAf0=
github.com/muir/sqltoken v0.3.0 h1:3xbcqr80f3IA4OlwkOpdIHC4DTu6gsi1TwMqgYL4Dpg=
github.com/muir/sqltoken v0.3.0/go.mod h1:+OSmbGI22QcVZ6DCzlHT8EAzEq/mqtqedtPP91Le+3A=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/redis/go-redis/v9 v9.18.0 h1:pMkxYPkEbMPwRdenAzUNyFNrDgHx9U+DrBabWNfSRQs=
github.com/redis/go-redis/v9 v9.18.0/go.mod h1:k3ufPphLU5YXwNTUcCRXGxUoF1fqxnhFQmscfkCoDA0=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/vinovest/sqlx v1.7.1 h1:kdq4v0N9kRLpytWGSWOw4aulOGdQPmIoMR6Y+cTBxow=
github.com/vinovest/sqlx v1.7.1/go.mod h1:3fAv74r4iDMv2PpFomADb+vex5ukzfYn4GseC9KngD8=
github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c=
github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI=
github.com/xdg-go/scram v1.2.0 h1:bYKF2AEwG5rqd1BumT4gAnvwU/M9nBp2pTSxeZw7Wvs=
github.com/xdg-go/scram v1.2.0/go.mod h1:3dlrS0iBaWKYVt2ZfA4cj48umJZ+cAEbR6/SjLA88I8=
github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8=
github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM=
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 h1:ilQV1hzziu+LLM3zUTJ0trRztfwgjqKnBWNtSRkbmwM=
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78/go.mod h1:aL8wCCfTfSfmXjznFBSZNN13rSJjlIOI1fUNAtF7rmI=
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
go.mongodb.org/mongo-driver/v2 v2.5.0 h1:yXUhImUjjAInNcpTcAlPHiT7bIXhshCTL3jVBkF3xaE=
go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0=
go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE=
go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts=
golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos=
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4=
golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k=
golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI=
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ=
golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk=
golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+26 -20
View File
@@ -8,20 +8,26 @@ import (
"git.nix13.pw/scuroneko/laniakea/tgapi" "git.nix13.pw/scuroneko/laniakea/tgapi"
) )
func (b *Bot) handle(u *tgapi.Update) { func (bot *Bot[T]) handle(u *tgapi.Update) {
ctx := &MsgContext{Bot: b, Update: *u, Api: b.api} ctx := &MsgContext{
for _, middleware := range b.middlewares { Update: *u, Api: bot.api,
middleware.Execute(ctx, b.dbContext) botLogger: bot.logger,
errorTemplate: bot.errorTemplate,
l10n: bot.l10n,
draftProvider: bot.draftProvider,
}
for _, middleware := range bot.middlewares {
middleware.Execute(ctx, bot.dbContext)
} }
if u.CallbackQuery != nil { if u.CallbackQuery != nil {
b.handleCallback(u, ctx) bot.handleCallback(u, ctx)
} else { } else {
b.handleMessage(u, ctx) bot.handleMessage(u, ctx)
} }
} }
func (b *Bot) handleMessage(update *tgapi.Update, ctx *MsgContext) { func (bot *Bot[T]) handleMessage(update *tgapi.Update, ctx *MsgContext) {
if update.Message == nil { if update.Message == nil {
return return
} }
@@ -34,7 +40,7 @@ func (b *Bot) handleMessage(update *tgapi.Update, ctx *MsgContext) {
} }
text = strings.TrimSpace(text) text = strings.TrimSpace(text)
prefix, hasPrefix := b.checkPrefixes(text) prefix, hasPrefix := bot.checkPrefixes(text)
if !hasPrefix { if !hasPrefix {
return return
} }
@@ -45,8 +51,8 @@ func (b *Bot) handleMessage(update *tgapi.Update, ctx *MsgContext) {
text = strings.TrimSpace(text[len(prefix):]) text = strings.TrimSpace(text[len(prefix):])
for _, plugin := range b.plugins { for _, plugin := range bot.plugins {
for cmd := range plugin.Commands { for cmd := range plugin.commands {
if !strings.HasPrefix(text, cmd) { if !strings.HasPrefix(text, cmd) {
continue continue
} }
@@ -71,20 +77,20 @@ func (b *Bot) handleMessage(update *tgapi.Update, ctx *MsgContext) {
ctx.Args = strings.Split(ctx.Text, " ") ctx.Args = strings.Split(ctx.Text, " ")
} }
if !plugin.executeMiddlewares(ctx, b.dbContext) { if !plugin.executeMiddlewares(ctx, bot.dbContext) {
return return
} }
go plugin.executeCmd(cmd, ctx, b.dbContext) go plugin.executeCmd(cmd, ctx, bot.dbContext)
return return
} }
} }
} }
func (b *Bot) handleCallback(update *tgapi.Update, ctx *MsgContext) { func (bot *Bot[T]) handleCallback(update *tgapi.Update, ctx *MsgContext) {
data := new(CallbackData) data := new(CallbackData)
err := json.Unmarshal([]byte(update.CallbackQuery.Data), data) err := json.Unmarshal([]byte(update.CallbackQuery.Data), data)
if err != nil { if err != nil {
b.logger.Errorln(err) bot.logger.Errorln(err)
return return
} }
@@ -95,22 +101,22 @@ func (b *Bot) handleCallback(update *tgapi.Update, ctx *MsgContext) {
ctx.CallbackQueryId = update.CallbackQuery.ID ctx.CallbackQueryId = update.CallbackQuery.ID
ctx.Args = data.Args ctx.Args = data.Args
for _, plugin := range b.plugins { for _, plugin := range bot.plugins {
_, ok := plugin.Payloads[data.Command] _, ok := plugin.payloads[data.Command]
if !ok { if !ok {
continue continue
} }
if !plugin.executeMiddlewares(ctx, b.dbContext) { if !plugin.executeMiddlewares(ctx, bot.dbContext) {
return return
} }
go plugin.executePayload(data.Command, ctx, b.dbContext) go plugin.executePayload(data.Command, ctx, bot.dbContext)
return return
} }
} }
func (b *Bot) checkPrefixes(text string) (string, bool) { func (bot *Bot[T]) checkPrefixes(text string) (string, bool) {
for _, prefix := range b.prefixes { for _, prefix := range bot.prefixes {
if strings.HasPrefix(text, prefix) { if strings.HasPrefix(text, prefix) {
return prefix, true return prefix, true
} }
-1
View File
@@ -17,7 +17,6 @@ func (l *L10n) AddDictEntry(key string, value DictEntry) *L10n {
func (l *L10n) GetFallbackLanguage() string { func (l *L10n) GetFallbackLanguage() string {
return l.fallbackLang return l.fallbackLang
} }
func (l *L10n) Translate(lang, key string) string { func (l *L10n) Translate(lang, key string) string {
s, ok := l.entries[key] s, ok := l.entries[key]
if !ok { if !ok {
+10 -26
View File
@@ -2,50 +2,34 @@ package laniakea
import ( import (
"encoding/json" "encoding/json"
"fmt"
"io"
"net/http"
"git.nix13.pw/scuroneko/laniakea/tgapi" "git.nix13.pw/scuroneko/laniakea/tgapi"
) )
func (b *Bot) Updates() ([]tgapi.Update, error) { func (bot *Bot[T]) Updates() ([]tgapi.Update, error) {
offset := b.GetUpdateOffset() offset := bot.GetUpdateOffset()
params := tgapi.UpdateParams{ params := tgapi.UpdateParams{
Offset: Ptr(offset), Offset: Ptr(offset),
Timeout: Ptr(30), Timeout: Ptr(30),
AllowedUpdates: b.GetUpdateTypes(), AllowedUpdates: bot.GetUpdateTypes(),
} }
updates, err := b.api.GetUpdates(params) updates, err := bot.api.GetUpdates(params)
if err != nil { if err != nil {
return nil, err return nil, err
} }
if bot.RequestLogger != nil {
for _, u := range updates { for _, u := range updates {
b.SetUpdateOffset(u.UpdateID + 1)
err = b.GetQueue().Enqueue(&u)
if err != nil {
return nil, err
}
if b.RequestLogger != nil {
j, err := json.Marshal(u) j, err := json.Marshal(u)
if err != nil { if err != nil {
b.Logger().Error(err) bot.GetLogger().Error(err)
} }
b.RequestLogger.Debugf("UPDATE %s\n", j) bot.RequestLogger.Debugf("UPDATE %s\n", j)
} }
} }
if len(updates) > 0 {
bot.SetUpdateOffset(updates[len(updates)-1].UpdateID + 1)
}
return updates, err return updates, err
} }
func (b *Bot) GetFileByLink(link string) ([]byte, error) {
u := fmt.Sprintf("https://api.telegram.org/file/bot%s/%s", b.token, link)
res, err := http.Get(u)
if err != nil {
return nil, err
}
defer res.Body.Close()
return io.ReadAll(res.Body)
}
+46 -37
View File
@@ -5,10 +5,10 @@ import (
"git.nix13.pw/scuroneko/laniakea/tgapi" "git.nix13.pw/scuroneko/laniakea/tgapi"
"git.nix13.pw/scuroneko/laniakea/utils" "git.nix13.pw/scuroneko/laniakea/utils"
"git.nix13.pw/scuroneko/slog"
) )
type MsgContext struct { type MsgContext struct {
Bot *Bot
Api *tgapi.API Api *tgapi.API
Msg *tgapi.Message Msg *tgapi.Message
@@ -20,6 +20,11 @@ type MsgContext struct {
Prefix string Prefix string
Text string Text string
Args []string Args []string
errorTemplate string
botLogger *slog.Logger
l10n *L10n
draftProvider *DraftProvider
} }
type AnswerMessage struct { type AnswerMessage struct {
@@ -41,7 +46,7 @@ func (ctx *MsgContext) edit(messageId int, text string, keyboard *InlineKeyboard
} }
msg, _, err := ctx.Api.EditMessageText(params) msg, _, err := ctx.Api.EditMessageText(params)
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
return nil return nil
} }
return &AnswerMessage{ return &AnswerMessage{
@@ -53,7 +58,7 @@ func (m *AnswerMessage) Edit(text string) *AnswerMessage {
} }
func (ctx *MsgContext) EditCallback(text string, keyboard *InlineKeyboard) *AnswerMessage { func (ctx *MsgContext) EditCallback(text string, keyboard *InlineKeyboard) *AnswerMessage {
if ctx.CallbackMsgId == 0 { if ctx.CallbackMsgId == 0 {
ctx.Api.Logger.Errorln("Can't edit non-callback update message") ctx.botLogger.Errorln("Can't edit non-callback update message")
return nil return nil
} }
@@ -73,9 +78,10 @@ func (ctx *MsgContext) editPhotoText(messageId int, text string, kb *InlineKeybo
if kb != nil { if kb != nil {
params.ReplyMarkup = kb.Get() params.ReplyMarkup = kb.Get()
} }
msg, _, err := ctx.Api.EditMessageCaption(params) msg, _, err := ctx.Api.EditMessageCaption(params)
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
} }
return &AnswerMessage{ return &AnswerMessage{
MessageID: msg.MessageID, ctx: ctx, Text: text, IsMedia: true, MessageID: msg.MessageID, ctx: ctx, Text: text, IsMedia: true,
@@ -83,7 +89,7 @@ func (ctx *MsgContext) editPhotoText(messageId int, text string, kb *InlineKeybo
} }
func (m *AnswerMessage) EditCaption(text string) *AnswerMessage { func (m *AnswerMessage) EditCaption(text string) *AnswerMessage {
if m.MessageID == 0 { if m.MessageID == 0 {
m.ctx.Api.Logger.Errorln("Can't edit caption message, message id is zero") m.ctx.botLogger.Errorln("Can't edit caption message, message id is zero")
return m return m
} }
return m.ctx.editPhotoText(m.MessageID, text, nil) return m.ctx.editPhotoText(m.MessageID, text, nil)
@@ -101,10 +107,13 @@ func (ctx *MsgContext) answer(text string, keyboard *InlineKeyboard) *AnswerMess
if keyboard != nil { if keyboard != nil {
params.ReplyMarkup = keyboard.Get() params.ReplyMarkup = keyboard.Get()
} }
if ctx.Msg.MessageThreadID > 0 {
params.MessageThreadID = ctx.Msg.MessageThreadID
}
msg, err := ctx.Api.SendMessage(params) msg, err := ctx.Api.SendMessage(params)
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
return nil return nil
} }
return &AnswerMessage{ return &AnswerMessage{
@@ -131,9 +140,13 @@ func (ctx *MsgContext) answerPhoto(photoId, text string, kb *InlineKeyboard) *An
if kb != nil { if kb != nil {
params.ReplyMarkup = kb.Get() params.ReplyMarkup = kb.Get()
} }
if ctx.Msg.MessageThreadID > 0 {
params.MessageThreadID = ctx.Msg.MessageThreadID
}
msg, err := ctx.Api.SendPhoto(params) msg, err := ctx.Api.SendPhoto(params)
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
return &AnswerMessage{ return &AnswerMessage{
ctx: ctx, Text: text, IsMedia: true, ctx: ctx, Text: text, IsMedia: true,
} }
@@ -155,15 +168,11 @@ func (ctx *MsgContext) delete(messageId int) {
MessageID: messageId, MessageID: messageId,
}) })
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
} }
} }
func (m *AnswerMessage) Delete() { func (m *AnswerMessage) Delete() { m.ctx.delete(m.MessageID) }
m.ctx.delete(m.MessageID) func (ctx *MsgContext) CallbackDelete() { ctx.delete(ctx.CallbackMsgId) }
}
func (ctx *MsgContext) CallbackDelete() {
ctx.delete(ctx.CallbackMsgId)
}
func (ctx *MsgContext) answerCallbackQuery(url, text string, showAlert bool) { func (ctx *MsgContext) answerCallbackQuery(url, text string, showAlert bool) {
if len(ctx.CallbackQueryId) == 0 { if len(ctx.CallbackQueryId) == 0 {
@@ -174,49 +183,49 @@ func (ctx *MsgContext) answerCallbackQuery(url, text string, showAlert bool) {
Text: text, ShowAlert: showAlert, URL: url, Text: text, ShowAlert: showAlert, URL: url,
}) })
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
} }
} }
func (ctx *MsgContext) AnswerCbQuery() { func (ctx *MsgContext) AnswerCbQuery() { ctx.answerCallbackQuery("", "", false) }
ctx.answerCallbackQuery("", "", false) func (ctx *MsgContext) AnswerCbQueryText(text string) { ctx.answerCallbackQuery("", text, false) }
} func (ctx *MsgContext) AnswerCbQueryAlert(text string) { ctx.answerCallbackQuery("", text, true) }
func (ctx *MsgContext) AnswerCbQueryText(text string) { func (ctx *MsgContext) AnswerCbQueryUrl(u string) { ctx.answerCallbackQuery(u, "", false) }
ctx.answerCallbackQuery("", text, false)
}
func (ctx *MsgContext) AnswerCbQueryAlert(text string) {
ctx.answerCallbackQuery("", text, true)
}
func (ctx *MsgContext) AnswerCbQueryUrl(u string) {
ctx.answerCallbackQuery(u, "", false)
}
func (ctx *MsgContext) SendAction(action tgapi.ChatActionType) { func (ctx *MsgContext) SendAction(action tgapi.ChatActionType) {
_, err := ctx.Api.SendChatAction(tgapi.SendChatActionP{ params := tgapi.SendChatActionP{
ChatID: ctx.Msg.Chat.ID, Action: action, ChatID: ctx.Msg.Chat.ID, Action: action,
}) }
if ctx.Msg.MessageThreadID > 0 {
params.MessageThreadID = ctx.Msg.MessageThreadID
}
_, err := ctx.Api.SendChatAction(params)
if err != nil { if err != nil {
ctx.Api.Logger.Errorln(err) ctx.botLogger.Errorln(err)
} }
} }
func (ctx *MsgContext) error(err error) { func (ctx *MsgContext) error(err error) {
text := fmt.Sprintf(ctx.Bot.errorTemplate, utils.EscapeMarkdown(err.Error())) text := fmt.Sprintf(ctx.errorTemplate, utils.EscapeMarkdown(err.Error()))
if ctx.CallbackQueryId != "" { if ctx.CallbackQueryId != "" {
ctx.answerCallbackQuery("", text, false) ctx.answerCallbackQuery("", text, false)
} else { } else {
ctx.answer(text, nil) ctx.answer(text, nil)
} }
ctx.Bot.Logger().Errorln(err) ctx.botLogger.Errorln(err)
}
func (ctx *MsgContext) Error(err error) {
ctx.error(err)
} }
func (ctx *MsgContext) Error(err error) { ctx.error(err) }
func (ctx *MsgContext) NewDraft() *Draft {
draft := ctx.draftProvider.NewDraft()
draft.chatID = ctx.Msg.Chat.ID
draft.messageThreadID = ctx.Msg.MessageThreadID
return draft
}
func (ctx *MsgContext) Translate(key string) string { func (ctx *MsgContext) Translate(key string) string {
if ctx.From == nil { if ctx.From == nil {
return key return key
} }
lang := Val(ctx.From.LanguageCode, ctx.Bot.l10n.GetFallbackLanguage()) lang := Val(ctx.From.LanguageCode, ctx.l10n.GetFallbackLanguage())
return ctx.Bot.L10n(lang, key) return ctx.l10n.Translate(lang, key)
} }
+46 -40
View File
@@ -45,32 +45,33 @@ func (c *CommandArg) SetRequired() *CommandArg {
return c return c
} }
type CommandExecutor func(ctx *MsgContext, dbContext *DatabaseContext) type CommandExecutor[T DbContext] func(ctx *MsgContext, dbContext *T)
type Command struct {
type Command[T DbContext] struct {
command string command string
description string description string
exec CommandExecutor exec CommandExecutor[T]
args extypes.Slice[CommandArg] args extypes.Slice[CommandArg]
middlewares extypes.Slice[Middleware] middlewares extypes.Slice[Middleware[T]]
skipAutoCmd bool skipAutoCmd bool
} }
func NewCommand(exec CommandExecutor, command string, args ...CommandArg) *Command { func NewCommand[T any](exec CommandExecutor[T], command string, args ...CommandArg) *Command[T] {
return &Command{command, "", exec, args, make(extypes.Slice[Middleware], 0), false} return &Command[T]{command, "", exec, args, make(extypes.Slice[Middleware[T]], 0), false}
} }
func (c *Command) Use(m Middleware) *Command { func (c *Command[T]) Use(m Middleware[T]) *Command[T] {
c.middlewares = c.middlewares.Push(m) c.middlewares = c.middlewares.Push(m)
return c return c
} }
func (c *Command) SetDescription(desc string) *Command { func (c *Command[T]) SetDescription(desc string) *Command[T] {
c.description = desc c.description = desc
return c return c
} }
func (c *Command) SkipCommandAutoGen() *Command { func (c *Command[T]) SkipCommandAutoGen() *Command[T] {
c.skipAutoCmd = true c.skipAutoCmd = true
return c return c
} }
func (c *Command) validateArgs(args []string) error { func (c *Command[T]) validateArgs(args []string) error {
cmdArgs := c.args.Filter(func(e CommandArg) bool { return !e.required }) cmdArgs := c.args.Filter(func(e CommandArg) bool { return !e.required })
if len(args) < cmdArgs.Len() { if len(args) < cmdArgs.Len() {
return ErrCmdArgCountMismatch return ErrCmdArgCountMismatch
@@ -91,54 +92,59 @@ func (c *Command) validateArgs(args []string) error {
return nil return nil
} }
type Plugin struct { type Plugin[T DbContext] struct {
Name string name string
Commands map[string]Command commands map[string]Command[T]
Payloads map[string]Command payloads map[string]Command[T]
Middlewares extypes.Slice[Middleware] middlewares extypes.Slice[Middleware[T]]
skipAutoCmd bool
} }
func NewPlugin(name string) *Plugin { func NewPlugin[T DbContext](name string) *Plugin[T] {
return &Plugin{ return &Plugin[T]{
name, map[string]Command{}, name, map[string]Command[T]{},
map[string]Command{}, extypes.Slice[Middleware]{}, map[string]Command[T]{}, extypes.Slice[Middleware[T]]{}, false,
} }
} }
func (p *Plugin) AddCommand(command *Command) *Plugin { func (p *Plugin[T]) AddCommand(command *Command[T]) *Plugin[T] {
p.Commands[command.command] = *command p.commands[command.command] = *command
return p return p
} }
func (p *Plugin) NewCommand(exec CommandExecutor, command string, args ...CommandArg) *Command { func (p *Plugin[T]) NewCommand(exec CommandExecutor[T], command string, args ...CommandArg) *Command[T] {
return NewCommand(exec, command, args...) return NewCommand(exec, command, args...)
} }
func (p *Plugin) AddPayload(command *Command) *Plugin { func (p *Plugin[T]) AddPayload(command *Command[T]) *Plugin[T] {
p.Payloads[command.command] = *command p.payloads[command.command] = *command
return p return p
} }
func (p *Plugin) AddMiddleware(middleware Middleware) *Plugin { func (p *Plugin[T]) AddMiddleware(middleware Middleware[T]) *Plugin[T] {
p.Middlewares = p.Middlewares.Push(middleware) p.middlewares = p.middlewares.Push(middleware)
return p
}
func (p *Plugin[T]) SkipCommandAutoGen() *Plugin[T] {
p.skipAutoCmd = true
return p return p
} }
func (p *Plugin) executeCmd(cmd string, ctx *MsgContext, dbContext *DatabaseContext) { func (p *Plugin[T]) executeCmd(cmd string, ctx *MsgContext, dbContext *T) {
command := p.Commands[cmd] command := p.commands[cmd]
if err := command.validateArgs(ctx.Args); err != nil { if err := command.validateArgs(ctx.Args); err != nil {
ctx.error(err) ctx.error(err)
return return
} }
command.exec(ctx, dbContext) command.exec(ctx, dbContext)
} }
func (p *Plugin) executePayload(payload string, ctx *MsgContext, dbContext *DatabaseContext) { func (p *Plugin[T]) executePayload(payload string, ctx *MsgContext, dbContext *T) {
pl := p.Payloads[payload] pl := p.payloads[payload]
if err := pl.validateArgs(ctx.Args); err != nil { if err := pl.validateArgs(ctx.Args); err != nil {
ctx.error(err) ctx.error(err)
return return
} }
pl.exec(ctx, dbContext) pl.exec(ctx, dbContext)
} }
func (p *Plugin) executeMiddlewares(ctx *MsgContext, db *DatabaseContext) bool { func (p *Plugin[T]) executeMiddlewares(ctx *MsgContext, db *T) bool {
for _, m := range p.Middlewares { for _, m := range p.middlewares {
if !m.Execute(ctx, db) { if !m.Execute(ctx, db) {
return false return false
} }
@@ -146,29 +152,29 @@ func (p *Plugin) executeMiddlewares(ctx *MsgContext, db *DatabaseContext) bool {
return true return true
} }
type MiddlewareExecutor func(ctx *MsgContext, db *DatabaseContext) bool type MiddlewareExecutor[T DbContext] func(ctx *MsgContext, db *T) bool
// Middleware // Middleware
// When async, returned value ignored // When async, returned value ignored
type Middleware struct { type Middleware[T DbContext] struct {
name string name string
executor MiddlewareExecutor executor MiddlewareExecutor[T]
order int order int
async bool async bool
} }
func NewMiddleware(name string, executor MiddlewareExecutor) *Middleware { func NewMiddleware[T DbContext](name string, executor MiddlewareExecutor[T]) *Middleware[T] {
return &Middleware{name, executor, 0, false} return &Middleware[T]{name, executor, 0, false}
} }
func (m *Middleware) SetOrder(order int) *Middleware { func (m *Middleware[T]) SetOrder(order int) *Middleware[T] {
m.order = order m.order = order
return m return m
} }
func (m *Middleware) SetAsync(async bool) *Middleware { func (m *Middleware[T]) SetAsync(async bool) *Middleware[T] {
m.async = async m.async = async
return m return m
} }
func (m *Middleware) Execute(ctx *MsgContext, db *DatabaseContext) bool { func (m *Middleware[T]) Execute(ctx *MsgContext, db *T) bool {
if m.async { if m.async {
go m.executor(ctx, db) go m.executor(ctx, db)
return true return true
+26 -38
View File
@@ -4,81 +4,69 @@ import (
"time" "time"
) )
type RunnerFn func(*Bot) error type RunnerFn[T DbContext] func(*Bot[T]) error
type RunnerBuilder struct { type Runner[T DbContext] struct {
name string name string
onetime bool onetime bool
async bool async bool
timeout time.Duration timeout time.Duration
fn RunnerFn fn RunnerFn[T]
}
type Runner struct {
Name string
Onetime bool
Async bool
Timeout time.Duration
Fn RunnerFn
} }
func NewRunner(name string, fn RunnerFn) *RunnerBuilder { func NewRunner[T DbContext](name string, fn RunnerFn[T]) *Runner[T] {
return &RunnerBuilder{ return &Runner[T]{
name: name, fn: fn, async: true, name: name, fn: fn, async: true,
} }
} }
func (b *RunnerBuilder) Onetime(onetime bool) *RunnerBuilder { func (b *Runner[T]) Onetime(onetime bool) *Runner[T] {
b.onetime = onetime b.onetime = onetime
return b return b
} }
func (b *RunnerBuilder) Async(async bool) *RunnerBuilder { func (b *Runner[T]) Async(async bool) *Runner[T] {
b.async = async b.async = async
return b return b
} }
func (b *RunnerBuilder) Timeout(timeout time.Duration) *RunnerBuilder { func (b *Runner[T]) Timeout(timeout time.Duration) *Runner[T] {
b.timeout = timeout b.timeout = timeout
return b return b
} }
func (b *RunnerBuilder) Build() Runner {
return Runner{
Name: b.name, Onetime: b.onetime, Async: b.async, Fn: b.fn, Timeout: b.timeout,
}
}
func (b *Bot) ExecRunners() { func (bot *Bot[T]) ExecRunners() {
b.logger.Infoln("Executing runners...") bot.logger.Infoln("Executing runners...")
for _, runner := range b.runners { for _, runner := range bot.runners {
if !runner.Onetime && !runner.Async { if !runner.onetime && !runner.async {
b.logger.Warnf("Runner %s not onetime, but sync\n", runner.Name) bot.logger.Warnf("Runner %s not onetime, but sync\n", runner.name)
continue continue
} }
if !runner.Onetime && runner.Async && runner.Timeout == (time.Second*0) { if !runner.onetime && runner.async && runner.timeout == (time.Second*0) {
b.logger.Warnf("Background runner \"%s\" should have timeout", runner.Name) bot.logger.Warnf("Background runner \"%s\" should have timeout", runner.name)
} }
if runner.Async && runner.Onetime { if runner.async && runner.onetime {
go func() { go func() {
err := runner.Fn(b) err := runner.fn(bot)
if err != nil { if err != nil {
b.logger.Warnf("Runner %s failed: %s\n", runner.Name, err) bot.logger.Warnf("Runner %s failed: %s\n", runner.name, err)
} }
}() }()
} else if !runner.Async && runner.Onetime { } else if !runner.async && runner.onetime {
t := time.Now() t := time.Now()
err := runner.Fn(b) err := runner.fn(bot)
if err != nil { if err != nil {
b.logger.Warnf("Runner %s failed: %s\n", runner.Name, err) bot.logger.Warnf("Runner %s failed: %s\n", runner.name, err)
} }
elapsed := time.Since(t) elapsed := time.Since(t)
if elapsed > time.Second*2 { if elapsed > time.Second*2 {
b.logger.Warnf("Runner %s too slow. Elapsed time %s>=2s", runner.Name, elapsed) bot.logger.Warnf("Runner %s too slow. Elapsed time %s>=2s", runner.name, elapsed)
} }
} else if !runner.Onetime { } else if !runner.onetime {
go func() { go func() {
for { for {
err := runner.Fn(b) err := runner.fn(bot)
if err != nil { if err != nil {
b.logger.Warnf("Runner %s failed: %s\n", runner.Name, err) bot.logger.Warnf("Runner %s failed: %s\n", runner.name, err)
} }
time.Sleep(runner.Timeout) time.Sleep(runner.timeout)
} }
}() }()
} }
+124 -21
View File
@@ -4,6 +4,7 @@ import (
"bytes" "bytes"
"context" "context"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"io" "io"
"net/http" "net/http"
@@ -11,23 +12,81 @@ import (
"git.nix13.pw/scuroneko/laniakea/utils" "git.nix13.pw/scuroneko/laniakea/utils"
"git.nix13.pw/scuroneko/slog" "git.nix13.pw/scuroneko/slog"
"golang.org/x/time/rate"
) )
type APIOpts struct {
token string
client *http.Client
useTestServer bool
apiUrl string
limiter *rate.Limiter
dropOverflowLimit bool
}
var ErrPoolUnexpected = errors.New("unexpected response from pool")
func NewAPIOpts(token string) *APIOpts {
return &APIOpts{token: token, client: nil, useTestServer: false, apiUrl: "https://api.telegram.org"}
}
func (opts *APIOpts) SetHTTPClient(client *http.Client) *APIOpts {
if client != nil {
opts.client = client
}
return opts
}
func (opts *APIOpts) UseTestServer(use bool) *APIOpts {
opts.useTestServer = use
return opts
}
func (opts *APIOpts) SetAPIUrl(apiUrl string) *APIOpts {
if apiUrl != "" {
opts.apiUrl = apiUrl
}
return opts
}
func (opts *APIOpts) SetLimiter(limiter *rate.Limiter) *APIOpts {
opts.limiter = limiter
return opts
}
func (opts *APIOpts) SetLimiterDrop(b bool) *APIOpts {
opts.dropOverflowLimit = b
return opts
}
type API struct { type API struct {
token string token string
client *http.Client client *http.Client
Logger *slog.Logger logger *slog.Logger
useTestServer bool
apiUrl string
pool *WorkerPool
limiter *rate.Limiter
dropOverflowLimit bool
} }
func NewAPI(token string) *API { func NewAPI(opts *APIOpts) *API {
l := slog.CreateLogger().Level(utils.GetLoggerLevel()).Prefix("API") l := slog.CreateLogger().Level(utils.GetLoggerLevel()).Prefix("API")
l.AddWriter(l.CreateJsonStdoutWriter()) l.AddWriter(l.CreateJsonStdoutWriter())
client := &http.Client{Timeout: time.Second * 45} client := opts.client
return &API{token, client, l} if client == nil {
client = &http.Client{Timeout: time.Second * 45}
}
pool := NewWorkerPool(16, 256)
pool.Start(context.Background())
return &API{
opts.token, client, l,
opts.useTestServer, opts.apiUrl,
pool, opts.limiter, opts.dropOverflowLimit,
}
} }
func (api *API) CloseApi() error { func (api *API) CloseApi() error {
return api.Logger.Close() api.pool.Stop()
return api.logger.Close()
} }
func (api *API) GetLogger() *slog.Logger { return api.logger }
type ApiResponse[R any] struct { type ApiResponse[R any] struct {
Ok bool `json:"ok"` Ok bool `json:"ok"`
@@ -35,7 +94,6 @@ type ApiResponse[R any] struct {
Result R `json:"result,omitempty"` Result R `json:"result,omitempty"`
ErrorCode int `json:"error_code,omitempty"` ErrorCode int `json:"error_code,omitempty"`
} }
type TelegramRequest[R, P any] struct { type TelegramRequest[R, P any] struct {
method string method string
params P params P
@@ -44,16 +102,32 @@ type TelegramRequest[R, P any] struct {
func NewRequest[R, P any](method string, params P) TelegramRequest[R, P] { func NewRequest[R, P any](method string, params P) TelegramRequest[R, P] {
return TelegramRequest[R, P]{method: method, params: params} return TelegramRequest[R, P]{method: method, params: params}
} }
func (r TelegramRequest[R, P]) DoWithContext(ctx context.Context, api *API) (R, error) { func (r TelegramRequest[R, P]) doRequest(ctx context.Context, api *API) (R, error) {
var zero R var zero R
if api.limiter != nil {
if api.dropOverflowLimit {
if !api.limiter.Allow() {
return zero, errors.New("rate limited")
}
} else {
if err := api.limiter.Wait(ctx); err != nil {
return zero, err
}
}
}
data, err := json.Marshal(r.params) data, err := json.Marshal(r.params)
if err != nil { if err != nil {
return zero, err return zero, err
} }
buf := bytes.NewBuffer(data) buf := bytes.NewBuffer(data)
u := fmt.Sprintf("https://api.telegram.org/bot%s/%s", api.token, r.method) methodPrefix := ""
req, err := http.NewRequestWithContext(ctx, "POST", u, buf) if api.useTestServer {
methodPrefix = "/test"
}
url := fmt.Sprintf("%s/bot%s%s/%s", api.apiUrl, api.token, methodPrefix, r.method)
req, err := http.NewRequestWithContext(ctx, "POST", url, buf)
if err != nil { if err != nil {
return zero, err return zero, err
} }
@@ -61,25 +135,59 @@ func (r TelegramRequest[R, P]) DoWithContext(ctx context.Context, api *API) (R,
req.Header.Set("Accept", "application/json") req.Header.Set("Accept", "application/json")
req.Header.Set("User-Agent", fmt.Sprintf("Laniakea/%s", utils.VersionString)) req.Header.Set("User-Agent", fmt.Sprintf("Laniakea/%s", utils.VersionString))
api.Logger.Debugln("REQ", r.method, buf.String()) api.logger.Debugln("REQ", api.apiUrl, r.method, buf.String())
res, err := api.client.Do(req) res, err := api.client.Do(req)
if err != nil { if err != nil {
return zero, err return zero, err
} }
defer res.Body.Close() defer func(Body io.ReadCloser) {
_ = Body.Close()
}(res.Body)
reader := io.LimitReader(res.Body, 10<<20) data, err = readBody(res.Body)
data, err = io.ReadAll(reader)
if err != nil { if err != nil {
return zero, err return zero, err
} }
api.Logger.Debugln("RES", r.method, string(data)) api.logger.Debugln("RES", r.method, string(data))
if res.StatusCode != http.StatusOK { if res.StatusCode != http.StatusOK {
return zero, fmt.Errorf("unexpected status code: %d", res.StatusCode) return zero, fmt.Errorf("unexpected status code: %d, %s", res.StatusCode, string(data))
}
return parseBody[R](data)
}
func (r TelegramRequest[R, P]) DoWithContext(ctx context.Context, api *API) (R, error) {
var zero R
result, err := api.pool.Submit(ctx, func(ctx context.Context) (any, error) {
return r.doRequest(ctx, api)
})
if err != nil {
return zero, err
} }
select {
case <-ctx.Done():
return zero, ctx.Err()
case res := <-result:
if res.Err != nil {
return zero, res.Err
}
if val, ok := res.Value.(R); ok {
return val, nil
}
return zero, ErrPoolUnexpected
}
}
func (r TelegramRequest[R, P]) Do(api *API) (R, error) {
return r.DoWithContext(context.Background(), api)
}
func readBody(body io.ReadCloser) ([]byte, error) {
reader := io.LimitReader(body, 10<<20)
return io.ReadAll(reader)
}
func parseBody[R any](data []byte) (R, error) {
var zero R
var resp ApiResponse[R] var resp ApiResponse[R]
err = json.Unmarshal(data, &resp) err := json.Unmarshal(data, &resp)
if err != nil { if err != nil {
return zero, err return zero, err
} }
@@ -87,9 +195,4 @@ func (r TelegramRequest[R, P]) DoWithContext(ctx context.Context, api *API) (R,
return zero, fmt.Errorf("[%d] %s", resp.ErrorCode, resp.Description) return zero, fmt.Errorf("[%d] %s", resp.ErrorCode, resp.Description)
} }
return resp.Result, nil return resp.Result, nil
}
func (r TelegramRequest[R, P]) Do(api *API) (R, error) {
ctx := context.Background()
return r.DoWithContext(ctx, api)
} }
+12
View File
@@ -56,6 +56,7 @@ type PromoteChatMember struct {
CanPinMessages bool `json:"can_pin_messages,omitempty"` CanPinMessages bool `json:"can_pin_messages,omitempty"`
CanManageTopics bool `json:"can_manage_topics,omitempty"` CanManageTopics bool `json:"can_manage_topics,omitempty"`
CanManageDirectMessages bool `json:"can_manage_direct_messages,omitempty"` CanManageDirectMessages bool `json:"can_manage_direct_messages,omitempty"`
CanManageTags bool `json:"can_manage_tags,omitempty"`
} }
func (api *API) PromoteChatMember(params PromoteChatMember) (bool, error) { func (api *API) PromoteChatMember(params PromoteChatMember) (bool, error) {
@@ -74,6 +75,17 @@ func (api *API) SetChatAdministratorCustomTitle(params SetChatAdministratorCusto
return req.Do(api) return req.Do(api)
} }
type SetChatMemberTagP struct {
ChatID int `json:"chat_id"`
UserID int `json:"user_id"`
Tag string `json:"tag,omitempty"`
}
func (api *API) SetChatMemberTag(params SetChatMemberTagP) (bool, error) {
req := NewRequest[bool]("setChatMemberTag", params)
return req.Do(api)
}
type BanChatSenderChatP struct { type BanChatSenderChatP struct {
ChatID int `json:"chat_id"` ChatID int `json:"chat_id"`
SenderChatID int `json:"sender_chat_id"` SenderChatID int `json:"sender_chat_id"`
+5
View File
@@ -99,6 +99,7 @@ type ChatPermissions struct {
CanSendPolls bool `json:"can_send_polls"` CanSendPolls bool `json:"can_send_polls"`
CanSendOtherMessages bool `json:"can_send_other_messages"` CanSendOtherMessages bool `json:"can_send_other_messages"`
CanAddWebPagePreview bool `json:"can_add_web_page_preview"` CanAddWebPagePreview bool `json:"can_add_web_page_preview"`
CatEditTag bool `json:"cat_edit_tag"`
CanChangeInfo bool `json:"can_change_info"` CanChangeInfo bool `json:"can_change_info"`
CanInviteUsers bool `json:"can_invite_users"` CanInviteUsers bool `json:"can_invite_users"`
CanPinMessages bool `json:"can_pin_messages"` CanPinMessages bool `json:"can_pin_messages"`
@@ -137,6 +138,7 @@ const (
type ChatMember struct { type ChatMember struct {
Status ChatMemberStatusType `json:"status"` Status ChatMemberStatusType `json:"status"`
User User `json:"user"` User User `json:"user"`
Tag string `json:"tag,omitempty"`
// Owner // Owner
IsAnonymous *bool `json:"is_anonymous"` IsAnonymous *bool `json:"is_anonymous"`
@@ -160,6 +162,7 @@ type ChatMember struct {
CanPinMessages *bool `json:"can_pin_messages,omitempty"` CanPinMessages *bool `json:"can_pin_messages,omitempty"`
CanManageTopics *bool `json:"can_manage_topics,omitempty"` CanManageTopics *bool `json:"can_manage_topics,omitempty"`
CanManageDirectMessages *bool `json:"can_manage_direct_messages,omitempty"` CanManageDirectMessages *bool `json:"can_manage_direct_messages,omitempty"`
CanManageTags *bool `json:"can_manage_tags,omitempty"`
// Member // Member
UntilDate *int `json:"until_date,omitempty"` UntilDate *int `json:"until_date,omitempty"`
@@ -175,6 +178,7 @@ type ChatMember struct {
CanSendPolls *bool `json:"can_send_polls,omitempty"` CanSendPolls *bool `json:"can_send_polls,omitempty"`
CanSendOtherMessages *bool `json:"can_send_other_messages,omitempty"` CanSendOtherMessages *bool `json:"can_send_other_messages,omitempty"`
CanAddWebPagePreview *bool `json:"can_add_web_page_preview,omitempty"` CanAddWebPagePreview *bool `json:"can_add_web_page_preview,omitempty"`
CanEditTag *bool `json:"can_edit_tag,omitempty"`
} }
type ChatBoostSource struct { type ChatBoostSource struct {
@@ -215,6 +219,7 @@ type ChatAdministratorRights struct {
CanPinMessages *bool `json:"can_pin_messages,omitempty"` CanPinMessages *bool `json:"can_pin_messages,omitempty"`
CanManageTopics *bool `json:"can_manage_topics,omitempty"` CanManageTopics *bool `json:"can_manage_topics,omitempty"`
CanManageDirectMessages *bool `json:"can_manage_direct_messages,omitempty"` CanManageDirectMessages *bool `json:"can_manage_direct_messages,omitempty"`
CanManageTags *bool `json:"can_manage_tags,omitempty"`
} }
type ChatBoostUpdated struct { type ChatBoostUpdated struct {
+1 -1
View File
@@ -268,7 +268,7 @@ func (api *API) SendDice(params SendDiceP) (Message, error) {
type SendMessageDraftP struct { type SendMessageDraftP struct {
ChatID int `json:"chat_id"` ChatID int `json:"chat_id"`
MessageThreadID int `json:"message_thread_id,omitempty"` MessageThreadID int `json:"message_thread_id,omitempty"`
DraftID int `json:"draft_id"` DraftID uint64 `json:"draft_id"`
Text string `json:"text"` Text string `json:"text"`
ParseMode ParseMode `json:"parse_mode,omitempty"` ParseMode ParseMode `json:"parse_mode,omitempty"`
Entities []MessageEntity `json:"entities,omitempty"` Entities []MessageEntity `json:"entities,omitempty"`
+5
View File
@@ -15,6 +15,7 @@ type Message struct {
SenderChat *Chat `json:"sender_chat,omitempty"` SenderChat *Chat `json:"sender_chat,omitempty"`
SenderBoostCount int `json:"sender_boost_count,omitempty"` SenderBoostCount int `json:"sender_boost_count,omitempty"`
SenderBusinessBot *User `json:"sender_business_bot,omitempty"` SenderBusinessBot *User `json:"sender_business_bot,omitempty"`
SenderTag string `json:"sender_tag,omitempty"`
Chat *Chat `json:"chat,omitempty"` Chat *Chat `json:"chat,omitempty"`
IsTopicMessage bool `json:"is_topic_message,omitempty"` IsTopicMessage bool `json:"is_topic_message,omitempty"`
@@ -74,6 +75,7 @@ const (
MessageEntityTextLink MessageEntityType = "text_link" MessageEntityTextLink MessageEntityType = "text_link"
MessageEntityTextMention MessageEntityType = "text_mention" MessageEntityTextMention MessageEntityType = "text_mention"
MessageEntityCustomEmoji MessageEntityType = "custom_emoji" MessageEntityCustomEmoji MessageEntityType = "custom_emoji"
MessageEntityDateTime MessageEntityType = "date_time"
) )
type MessageEntity struct { type MessageEntity struct {
@@ -85,6 +87,9 @@ type MessageEntity struct {
User *User `json:"user,omitempty"` User *User `json:"user,omitempty"`
Language string `json:"language,omitempty"` Language string `json:"language,omitempty"`
CustomEmojiID string `json:"custom_emoji_id,omitempty"` CustomEmojiID string `json:"custom_emoji_id,omitempty"`
UnixTime int `json:"unix_time,omitempty"`
DateTimeFormat string `json:"date_time_format,omitempty"`
} }
type ReplyParameters struct { type ReplyParameters struct {
+16
View File
@@ -1,5 +1,11 @@
package tgapi package tgapi
import (
"fmt"
"io"
"net/http"
)
type ParseMode string type ParseMode string
const ( const (
@@ -44,3 +50,13 @@ func (api *API) GetFile(params GetFileP) (File, error) {
req := NewRequest[File]("getFile", params) req := NewRequest[File]("getFile", params)
return req.Do(api) return req.Do(api)
} }
func (api *API) GetFileByLink(link string) ([]byte, error) {
u := fmt.Sprintf("https://api.telegram.org/file/bot%s/%s", api.token, link)
res, err := http.Get(u)
if err != nil {
return nil, err
}
defer res.Body.Close()
return io.ReadAll(res.Body)
}
+92
View File
@@ -0,0 +1,92 @@
package tgapi
import (
"context"
"errors"
"sync"
)
var ErrPoolQueueFull = errors.New("worker pool queue full")
type RequestEnvelope struct {
DoFunc func(context.Context) (any, error) // функция, которая выполнит запрос и вернет any
ResultCh chan RequestResult // канал для результата
}
type RequestResult struct {
Value any
Err error
}
// WorkerPool управляет воркерами и очередью
type WorkerPool struct {
taskCh chan RequestEnvelope
queueSize int
workers int
wg sync.WaitGroup
quit chan struct{}
started bool
startedMu sync.Mutex
}
func NewWorkerPool(workers int, queueSize int) *WorkerPool {
return &WorkerPool{
taskCh: make(chan RequestEnvelope, queueSize),
queueSize: queueSize,
workers: workers,
quit: make(chan struct{}),
}
}
// Start запускает воркеров
func (p *WorkerPool) Start(ctx context.Context) {
p.startedMu.Lock()
defer p.startedMu.Unlock()
if p.started {
return
}
p.started = true
for i := 0; i < p.workers; i++ {
p.wg.Add(1)
go p.worker(ctx)
}
}
// Stop останавливает пул (ждет завершения текущих задач)
func (p *WorkerPool) Stop() {
close(p.quit)
p.wg.Wait()
}
// Submit отправляет задачу в очередь и возвращает канал для результата
func (p *WorkerPool) Submit(ctx context.Context, do func(context.Context) (any, error)) (<-chan RequestResult, error) {
if len(p.taskCh) >= p.queueSize {
return nil, ErrPoolQueueFull
}
resultCh := make(chan RequestResult, 1) // буфер 1, чтобы не блокировать воркера
envelope := RequestEnvelope{do, resultCh}
select {
case <-ctx.Done():
return nil, ctx.Err()
case p.taskCh <- envelope:
return resultCh, nil
default:
return nil, ErrPoolQueueFull
}
}
// worker выполняет задачи
func (p *WorkerPool) worker(ctx context.Context) {
defer p.wg.Done()
for {
select {
case <-p.quit:
return
case envelope := <-p.taskCh:
// Выполняем задачу с переданным контекстом (или можно использовать свой)
val, err := envelope.DoFunc(ctx)
envelope.ResultCh <- RequestResult{Value: val, Err: err}
close(envelope.ResultCh)
}
}
}
+75 -43
View File
@@ -3,9 +3,8 @@ package tgapi
import ( import (
"bytes" "bytes"
"context" "context"
"encoding/json" "errors"
"fmt" "fmt"
"io"
"mime/multipart" "mime/multipart"
"net/http" "net/http"
"path/filepath" "path/filepath"
@@ -54,6 +53,7 @@ func NewUploader(api *API) *Uploader {
return &Uploader{api, logger} return &Uploader{api, logger}
} }
func (u *Uploader) Close() error { return u.logger.Close() } func (u *Uploader) Close() error { return u.logger.Close() }
func (u *Uploader) GetLogger() *slog.Logger { return u.logger }
type UploaderRequest[R, P any] struct { type UploaderRequest[R, P any] struct {
method string method string
@@ -64,73 +64,105 @@ type UploaderRequest[R, P any] struct {
func NewUploaderRequest[R, P any](method string, params P, files ...UploaderFile) UploaderRequest[R, P] { func NewUploaderRequest[R, P any](method string, params P, files ...UploaderFile) UploaderRequest[R, P] {
return UploaderRequest[R, P]{method, files, params} return UploaderRequest[R, P]{method, files, params}
} }
func (u UploaderRequest[R, P]) DoWithContext(ctx context.Context, up *Uploader) (R, error) { func (r UploaderRequest[R, P]) doRequest(ctx context.Context, up *Uploader) (R, error) {
var zero R var zero R
url := fmt.Sprintf("https://api.telegram.org/bot%s/%s", up.api.token, u.method) if up.api.limiter != nil {
if up.api.dropOverflowLimit {
buf := bytes.NewBuffer(nil) if !up.api.limiter.Allow() {
w := multipart.NewWriter(buf) return zero, errors.New("rate limited")
for _, file := range u.files {
fw, err := w.CreateFormFile(string(file.field), file.filename)
if err != nil {
_ = w.Close()
return zero, err
} }
_, err = fw.Write(file.data) } else {
if err != nil { if err := up.api.limiter.Wait(ctx); err != nil {
_ = w.Close()
return zero, err return zero, err
} }
} }
err := utils.Encode(w, u.params)
if err != nil {
_ = w.Close()
return zero, err
} }
err = w.Close()
buf, contentType, err := prepareMultipart(r.files, r.params)
if err != nil { if err != nil {
return zero, err return zero, err
} }
methodPrefix := ""
if up.api.useTestServer {
methodPrefix = "/test"
}
url := fmt.Sprintf("%s/bot%s%s/%s", up.api.apiUrl, up.api.token, methodPrefix, r.method)
req, err := http.NewRequestWithContext(ctx, "POST", url, buf) req, err := http.NewRequestWithContext(ctx, "POST", url, buf)
if err != nil { if err != nil {
return zero, err return zero, err
} }
req.Header.Set("Content-Type", w.FormDataContentType()) req.Header.Set("Content-Type", contentType)
req.Header.Set("Accept", "application/json") req.Header.Set("Accept", "application/json")
req.Header.Set("User-Agent", fmt.Sprintf("Laniakea/%s", utils.VersionString)) req.Header.Set("User-Agent", fmt.Sprintf("Laniakea/%s", utils.VersionString))
up.logger.Debugln("UPLOADER REQ", u.method) up.logger.Debugln("UPLOADER REQ", r.method)
res, err := up.api.client.Do(req) res, err := up.api.client.Do(req)
if err != nil { if err != nil {
return zero, err return zero, err
} }
defer res.Body.Close() defer res.Body.Close()
body, err := readBody(res.Body)
up.logger.Debugln("UPLOADER RES", r.method, string(body))
if res.StatusCode != http.StatusOK { if res.StatusCode != http.StatusOK {
return zero, fmt.Errorf("unexpected status code: %d", res.StatusCode) return zero, fmt.Errorf("unexpected status code: %d, %s", res.StatusCode, string(body))
} }
reader := io.LimitReader(res.Body, 10<<20) return parseBody[R](body)
body, err := io.ReadAll(reader)
if err != nil {
return zero, err
}
up.logger.Debugln("UPLOADER RES", u.method, string(body))
var resp ApiResponse[R]
err = json.Unmarshal(body, &resp)
if err != nil {
return zero, err
}
if !resp.Ok {
return zero, fmt.Errorf("[%d] %s", resp.ErrorCode, resp.Description)
}
return resp.Result, nil
} }
func (u UploaderRequest[R, P]) Do(up *Uploader) (R, error) { func (r UploaderRequest[R, P]) DoWithContext(ctx context.Context, up *Uploader) (R, error) {
return u.DoWithContext(context.Background(), up) var zero R
result, err := up.api.pool.Submit(ctx, func(ctx context.Context) (any, error) {
return r.doRequest(ctx, up)
})
if err != nil {
return zero, err
}
select {
case <-ctx.Done():
return zero, ctx.Err()
case res := <-result:
if res.Err != nil {
return zero, res.Err
}
if val, ok := res.Value.(R); ok {
return val, nil
}
return zero, ErrPoolUnexpected
}
}
func (r UploaderRequest[R, P]) Do(up *Uploader) (R, error) {
return r.DoWithContext(context.Background(), up)
}
func prepareMultipart[P any](files []UploaderFile, params P) (*bytes.Buffer, string, error) {
buf := bytes.NewBuffer(nil)
w := multipart.NewWriter(buf)
for _, file := range files {
fw, err := w.CreateFormFile(string(file.field), file.filename)
if err != nil {
_ = w.Close()
return buf, w.FormDataContentType(), err
}
_, err = fw.Write(file.data)
if err != nil {
_ = w.Close()
return buf, w.FormDataContentType(), err
}
}
err := utils.Encode(w, params)
if err != nil {
_ = w.Close()
return buf, w.FormDataContentType(), err
}
err = w.Close()
return buf, w.FormDataContentType(), err
} }
func uploaderTypeByExt(filename string) UploaderFileType { func uploaderTypeByExt(filename string) UploaderFileType {
+2 -3
View File
@@ -3,14 +3,13 @@ package laniakea
import "git.nix13.pw/scuroneko/laniakea/utils" import "git.nix13.pw/scuroneko/laniakea/utils"
func Ptr[T any](v T) *T { return &v } func Ptr[T any](v T) *T { return &v }
func Val[T any](p *T, def T) T { func Val[T any](p *T, def T) T {
if p != nil { if p != nil {
return *p return *p
} }
return def return def
} }
func EscapeMarkdown(s string) string { return utils.EscapeMarkdown(s) }
func EscapeMarkdownV2(s string) string { return utils.EscapeMarkdownV2(s) }
const VersionString = utils.VersionString const VersionString = utils.VersionString
var EscapeMarkdown = utils.EscapeMarkdown
+5 -4
View File
@@ -1,8 +1,9 @@
package utils package utils
const ( const (
VersionString = "0.6.1" VersionString = "1.0.0-beta.4"
VersionMajor = 0 VersionMajor = 1
VersionMinor = 6 VersionMinor = 0
VersionPatch = 1 VersionPatch = 0
Beta = 4
) )