main.go 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121
  1. // Command ingestd receives alerts via HTTP POST / WebSocket / MQTT / gRPC,
  2. // validates, rate-limits, dedupes, and publishes to NATS JetStream.
  3. //
  4. // M0: HTTP POST endpoint only. Other transports land in M5 / M4 / M11.
  5. package main
  6. import (
  7. "context"
  8. "os"
  9. "os/signal"
  10. "syscall"
  11. "git3.techno-world.net/lrosales/broad-announce/internal/broker"
  12. "git3.techno-world.net/lrosales/broad-announce/internal/config"
  13. "git3.techno-world.net/lrosales/broad-announce/internal/dedupe"
  14. "git3.techno-world.net/lrosales/broad-announce/internal/httpserver"
  15. "git3.techno-world.net/lrosales/broad-announce/internal/observability"
  16. "git3.techno-world.net/lrosales/broad-announce/internal/ratelimit"
  17. "git3.techno-world.net/lrosales/broad-announce/internal/store"
  18. )
  19. func main() {
  20. cfg, err := config.LoadIngestd()
  21. if err != nil {
  22. // Logger isn't up yet; stderr is the only thing we have.
  23. os.Stderr.WriteString("config: " + err.Error() + "\n")
  24. os.Exit(1)
  25. }
  26. logger := observability.Init(cfg.Env, cfg.LogLevel, "ingestd")
  27. logger.Info("starting", "env", cfg.Env, "addr", cfg.HTTPAddr)
  28. ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
  29. defer stop()
  30. // Redis
  31. r, err := store.ConnectRedis(ctx, cfg.RedisURL)
  32. if err != nil {
  33. logger.Error("redis connect", "err", err)
  34. os.Exit(1)
  35. }
  36. defer r.Close()
  37. logger.Info("redis connected")
  38. // NATS JetStream
  39. br, err := broker.Connect(ctx, cfg.NATSURL)
  40. if err != nil {
  41. logger.Error("nats connect", "err", err)
  42. os.Exit(1)
  43. }
  44. defer br.Close()
  45. logger.Info("nats connected", "url", cfg.NATSURL)
  46. // Metrics
  47. reg, m := observability.NewRegistry("ingestd")
  48. limiter := ratelimit.New(r.Client)
  49. ded := dedupe.New(r.Client, dedupe.DefaultWindow)
  50. // M0 source registry: loaded from env. M2 replaces with DB.
  51. sources := loadSourcesFromEnv(logger)
  52. js, err := br.NC().JetStream()
  53. if err != nil {
  54. logger.Error("nats jetstream context", "err", err)
  55. os.Exit(1)
  56. }
  57. deps := &httpDeps{
  58. processDeps: processDeps{
  59. Logger: logger.With("component", "http"),
  60. Metrics: m,
  61. Limiter: limiter,
  62. Deduper: ded,
  63. Sources: sources,
  64. JetStream: newNatsPublisher(js),
  65. CompanyRatePerSec: cfg.RateLimitPerCompany,
  66. },
  67. MaxBytes: cfg.MaxPayloadBytes,
  68. }
  69. // HTTP server
  70. srv := httpserver.New(httpserver.Config{
  71. Addr: cfg.HTTPAddr,
  72. ServiceName: "ingestd",
  73. ShutdownGrace: cfg.ShutdownGrace,
  74. }, logger, observability.MetricsHandler(reg))
  75. RegisterRoutes(srv.Mux(), deps)
  76. // MQTT subscriber (M4). Disabled if BA_INGESTD_MQTT_BROKER is
  77. // empty. The subscriber shares the processDeps with the HTTP
  78. // handler so the dedupe window, rate limits, and metrics are
  79. // per-source exactly once across both transports.
  80. mqttCfg := loadMQTTConfig(logger)
  81. mqttErrCh := make(chan error, 1)
  82. go func() {
  83. mqttErrCh <- startMQTT(ctx, mqttCfg, &deps.processDeps, logger, m)
  84. }()
  85. // Run + graceful shutdown
  86. errCh := make(chan error, 1)
  87. go func() { errCh <- srv.Start() }()
  88. select {
  89. case <-ctx.Done():
  90. logger.Info("shutdown signal received")
  91. case err := <-errCh:
  92. if err != nil {
  93. logger.Error("http server", "err", err)
  94. os.Exit(1)
  95. }
  96. case err := <-mqttErrCh:
  97. if err != nil {
  98. logger.Error("mqtt subscriber", "err", err)
  99. os.Exit(1)
  100. }
  101. }
  102. if err := srv.Shutdown(ctx); err != nil {
  103. logger.Warn("graceful shutdown", "err", err)
  104. }
  105. logger.Info("bye")
  106. }