diff --git a/README.md b/README.md index 3086ee343..3d29e6c65 100644 --- a/README.md +++ b/README.md @@ -71,6 +71,9 @@ around the opposite principle — **prove it is safe, then act**: config validation, locks, processes, preflight, inventory, SLA and events. - **Notifications** to email, Slack, Teams and webhook sinks (ntfy/Telegram/ Gotify) with a templated default message. +- An optional **interactive Telegram report bot** (read-only): ask it `/status`, + `/services`, `/sla` and it replies with live reports — long polling only, no + inbound port, answering allow-listed chats. - A **daemon-wide panic switch** to pause all automatic remediation instantly. - **Guided wizards** for common setups (service, docker, vm, mount, volume, net, uplink). diff --git a/cmd/sermod/main.go b/cmd/sermod/main.go index 0d07c129c..1dd0ab243 100644 --- a/cmd/sermod/main.go +++ b/cmd/sermod/main.go @@ -34,6 +34,7 @@ import ( "sermo/internal/rules" "sermo/internal/servicemgr" "sermo/internal/state" + "sermo/internal/telegrambot" "sermo/internal/web" ) @@ -368,13 +369,20 @@ func run(args []string) int { } } + botCfg := telegrambot.ParseConfig(config.SectionMap(cfg.Global.Raw, config.SectionTelegramBot)) + var webHolder *app.WebBackendHolder - var webDone chan struct{} addr, webDisabledReason := webListenAddr(cfg) - if addr != "" { + // The web backend feeds both the dashboard and the report bot; build it when + // either is enabled, even if the HTTP server itself stays off. + if addr != "" || botCfg.Enabled { var webWarnings []string webHolder, webWarnings = app.NewWebBackendHolder(ctx, cfg, deps) app.LogBuildNotices(logger, "build web backend", webWarnings) + } + + var webDone chan struct{} + if addr != "" { auth := webAuth(cfg) server := &web.Server{ Addr: addr, @@ -410,6 +418,21 @@ func run(args []string) int { logger.Warn("web ui disabled; no port will be opened", logFieldReason, webDisabledReason) } + // Interactive read-only report bot (long polling; no inbound socket). It + // reads the same web backend the dashboard serves and replies to commands + // from allow-listed chats only. + var botDone chan struct{} + if botCfg.Enabled { + bot := telegrambot.New(app.NewTelegramReporter(webHolder, store, time.Now), botCfg, logger) + deps.TelegramBot = bot + botDone = make(chan struct{}) + go func() { + defer close(botDone) + bot.Run(ctx) + }() + logger.Info("telegram report bot enabled", "allowed_chats", len(botCfg.AllowedChats)) + } + pruneDone := startOldHistoryPrune(ctx, logger, store, time.Now().Add(-state.DefaultHistoryRetention)) logger.Info("sermod starting", logFieldBackend, detection.Backend, logFieldServices, len(workers), logFieldWatches, len(watches)) @@ -453,6 +476,9 @@ func run(args []string) int { if webDone != nil { <-webDone } + if botDone != nil { + <-botDone + } if !drainOrTimeout(pruneDone, shutdownPruneDrainTimeout) { logger.Warn("history prune still running at shutdown; closing the store without it") } diff --git a/docs/configuration.es.md b/docs/configuration.es.md index 35999c756..228529c1c 100644 --- a/docs/configuration.es.md +++ b/docs/configuration.es.md @@ -59,6 +59,7 @@ directorio equivocado. La configuración distribuida la omite. - [Notificaciones](#notificaciones) - [Plantillas de notificación](#plantillas-de-notificación) - [Selección por defecto y precedencia](#selección-por-defecto-y-precedencia) +- [Bot de informes de Telegram](#bot-de-informes-de-telegram) - [Host watches](#host-watches) - [then.expand — crecimiento de volumen (watch de storage)](#thenexpand--crecimiento-de-volumen-watch-de-storage) - [Control manual de reconstrucción RAID](#control-manual-de-reconstrucción-raid) @@ -1200,10 +1201,19 @@ notifiers: - **`telegram`** — envía a través de un **bot de Telegram** (`sendMessage`). - **`token`** — el token del bot de `@BotFather`. Queda dentro de la URL de la - API y nunca aparece en el dashboard. + API y nunca aparece en el dashboard. Prefiere `${env:...}`; si queda vacío + (la variable no está definida) el notifier queda inactivo en lugar de fallar + la carga de configuración. - **`chat_id`** — el id numérico del chat/grupo o un nombre `@canal`. El asunto es la línea principal y el detalle (los campos `SERMO_*`) sigue como texto plano. + - **`parse_mode`** *(opcional)* — `MarkdownV2`, `Markdown` o `HTML` para + renderizar el mensaje como texto con formato (negrita, código, enlaces) en + lugar de texto plano. Omítelo para texto plano. + - **`silent`** *(opcional)* — `true` entrega el mensaje en silencio, sin sonido + ni vibración (`disable_notification` de la Bot API). + - **`message_thread_id`** *(opcional)* — un id entero de tema (forum topic), + para publicar en un tema concreto de un grupo. ```yaml # /etc/sermo/notifiers/telegram.yml @@ -1212,6 +1222,9 @@ notifiers: type: telegram token: "123456789:AAF...XXXX" chat_id: -1001234567890 + # parse_mode: MarkdownV2 # opcional: formatea el texto del mensaje + # silent: true # opcional: entrega sin sonido + # message_thread_id: 42 # opcional: publica en un tema del grupo ``` - **`tty`** — escribe directamente en las sesiones de terminal Linux activas, similar a @@ -1346,6 +1359,60 @@ solo-alerta (estado de disparo + eventos en la interfaz y el log, pero sin accio herencia de los globales). Consulta la sección de host watches a continuación para el ejemplo de `check` + `for` desnudo. +## Bot de informes de Telegram + +La sección opcional de nivel superior **`telegram_bot`** ejecuta un bot de Telegram +interactivo y de **solo lectura** dentro de `sermod`. Mientras que un notifier +`telegram` *empuja* alertas, el bot permite al operador *pedir* informes bajo demanda: +envíale `/status` y responde con el resumen actual del fleet. Nunca puede cambiar el +host — solo lee el mismo estado que sirve el dashboard web. + +Recibe comandos mediante **long polling** de la Bot API (`getUpdates`), así que no +necesita puerto entrante, ni exposición pública, ni proxy inverso — en línea con la +postura de solo-salida de Sermo. El token del bot queda dentro de la URL de la API y se +depura de logs y errores, igual que el notifier `telegram`. + +```yaml +# /etc/sermo/sermo.yml (o un fragmento drop-in) +telegram_bot: + token: "${env:TELEGRAM_BOT_TOKEN}" # token del bot de @BotFather + allowed_chats: # obligatorio: solo se responde a estos chats + - 123456789 + - -1001234567890 + # poll_interval: 30s # timeout opcional del long-poll getUpdates + # enabled: false # opcional; omitido => habilitado +``` + +- **`token`** — el token del bot de `@BotFather` (usado tanto para `getUpdates` como + para las respuestas). Prefiere `${env:...}` para no escribirlo en un archivo. Opcional: + si queda vacío (la variable de entorno no está definida) el bot simplemente queda + inactivo y el resto de la configuración se carga igual. +- **`allowed_chats`** — la lista de ids numéricos de chat que el bot responde. Un + mensaje de cualquier otro chat se ignora, nunca se responde. Este es el control de + acceso: limítalo a los operadores/grupos que pueden consultar el daemon. +- **`poll_interval`** *(opcional)* — el timeout del long-poll, por defecto `30s`, + acotado a `1s`–`10m`. +- **`enabled`** *(opcional)* — ponlo en `false` para conservar la sección pero detener + el polling. + +Comandos (todos de solo lectura): + +| Comando | Respuesta | +| --- | --- | +| `/status` | resumen del fleet: services ok/fallando, monitorizados/pausados, errores recientes, uptime del host | +| `/services [name]` | la lista de services, o el estado y salud de un service concreto | +| `/watches` | estados de host watches y watches de servicio | +| `/sla ` | ventanas de disponibilidad (hora…año) de un service | +| `/events [count]` | los eventos más recientes (por defecto 10, máximo 50) | +| `/help` | la lista de comandos | + +Recargar (`sermoctl daemon reload` / `SIGHUP`) aplica cambios en `token`, +`allowed_chats` y `poll_interval` sin reiniciar. Como la goroutine solo se arranca al +inicio cuando la sección está presente, **habilitar el bot por primera vez requiere un +reinicio** (la misma regla que sigue la interfaz web para su puerto). Al arrancar, el bot +descarta cualquier comando encolado mientras estuvo caído, así que un reinicio nunca +reproduce solicitudes antiguas. + ## Host watches Los `watches` monitorizan recursos a nivel de host independientemente de cualquier diff --git a/docs/configuration.md b/docs/configuration.md index 752c378cf..e4314a389 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -57,6 +57,7 @@ configuration omits it. - [Notifications](#notifications) - [Notification templates](#notification-templates) - [Default selection and precedence](#default-selection-and-precedence) +- [Telegram report bot](#telegram-report-bot) - [Host watches](#host-watches) - [then.expand — volume growth (storage watch)](#thenexpand--volume-growth-storage-watch) - [Manual RAID reconstruction control](#manual-raid-reconstruction-control) @@ -1155,10 +1156,19 @@ notifiers: - **`telegram`** — sends through a **Telegram bot** (`sendMessage`). - **`token`** — the bot token from `@BotFather`. It stays inside the API URL - and never appears on the dashboard. + and never appears on the dashboard. Prefer `${env:...}`; when it resolves + empty (the variable is unset) the notifier stays inactive instead of failing + config load. - **`chat_id`** — the numeric chat/group id or an `@channel` name. The subject is the lead line and the detail (the `SERMO_*` fields) follows as plain text. + - **`parse_mode`** *(optional)* — `MarkdownV2`, `Markdown` or `HTML` to render + the message as formatted text (bold, code, links) instead of plain text. Omit + for plain text. + - **`silent`** *(optional)* — `true` delivers the message quietly, with no + sound or vibration (Bot API `disable_notification`). + - **`message_thread_id`** *(optional)* — an integer forum-topic id, to post + into a specific topic within a group. ```yaml # /etc/sermo/notifiers/telegram.yml @@ -1167,6 +1177,9 @@ notifiers: type: telegram token: "123456789:AAF...XXXX" chat_id: -1001234567890 + # parse_mode: MarkdownV2 # optional: format the message text + # silent: true # optional: deliver without sound + # message_thread_id: 42 # optional: post into a forum topic ``` - **`tty`** — writes directly to active Linux terminal sessions, similar to @@ -1296,6 +1309,59 @@ alert-only behaviour (firing state + events in the UI and log, but no actions and no inheritance of globals). See the host watches section below for the bare `check` + `for` example. +## Telegram report bot + +The optional top-level **`telegram_bot`** section runs an interactive, +**read-only** Telegram bot inside `sermod`. Where a `telegram` notifier *pushes* +alerts, the bot lets an operator *ask* for reports on demand: send it `/status` +and it replies with the current fleet summary. It can never change the host — it +only reads the same state the web dashboard serves. + +It receives commands over Bot API **long polling** (`getUpdates`), so it needs +no inbound port, no public exposure and no reverse proxy — matching Sermo's +outbound-only posture. The bot token stays inside the API URL and is scrubbed +from logs and errors, exactly like the `telegram` notifier. + +```yaml +# /etc/sermo/sermo.yml (or a drop-in fragment) +telegram_bot: + token: "${env:TELEGRAM_BOT_TOKEN}" # bot token from @BotFather + allowed_chats: # required: only these chats are answered + - 123456789 + - -1001234567890 + # poll_interval: 30s # optional getUpdates long-poll timeout + # enabled: false # optional; omitted => enabled +``` + +- **`token`** — the bot token from `@BotFather` (used for both `getUpdates` and + replies). Prefer `${env:...}` so it is not written in a file. Optional: when it + resolves empty (the env var is unset) the bot simply stays inactive and the + rest of the config still loads. +- **`allowed_chats`** — the list of numeric chat ids the bot answers. A message + from any other chat is ignored, never answered. This is the access control: + keep it to the operators/groups that may query the daemon. +- **`poll_interval`** *(optional)* — the long-poll timeout, default `30s`, + clamped to `1s`–`10m`. +- **`enabled`** *(optional)* — set `false` to keep the section but stop polling. + +Commands (all read-only): + +| Command | Reply | +| --- | --- | +| `/status` | fleet summary: services ok/failing, monitored/paused, recent errors, host uptime | +| `/services [name]` | the service list, or one named service's state and health | +| `/watches` | host and service watch states | +| `/sla ` | availability windows (hour…year) for a service | +| `/events [count]` | the most recent events (default 10, max 50) | +| `/help` | the command list | + +Reloading (`sermoctl daemon reload` / `SIGHUP`) applies changes to `token`, +`allowed_chats` and `poll_interval` without a restart. Because the goroutine is +only started at boot when the section is present, **enabling the bot for the +first time requires a restart** (the same rule the web UI follows for its port). +On startup the bot discards any commands queued while it was down, so a restart +never replays old requests. + ## Host watches `watches` monitor host-level resources independently of any service and run a diff --git a/docs/sermo-all.yml b/docs/sermo-all.yml index e9a7bc47a..8becc7442 100644 --- a/docs/sermo-all.yml +++ b/docs/sermo-all.yml @@ -120,6 +120,13 @@ notifiers: enabled: false # keep configured but skip delivery type: teams webhook: "https://prod-01.westeurope.logic.azure.com:443/workflows/x" + ops-telegram: + type: telegram + token: "${env:TELEGRAM_TOKEN:-123456789:AAF...XXXX}" + chat_id: "-1001234567890" + parse_mode: MarkdownV2 # optional: HTML, Markdown or MarkdownV2 + silent: true # optional: deliver without sound + message_thread_id: 42 # optional: post into a forum topic tty: type: tty users: [root] # optional; omit to notify every active terminal @@ -130,6 +137,15 @@ notifiers: # delivery per site (a watch whose only action is [none] is monitor-only). notify: [ops-email] +# Interactive read-only Telegram report bot (long polling; no inbound port). +# Distinct from a `telegram` notifier: it answers /status, /services, /sla, ... +# from allow-listed chats only, and can never change the host. +telegram_bot: + enabled: true + token: "${env:TELEGRAM_BOT_TOKEN:-123456789:AAF...XXXX}" + allowed_chats: [123456789, -1001234567890] + poll_interval: 30s + # Host watches: one entry per concern. This reference groups several examples # only to keep the schema readable; real config stores each watch as its own # `name:` document under any directory listed in paths.watches. Storage checks diff --git a/examples/notifiers/ops-telegram.yml b/examples/notifiers/ops-telegram.yml new file mode 100644 index 000000000..b8df85782 --- /dev/null +++ b/examples/notifiers/ops-telegram.yml @@ -0,0 +1,10 @@ +notifiers: + ops-telegram: + # enabled: false # keep configured but skip delivery attempts + type: telegram + # template: default-alert + token: "${env:TELEGRAM_TOKEN}" # bot token from @BotFather; kept inside the API URL, never surfaced + chat_id: "-1001234567890" # numeric chat/channel id or @channelname + # parse_mode: MarkdownV2 # HTML, Markdown or MarkdownV2; omit for plain text + # silent: true # deliver quietly, without sound or vibration + # message_thread_id: 42 # post into a specific forum topic within a group diff --git a/internal/app/daemon.go b/internal/app/daemon.go index 2a940ea5e..fdb587955 100644 --- a/internal/app/daemon.go +++ b/internal/app/daemon.go @@ -22,6 +22,7 @@ import ( "sermo/internal/rules" "sermo/internal/servicemgr" "sermo/internal/state" + "sermo/internal/telegrambot" "sermo/internal/web" ) @@ -288,6 +289,9 @@ type Deps struct { Events *EventLog // DiagnosticLog exports scheduled diagnostics to engine.diagnostics when set. DiagnosticLog *DiagnosticLog + // TelegramBot is the interactive read-only report bot. Optional: nil when the + // `telegram_bot` section is absent. It is refreshed on reload via UpdateConfig. + TelegramBot *telegrambot.Bot // SystemFreshness caches system metrics so concurrent workers in one cycle // share a computation; it must be below the scheduler interval. SystemFreshness time.Duration diff --git a/internal/app/monitor.go b/internal/app/monitor.go index e6602dd34..b8b81781c 100644 --- a/internal/app/monitor.go +++ b/internal/app/monitor.go @@ -14,6 +14,7 @@ import ( "sermo/internal/notify" "sermo/internal/process" "sermo/internal/rules" + "sermo/internal/telegrambot" ) // Monitor runs service workers and host watches in reloadable generations. @@ -196,6 +197,9 @@ func (m *Monitor) installGenerationLocked(ctx context.Context, newCfg *config.Co m.deps.DiagnosticLog.UpdateConfig(newCfg) go m.deps.DiagnosticLog.Export() } + if m.deps.TelegramBot != nil { + m.deps.TelegramBot.UpdateConfig(telegrambot.ParseConfig(config.SectionMap(newCfg.Global.Raw, config.SectionTelegramBot))) + } LogBuildNotices(m.Logger, "reload build", warnings) m.startGenerationLocked(ctx, false) diff --git a/internal/app/telegramreport.go b/internal/app/telegramreport.go new file mode 100644 index 000000000..ec590a8e0 --- /dev/null +++ b/internal/app/telegramreport.go @@ -0,0 +1,151 @@ +package app + +import ( + "context" + "fmt" + "strings" + "time" + + "sermo/internal/state" + "sermo/internal/telegrambot" + "sermo/internal/web" +) + +// Service check-health values surfaced by the web backend that the status +// rollup counts. +const ( + webCheckHealthOK = "ok" + webCheckHealthFailing = "failing" +) + +// telegramSLAReader is the slice of the state store the bot needs for /sla. +type telegramSLAReader interface { + SLAReport(service string, now time.Time) ([]state.SLAValue, error) +} + +// telegramReporter adapts the reload-safe web backend and the state store into +// the read-only view the Telegram bot renders. It reuses exactly the data the +// web dashboard serves, so bot reports never diverge from the UI. +type telegramReporter struct { + web *WebBackendHolder + sla telegramSLAReader + now func() time.Time +} + +// NewTelegramReporter builds a telegrambot.Reporter over the reload-safe web +// backend holder and an SLA source (the state store). now defaults to time.Now. +func NewTelegramReporter(webHolder *WebBackendHolder, sla telegramSLAReader, now func() time.Time) telegrambot.Reporter { + if now == nil { + now = time.Now + } + return &telegramReporter{web: webHolder, sla: sla, now: now} +} + +func (r *telegramReporter) Status(ctx context.Context) (telegrambot.StatusReport, error) { + snap := r.web.DashboardSnapshot(ctx, 0) + rep := telegrambot.StatusReport{ + Host: snap.Daemon.Hostname, + Services: len(snap.Services), + Monitored: snap.Monitoring.Monitored, + Paused: snap.Monitoring.Paused, + Errors: snap.Activity.Errors, + LastEvent: snap.Activity.LastEventKind, + HostUptime: snap.Daemon.HostUptime, + } + for _, s := range snap.Services { + switch s.CheckHealth { + case webCheckHealthOK: + rep.OK++ + case webCheckHealthFailing: + rep.Failing++ + } + } + return rep, nil +} + +func (r *telegramReporter) Services(ctx context.Context) ([]telegrambot.ServiceLine, error) { + services := r.web.Services(ctx) + lines := make([]telegrambot.ServiceLine, 0, len(services)) + for _, s := range services { + lines = append(lines, telegrambot.ServiceLine{ + Name: s.Name, + State: s.State, + Health: s.CheckHealth, + Monitored: s.Monitored, + }) + } + return lines, nil +} + +func (r *telegramReporter) Watches(ctx context.Context) ([]telegrambot.WatchLine, error) { + watches := r.web.Watches(ctx) + lines := make([]telegrambot.WatchLine, 0, len(watches)) + for _, w := range watches { + lines = append(lines, telegrambot.WatchLine{ + Name: w.Name, + Scope: w.Scope, + State: w.State, + Monitored: w.Monitored, + }) + } + return lines, nil +} + +func (r *telegramReporter) SLA(ctx context.Context, service string) ([]telegrambot.SLAWindow, bool, error) { + if !r.serviceExists(ctx, service) { + return nil, false, nil + } + values, err := r.sla.SLAReport(service, r.now()) + if err != nil { + return nil, true, err + } + windows := make([]telegrambot.SLAWindow, 0, len(values)) + for _, v := range values { + windows = append(windows, telegrambot.SLAWindow{Window: v.Window, Ratio: formatSLARatio(v)}) + } + return windows, true, nil +} + +func (r *telegramReporter) Events(ctx context.Context, limit int) ([]telegrambot.EventLine, error) { + events := r.web.Events(ctx, limit) + lines := make([]telegrambot.EventLine, 0, len(events)) + for _, e := range events { + lines = append(lines, telegrambot.EventLine{ + Time: e.Time, + Target: eventTarget(e), + Kind: e.Kind, + Message: e.Message, + }) + } + return lines, nil +} + +func (r *telegramReporter) serviceExists(ctx context.Context, name string) bool { + for _, s := range r.web.Services(ctx) { + if strings.EqualFold(s.Name, name) { + return true + } + } + return false +} + +func formatSLARatio(v state.SLAValue) string { + ratio, ok := v.Ratio() + if !ok { + return "n/a" + } + return fmt.Sprintf("%.2f%%", ratio*100) +} + +func eventTarget(e web.Event) string { + switch { + case e.Service != "": + return e.Service + case e.Watch != "": + return e.Watch + case e.App != "": + return e.App + default: + return "" + } +} diff --git a/internal/checks/firewall_nftables.go b/internal/checks/firewall_nftables.go index 7bb44fcf5..4814f9299 100644 --- a/internal/checks/firewall_nftables.go +++ b/internal/checks/firewall_nftables.go @@ -1,3 +1,5 @@ +//go:build linux + package checks import ( diff --git a/internal/checks/firewall_nftables_other.go b/internal/checks/firewall_nftables_other.go new file mode 100644 index 000000000..f15269c23 --- /dev/null +++ b/internal/checks/firewall_nftables_other.go @@ -0,0 +1,24 @@ +//go:build !linux + +package checks + +import ( + "context" + "errors" + + "sermo/internal/execx" +) + +// nftablesRuleCounter reads the loaded nftables rule count. Tests override it to +// avoid real netlink I/O. +var nftablesRuleCounter = countLoadedNftablesRules + +// countLoadedNftablesRules reports that nftables rule counting is unavailable off +// Linux, where the netlink API it relies on does not exist. It still honors a +// cancelled context so callers observe the same cancellation signal as on Linux. +func countLoadedNftablesRules(ctx context.Context) (uint64, error) { + if err := ctx.Err(); err != nil { + return 0, execx.ContextError(err) + } + return 0, errors.New("nftables rule counting is only supported on Linux") +} diff --git a/internal/checks/route.go b/internal/checks/route.go index cf4f0180b..9068cc8a7 100644 --- a/internal/checks/route.go +++ b/internal/checks/route.go @@ -118,9 +118,9 @@ func SampleRoutes(family string) ([]DefaultRoute, error) { return defaultRouteSa func netlinkFamily(family string) (int, error) { switch family { case FamilyIPv4, "": - return netlink.FAMILY_V4, nil + return familyV4, nil case FamilyIPv6: - return netlink.FAMILY_V6, nil + return familyV6, nil default: return 0, fmt.Errorf("unknown route family %q", family) } diff --git a/internal/checks/route_family_linux.go b/internal/checks/route_family_linux.go new file mode 100644 index 000000000..38071f862 --- /dev/null +++ b/internal/checks/route_family_linux.go @@ -0,0 +1,12 @@ +//go:build linux + +package checks + +import "github.com/vishvananda/netlink" + +// familyV4 and familyV6 are the netlink address-family selectors for the route +// query. On Linux they are the kernel constants netlink itself uses. +const ( + familyV4 = netlink.FAMILY_V4 + familyV6 = netlink.FAMILY_V6 +) diff --git a/internal/checks/route_family_other.go b/internal/checks/route_family_other.go new file mode 100644 index 000000000..69f0a88d1 --- /dev/null +++ b/internal/checks/route_family_other.go @@ -0,0 +1,11 @@ +//go:build !linux + +package checks + +// familyV4 and familyV6 mirror netlink.FAMILY_V4 / FAMILY_V6 (AF_INET / +// AF_INET6) so off-Linux builds compile. The netlink route query is a no-op stub +// off Linux, so these values are never used against a live kernel. +const ( + familyV4 = 2 // AF_INET + familyV6 = 10 // AF_INET6 on Linux +) diff --git a/internal/checks/route_test.go b/internal/checks/route_test.go index fda2ad62c..64a46dd69 100644 --- a/internal/checks/route_test.go +++ b/internal/checks/route_test.go @@ -53,10 +53,10 @@ func TestDefaultRoutesFromNetlinkIPv6SkipsLoopback(t *testing.T) { } func TestNetlinkFamily(t *testing.T) { - if family, err := netlinkFamily("ipv4"); err != nil || family != netlink.FAMILY_V4 { + if family, err := netlinkFamily("ipv4"); err != nil || family != familyV4 { t.Fatalf("ipv4 family = %d, %v", family, err) } - if family, err := netlinkFamily("ipv6"); err != nil || family != netlink.FAMILY_V6 { + if family, err := netlinkFamily("ipv6"); err != nil || family != familyV6 { t.Fatalf("ipv6 family = %d, %v", family, err) } if _, err := netlinkFamily("ipx"); err == nil { diff --git a/internal/config/model.go b/internal/config/model.go index b764580c6..858f86b4f 100644 --- a/internal/config/model.go +++ b/internal/config/model.go @@ -135,6 +135,18 @@ const ( WebKeyPort = "port" ) +// SectionTelegramBot is the top-level interactive Telegram report-bot block. It +// is distinct from a `telegram` notifier under `notifiers`: the bot receives +// commands (long polling) and answers with read-only reports. +const SectionTelegramBot = "telegram_bot" + +// SectionMap returns the named top-level section as a mapping, or nil when the +// section is absent or is not a mapping. +func SectionMap(raw map[string]any, section string) map[string]any { + m, _ := raw[section].(map[string]any) + return m +} + // stop_policy timeout and kill-guard field keys. const ( keyGracefulTimeout = process.StopPolicyKeyGracefulTimeout diff --git a/internal/config/validate.go b/internal/config/validate.go index 5bbe7230d..d34b2c258 100644 --- a/internal/config/validate.go +++ b/internal/config/validate.go @@ -124,6 +124,7 @@ func validateGlobal(cfg *Config) []Issue { validateGlobalPaths(cfg, raw, add) validateGlobalSecurity(raw, add) validateGlobalWebAndEmission(raw, add) + validateTelegramBot(raw, add) validateGlobalDefaults(cfg, raw, add) return issues diff --git a/internal/config/validate_global.go b/internal/config/validate_global.go index 99cf2f7d5..bc7a3fadd 100644 --- a/internal/config/validate_global.go +++ b/internal/config/validate_global.go @@ -8,6 +8,7 @@ import ( "sermo/internal/cfgval" "sermo/internal/notify" "sermo/internal/rules" + "sermo/internal/telegrambot" ) // validateWatches checks each host-watch entry: a known check type with valid @@ -115,12 +116,7 @@ func validateNotifierType(name string, entry map[string]any, add func(string, .. case notify.TypeSlack, notify.TypeTeams: validateWebhookNotifier(name, typ, entry, add) case notify.TypeTelegram: - if cfgval.String(entry[notify.KeyToken]) == "" { - add("%s is required for a telegram notifier", notifierFieldPath(name, notify.KeyToken)) - } - if cfgval.String(entry[notify.KeyChatID]) == "" { - add("%s is required for a telegram notifier", notifierFieldPath(name, notify.KeyChatID)) - } + validateTelegramNotifier(name, entry, add) case notify.TypeTTY: if users, present := entry[notify.KeyUsers]; present && !cfgval.IsStringOrStringList(users) { add(validationStringListFormat, notifierFieldPath(name, notify.KeyUsers)) @@ -153,6 +149,68 @@ func validateEmailNotifier(name string, entry map[string]any, add func(string, . } } +// validateTelegramBot checks the optional top-level `telegram_bot` section. The +// token is usually sourced from ${env:...}, so an empty token (an unset +// variable) leaves the bot inactive rather than failing config load — mirroring +// telegrambot.Config.active(). When a token is present the section requires at +// least one allowed chat id; a poll interval, when set, must be positive. +func validateTelegramBot(raw map[string]any, add func(string, ...any)) { + section, ok := raw[SectionTelegramBot].(map[string]any) + if !ok { + return + } + field := func(key string) string { return SectionTelegramBot + "." + key } + if v, present := section[telegrambot.KeyEnabled]; present { + if _, isBool := v.(bool); !isBool { + add(validationBooleanFormat, field(telegrambot.KeyEnabled)) + } + } + if enabled, ok := section[telegrambot.KeyEnabled].(bool); ok && !enabled { + return + } + // An empty token (typically an unset ${env:...} secret) leaves the bot + // inactive instead of failing validation, so a host without the token still + // loads its config. Mirrors telegrambot.Config.active(). + if cfgval.String(section[telegrambot.KeyToken]) == "" { + return + } + if ids, ok := cfgval.IntList(section[telegrambot.KeyAllowedChats]); !ok || len(ids) == 0 { + add("%s must list at least one chat id", field(telegrambot.KeyAllowedChats)) + } + if v, present := section[telegrambot.KeyPollInterval]; present && cfgval.Duration(v) <= 0 { + add("%s must be a positive duration", field(telegrambot.KeyPollInterval)) + } +} + +func validateTelegramNotifier(name string, entry map[string]any, add func(string, ...any)) { + // The token is usually sourced from ${env:...}; an empty token (an unset + // variable) leaves the notifier inactive rather than failing config load. + // Build() skips a tokenless telegram notifier with a warning, and its name + // stays defined so `notify` references to it still resolve. + if cfgval.String(entry[notify.KeyToken]) == "" { + return + } + if cfgval.String(entry[notify.KeyChatID]) == "" { + add("%s is required for a telegram notifier", notifierFieldPath(name, notify.KeyChatID)) + } + if v, present := entry[notify.KeyParseMode]; present { + mode, ok := v.(string) + if !ok || !notify.ValidTelegramParseMode(mode) { + add("%s must be one of %s", notifierFieldPath(name, notify.KeyParseMode), strings.Join(notify.TelegramParseModes(), ", ")) + } + } + if v, present := entry[notify.KeySilent]; present { + if _, ok := v.(bool); !ok { + add(validationBooleanFormat, notifierFieldPath(name, notify.KeySilent)) + } + } + if v, present := entry[notify.KeyMessageThreadID]; present { + if _, ok := cfgval.Int(v); !ok { + add("%s must be an integer", notifierFieldPath(name, notify.KeyMessageThreadID)) + } + } +} + func validateWebhookNotifier(name, typ string, entry map[string]any, add func(string, ...any)) { webhook := cfgval.String(entry[notify.KeyWebhook]) if webhook == "" { diff --git a/internal/config/validate_notify_test.go b/internal/config/validate_notify_test.go index 96db6cc20..c9eef5acc 100644 --- a/internal/config/validate_notify_test.go +++ b/internal/config/validate_notify_test.go @@ -117,6 +117,69 @@ func TestValidateTeamsNotifier(t *testing.T) { } } +func TestValidateTelegramNotifier(t *testing.T) { + cases := []struct { + name string + entry map[string]any + wantSub string // "" = expect no issue + }{ + {"valid minimal", map[string]any{"type": "telegram", "token": "t", "chat_id": "1"}, ""}, + {"valid full", map[string]any{"type": "telegram", "token": "t", "chat_id": "1", "parse_mode": "HTML", "silent": true, "message_thread_id": 7}, ""}, + {"empty token leaves notifier inactive", map[string]any{"type": "telegram", "chat_id": "1"}, ""}, + {"missing chat", map[string]any{"type": "telegram", "token": "t"}, "chat_id is required for a telegram notifier"}, + {"bad parse_mode", map[string]any{"type": "telegram", "token": "t", "chat_id": "1", "parse_mode": "rtf"}, "parse_mode must be one of"}, + {"non-bool silent", map[string]any{"type": "telegram", "token": "t", "chat_id": "1", "silent": "yes"}, "silent must be a boolean"}, + {"non-int thread", map[string]any{"type": "telegram", "token": "t", "chat_id": "1", "message_thread_id": "nope"}, "message_thread_id must be an integer"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + issues := collect(func(add func(string, ...any)) { + validateNotifiers(map[string]any{"tg": c.entry}, t.TempDir(), add) + }) + joined := strings.Join(issues, "\n") + if c.wantSub == "" { + if len(issues) != 0 { + t.Errorf("expected no issues, got: %v", issues) + } + } else if !strings.Contains(joined, c.wantSub) { + t.Errorf("expected %q, got: %v", c.wantSub, issues) + } + }) + } +} + +func TestValidateTelegramBot(t *testing.T) { + cases := []struct { + name string + section map[string]any + wantSub string // "" = expect no issue + }{ + {"valid", map[string]any{"token": "t", "allowed_chats": []any{123, -1001234567890}}, ""}, + {"valid full", map[string]any{"token": "t", "allowed_chats": []any{1}, "poll_interval": "45s"}, ""}, + {"disabled skips checks", map[string]any{"enabled": false}, ""}, + {"empty token leaves bot inactive", map[string]any{"allowed_chats": []any{123}}, ""}, + {"enabled without token is optional", map[string]any{"enabled": true, "allowed_chats": []any{123}}, ""}, + {"no chats", map[string]any{"token": "t"}, "telegram_bot.allowed_chats must list at least one chat id"}, + {"bad interval", map[string]any{"token": "t", "allowed_chats": []any{1}, "poll_interval": "nope"}, "telegram_bot.poll_interval must be a positive duration"}, + {"non-bool enabled", map[string]any{"enabled": "yes", "token": "t", "allowed_chats": []any{1}}, "telegram_bot.enabled must be a boolean"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + issues := collect(func(add func(string, ...any)) { + validateTelegramBot(map[string]any{SectionTelegramBot: c.section}, add) + }) + joined := strings.Join(issues, "\n") + if c.wantSub == "" { + if len(issues) != 0 { + t.Errorf("expected no issues, got: %v", issues) + } + } else if !strings.Contains(joined, c.wantSub) { + t.Errorf("expected %q, got: %v", c.wantSub, issues) + } + }) + } +} + func TestValidateNotifierTemplate(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "default-alert.yml"), []byte("subject: '{{ .Subject }}'\nbody: '{{ .Body }}'\n"), 0o644); err != nil { diff --git a/internal/notify/notify.go b/internal/notify/notify.go index 577860cb7..fca92a6c3 100644 --- a/internal/notify/notify.go +++ b/internal/notify/notify.go @@ -39,15 +39,18 @@ type Notifier interface { // Config key constants name fields inside one notifier configuration entry. const ( - KeyDSN = "dsn" - KeyFrom = "from" - KeyTemplate = "template" - KeyTo = "to" - KeyChatID = "chat_id" - KeyToken = "token" - KeyType = "type" - KeyUsers = "users" - KeyWebhook = "webhook" + KeyDSN = "dsn" + KeyFrom = "from" + KeyTemplate = "template" + KeyTo = "to" + KeyChatID = "chat_id" + KeyToken = "token" + KeyType = "type" + KeyUsers = "users" + KeyWebhook = "webhook" + KeyParseMode = "parse_mode" + KeySilent = "silent" + KeyMessageThreadID = "message_thread_id" ) // Option customizes notifier construction. diff --git a/internal/notify/telegram.go b/internal/notify/telegram.go index de89b9471..1ea5cd6dd 100644 --- a/internal/notify/telegram.go +++ b/internal/notify/telegram.go @@ -2,6 +2,7 @@ package notify import ( "errors" + "slices" "sermo/internal/cfgval" ) @@ -11,11 +12,33 @@ const ( telegramSendMessagePath = "/sendMessage" telegramChatIDKey = "chat_id" telegramTextKey = "text" + telegramParseModeKey = "parse_mode" + telegramSilentKey = "disable_notification" + telegramThreadIDKey = "message_thread_id" ) +// telegramParseModes are the Bot API `parse_mode` values Sermo accepts, sorted. +var telegramParseModes = []string{"HTML", "Markdown", "MarkdownV2"} + +// TelegramParseModes returns the accepted `parse_mode` values, for validation +// and docs. +func TelegramParseModes() []string { return slices.Clone(telegramParseModes) } + +// ValidTelegramParseMode reports whether s is an accepted `parse_mode`. +func ValidTelegramParseMode(s string) bool { return slices.Contains(telegramParseModes, s) } + +// telegramOptions carries the optional sendMessage tuning read from config. +type telegramOptions struct { + parseMode string // Bot API parse_mode; empty means plain text + silent bool // disable_notification: deliver without sound + threadID int // message_thread_id: target a forum topic + hasThread bool // whether a thread id was configured +} + // buildTelegram constructs a Telegram bot notifier from a config entry: // `token` is the bot token (kept inside the API URL, never surfaced) and -// `chat_id` the numeric chat or `@channel` target. +// `chat_id` the numeric chat or `@channel` target. Optional `parse_mode`, +// `silent` and `message_thread_id` tune the sendMessage delivery. func buildTelegram(name string, entry map[string]any) (Notifier, error) { token := cfgval.String(entry[KeyToken]) if token == "" { @@ -25,20 +48,41 @@ func buildTelegram(name string, entry map[string]any) (Notifier, error) { if chatID == "" { return nil, errors.New("telegram notifier requires a chat_id") } + opts := telegramOptions{ + parseMode: cfgval.String(entry[KeyParseMode]), + silent: cfgval.Bool(entry[KeySilent]), + } + if _, present := entry[KeyMessageThreadID]; present { + if id, ok := cfgval.Int(entry[KeyMessageThreadID]); ok { + opts.threadID, opts.hasThread = id, true + } + } return &webhookNotifier{ name: name, typ: TypeTelegram, webhook: telegramAPIBase + token + telegramSendMessagePath, - payload: func(msg Message) []byte { return telegramPayload(chatID, msg) }, + payload: func(msg Message) []byte { return telegramPayload(chatID, opts, msg) }, }, nil } // telegramPayload renders the sendMessage body: the subject as the lead line -// and the detail (the SERMO_* fields) below it, as plain text. -func telegramPayload(chatID string, msg Message) []byte { +// and the detail (the SERMO_* fields) below it. Optional tuning fields are +// added only when configured, so an unconfigured notifier posts exactly the +// plain `chat_id`+`text` body it always did. +func telegramPayload(chatID string, opts telegramOptions, msg Message) []byte { text := msg.Subject if msg.Body != "" { text = msg.Subject + notifyLF + msg.Body } - return webhookPayload(map[string]string{telegramChatIDKey: chatID, telegramTextKey: text}) + body := map[string]any{telegramChatIDKey: chatID, telegramTextKey: text} + if opts.parseMode != "" { + body[telegramParseModeKey] = opts.parseMode + } + if opts.silent { + body[telegramSilentKey] = true + } + if opts.hasThread { + body[telegramThreadIDKey] = opts.threadID + } + return webhookPayload(body) } diff --git a/internal/notify/telegram_test.go b/internal/notify/telegram_test.go index e460c4d99..2e9348590 100644 --- a/internal/notify/telegram_test.go +++ b/internal/notify/telegram_test.go @@ -48,3 +48,55 @@ func TestTelegramSendPostsSendMessage(t *testing.T) { t.Fatalf("unexpected telegram body: %+v", body) } } + +func TestTelegramSendIncludesConfiguredOptions(t *testing.T) { + n, err := buildTelegram("tg", map[string]any{ + "type": "telegram", "token": "123:abc", "chat_id": "1", + "parse_mode": "MarkdownV2", "silent": true, "message_thread_id": 42, + }) + if err != nil { + t.Fatal(err) + } + wn := n.(*webhookNotifier) + var gotURL string + var gotPayload []byte + wn.post = capturingPost(t, TypeTelegram, &gotURL, &gotPayload) + if err := n.Send(context.Background(), Message{Subject: "s", Body: "b"}); err != nil { + t.Fatal(err) + } + var body struct { + ParseMode string `json:"parse_mode"` + DisableNotification bool `json:"disable_notification"` + MessageThreadID int `json:"message_thread_id"` + } + if err := json.Unmarshal(gotPayload, &body); err != nil { + t.Fatalf("payload not JSON: %v (%s)", err, gotPayload) + } + if body.ParseMode != "MarkdownV2" || !body.DisableNotification || body.MessageThreadID != 42 { + t.Fatalf("unexpected telegram options: %+v", body) + } +} + +func TestTelegramSendOmitsUnsetOptions(t *testing.T) { + n, err := buildTelegram("tg", map[string]any{"type": "telegram", "token": "123:abc", "chat_id": "1"}) + if err != nil { + t.Fatal(err) + } + wn := n.(*webhookNotifier) + var gotURL string + var gotPayload []byte + wn.post = capturingPost(t, TypeTelegram, &gotURL, &gotPayload) + if err := n.Send(context.Background(), Message{Subject: "s"}); err != nil { + t.Fatal(err) + } + var body map[string]any + if err := json.Unmarshal(gotPayload, &body); err != nil { + t.Fatalf("payload not JSON: %v (%s)", err, gotPayload) + } + // A plain notifier must post exactly the fields it always did. + for _, k := range []string{telegramParseModeKey, telegramSilentKey, telegramThreadIDKey} { + if _, ok := body[k]; ok { + t.Fatalf("unset option %q should be omitted, got %v", k, body) + } + } +} diff --git a/internal/telegrambot/bot.go b/internal/telegrambot/bot.go new file mode 100644 index 000000000..e8a58e206 --- /dev/null +++ b/internal/telegrambot/bot.go @@ -0,0 +1,167 @@ +package telegrambot + +import ( + "context" + "log/slog" + "strings" + "sync" + "time" +) + +const ( + // pollErrorBackoff paces retries after a getUpdates error so a persistent + // failure (bad token, outage) does not spin. + pollErrorBackoff = 5 * time.Second + // idlePollInterval is how long Run waits between checks while disabled or + // tokenless, so a reload can enable it without a restart. + idlePollInterval = 5 * time.Second +) + +// Bot is a read-only Telegram command bot driven by long polling. Construct it +// once and reconfigure it on SIGHUP reload via UpdateConfig. +type Bot struct { + reporter Reporter + log *slog.Logger + + mu sync.Mutex + cfg Config + client *client + + // offset is touched only by the Run goroutine, so it needs no lock. + offset int64 +} + +// New builds a bot from the initial config. reporter supplies report data; +// logger may be nil. +func New(reporter Reporter, cfg Config, logger *slog.Logger) *Bot { + if logger == nil { + logger = slog.Default() + } + b := &Bot{reporter: reporter, log: logger} + b.UpdateConfig(cfg) + return b +} + +// UpdateConfig swaps the configuration used from the next poll (config reload). +// It rebuilds the API client when the token or poll interval changes, and drops +// it when the token is cleared. +func (b *Bot) UpdateConfig(cfg Config) { + if b == nil { + return + } + b.mu.Lock() + defer b.mu.Unlock() + rebuild := b.client == nil || cfg.Token != b.cfg.Token || cfg.PollInterval != b.cfg.PollInterval + b.cfg = cfg + switch { + case cfg.Token == "": + b.client = nil + case rebuild: + b.client = newClient(cfg.Token, cfg.PollInterval+pollClientMargin) + } +} + +// snapshot returns the current config and client under the lock. +func (b *Bot) snapshot() (Config, *client) { + b.mu.Lock() + defer b.mu.Unlock() + return b.cfg, b.client +} + +// Run polls Telegram for commands until ctx is cancelled. It first discards any +// backlog queued before startup so a restart does not replay old commands. +func (b *Bot) Run(ctx context.Context) { + if b == nil { + return + } + b.skipBacklog(ctx) + for { + if ctx.Err() != nil { + return + } + cfg, cl := b.snapshot() + if !cfg.active() || cl == nil { + // Disabled or tokenless (possibly after a reload): idle rather than + // busy-loop, and re-check on the next tick. + if !sleepCtx(ctx, idlePollInterval) { + return + } + continue + } + updates, err := cl.getUpdates(ctx, b.offset, cfg.PollInterval) + if err != nil { + if ctx.Err() != nil { + return + } + b.log.Warn("telegram getUpdates failed", "error", err) + if !sleepCtx(ctx, pollErrorBackoff) { + return + } + continue + } + for _, u := range updates { + b.offset = u.UpdateID + 1 + b.handleUpdate(ctx, cfg, cl, u) + } + } +} + +// skipBacklog advances the offset past updates already queued at startup +// without acting on them. Best effort: on error the offset stays at zero and +// the main loop proceeds. +func (b *Bot) skipBacklog(ctx context.Context) { + cfg, cl := b.snapshot() + if !cfg.active() || cl == nil { + return + } + updates, err := cl.getUpdates(ctx, 0, 0) + if err != nil { + return + } + for _, u := range updates { + if next := u.UpdateID + 1; next > b.offset { + b.offset = next + } + } +} + +// handleUpdate authorizes and dispatches one update, then replies. A panic in a +// handler is recovered so one bad command cannot stop the poll loop. +func (b *Bot) handleUpdate(ctx context.Context, cfg Config, cl *client, u update) { + defer func() { + if r := recover(); r != nil { + b.log.Error("telegram command panic", "recover", r) + } + }() + msg := u.Message + if msg == nil || strings.TrimSpace(msg.Text) == "" { + return + } + if !cfg.allows(msg.Chat.ID) { + // Never reply to a chat that is not on the allow-list. + b.log.Warn("telegram command from unauthorized chat ignored", "chat_id", msg.Chat.ID) + return + } + reply, err := b.dispatch(ctx, msg.Text) + if err != nil { + reply = "Error: " + err.Error() + } + if reply == "" { + return + } + if err := cl.sendMessage(ctx, msg.Chat.ID, msg.MessageThreadID, reply); err != nil { + b.log.Warn("telegram sendMessage failed", "error", err) + } +} + +// sleepCtx waits d or returns false if ctx is cancelled first. +func sleepCtx(ctx context.Context, d time.Duration) bool { + timer := time.NewTimer(d) + defer timer.Stop() + select { + case <-ctx.Done(): + return false + case <-timer.C: + return true + } +} diff --git a/internal/telegrambot/bot_test.go b/internal/telegrambot/bot_test.go new file mode 100644 index 000000000..71b5fc1f4 --- /dev/null +++ b/internal/telegrambot/bot_test.go @@ -0,0 +1,161 @@ +package telegrambot + +import ( + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" +) + +// fakeReporter returns canned data for command tests. +type fakeReporter struct { + status StatusReport + services []ServiceLine + watches []WatchLine + sla []SLAWindow + slaOK bool + events []EventLine + lastN int + err error +} + +func (f *fakeReporter) Status(context.Context) (StatusReport, error) { return f.status, f.err } +func (f *fakeReporter) Services(context.Context) ([]ServiceLine, error) { + return f.services, f.err +} +func (f *fakeReporter) Watches(context.Context) ([]WatchLine, error) { return f.watches, f.err } +func (f *fakeReporter) SLA(_ context.Context, _ string) ([]SLAWindow, bool, error) { + return f.sla, f.slaOK, f.err +} +func (f *fakeReporter) Events(_ context.Context, limit int) ([]EventLine, error) { + f.lastN = limit + return f.events, f.err +} + +func discardLogger() *slog.Logger { return slog.New(slog.NewTextHandler(io.Discard, nil)) } + +func TestParseCommand(t *testing.T) { + cases := []struct { + in string + wantName string + wantArgs []string + }{ + {"/status", "/status", nil}, + {" /Status ", "/status", nil}, + {"/status@SermoBot", "/status", nil}, + {"/sla web", "/sla", []string{"web"}}, + {"status", "/status", nil}, + {"", "", nil}, + {" ", "", nil}, + } + for _, c := range cases { + name, args := parseCommand(c.in) + if name != c.wantName { + t.Errorf("parseCommand(%q) name = %q, want %q", c.in, name, c.wantName) + } + if strings.Join(args, ",") != strings.Join(c.wantArgs, ",") { + t.Errorf("parseCommand(%q) args = %v, want %v", c.in, args, c.wantArgs) + } + } +} + +func TestDispatch(t *testing.T) { + rep := &fakeReporter{ + status: StatusReport{Host: "srv1", Services: 3, OK: 2, Failing: 1}, + services: []ServiceLine{{Name: "web", State: "running", Health: "ok", Monitored: true}}, + watches: []WatchLine{{Name: "disk", Scope: "host", State: "ok", Monitored: true}}, + sla: []SLAWindow{{Window: "day", Ratio: "99.9%"}}, + slaOK: true, + events: []EventLine{{Time: "t", Kind: "firing", Message: "down"}}, + } + b := &Bot{reporter: rep, log: discardLogger()} + ctx := context.Background() + + cases := []struct { + in string + wantSub string + }{ + {"/status", "Sermo status — srv1"}, + {"/services", "web: running / ok"}, + {"/services web", "State: running"}, + {"/services ghost", `No service named "ghost"`}, + {"/watches", "disk (host): ok"}, + {"/sla web", "SLA — web"}, + {"/sla", "Usage: /sla "}, + {"/events", "Recent events (1)"}, + {"/help", "read-only commands"}, + {"", "read-only commands"}, + {"/bogus", "Unknown command /bogus"}, + {"/status@SomeBot", "Sermo status"}, + } + for _, c := range cases { + got, err := b.dispatch(ctx, c.in) + if err != nil { + t.Fatalf("dispatch(%q): %v", c.in, err) + } + if !strings.Contains(got, c.wantSub) { + t.Errorf("dispatch(%q) = %q, want substring %q", c.in, got, c.wantSub) + } + } +} + +func TestDispatchEventsLimitCapped(t *testing.T) { + rep := &fakeReporter{} + b := &Bot{reporter: rep, log: discardLogger()} + if _, err := b.dispatch(context.Background(), "/events 9999"); err != nil { + t.Fatal(err) + } + if rep.lastN != EventsMaxLimit { + t.Fatalf("events limit = %d, want cap %d", rep.lastN, EventsMaxLimit) + } +} + +func TestDispatchPropagatesReporterError(t *testing.T) { + b := &Bot{reporter: &fakeReporter{err: errors.New("backend down")}, log: discardLogger()} + if _, err := b.dispatch(context.Background(), "/status"); err == nil { + t.Fatal("expected reporter error to propagate") + } +} + +func TestHandleUpdateAuthorization(t *testing.T) { + var sends atomic.Int32 + var lastText string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/"+telegramSendMessageMethod) { + sends.Add(1) + var body struct { + Text string `json:"text"` + } + _ = json.NewDecoder(r.Body).Decode(&body) + lastText = body.Text + } + io.WriteString(w, `{"ok":true}`) + })) + defer srv.Close() + + cfg := Config{Enabled: true, Token: "t", AllowedChats: []int64{42}} + b := &Bot{reporter: &fakeReporter{status: StatusReport{Host: "h"}}, log: discardLogger()} + cl := testClient(srv.URL, "t") + ctx := context.Background() + + // Unauthorized chat: no reply is ever sent. + b.handleUpdate(ctx, cfg, cl, update{UpdateID: 1, Message: &message{Chat: chat{ID: 99}, Text: "/status"}}) + if sends.Load() != 0 { + t.Fatalf("unauthorized chat must not get a reply, got %d sends", sends.Load()) + } + + // Authorized chat: a reply is sent. + b.handleUpdate(ctx, cfg, cl, update{UpdateID: 2, Message: &message{Chat: chat{ID: 42}, Text: "/status"}}) + if sends.Load() != 1 { + t.Fatalf("authorized chat should get one reply, got %d sends", sends.Load()) + } + if !strings.Contains(lastText, "Sermo status") { + t.Fatalf("unexpected reply text: %q", lastText) + } +} diff --git a/internal/telegrambot/client.go b/internal/telegrambot/client.go new file mode 100644 index 000000000..b35648d48 --- /dev/null +++ b/internal/telegrambot/client.go @@ -0,0 +1,147 @@ +package telegrambot + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "time" + + "sermo/internal/httpx" + "sermo/internal/netutil" +) + +const ( + telegramAPIBase = "https://api.telegram.org/bot" + telegramGetUpdatesMethod = "getUpdates" + telegramSendMessageMethod = "sendMessage" + // pollClientMargin extends the HTTP timeout past the long-poll timeout so a + // legitimately held-open getUpdates is not cut off by the client. + pollClientMargin = 10 * time.Second + // errorBodySnippetLimit bounds the API error body captured into an error. + errorBodySnippetLimit = 256 + httpStatusClassDivisor = 100 + httpStatusClassSuccess = 2 +) + +// httpDoer performs an HTTP request; *http.Client satisfies it. Injected so +// tests exercise the client without real network I/O. +type httpDoer interface { + Do(*http.Request) (*http.Response, error) +} + +// client talks to the Telegram Bot API for one bot token. The token lives only +// inside the request URL and is scrubbed from any surfaced error. +type client struct { + base string // API base, up to but excluding the token; overridable in tests + token string + http httpDoer +} + +func newClient(token string, timeout time.Duration) *client { + return &client{ + base: telegramAPIBase, + token: token, + http: &http.Client{Timeout: timeout, Transport: httpx.CloneDefaultTransport()}, + } +} + +func (c *client) methodURL(method string) string { + return c.base + c.token + "/" + method +} + +// update is one getUpdates result item; only message updates are requested. +// Only the fields the bot acts on are decoded. +type update struct { + UpdateID int64 `json:"update_id"` + Message *message `json:"message"` +} + +type message struct { + Chat chat `json:"chat"` + Text string `json:"text"` + MessageThreadID int `json:"message_thread_id"` +} + +type chat struct { + ID int64 `json:"id"` +} + +// getUpdates long-polls for updates newer than offset, holding the connection +// up to timeout. Only message updates are requested. +func (c *client) getUpdates(ctx context.Context, offset int64, timeout time.Duration) ([]update, error) { + body := map[string]any{ + "offset": offset, + "timeout": int(timeout / time.Second), + "allowed_updates": []string{"message"}, + } + var resp struct { + OK bool `json:"ok"` + Result []update `json:"result"` + Description string `json:"description"` + } + if err := c.call(ctx, telegramGetUpdatesMethod, body, &resp); err != nil { + return nil, err + } + if !resp.OK { + return nil, fmt.Errorf("telegram getUpdates: %s", resp.Description) + } + return resp.Result, nil +} + +// sendMessage posts a plain-text reply to chatID, optionally within a forum +// topic thread (threadID 0 means the chat's main timeline). +func (c *client) sendMessage(ctx context.Context, chatID int64, threadID int, text string) error { + body := map[string]any{"chat_id": chatID, "text": text} + if threadID != 0 { + body["message_thread_id"] = threadID + } + var resp struct { + OK bool `json:"ok"` + Description string `json:"description"` + } + if err := c.call(ctx, telegramSendMessageMethod, body, &resp); err != nil { + return err + } + if !resp.OK { + return fmt.Errorf("telegram sendMessage: %s", resp.Description) + } + return nil +} + +// call POSTs a JSON body to a Bot API method and decodes the JSON response. +// The bot token is embedded in the URL, so every error is scrubbed of the URL +// via netutil.URLErrorCause before it is returned — no credential ever reaches +// a log line. +func (c *client) call(ctx context.Context, method string, body map[string]any, out any) error { + payload, err := json.Marshal(body) + if err != nil { + return fmt.Errorf("encode %s request: %w", method, err) + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.methodURL(method), bytes.NewReader(payload)) + if err != nil { + return fmt.Errorf("build %s request: %w", method, netutil.URLErrorCause(err)) + } + req.Header.Set(httpx.HeaderContentType, httpx.ContentTypeJSON) + + resp, err := c.http.Do(req) + if err != nil { + return fmt.Errorf("%s request: %w", method, netutil.URLErrorCause(err)) + } + defer resp.Body.Close() + + if resp.StatusCode/httpStatusClassDivisor != httpStatusClassSuccess { + snippet, _ := io.ReadAll(io.LimitReader(resp.Body, errorBodySnippetLimit)) + return fmt.Errorf("telegram %s returned %s: %s", method, resp.Status, strings.TrimSpace(string(snippet))) + } + if out == nil { + return nil + } + if err := json.NewDecoder(resp.Body).Decode(out); err != nil { + return fmt.Errorf("decode %s response: %w", method, err) + } + return nil +} diff --git a/internal/telegrambot/client_test.go b/internal/telegrambot/client_test.go new file mode 100644 index 000000000..1cf3f9d4c --- /dev/null +++ b/internal/telegrambot/client_test.go @@ -0,0 +1,124 @@ +package telegrambot + +import ( + "context" + "encoding/json" + "errors" + "io" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "testing" + "time" +) + +// doerFunc adapts a function to httpDoer. +type doerFunc func(*http.Request) (*http.Response, error) + +func (f doerFunc) Do(r *http.Request) (*http.Response, error) { return f(r) } + +// testClient points a client at a test server base and token. +func testClient(base, token string) *client { + c := newClient(token, time.Second) + c.base = base + "/bot" + return c +} + +func TestGetUpdatesParsesMessages(t *testing.T) { + const token = "123:abc" + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !strings.HasSuffix(r.URL.Path, "/bot"+token+"/"+telegramGetUpdatesMethod) { + t.Errorf("unexpected path %q", r.URL.Path) + } + var body struct { + Offset int64 `json:"offset"` + Timeout int `json:"timeout"` + AllowedUpdates []string `json:"allowed_updates"` + } + _ = json.NewDecoder(r.Body).Decode(&body) + if body.Offset != 7 || body.Timeout != 30 || len(body.AllowedUpdates) != 1 || body.AllowedUpdates[0] != "message" { + t.Errorf("unexpected getUpdates request: %+v", body) + } + io.WriteString(w, `{"ok":true,"result":[{"update_id":9,"message":{"message_id":1,"chat":{"id":42,"type":"private"},"text":"/status"}}]}`) + })) + defer srv.Close() + + updates, err := testClient(srv.URL, token).getUpdates(context.Background(), 7, 30*time.Second) + if err != nil { + t.Fatal(err) + } + if len(updates) != 1 || updates[0].UpdateID != 9 || updates[0].Message == nil || + updates[0].Message.Chat.ID != 42 || updates[0].Message.Text != "/status" { + t.Fatalf("unexpected updates: %+v", updates) + } +} + +func TestSendMessagePostsChatAndText(t *testing.T) { + var got struct { + ChatID int64 `json:"chat_id"` + Text string `json:"text"` + MessageThreadID int `json:"message_thread_id"` + } + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !strings.HasSuffix(r.URL.Path, "/"+telegramSendMessageMethod) { + t.Errorf("unexpected path %q", r.URL.Path) + } + _ = json.NewDecoder(r.Body).Decode(&got) + io.WriteString(w, `{"ok":true}`) + })) + defer srv.Close() + + if err := testClient(srv.URL, "t").sendMessage(context.Background(), 42, 5, "hello"); err != nil { + t.Fatal(err) + } + if got.ChatID != 42 || got.Text != "hello" || got.MessageThreadID != 5 { + t.Fatalf("unexpected sendMessage body: %+v", got) + } +} + +func TestSendMessageOmitsThreadWhenZero(t *testing.T) { + var raw map[string]any + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _ = json.NewDecoder(r.Body).Decode(&raw) + io.WriteString(w, `{"ok":true}`) + })) + defer srv.Close() + + if err := testClient(srv.URL, "t").sendMessage(context.Background(), 42, 0, "hi"); err != nil { + t.Fatal(err) + } + if _, ok := raw["message_thread_id"]; ok { + t.Fatalf("thread id 0 should be omitted, got %v", raw) + } +} + +func TestCallReportsAPIError(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusBadRequest) + io.WriteString(w, `{"ok":false,"description":"chat not found"}`) + })) + defer srv.Close() + + err := testClient(srv.URL, "t").sendMessage(context.Background(), 1, 0, "x") + if err == nil || !strings.Contains(err.Error(), "chat not found") { + t.Fatalf("want API error surfaced, got %v", err) + } +} + +func TestCallScrubsTokenFromTransportError(t *testing.T) { + const token = "SUPERSECRETTOKEN" + c := newClient(token, time.Second) + c.http = doerFunc(func(r *http.Request) (*http.Response, error) { + // A real transport error is a *url.Error whose text embeds the full + // request URL — for Telegram that URL carries the bot token. + return nil, &url.Error{Op: "Post", URL: r.URL.String(), Err: errors.New("dial tcp: refused")} + }) + err := c.call(context.Background(), telegramGetUpdatesMethod, map[string]any{}, nil) + if err == nil { + t.Fatal("expected an error") + } + if strings.Contains(err.Error(), token) { + t.Fatalf("token leaked into error: %q", err.Error()) + } +} diff --git a/internal/telegrambot/commands.go b/internal/telegrambot/commands.go new file mode 100644 index 000000000..3bf69ac0d --- /dev/null +++ b/internal/telegrambot/commands.go @@ -0,0 +1,143 @@ +package telegrambot + +import ( + "context" + "fmt" + "maps" + "slices" + "strconv" + "strings" +) + +// command is one read-only bot command. +type command struct { + help string + handler func(ctx context.Context, b *Bot, args []string) (string, error) +} + +// commands is the read-only command registry, mirroring sermoctl's dispatch +// table. Every handler only reads state through the Reporter. It is populated +// in init to avoid a self-reference cycle (cmdHelp -> helpText -> commands). +var commands map[string]command + +func init() { + commands = map[string]command{ + "/status": {help: "daemon and service health summary", handler: cmdStatus}, + "/services": {help: "list services, or detail one: /services [name]", handler: cmdServices}, + "/watches": {help: "host and service watch states", handler: cmdWatches}, + "/sla": {help: "availability windows for a service: /sla ", handler: cmdSLA}, + "/events": {help: "recent events: /events [count]", handler: cmdEvents}, + "/help": {help: "show this help", handler: cmdHelp}, + } +} + +// dispatch parses one message and returns the reply text. An empty or unknown +// command yields help. +func (b *Bot) dispatch(ctx context.Context, text string) (string, error) { + name, args := parseCommand(text) + if name == "" { + return helpText(), nil + } + cmd, ok := commands[name] + if !ok { + return "Unknown command " + name + ".\n\n" + helpText(), nil + } + return cmd.handler(ctx, b, args) +} + +// parseCommand splits a message into a normalized command name and its +// arguments. A bot-mention suffix (/status@MyBot) is stripped, the name is +// lowercased, and a leading slash is ensured. +func parseCommand(text string) (name string, args []string) { + fields := strings.Fields(strings.TrimSpace(text)) + if len(fields) == 0 { + return "", nil + } + name = fields[0] + if i := strings.IndexByte(name, '@'); i >= 0 { + name = name[:i] + } + name = strings.ToLower(name) + if !strings.HasPrefix(name, "/") { + name = "/" + name + } + return name, fields[1:] +} + +func cmdStatus(ctx context.Context, b *Bot, _ []string) (string, error) { + rep, err := b.reporter.Status(ctx) + if err != nil { + return "", err + } + return formatStatus(rep), nil +} + +func cmdServices(ctx context.Context, b *Bot, args []string) (string, error) { + lines, err := b.reporter.Services(ctx) + if err != nil { + return "", err + } + if len(args) > 0 { + name := args[0] + for _, s := range lines { + if strings.EqualFold(s.Name, name) { + return formatServiceDetail(s), nil + } + } + return fmt.Sprintf("No service named %q.", name), nil + } + return formatServices(lines), nil +} + +func cmdWatches(ctx context.Context, b *Bot, _ []string) (string, error) { + lines, err := b.reporter.Watches(ctx) + if err != nil { + return "", err + } + return formatWatches(lines), nil +} + +func cmdSLA(ctx context.Context, b *Bot, args []string) (string, error) { + if len(args) == 0 { + return "Usage: /sla ", nil + } + service := args[0] + windows, ok, err := b.reporter.SLA(ctx, service) + if err != nil { + return "", err + } + if !ok { + return fmt.Sprintf("No service named %q.", service), nil + } + return formatSLA(service, windows), nil +} + +func cmdEvents(ctx context.Context, b *Bot, args []string) (string, error) { + limit := EventsDefaultLimit + if len(args) > 0 { + if n, err := strconv.Atoi(args[0]); err == nil && n > 0 { + limit = n + } + } + if limit > EventsMaxLimit { + limit = EventsMaxLimit + } + lines, err := b.reporter.Events(ctx, limit) + if err != nil { + return "", err + } + return formatEvents(lines), nil +} + +func cmdHelp(_ context.Context, _ *Bot, _ []string) (string, error) { + return helpText(), nil +} + +func helpText() string { + var b strings.Builder + b.WriteString("Sermo bot — read-only commands:\n") + for _, name := range slices.Sorted(maps.Keys(commands)) { + fmt.Fprintf(&b, "%s — %s\n", name, commands[name].help) + } + return strings.TrimRight(b.String(), "\n") +} diff --git a/internal/telegrambot/config.go b/internal/telegrambot/config.go new file mode 100644 index 000000000..c7fb8bd08 --- /dev/null +++ b/internal/telegrambot/config.go @@ -0,0 +1,86 @@ +// Package telegrambot implements an interactive, read-only Telegram bot that +// answers operator commands (/status, /services, /sla, ...) with Sermo reports. +// +// It receives commands over Bot API long polling (getUpdates), so it needs no +// inbound socket and no public exposure — matching Sermo's outbound-only +// posture. It is configured under the top-level `telegram_bot` section and, +// like the notifier transports, only ever reads state: it can never change the +// host. Report data is supplied through the narrow Reporter interface so this +// package builds and tests free of the daemon's platform-specific +// dependencies; the daemon wires an adapter over its web backend and store. +package telegrambot + +import ( + "time" + + "sermo/internal/cfgval" +) + +// Config field keys inside the `telegram_bot` section. +const ( + KeyEnabled = "enabled" + KeyToken = "token" + KeyAllowedChats = "allowed_chats" + KeyPollInterval = "poll_interval" +) + +const ( + // DefaultPollInterval is the getUpdates long-poll timeout when unset. + DefaultPollInterval = 30 * time.Second + // MinPollInterval and MaxPollInterval bound the configured poll interval. + MinPollInterval = 1 * time.Second + MaxPollInterval = 10 * time.Minute +) + +// Config is the parsed `telegram_bot` section. +type Config struct { + Enabled bool + Token string + AllowedChats []int64 + PollInterval time.Duration +} + +// ParseConfig reads a `telegram_bot` section map into a Config. A nil/absent +// section yields a disabled Config. Parsing is lenient: schema problems are +// reported by config validation, not here. +func ParseConfig(raw map[string]any) Config { + if raw == nil { + return Config{} + } + cfg := Config{ + Enabled: !cfgval.Disabled(raw), + Token: cfgval.String(raw[KeyToken]), + PollInterval: clampInterval(cfgval.DurationOr(raw[KeyPollInterval], DefaultPollInterval)), + } + if ids, ok := cfgval.IntList(raw[KeyAllowedChats]); ok { + cfg.AllowedChats = make([]int64, len(ids)) + for i, id := range ids { + cfg.AllowedChats[i] = int64(id) + } + } + return cfg +} + +func clampInterval(d time.Duration) time.Duration { + switch { + case d < MinPollInterval: + return MinPollInterval + case d > MaxPollInterval: + return MaxPollInterval + default: + return d + } +} + +// active reports whether the bot should poll: enabled and holding a token. +func (c Config) active() bool { return c.Enabled && c.Token != "" } + +// allows reports whether chat id may command the bot. +func (c Config) allows(chatID int64) bool { + for _, id := range c.AllowedChats { + if id == chatID { + return true + } + } + return false +} diff --git a/internal/telegrambot/config_test.go b/internal/telegrambot/config_test.go new file mode 100644 index 000000000..a52c6b54d --- /dev/null +++ b/internal/telegrambot/config_test.go @@ -0,0 +1,55 @@ +package telegrambot + +import ( + "testing" + "time" +) + +func TestParseConfigDefaultsAndClamp(t *testing.T) { + if got := ParseConfig(nil); got.Enabled || got.Token != "" || got.AllowedChats != nil { + t.Fatalf("nil section should be a disabled empty config, got %+v", got) + } + + cfg := ParseConfig(map[string]any{ + "token": "123:abc", + "allowed_chats": []any{123, -1001234567890}, + }) + if !cfg.Enabled { + t.Fatalf("present section without enabled:false should be enabled") + } + if cfg.PollInterval != DefaultPollInterval { + t.Fatalf("poll interval = %s, want default %s", cfg.PollInterval, DefaultPollInterval) + } + if len(cfg.AllowedChats) != 2 || cfg.AllowedChats[0] != 123 || cfg.AllowedChats[1] != -1001234567890 { + t.Fatalf("allowed chats = %v", cfg.AllowedChats) + } + + if got := ParseConfig(map[string]any{"poll_interval": "1ms"}).PollInterval; got != MinPollInterval { + t.Fatalf("tiny interval should clamp to %s, got %s", MinPollInterval, got) + } + if got := ParseConfig(map[string]any{"poll_interval": "24h"}).PollInterval; got != MaxPollInterval { + t.Fatalf("huge interval should clamp to %s, got %s", MaxPollInterval, got) + } + if got := ParseConfig(map[string]any{"poll_interval": "45s"}).PollInterval; got != 45*time.Second { + t.Fatalf("in-range interval should pass through, got %s", got) + } +} + +func TestConfigActiveAndAllows(t *testing.T) { + if (Config{Enabled: true}).active() { + t.Fatalf("enabled without token must not be active") + } + if (Config{Token: "t"}).active() { + t.Fatalf("token without enabled must not be active") + } + c := Config{Enabled: true, Token: "t", AllowedChats: []int64{42}} + if !c.active() { + t.Fatalf("enabled+token must be active") + } + if !c.allows(42) || c.allows(43) { + t.Fatalf("allow-list check wrong for %+v", c) + } + if (Config{}).allows(0) { + t.Fatalf("empty allow-list must reject everyone") + } +} diff --git a/internal/telegrambot/report.go b/internal/telegrambot/report.go new file mode 100644 index 000000000..003b26076 --- /dev/null +++ b/internal/telegrambot/report.go @@ -0,0 +1,160 @@ +package telegrambot + +import ( + "context" + "fmt" + "strings" +) + +// Reporter is the read-only view of Sermo state the bot renders into replies. +// It is deliberately narrow and free of daemon/web types so this package builds +// and tests without the daemon's platform-specific dependencies; the daemon +// supplies an adapter over its web backend and state store. +type Reporter interface { + Status(ctx context.Context) (StatusReport, error) + Services(ctx context.Context) ([]ServiceLine, error) + Watches(ctx context.Context) ([]WatchLine, error) + // SLA returns availability windows for a service; ok is false when no such + // service is configured. + SLA(ctx context.Context, service string) (windows []SLAWindow, ok bool, err error) + Events(ctx context.Context, limit int) ([]EventLine, error) +} + +// StatusReport is the /status rollup. +type StatusReport struct { + Host string + Services int + OK int + Failing int + Monitored int + Paused int + Errors int + LastEvent string + HostUptime string +} + +// ServiceLine is one service's summary for /services. +type ServiceLine struct { + Name string + State string + Health string + Monitored bool +} + +// WatchLine is one watch's summary for /watches. +type WatchLine struct { + Name string + Scope string + State string + Monitored bool +} + +// EventLine is one entry of the /events feed. +type EventLine struct { + Time string + Target string // the service/watch/app the event concerns, if any + Kind string + Message string +} + +// SLAWindow is one availability window for /sla. +type SLAWindow struct { + Window string + Ratio string // formatted percentage, or "n/a" +} + +const ( + // EventsDefaultLimit and EventsMaxLimit bound the /events reply size. + EventsDefaultLimit = 10 + EventsMaxLimit = 50 +) + +func formatStatus(r StatusReport) string { + var b strings.Builder + fmt.Fprintf(&b, "Sermo status — %s\n", nonEmpty(r.Host, "host")) + fmt.Fprintf(&b, "Services: %d (ok %d, failing %d)\n", r.Services, r.OK, r.Failing) + fmt.Fprintf(&b, "Monitoring: %d monitored, %d paused\n", r.Monitored, r.Paused) + fmt.Fprintf(&b, "Recent errors: %d\n", r.Errors) + if r.LastEvent != "" { + fmt.Fprintf(&b, "Last event: %s\n", r.LastEvent) + } + if r.HostUptime != "" { + fmt.Fprintf(&b, "Host uptime: %s\n", r.HostUptime) + } + return strings.TrimRight(b.String(), "\n") +} + +func formatServices(lines []ServiceLine) string { + if len(lines) == 0 { + return "No services configured." + } + var b strings.Builder + fmt.Fprintf(&b, "Services (%d):\n", len(lines)) + for _, s := range lines { + fmt.Fprintf(&b, "- %s: %s / %s%s\n", s.Name, nonEmpty(s.State, "?"), nonEmpty(s.Health, "unknown"), monitorSuffix(s.Monitored)) + } + return strings.TrimRight(b.String(), "\n") +} + +func formatServiceDetail(s ServiceLine) string { + mon := "monitored" + if !s.Monitored { + mon = "not monitored" + } + return fmt.Sprintf("%s\nState: %s\nHealth: %s\nMonitoring: %s", + s.Name, nonEmpty(s.State, "?"), nonEmpty(s.Health, "unknown"), mon) +} + +func formatWatches(lines []WatchLine) string { + if len(lines) == 0 { + return "No watches configured." + } + var b strings.Builder + fmt.Fprintf(&b, "Watches (%d):\n", len(lines)) + for _, w := range lines { + fmt.Fprintf(&b, "- %s (%s): %s%s\n", w.Name, nonEmpty(w.Scope, "?"), nonEmpty(w.State, "?"), monitorSuffix(w.Monitored)) + } + return strings.TrimRight(b.String(), "\n") +} + +func formatSLA(service string, windows []SLAWindow) string { + if len(windows) == 0 { + return fmt.Sprintf("No SLA data for %s.", service) + } + var b strings.Builder + fmt.Fprintf(&b, "SLA — %s\n", service) + for _, w := range windows { + fmt.Fprintf(&b, "- %s: %s\n", w.Window, w.Ratio) + } + return strings.TrimRight(b.String(), "\n") +} + +func formatEvents(lines []EventLine) string { + if len(lines) == 0 { + return "No recent events." + } + var b strings.Builder + fmt.Fprintf(&b, "Recent events (%d):\n", len(lines)) + for _, e := range lines { + target := "" + if e.Target != "" { + target = " " + e.Target + } + fmt.Fprintf(&b, "- %s [%s]%s: %s\n", e.Time, nonEmpty(e.Kind, "?"), target, e.Message) + } + return strings.TrimRight(b.String(), "\n") +} + +func monitorSuffix(monitored bool) string { + if monitored { + return "" + } + return " [not monitored]" +} + +func nonEmpty(s, fallback string) string { + if strings.TrimSpace(s) == "" { + return fallback + } + return s +}