main.go 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218
  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. errCh := make(chan error, 1)
  103. go func() { errCh <- srv.Start() }()
  104. select {
  105. case <-ctx.Done():
  106. logger.Info("shutdown signal received")
  107. case err := <-errCh:
  108. if err != nil {
  109. logger.Error("http server", "err", err)
  110. os.Exit(1)
  111. }
  112. }
  113. runCancel()
  114. time.Sleep(500 * time.Millisecond) // let consumer drain
  115. // M6.5: drain any pending collapses so we don't lose
  116. // the last few alerts of a graceful shutdown.
  117. collapser.FlushAll()
  118. if err := srv.Shutdown(ctx); err != nil {
  119. logger.Warn("graceful shutdown", "err", err)
  120. }
  121. logger.Info("bye")
  122. }
  123. func consume(
  124. ctx context.Context,
  125. logger *slog.Logger,
  126. c jetstream.Consumer,
  127. br *broker.Client,
  128. r *routing.Resolver,
  129. collapser *dedupe.Collapser,
  130. state *fanoutState,
  131. ) {
  132. // M6.5: start the Collapser's flush loop. It exits
  133. // when ctx is canceled.
  134. collapseCtx, collapseCancel := context.WithCancel(ctx)
  135. defer collapseCancel()
  136. go collapser.Run(collapseCtx)
  137. for {
  138. if ctx.Err() != nil {
  139. return
  140. }
  141. batch, err := c.Fetch(16, jetstream.FetchMaxWait(2*time.Second))
  142. if err != nil {
  143. if ctx.Err() != nil {
  144. return
  145. }
  146. logger.Warn("nats fetch", "err", err)
  147. time.Sleep(500 * time.Millisecond)
  148. continue
  149. }
  150. for m := range batch.Messages() {
  151. handleOne(ctx, logger, m, br, r, collapser, state)
  152. if batch.Error() != nil {
  153. logger.Warn("batch error", "err", batch.Error())
  154. break
  155. }
  156. }
  157. }
  158. }
  159. func handleOne(
  160. ctx context.Context,
  161. logger *slog.Logger,
  162. m jetstream.Msg,
  163. br *broker.Client,
  164. r *routing.Resolver,
  165. collapser *dedupe.Collapser,
  166. state *fanoutState,
  167. ) {
  168. var a alert.Alert
  169. if err := json.Unmarshal(m.Data(), &a); err != nil {
  170. logger.Warn("malformed alert payload", "err", err, "subject", m.Subject())
  171. _ = m.Ack()
  172. return
  173. }
  174. // Extract company_id from subject "alerts.<company_id>".
  175. parts := strings.SplitN(m.Subject(), ".", 2)
  176. if len(parts) != 2 {
  177. logger.Warn("bad subject", "subject", m.Subject())
  178. _ = m.Ack()
  179. return
  180. }
  181. companyID := parts[1]
  182. if a.CompanyID != "" && a.CompanyID != companyID {
  183. logger.Warn("company_id mismatch", "subject", companyID, "body", a.CompanyID)
  184. }
  185. a.CompanyID = companyID
  186. // M6.5: route through the Collapser-aware handler. The
  187. // collapse decision is made here, the actual delivery
  188. // publish may happen later (on flush) — but the NATS
  189. // message is always Ack'd now.
  190. _, shouldAck, shouldNak := observeAndFanout(ctx, logger, collapser, state, r, br, companyID, a)
  191. switch {
  192. case shouldNak:
  193. _ = m.Nak()
  194. case shouldAck:
  195. _ = m.Ack()
  196. }
  197. }
  198. var _ = fmt.Sprintf // keep import