main.go 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. // Command routerd consumes alerts from NATS JetStream, resolves
  2. // recipients (companies → sources → allowed_targets ∪ routing_rules
  3. // → groups → individuals ∩ subscriptions), and enqueues one
  4. // delivery per (individual, channel) to
  5. // deliveries.<channel>.<company_id> subjects.
  6. //
  7. // M0: connects to NATS, /health, /metrics. No business logic.
  8. // M1: broadcast mode — every active fcm_token in the company.
  9. // M2: rules engine from SPEC §6 — uses internal/routing.Resolver
  10. // which honors source.allowed_targets, routing_rules, and
  11. // subscriptions (with min_severity, quiet hours, channel_mask).
  12. // Hard-fails on zero recipients (no silent broadcast).
  13. package main
  14. import (
  15. "context"
  16. "encoding/json"
  17. "fmt"
  18. "log/slog"
  19. "os"
  20. "os/signal"
  21. "strings"
  22. "syscall"
  23. "time"
  24. "git3.techno-world.net/lrosales/broad-announce/internal/alert"
  25. "git3.techno-world.net/lrosales/broad-announce/internal/broker"
  26. "git3.techno-world.net/lrosales/broad-announce/internal/config"
  27. "git3.techno-world.net/lrosales/broad-announce/internal/dedupe"
  28. "git3.techno-world.net/lrosales/broad-announce/internal/httpserver"
  29. "git3.techno-world.net/lrosales/broad-announce/internal/observability"
  30. "git3.techno-world.net/lrosales/broad-announce/internal/postgres"
  31. "git3.techno-world.net/lrosales/broad-announce/internal/routing"
  32. "github.com/nats-io/nats.go/jetstream"
  33. )
  34. // routerdMetrics is the package-level routerd metrics instance.
  35. // Set in main() before the consume goroutine starts.
  36. var routerdMetrics *observability.RouterdMetrics
  37. func main() {
  38. cfg, err := config.LoadRouterd()
  39. if err != nil {
  40. os.Stderr.WriteString("config: " + err.Error() + "\n")
  41. os.Exit(1)
  42. }
  43. logger := observability.Init(cfg.Env, cfg.LogLevel, "routerd")
  44. logger.Info("starting",
  45. "env", cfg.Env,
  46. "addr", cfg.HTTPAddr,
  47. "dedupe_flush_ms", cfg.DedupeFlushMs,
  48. )
  49. ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
  50. defer stop()
  51. br, err := broker.Connect(ctx, cfg.NATSURL)
  52. if err != nil {
  53. logger.Error("nats connect", "err", err)
  54. os.Exit(1)
  55. }
  56. defer br.Close()
  57. pool, err := postgres.Connect(ctx, cfg.PostgresDSN)
  58. if err != nil {
  59. logger.Error("postgres connect", "err", err)
  60. os.Exit(1)
  61. }
  62. defer pool.Close()
  63. resolver := routing.New(pool, logger.With("subsystem", "routing"))
  64. // M6.5: build the router-level dedupe Collapser. It collapses
  65. // bursts of identical alerts into 1 delivery per (source,
  66. // dedupe_key) pair, with a max-wait debounce so a continuous
  67. // stream re-flushes every flush window.
  68. collapseState := newFanoutState()
  69. collapser := dedupe.NewCollapser(
  70. time.Duration(cfg.DedupeFlushMs)*time.Millisecond,
  71. onFlushWithKey(logger, collapseState, br),
  72. )
  73. // Subscribe to all alerts.* subjects.
  74. js := br.JS()
  75. stream, err := js.Stream(ctx, "ALERTS")
  76. if err != nil {
  77. logger.Error("nats stream ALERTS", "err", err)
  78. os.Exit(1)
  79. }
  80. consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
  81. Name: "routerd",
  82. Durable: "routerd",
  83. FilterSubjects: []string{"alerts.>"},
  84. AckPolicy: jetstream.AckExplicitPolicy,
  85. })
  86. if err != nil {
  87. logger.Error("nats consumer", "err", err)
  88. os.Exit(1)
  89. }
  90. runCtx, runCancel := context.WithCancel(ctx)
  91. defer runCancel()
  92. // M9 L1: routerd metrics. Create before the consumer starts
  93. // so the /metrics endpoint is valid immediately.
  94. reg, _ := observability.NewRegistry("routerd")
  95. routerdMetrics = observability.NewRouterdMetrics(reg, "routerd")
  96. go consume(runCtx, logger, consumer, br, resolver, collapser, collapseState)
  97. srv := httpserver.New(httpserver.Config{
  98. Addr: cfg.HTTPAddr,
  99. ServiceName: "routerd",
  100. ShutdownGrace: cfg.ShutdownGrace,
  101. }, logger, observability.MetricsHandler(reg))
  102. // M13a W5: admin routes (JWT-gated). When BA_AUTHD_JWT_SECRET
  103. // is unset, wireAdminRoutes is a no-op so the LAN deploy
  104. // path keeps working unchanged.
  105. wireAdminRoutes(srv.Mux(), collapser, collapseState, logger)
  106. errCh := make(chan error, 1)
  107. go func() { errCh <- srv.Start() }()
  108. select {
  109. case <-ctx.Done():
  110. logger.Info("shutdown signal received")
  111. case err := <-errCh:
  112. if err != nil {
  113. logger.Error("http server", "err", err)
  114. os.Exit(1)
  115. }
  116. }
  117. runCancel()
  118. time.Sleep(500 * time.Millisecond) // let consumer drain
  119. // M6.5: drain any pending collapses so we don't lose
  120. // the last few alerts of a graceful shutdown.
  121. collapser.FlushAll()
  122. if err := srv.Shutdown(ctx); err != nil {
  123. logger.Warn("graceful shutdown", "err", err)
  124. }
  125. logger.Info("bye")
  126. }
  127. func consume(
  128. ctx context.Context,
  129. logger *slog.Logger,
  130. c jetstream.Consumer,
  131. br *broker.Client,
  132. r *routing.Resolver,
  133. collapser *dedupe.Collapser,
  134. state *fanoutState,
  135. ) {
  136. // M6.5: start the Collapser's flush loop. It exits
  137. // when ctx is canceled.
  138. collapseCtx, collapseCancel := context.WithCancel(ctx)
  139. defer collapseCancel()
  140. go collapser.Run(collapseCtx)
  141. for {
  142. if ctx.Err() != nil {
  143. return
  144. }
  145. batch, err := c.Fetch(16, jetstream.FetchMaxWait(2*time.Second))
  146. if err != nil {
  147. if ctx.Err() != nil {
  148. return
  149. }
  150. logger.Warn("nats fetch", "err", err)
  151. time.Sleep(500 * time.Millisecond)
  152. continue
  153. }
  154. for m := range batch.Messages() {
  155. handleOne(ctx, logger, m, br, r, collapser, state)
  156. if batch.Error() != nil {
  157. logger.Warn("batch error", "err", batch.Error())
  158. break
  159. }
  160. }
  161. }
  162. }
  163. func handleOne(
  164. ctx context.Context,
  165. logger *slog.Logger,
  166. m jetstream.Msg,
  167. br *broker.Client,
  168. r *routing.Resolver,
  169. collapser *dedupe.Collapser,
  170. state *fanoutState,
  171. ) {
  172. var a alert.Alert
  173. if err := json.Unmarshal(m.Data(), &a); err != nil {
  174. logger.Warn("malformed alert payload", "err", err, "subject", m.Subject())
  175. _ = m.Ack()
  176. return
  177. }
  178. // Extract company_id from subject "alerts.<company_id>".
  179. parts := strings.SplitN(m.Subject(), ".", 2)
  180. if len(parts) != 2 {
  181. logger.Warn("bad subject", "subject", m.Subject())
  182. _ = m.Ack()
  183. return
  184. }
  185. companyID := parts[1]
  186. if a.CompanyID != "" && a.CompanyID != companyID {
  187. logger.Warn("company_id mismatch", "subject", companyID, "body", a.CompanyID)
  188. }
  189. a.CompanyID = companyID
  190. // M6.5: route through the Collapser-aware handler. The
  191. // collapse decision is made here, the actual delivery
  192. // publish may happen later (on flush) — but the NATS
  193. // message is always Ack'd now.
  194. _, shouldAck, shouldNak := observeAndFanout(ctx, logger, collapser, state, r, br, companyID, a)
  195. switch {
  196. case shouldNak:
  197. _ = m.Nak()
  198. case shouldAck:
  199. _ = m.Ack()
  200. }
  201. }
  202. var _ = fmt.Sprintf // keep import