| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223 |
- // 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.<channel>.<company_id> 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))
- // M13a W5: admin routes (JWT-gated). When BA_AUTHD_JWT_SECRET
- // is unset, wireAdminRoutes is a no-op so the LAN deploy
- // path keeps working unchanged.
- wireAdminRoutes(srv.Mux(), collapser, collapseState, logger)
- 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.<company_id>".
- 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
|