| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195 |
- // Command telegramd is the long-polling bot loop for
- // broad-announce M3 (SPEC §8). It pulls the list of
- // active telegram_bots from Postgres, then for each bot
- // long-polls getUpdates, parses commands, applies them
- // via internal/telegram.Handler, and replies via
- // SendMessage.
- //
- // M3 ships long-polling only. Webhook mode is M5/M9.
- //
- // Configuration:
- // BA_TELEGRAM_BOT_TOKEN — fallback if DB has no row
- // BA_TELEGRAM_BASE_URL — override for tests (default
- // https://api.telegram.org).
- // M3 dev: point at faketgmd.
- //
- // One process polls all bots. With M3's one-bot-per-company
- // this is a single update loop; M3+ can shard by company
- // hash if the volume warrants.
- package main
- import (
- "context"
- "log/slog"
- "net/http"
- "os"
- "os/signal"
- "syscall"
- "time"
- "git3.techno-world.net/lrosales/broad-announce/internal/broker"
- "git3.techno-world.net/lrosales/broad-announce/internal/config"
- "git3.techno-world.net/lrosales/broad-announce/internal/httpserver"
- "git3.techno-world.net/lrosales/broad-announce/internal/observability"
- "git3.techno-world.net/lrosales/broad-announce/internal/postgres"
- "git3.techno-world.net/lrosales/broad-announce/internal/telegram"
- )
- type botConfig struct {
- BotID string
- CompanyID string
- BotToken string
- }
- func main() {
- cfg, err := config.LoadCommon("telegramd")
- if err != nil {
- os.Stderr.WriteString("config: " + err.Error() + "\n")
- os.Exit(1)
- }
- logger := observability.Init(cfg.Env, cfg.LogLevel, "telegramd")
- logger.Info("starting", "env", cfg.Env, "addr", cfg.HTTPAddr)
- ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
- defer stop()
- br, err := broker.Connect(ctx, cfg.NATSURL)
- if err != nil {
- logger.Error("nats connect", "err", err)
- os.Exit(1)
- }
- defer br.Close()
- _ = br // M3: long-polling only; NATS used later for ack fan-out
- pool, err := postgres.Connect(ctx, cfg.PostgresDSN)
- if err != nil {
- logger.Error("postgres connect", "err", err)
- os.Exit(1)
- }
- defer pool.Close()
- baseURL := os.Getenv("BA_TELEGRAM_FAKE_URL")
- if baseURL == "" {
- baseURL = os.Getenv("BA_TELEGRAM_BASE_URL")
- }
- if baseURL == "" {
- baseURL = "http://faketgmd:8830" // M3 docker-compose default
- }
- logger.Info("telegram target", "base_url", baseURL)
- client := &telegram.HTTPBotClient{
- BaseURL: baseURL,
- HTTP: &http.Client{Timeout: 60 * time.Second},
- }
- handler := telegram.NewHandler(pool, logger.With("subsystem", "telegram"))
- runCtx, runCancel := context.WithCancel(ctx)
- defer runCancel()
- go pollAllBots(runCtx, logger, pool, client, handler)
- reg, _ := observability.NewRegistry("telegramd")
- srv := httpserver.New(httpserver.Config{
- Addr: cfg.HTTPAddr,
- ServiceName: "telegramd",
- ShutdownGrace: cfg.ShutdownGrace,
- }, logger, observability.MetricsHandler(reg))
- errCh := make(chan error, 1)
- go func() { errCh <- srv.Start() }()
- select {
- case <-ctx.Done():
- logger.Info("shutdown signal received")
- case err := <-errCh:
- if err != nil {
- logger.Error("http server", "err", err)
- os.Exit(1)
- }
- }
- runCancel()
- time.Sleep(500 * time.Millisecond)
- if err := srv.Shutdown(ctx); err != nil {
- logger.Warn("graceful shutdown", "err", err)
- }
- logger.Info("bye")
- }
- // pollAllBots loads the bot list once at startup, then
- // long-polls each in a goroutine. M3: a single reload on
- // SIGUSR1 is overkill; restarting the binary picks up new
- // bots. M5 can add a watch on the table.
- func pollAllBots(ctx context.Context, logger *slog.Logger, pool *postgres.Pool, client telegram.BotClient, h *telegram.Handler) {
- bots, err := loadBots(ctx, pool)
- if err != nil {
- logger.Error("load bots", "err", err)
- return
- }
- logger.Info("loaded bots", "count", len(bots))
- if len(bots) == 0 {
- // No bots yet. Block on ctx so the process stays up.
- <-ctx.Done()
- return
- }
- for _, b := range bots {
- go pollOneBot(ctx, logger.With("bot", b.BotID, "company", b.CompanyID), client, h, b)
- }
- <-ctx.Done()
- }
- func loadBots(ctx context.Context, pool *postgres.Pool) ([]botConfig, error) {
- rows, err := pool.Query(ctx, `
- SELECT bot_id, company_id, bot_token
- FROM telegram_bots
- WHERE status = 'active'
- ORDER BY company_id, bot_id
- `)
- if err != nil {
- return nil, err
- }
- defer rows.Close()
- var out []botConfig
- for rows.Next() {
- var b botConfig
- if err := rows.Scan(&b.BotID, &b.CompanyID, &b.BotToken); err != nil {
- return nil, err
- }
- out = append(out, b)
- }
- return out, rows.Err()
- }
- func pollOneBot(ctx context.Context, logger *slog.Logger, client telegram.BotClient, h *telegram.Handler, b botConfig) {
- var offset int64
- for {
- if ctx.Err() != nil {
- return
- }
- updates, err := client.GetUpdates(ctx, b.BotToken, offset, 25)
- if err != nil {
- if ctx.Err() != nil {
- return
- }
- logger.Warn("getUpdates", "err", err)
- time.Sleep(2 * time.Second)
- continue
- }
- for _, u := range updates {
- if u.UpdateID >= offset {
- offset = u.UpdateID + 1
- }
- if u.Message == nil {
- continue
- }
- reply, err := h.Handle(ctx, u.Message)
- if err != nil {
- logger.Warn("command handle", "err", err, "update_id", u.UpdateID)
- continue
- }
- if reply == "" {
- continue
- }
- if _, err := client.SendMessage(ctx, b.BotToken, u.Message.Chat.ID, reply); err != nil {
- logger.Warn("sendMessage reply", "err", err, "chat_id", u.Message.Chat.ID)
- }
- }
- }
- }
|