package alerting import ( "bytes" "context" "crypto/sha256" "encoding/json" "errors" "fmt" "log/slog" "net/http" "sync" "time" "github.com/jackc/pgx/v5/pgxpool" "github.com/netpulse/netpulse/server/internal/crypto" "github.com/netpulse/netpulse/server/internal/store" ) // Приймач натискань кнопок Telegram. // // ЧОМУ ДОВГЕ ОПИТУВАННЯ, А НЕ ВЕБХУК // // Bot API дає два способи отримувати оновлення, і вибір тут зробило // саме розгортання, а не смак. // // Вебхук вимагає, щоб Telegram МІГ ДО НАС ДОСТУКАТИСЬ: публічний // порт із переліку 443/80/88/8443 і TLS-сертифікат, якому довіряє // їхній бік. Самопідписаний приймається лише як завантажений у // setWebhook файл, і навіть тоді потрібне ім'я, на яке він виданий. // Наш стенд — самопідписаний TLS на голій IP-адресі без домену. Це не // «поки не налаштували», а стан, у якому продукт живе: self-hosted // інсталяція в мережі оператора зазвичай узагалі не має входу ззовні. // Вебхук там не запрацює ніколи, і код, написаний під нього, був би // кодом, який не працює в жодній наявній інсталяції. // // Довге опитування не вимагає від нас ані вхідного порту, ані імені, // ані сертифіката: з'єднання ініціює сервер, TLS перевіряється в бік // api.telegram.org, тобто в той бік, де сертифікат справжній. Ціна — // одне висяче HTTP-з'єднання на бота й курсор у базі (0061). // // Секретний токен у заголовку X-Telegram-Bot-Api-Secret-Token — це // захист вебхука від сторонніх POST-ів на наш відкритий шлях. Тут // відкритого шляху немає взагалі: приймати нема чого, ми самі ходимо // по оновлення й показуємо в URL токен бота. Отвору, який той заголовок // затуляє, у цій схемі не існує. // // Якщо колись з'явиться домен і справжній сертифікат, вебхук стане // кращим (менше з'єднань, менша затримка) — і перевірка натискання // (chatMatch → прив'язка → права → дія) переїде в нього без змін: вона // навмисно не знає, звідки прийшло оновлення. // Bot читає оновлення ботів усіх кабінетів і виконує натиснуте. type Bot struct { st *store.Store ring *crypto.Keyring log *slog.Logger // hc — для довгого опитування. Таймаут свідомо більший за // pollTimeout: getUpdates мовчить рівно стільки, скільки просили, // і клієнт, який рветься раніше, перетворював би штатне очікування // на потік помилок. hc *http.Client // action — для коротких викликів (відповідь на натискання, // редагування повідомлення). Окремий клієнт, бо чекати на них 40 // секунд немає жодного сенсу. action *http.Client pollTimeout time.Duration } // NewBot створює приймач. ring обов'язковий: без ключів шифрування // токени ботів не розшифрувати, а отже й опитувати нікого. func NewBot(st *store.Store, ring *crypto.Keyring, log *slog.Logger) *Bot { poll := 25 * time.Second return &Bot{ st: st, ring: ring, log: log.With("component", "telegram"), hc: &http.Client{Timeout: poll + 15*time.Second}, action: &http.Client{Timeout: 15 * time.Second}, pollTimeout: poll, } } // telegramLockKey — довільна стала, аби її не займав ніхто інший у цій // же базі. Сусідня з ключем движка алертів (див. engine.go). const telegramLockKey = 0x6e70_7467 // "nptg" // Run тримає опитування до скасування контексту. // // Advisory-блокування береться на ВЕСЬ час роботи, а не на такт, як у // движка алертів. Причина в тому, що getUpdates ексклюзивний: вибране // оновлення другому читачеві вже не дістанеться, і два процеси на // одному боті ділили б натискання між собою навпіл. Блокування живе // разом із з'єднанням, тож падіння процесу звільняє його само — сусід // підхопить опитування за пів хвилини. func (b *Bot) Run(ctx context.Context) { if b.ring == nil { b.log.Info("приймач Telegram вимкнено: немає ключів шифрування") return } b.log.Info("приймач натискань Telegram запущено") for ctx.Err() == nil { conn, held := b.acquire(ctx) if !held { select { case <-ctx.Done(): return case <-time.After(30 * time.Second): continue } } b.serve(ctx, conn) // context.WithoutCancel: під час зупинки контекст уже мертвий, // а блокування зняти все одно треба — інакше сусідній процес // чекатиме на нього до розриву з'єднання. _, _ = conn.Exec(context.WithoutCancel(ctx), `SELECT pg_advisory_unlock($1)`, int64(telegramLockKey)) conn.Release() } } func (b *Bot) acquire(ctx context.Context) (*pgxpool.Conn, bool) { // WorkerPool, а не Pool: блокування має жити стільки ж, скільки // саме опитування, а опитування за побудовою ходить поверх усіх // кабінетів — це та сама роль, що й у решти фонових тактів. conn, err := b.st.WorkerPool().Acquire(ctx) if err != nil { return nil, false } var got bool if err := conn.QueryRow(ctx, `SELECT pg_try_advisory_lock($1)`, int64(telegramLockKey)).Scan(&got); err != nil || !got { conn.Release() return nil, false } return conn, true } // serve крутить такти, доки тримається блокування. func (b *Bot) serve(ctx context.Context, conn *pgxpool.Conn) { for ctx.Err() == nil { // Перелік ботів перечитується щотакту. Такт — це майже завжди // очікування на getUpdates, тобто раз на ~25 секунд, і за цю // ціну щойно доданий канал починає слухати кнопки сам, без // перезапуску процесу. groups, err := b.groups(ctx) if err != nil { b.log.Error("читання каналів Telegram", "помилка", err) select { case <-ctx.Done(): return case <-time.After(30 * time.Second): continue } } if len(groups) == 0 { // Жодного telegram-каналу: спати довше, ніж такт опитування. select { case <-ctx.Done(): return case <-time.After(60 * time.Second): continue } } var wg sync.WaitGroup for _, g := range groups { wg.Add(1) go func(g botGroup) { defer wg.Done() b.pollOnce(ctx, g) }(g) } wg.Wait() // Живе з'єднання — доказ, що блокування ще наше. Мертве означає, // що Postgres його вже зняв і опитувати далі не можна: сусідній // процес міг узяти бота собі. if err := conn.Ping(ctx); err != nil { b.log.Warn("з'єднання з блокуванням втрачено", "помилка", err) return } } } // botGroup — один бот і всі канали, які через нього шлють. // // Групування саме за токеном, а не за каналом: один бот цілком може // обслуговувати кілька чатів і навіть кілька кабінетів, а getUpdates // у нього одна черга на всіх. type botGroup struct { token string hash []byte chans []store.Channel } func (b *Bot) groups(ctx context.Context) ([]botGroup, error) { channels, err := b.st.TelegramChannels(ctx, b.ring) if err != nil { return nil, err } byToken := map[string]*botGroup{} var out []botGroup for _, c := range channels { if c.Secret == "" { continue } g, ok := byToken[c.Secret] if !ok { sum := sha256.Sum256([]byte(c.Secret)) g = &botGroup{token: c.Secret, hash: sum[:]} byToken[c.Secret] = g } g.chans = append(g.chans, c) } for _, g := range byToken { out = append(out, *g) } return out, nil } // pollOnce робить один getUpdates і обробляє все, що прийшло. func (b *Bot) pollOnce(ctx context.Context, g botGroup) { offset, err := b.st.TelegramCursor(ctx, g.hash) if err != nil { b.log.Error("читання курсора", "помилка", err) return } updates, err := b.getUpdates(ctx, g, offset) if err != nil { if ctx.Err() != nil { return } b.log.Warn("getUpdates", "бот", g.chans[0].Name, "помилка", err) // Пауза після помилки: без неї недоступний api.telegram.org // перетворював би такт на щільний цикл запитів. select { case <-ctx.Done(): case <-time.After(10 * time.Second): } return } var next int64 for _, u := range updates { if u.UpdateID >= next { next = u.UpdateID + 1 } b.handle(ctx, g, u) } if next == 0 { return } // Курсор посувається НЕЗАЛЕЖНО від того, чи вдалася сама дія. // Оновлення, на якому обробник спіткнувся, інакше приходило б знову // й знову, і одна крива кнопка глушила б усі наступні назавжди. // Людина при цьому не лишається без відповіді: невдача їй сказана // текстом у answerCallbackQuery. if err := b.st.SaveTelegramCursor(context.WithoutCancel(ctx), g.hash, next); err != nil { b.log.Error("збереження курсора", "помилка", err) } } func (b *Bot) getUpdates(ctx context.Context, g botGroup, offset int64) ([]tgUpdate, error) { body := map[string]any{ "timeout": int(b.pollTimeout.Seconds()), // Просимо рівно ті два типи, які вміємо: натискання кнопок і // повідомлення з командою прив'язки. Решта (правки, реакції, // вступи в чат) не має навіть потрапляти в чергу — вона займала // б місце й змушувала б нас її вичитувати. "allowed_updates": []string{"callback_query", "message"}, } if offset > 0 { body["offset"] = offset } var out struct { OK bool `json:"ok"` Description string `json:"description"` Result []tgUpdate `json:"result"` } if err := b.call(ctx, b.hc, g.token, "getUpdates", body, &out); err != nil { return nil, err } if !out.OK { return nil, fmt.Errorf("%s", out.Description) } return out.Result, nil } // handle розводить оновлення по обробниках. func (b *Bot) handle(ctx context.Context, g botGroup, u tgUpdate) { switch { case u.CallbackQuery != nil: b.handleCallback(ctx, g, u.CallbackQuery) case u.Message != nil: b.handleMessage(ctx, g, u.Message) } } // --------------------------------------------------------------------- // Натискання кнопки // --------------------------------------------------------------------- // handleCallback виконує натиснуте. // // Порядок перевірок навмисно такий: спершу «чий це чат» (звідси // кабінет), потім «що просять» (розбір callback_data), потім «хто саме // натиснув» (прив'язка), потім «чи можна йому» (права й доступ до // хоста) — і лише тоді дія. Жоден крок не бере кабінет чи особу з // вмісту кнопки: підробити її може будь-хто, хто бачив формат. func (b *Bot) handleCallback(ctx context.Context, g botGroup, cq *tgCallbackQuery) { answer := "Не вдалося обробити" // Відповідь на натискання обов'язкова й безумовна. Доки її немає, // Telegram крутить на кнопці годинник — і людина бачить не // «відмовлено», а «зламалось». Тому вона в defer, а не в кінці // щасливого шляху, і йде з власним контекстом: під час зупинки // процесу натискання все одно має отримати відповідь. defer func() { ansCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) defer cancel() if err := b.answerCallback(ansCtx, g.token, cq.ID, answer); err != nil { b.log.Warn("answerCallbackQuery", "помилка", err) } }() if cq.Message == nil { answer = "Повідомлення застаріле — відкрийте алерт у NetPulse" return } ch, ok := matchChannel(g.chans, cq.Message.Chat.ID, cq.Message.ThreadID) if !ok { // Бот стоїть у чаті, якого немає в жодному каналі. Кабінет із // такого натискання не виводиться ніяк, і вгадувати його за // вмістом кнопки — рівно те, чого робити не можна. answer = "Цей чат не налаштовано в NetPulse" b.log.Warn("натискання з невідомого чату", "chat", cq.Message.Chat.ID) return } act, err := parseCallbackData(cq.Data) if err != nil { answer = "Кнопка застаріла або невідома" b.log.Warn("розбір callback_data", "канал", ch.Name, "помилка", err) return } acc, err := b.st.TelegramAccountByTgID(ctx, ch.TenantID, cq.From.ID) if errors.Is(err, store.ErrTelegramNotLinked) { // Текст називає те, що людина побачить на екрані, а не назву // сторінки в коді. Попередній варіант відсилав до «Профіль → // Telegram», і саме так це вперше й не спрацювало: пункту // «Профіль» у меню не було, посилання в шапці не малювалось // (власник входить іменем, а не поштою, а /me імені не // віддавало), і людина, яка сумлінно виконала інструкцію, // доходила до висновку, що зламалось. answer = "Ваш Telegram не прив'язано до NetPulse.\n" + "У NetPulse: меню → Мій профіль → картка «Telegram» → " + "«Отримати код», далі надішліть цьому боту: /link КОД" return } if err != nil { b.log.Error("пошук прив'язки", "канал", ch.Name, "помилка", err) return } perms, err := b.st.UserPermissions(ctx, acc.UserID, ch.TenantID) if err != nil { b.log.Error("права користувача", "помилка", err) return } if !hasPerm(perms, "alerts:ack") { // Порожній набір прав означає ще й відкликане членство: людину // прибрали з кабінету, а прив'язка лишилась. Відповідь однакова // навмисно — з боку Telegram це та сама відмова. answer = "Немає права підтверджувати алерти" return } sc, err := b.st.LoadScope(ctx, ch.TenantID, acc.UserID) if err != nil { b.log.Error("доступ до хостів", "помилка", err) return } var line string switch act.Kind { case "ack": answer, line = b.doAck(ctx, ch, acc, sc, act.ID) case "mute": answer, line = b.doMute(ctx, ch, acc, sc, act.ID) } if line == "" { return } _ = b.st.TouchTelegramAccount(ctx, ch.TenantID, acc.ID) // Правка самого повідомлення — не прикраса. answerCallbackQuery // показує спливаючий рядок на кілька секунд і тому, хто натиснув; // у чат він не потрапляє, а чат читає вся зміна. Без правки // повідомлення про аварію так і лишається з живими кнопками, і // наступний черговий натискає їх ще раз. if err := b.editMessage(ctx, g.token, cq.Message, line); err != nil { b.log.Warn("правка повідомлення", "помилка", err) } } // doAck підтверджує алерт від імені прив'язаного користувача. // // Викликає той самий store.AckAlert, що й POST /api/v1/alerts/{id}/ack: // підтвердження з телефона й підтвердження з браузера мають лишати в // базі однаковий слід, а власна копія логіки розійшлася б із оригіналом // на першій же зміні — і розбіжність побачили б не тут, а в звіті. func (b *Bot) doAck(ctx context.Context, ch store.Channel, acc store.TelegramAccount, sc store.Scope, alertID string) (answer, line string) { cur, err := b.st.AlertAckState(ctx, ch.TenantID, alertID) if errors.Is(err, store.ErrAlertNotFound) { // Кабінет узято з чату, тож «не знайдено» тут означає саме // «немає в цьому кабінеті» — зокрема й тоді, коли алерт із // таким id є в чужому. return "Алерт не знайдено", "" } if err != nil { b.log.Error("читання алерту", "помилка", err) return "Не вдалося прочитати алерт", "" } if cur.DeviceID != "" && !sc.CanWrite(cur.DeviceID) { return "Немає доступу до цього хоста", "" } // Ідемпотентність. Друге натискання не має ні падати помилкою, ні // переписувати автора: перший, хто взявся, лишається першим. if cur.State == "acknowledged" { who := cur.AckedByEmail if who == "" { who = "невідомо ким" } at := time.Now() if cur.AckedAt != nil { at = *cur.AckedAt } return "Уже підтверджено: " + who, ackLine(who, at) } if cur.State == "resolved" || cur.State == "expired" { return "Алерт уже закрито", "" } a, err := b.st.AckAlert(ctx, ch.TenantID, alertID, acc.UserID, "підтверджено з Telegram") if errors.Is(err, store.ErrAlertNotFound) { // Хтось встиг підтвердити між читанням і записом — для людини // це той самий результат, що й гілка вище. return "Уже підтверджено", "" } if err != nil { b.log.Error("підтвердження алерту", "помилка", err) return "Не вдалося підтвердити", "" } at := time.Now() if a.AckedAt != nil { at = *a.AckedAt } return "Підтверджено", ackLine(acc.Email, at) } // doMute глушить хост на годину — тією ж дією, що й POST /api/v1/mutes. func (b *Bot) doMute(ctx context.Context, ch store.Channel, acc store.TelegramAccount, sc store.Scope, deviceID string) (answer, line string) { name, err := b.st.DeviceNameInTenant(ctx, ch.TenantID, deviceID) if errors.Is(err, store.ErrNotFound) { return "Хост не знайдено", "" } if err != nil { b.log.Error("пошук хоста", "помилка", err) return "Не вдалося знайти хост", "" } if !sc.CanWrite(deviceID) { return "Немає доступу до хоста " + name, "" } // Ідемпотентність: уже заглушений хост не глушиться вдруге. // Інакше подвійне натискання мовчки подвоювало б час тиші, і // дізнались би про це аж тоді, коли алерт не прийшов. if until, muted, err := b.st.ActiveMute(ctx, ch.TenantID, deviceID); err == nil && muted { return "Уже заглушено до " + until.In(tgLocation).Format("15:04"), muteLine(acc.Email, until) } until := time.Now().Add(time.Hour) if max := time.Now().Add(store.MaxMute); until.After(max) { until = max } if err := b.st.MuteDevice(ctx, ch.TenantID, deviceID, acc.UserID, "заглушено з Telegram", until); err != nil { b.log.Error("заглушення хоста", "помилка", err) return "Не вдалося заглушити", "" } return "Заглушено до " + until.In(tgLocation).Format("15:04"), muteLine(acc.Email, until) } func hasPerm(perms []string, want string) bool { for _, p := range perms { if p == want || p == "*" { return true } } return false } // --------------------------------------------------------------------- // Прив'язка акаунта // --------------------------------------------------------------------- // handleMessage відповідає лише на дві команди й мовчить на решту. // // Бот часто стоїть у робочому груповому чаті. Відповідь на кожне // повідомлення зробила б його джерелом шуму — і першою реакцією // команди стало б вимкнути сповіщення того чату, тобто рівно те, чому // продукт має запобігати. func (b *Bot) handleMessage(ctx context.Context, g botGroup, m *tgMessage) { if m.From == nil || m.From.IsBot { return } code, isLink := parseLinkCommand(m.Text) if !isLink { return } if code == "" { b.reply(ctx, g.token, m, "Надішліть код із профілю NetPulse: /link КОД") return } // Кабінети, яким належить цей бот. Без цього переліку код був би // універсальним: чинний код кабінету А, надісланий боту кабінету Б, // прив'язав би людину туди, де її бот навіть не стоїть. seen := map[string]bool{} var tenants []string for _, c := range g.chans { if !seen[c.TenantID] { seen[c.TenantID] = true tenants = append(tenants, c.TenantID) } } acc, err := b.st.RedeemTelegramLinkCode(ctx, code, tenants, m.From.ID, m.From.Username, tgDisplayName(*m.From)) if errors.Is(err, store.ErrTelegramLinkInvalid) { b.reply(ctx, g.token, m, "Код недійсний, вже використаний або прострочений. "+ "Візьміть новий у профілі NetPulse.") return } if err != nil { b.log.Error("прив'язка telegram", "помилка", err) b.reply(ctx, g.token, m, "Не вдалося прив'язати. Спробуйте пізніше.") return } _ = b.st.WriteAudit(ctx, acc.TenantID, store.AuditEntry{ ActorUserID: acc.UserID, Action: store.AuditActionTelegramLink, ObjectType: store.AuditObjectTelegram, ObjectID: acc.ID, Meta: map[string]any{"tg_user_id": acc.TgUserID, "tg_username": acc.TgUsername}, }) b.reply(ctx, g.token, m, "Готово: кнопки під алертами тепер працюють від вашого імені. "+ "Повідомлення з кодом можна видалити — код уже зужито.") } // --------------------------------------------------------------------- // Виклики Bot API // --------------------------------------------------------------------- func (b *Bot) answerCallback(ctx context.Context, token, queryID, text string) error { // show_alert=false: спливаючий рядок замість вікна з кнопкою «ОК». // Черговий тримає телефон однією рукою, і зайве підтвердження на // кожне натискання коштувало б рівно стільки ж, скільки економить // сама кнопка. body := map[string]any{"callback_query_id": queryID, "text": text} return b.call(ctx, b.action, token, "answerCallbackQuery", body, nil) } // editMessage дописує підсумок у повідомлення й прибирає кнопки. func (b *Bot) editMessage(ctx context.Context, token string, m *tgMessage, line string) error { text := withStatus(m.Text, line) if m.Text == "" { text = line } body := map[string]any{ "chat_id": m.Chat.ID, "message_id": m.MessageID, "text": text, // parse_mode навмисно НЕ задається, хоч надсилали ми з HTML. // Telegram віддає в message.text уже готовий текст без розмітки, // і повторна відправка його як HTML або зламалася б на першому // «<» у назві інтерфейсу, або перетворила б частину тексту // алерту на теги. // // Порожній inline_keyboard замість пропуску поля: так кнопки // зникають гарантовано, а не за замовчуванням, на яке довелось // би покладатися. "reply_markup": map[string]any{"inline_keyboard": [][]any{}}, } return b.call(ctx, b.action, token, "editMessageText", body, nil) } func (b *Bot) reply(ctx context.Context, token string, m *tgMessage, text string) { body := map[string]any{"chat_id": m.Chat.ID, "text": text} if m.ThreadID != 0 { body["message_thread_id"] = m.ThreadID } if err := b.call(ctx, b.action, token, "sendMessage", body, nil); err != nil { b.log.Warn("відповідь боту", "помилка", err) } } // call — один виклик Bot API. // // Токен іде в шляху URL (так вимагає Bot API), тому будь-яка помилка // транспорту проходить через scrubToken: http.Client вкладає в її текст // повний URL, а в журналі токен бота — це чинний доступ. func (b *Bot) call(ctx context.Context, hc *http.Client, token, method string, body, out any) error { payload, err := json.Marshal(body) if err != nil { return err } req, err := http.NewRequestWithContext(ctx, http.MethodPost, "https://api.telegram.org/bot"+token+"/"+method, bytes.NewReader(payload)) if err != nil { return scrubToken(err, token) } req.Header.Set("Content-Type", "application/json") res, err := hc.Do(req) if err != nil { return scrubToken(err, token) } defer res.Body.Close() if out == nil { // Тіло відповіді нікому не потрібне, але прочитати його треба: // недочитане з'єднання не повертається в keep-alive, а на // довгому опитуванні це нове TLS-рукостискання щохвилини. var sink struct { OK bool `json:"ok"` Description string `json:"description"` } if err := json.NewDecoder(res.Body).Decode(&sink); err != nil { return fmt.Errorf("%s: відповідь %d нерозбірлива", method, res.StatusCode) } if !sink.OK { return fmt.Errorf("%s: %s", method, sink.Description) } return nil } if err := json.NewDecoder(res.Body).Decode(out); err != nil { return fmt.Errorf("%s: відповідь %d нерозбірлива", method, res.StatusCode) } return nil }