// 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") }