// Command routerd consumes alerts from NATS JetStream, resolves // recipients (companies → sources → allowed_targets ∪ routing_rules // → groups → individuals ∩ subscriptions), and enqueues one // delivery per (individual, channel) to // deliveries.. subjects. // // M0: connects to NATS, /health, /metrics. No business logic. // M1: broadcast mode — every active fcm_token in the company. // M2: rules engine from SPEC §6 — uses internal/routing.Resolver // which honors source.allowed_targets, routing_rules, and // subscriptions (with min_severity, quiet hours, channel_mask). // Hard-fails on zero recipients (no silent broadcast). package main import ( "context" "encoding/json" "fmt" "log/slog" "os" "os/signal" "strings" "syscall" "time" "git3.techno-world.net/lrosales/broad-announce/internal/alert" "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/postgres" "git3.techno-world.net/lrosales/broad-announce/internal/routing" "github.com/nats-io/nats.go/jetstream" ) // routerdMetrics is the package-level routerd metrics instance. // Set in main() before the consume goroutine starts. var routerdMetrics *observability.RouterdMetrics func main() { cfg, err := config.LoadRouterd() if err != nil { os.Stderr.WriteString("config: " + err.Error() + "\n") os.Exit(1) } logger := observability.Init(cfg.Env, cfg.LogLevel, "routerd") logger.Info("starting", "env", cfg.Env, "addr", cfg.HTTPAddr, "dedupe_flush_ms", cfg.DedupeFlushMs, ) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() br, err := broker.Connect(ctx, cfg.NATSURL) if err != nil { logger.Error("nats connect", "err", err) os.Exit(1) } defer br.Close() pool, err := postgres.Connect(ctx, cfg.PostgresDSN) if err != nil { logger.Error("postgres connect", "err", err) os.Exit(1) } defer pool.Close() resolver := routing.New(pool, logger.With("subsystem", "routing")) // M6.5: build the router-level dedupe Collapser. It collapses // bursts of identical alerts into 1 delivery per (source, // dedupe_key) pair, with a max-wait debounce so a continuous // stream re-flushes every flush window. collapseState := newFanoutState() collapser := dedupe.NewCollapser( time.Duration(cfg.DedupeFlushMs)*time.Millisecond, onFlushWithKey(logger, collapseState, br), ) // Subscribe to all alerts.* subjects. js := br.JS() stream, err := js.Stream(ctx, "ALERTS") if err != nil { logger.Error("nats stream ALERTS", "err", err) os.Exit(1) } consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{ Name: "routerd", Durable: "routerd", FilterSubjects: []string{"alerts.>"}, AckPolicy: jetstream.AckExplicitPolicy, }) if err != nil { logger.Error("nats consumer", "err", err) os.Exit(1) } runCtx, runCancel := context.WithCancel(ctx) defer runCancel() // M9 L1: routerd metrics. Create before the consumer starts // so the /metrics endpoint is valid immediately. reg, _ := observability.NewRegistry("routerd") routerdMetrics = observability.NewRouterdMetrics(reg, "routerd") go consume(runCtx, logger, consumer, br, resolver, collapser, collapseState) srv := httpserver.New(httpserver.Config{ Addr: cfg.HTTPAddr, ServiceName: "routerd", ShutdownGrace: cfg.ShutdownGrace, }, logger, observability.MetricsHandler(reg)) 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) } } runCancel() time.Sleep(500 * time.Millisecond) // let consumer drain // M6.5: drain any pending collapses so we don't lose // the last few alerts of a graceful shutdown. collapser.FlushAll() if err := srv.Shutdown(ctx); err != nil { logger.Warn("graceful shutdown", "err", err) } logger.Info("bye") } func consume( ctx context.Context, logger *slog.Logger, c jetstream.Consumer, br *broker.Client, r *routing.Resolver, collapser *dedupe.Collapser, state *fanoutState, ) { // M6.5: start the Collapser's flush loop. It exits // when ctx is canceled. collapseCtx, collapseCancel := context.WithCancel(ctx) defer collapseCancel() go collapser.Run(collapseCtx) for { if ctx.Err() != nil { return } batch, err := c.Fetch(16, jetstream.FetchMaxWait(2*time.Second)) if err != nil { if ctx.Err() != nil { return } logger.Warn("nats fetch", "err", err) time.Sleep(500 * time.Millisecond) continue } for m := range batch.Messages() { handleOne(ctx, logger, m, br, r, collapser, state) if batch.Error() != nil { logger.Warn("batch error", "err", batch.Error()) break } } } } func handleOne( ctx context.Context, logger *slog.Logger, m jetstream.Msg, br *broker.Client, r *routing.Resolver, collapser *dedupe.Collapser, state *fanoutState, ) { var a alert.Alert if err := json.Unmarshal(m.Data(), &a); err != nil { logger.Warn("malformed alert payload", "err", err, "subject", m.Subject()) _ = m.Ack() return } // Extract company_id from subject "alerts.". parts := strings.SplitN(m.Subject(), ".", 2) if len(parts) != 2 { logger.Warn("bad subject", "subject", m.Subject()) _ = m.Ack() return } companyID := parts[1] if a.CompanyID != "" && a.CompanyID != companyID { logger.Warn("company_id mismatch", "subject", companyID, "body", a.CompanyID) } a.CompanyID = companyID // M6.5: route through the Collapser-aware handler. The // collapse decision is made here, the actual delivery // publish may happen later (on flush) — but the NATS // message is always Ack'd now. _, shouldAck, shouldNak := observeAndFanout(ctx, logger, collapser, state, r, br, companyID, a) switch { case shouldNak: _ = m.Nak() case shouldAck: _ = m.Ack() } } var _ = fmt.Sprintf // keep import