// 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) } } } }