main.go 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  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. func main() {
  35. cfg, err := config.LoadRouterd()
  36. if err != nil {
  37. os.Stderr.WriteString("config: " + err.Error() + "\n")
  38. os.Exit(1)
  39. }
  40. logger := observability.Init(cfg.Env, cfg.LogLevel, "routerd")
  41. logger.Info("starting",
  42. "env", cfg.Env,
  43. "addr", cfg.HTTPAddr,
  44. "dedupe_flush_ms", cfg.DedupeFlushMs,
  45. )
  46. ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
  47. defer stop()
  48. br, err := broker.Connect(ctx, cfg.NATSURL)
  49. if err != nil {
  50. logger.Error("nats connect", "err", err)
  51. os.Exit(1)
  52. }
  53. defer br.Close()
  54. pool, err := postgres.Connect(ctx, cfg.PostgresDSN)
  55. if err != nil {
  56. logger.Error("postgres connect", "err", err)
  57. os.Exit(1)
  58. }
  59. defer pool.Close()
  60. resolver := routing.New(pool, logger.With("subsystem", "routing"))
  61. // M6.5: build the router-level dedupe Collapser. It collapses
  62. // bursts of identical alerts into 1 delivery per (source,
  63. // dedupe_key) pair, with a max-wait debounce so a continuous
  64. // stream re-flushes every flush window.
  65. collapseState := newFanoutState()
  66. collapser := dedupe.NewCollapser(
  67. time.Duration(cfg.DedupeFlushMs)*time.Millisecond,
  68. onFlushWithKey(logger, collapseState, br),
  69. )
  70. // Subscribe to all alerts.* subjects.
  71. js := br.JS()
  72. stream, err := js.Stream(ctx, "ALERTS")
  73. if err != nil {
  74. logger.Error("nats stream ALERTS", "err", err)
  75. os.Exit(1)
  76. }
  77. consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
  78. Name: "routerd",
  79. Durable: "routerd",
  80. FilterSubjects: []string{"alerts.>"},
  81. AckPolicy: jetstream.AckExplicitPolicy,
  82. })
  83. if err != nil {
  84. logger.Error("nats consumer", "err", err)
  85. os.Exit(1)
  86. }
  87. runCtx, runCancel := context.WithCancel(ctx)
  88. defer runCancel()
  89. go consume(runCtx, logger, consumer, br, resolver, collapser, collapseState)
  90. reg, _ := observability.NewRegistry("routerd")
  91. srv := httpserver.New(httpserver.Config{
  92. Addr: cfg.HTTPAddr,
  93. ServiceName: "routerd",
  94. ShutdownGrace: cfg.ShutdownGrace,
  95. }, logger, observability.MetricsHandler(reg))
  96. errCh := make(chan error, 1)
  97. go func() { errCh <- srv.Start() }()
  98. select {
  99. case <-ctx.Done():
  100. logger.Info("shutdown signal received")
  101. case err := <-errCh:
  102. if err != nil {
  103. logger.Error("http server", "err", err)
  104. os.Exit(1)
  105. }
  106. }
  107. runCancel()
  108. time.Sleep(500 * time.Millisecond) // let consumer drain
  109. // M6.5: drain any pending collapses so we don't lose
  110. // the last few alerts of a graceful shutdown.
  111. collapser.FlushAll()
  112. if err := srv.Shutdown(ctx); err != nil {
  113. logger.Warn("graceful shutdown", "err", err)
  114. }
  115. logger.Info("bye")
  116. }
  117. func consume(
  118. ctx context.Context,
  119. logger *slog.Logger,
  120. c jetstream.Consumer,
  121. br *broker.Client,
  122. r *routing.Resolver,
  123. collapser *dedupe.Collapser,
  124. state *fanoutState,
  125. ) {
  126. // M6.5: start the Collapser's flush loop. It exits
  127. // when ctx is canceled.
  128. collapseCtx, collapseCancel := context.WithCancel(ctx)
  129. defer collapseCancel()
  130. go collapser.Run(collapseCtx)
  131. for {
  132. if ctx.Err() != nil {
  133. return
  134. }
  135. batch, err := c.Fetch(16, jetstream.FetchMaxWait(2*time.Second))
  136. if err != nil {
  137. if ctx.Err() != nil {
  138. return
  139. }
  140. logger.Warn("nats fetch", "err", err)
  141. time.Sleep(500 * time.Millisecond)
  142. continue
  143. }
  144. for m := range batch.Messages() {
  145. handleOne(ctx, logger, m, br, r, collapser, state)
  146. if batch.Error() != nil {
  147. logger.Warn("batch error", "err", batch.Error())
  148. break
  149. }
  150. }
  151. }
  152. }
  153. func handleOne(
  154. ctx context.Context,
  155. logger *slog.Logger,
  156. m jetstream.Msg,
  157. br *broker.Client,
  158. r *routing.Resolver,
  159. collapser *dedupe.Collapser,
  160. state *fanoutState,
  161. ) {
  162. var a alert.Alert
  163. if err := json.Unmarshal(m.Data(), &a); err != nil {
  164. logger.Warn("malformed alert payload", "err", err, "subject", m.Subject())
  165. _ = m.Ack()
  166. return
  167. }
  168. // Extract company_id from subject "alerts.<company_id>".
  169. parts := strings.SplitN(m.Subject(), ".", 2)
  170. if len(parts) != 2 {
  171. logger.Warn("bad subject", "subject", m.Subject())
  172. _ = m.Ack()
  173. return
  174. }
  175. companyID := parts[1]
  176. if a.CompanyID != "" && a.CompanyID != companyID {
  177. logger.Warn("company_id mismatch", "subject", companyID, "body", a.CompanyID)
  178. }
  179. a.CompanyID = companyID
  180. // M6.5: route through the Collapser-aware handler. The
  181. // collapse decision is made here, the actual delivery
  182. // publish may happen later (on flush) — but the NATS
  183. // message is always Ack'd now.
  184. _, shouldAck, shouldNak := observeAndFanout(ctx, logger, collapser, state, r, br, companyID, a)
  185. switch {
  186. case shouldNak:
  187. _ = m.Nak()
  188. case shouldAck:
  189. _ = m.Ack()
  190. }
  191. }
  192. var _ = fmt.Sprintf // keep import