| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121 |
- // Command ingestd receives alerts via HTTP POST / WebSocket / MQTT / gRPC,
- // validates, rate-limits, dedupes, and publishes to NATS JetStream.
- //
- // M0: HTTP POST endpoint only. Other transports land in M5 / M4 / M11.
- package main
- import (
- "context"
- "os"
- "os/signal"
- "syscall"
- "git3.techno-world.net/lrosales/broad-announce/internal/broker"
- "git3.techno-world.net/lrosales/broad-announce/internal/config"
- "git3.techno-world.net/lrosales/broad-announce/internal/dedupe"
- "git3.techno-world.net/lrosales/broad-announce/internal/httpserver"
- "git3.techno-world.net/lrosales/broad-announce/internal/observability"
- "git3.techno-world.net/lrosales/broad-announce/internal/ratelimit"
- "git3.techno-world.net/lrosales/broad-announce/internal/store"
- )
- func main() {
- cfg, err := config.LoadIngestd()
- if err != nil {
- // Logger isn't up yet; stderr is the only thing we have.
- os.Stderr.WriteString("config: " + err.Error() + "\n")
- os.Exit(1)
- }
- logger := observability.Init(cfg.Env, cfg.LogLevel, "ingestd")
- logger.Info("starting", "env", cfg.Env, "addr", cfg.HTTPAddr)
- ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
- defer stop()
- // Redis
- r, err := store.ConnectRedis(ctx, cfg.RedisURL)
- if err != nil {
- logger.Error("redis connect", "err", err)
- os.Exit(1)
- }
- defer r.Close()
- logger.Info("redis connected")
- // NATS JetStream
- br, err := broker.Connect(ctx, cfg.NATSURL)
- if err != nil {
- logger.Error("nats connect", "err", err)
- os.Exit(1)
- }
- defer br.Close()
- logger.Info("nats connected", "url", cfg.NATSURL)
- // Metrics
- reg, m := observability.NewRegistry("ingestd")
- limiter := ratelimit.New(r.Client)
- ded := dedupe.New(r.Client, dedupe.DefaultWindow)
- // M0 source registry: loaded from env. M2 replaces with DB.
- sources := loadSourcesFromEnv(logger)
- js, err := br.NC().JetStream()
- if err != nil {
- logger.Error("nats jetstream context", "err", err)
- os.Exit(1)
- }
- deps := &httpDeps{
- processDeps: processDeps{
- Logger: logger.With("component", "http"),
- Metrics: m,
- Limiter: limiter,
- Deduper: ded,
- Sources: sources,
- JetStream: newNatsPublisher(js),
- CompanyRatePerSec: cfg.RateLimitPerCompany,
- },
- MaxBytes: cfg.MaxPayloadBytes,
- }
- // HTTP server
- srv := httpserver.New(httpserver.Config{
- Addr: cfg.HTTPAddr,
- ServiceName: "ingestd",
- ShutdownGrace: cfg.ShutdownGrace,
- }, logger, observability.MetricsHandler(reg))
- RegisterRoutes(srv.Mux(), deps)
- // MQTT subscriber (M4). Disabled if BA_INGESTD_MQTT_BROKER is
- // empty. The subscriber shares the processDeps with the HTTP
- // handler so the dedupe window, rate limits, and metrics are
- // per-source exactly once across both transports.
- mqttCfg := loadMQTTConfig(logger)
- mqttErrCh := make(chan error, 1)
- go func() {
- mqttErrCh <- startMQTT(ctx, mqttCfg, &deps.processDeps, logger, m)
- }()
- // Run + graceful shutdown
- errCh := make(chan error, 1)
- go func() { errCh <- srv.Start() }()
- select {
- case <-ctx.Done():
- logger.Info("shutdown signal received")
- case err := <-errCh:
- if err != nil {
- logger.Error("http server", "err", err)
- os.Exit(1)
- }
- case err := <-mqttErrCh:
- if err != nil {
- logger.Error("mqtt subscriber", "err", err)
- os.Exit(1)
- }
- }
- if err := srv.Shutdown(ctx); err != nil {
- logger.Warn("graceful shutdown", "err", err)
- }
- logger.Info("bye")
- }
|