// 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/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" ) func main() { cfg, err := config.LoadCommon("routerd") 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) 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")) // 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() go consume(runCtx, logger, consumer, br, resolver) reg, _ := observability.NewRegistry("routerd") 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 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) { 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) 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) { 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 targets, err := r.ResolveTargets(ctx, &a) if err != nil { logger.Error("resolve targets", "err", err, "alert_id", a.ID, "company", companyID) // Nack so the message is redelivered. In M9 we add the // circuit breaker; for M2 we just retry. _ = m.Nak() return } if len(targets) == 0 { // M2 hard-fail. The alert is acked (it can't be delivered, // retrying won't help) and the warning is logged. M3+ will // route these to dlq.no_recipients. logger.Warn("zero recipients, dropping", "alert_id", a.ID, "company", companyID, "source_id", a.SourceID, "severity", string(a.Severity), "category", a.Category, ) _ = m.Ack() return } // Enqueue one deliveries.. per (target). js, err := br.NC().JetStream() if err != nil { logger.Error("js ctx", "err", err) _ = m.Nak() return } enqueued := 0 for _, t := range targets { envelope := deliveryEnvelope{ Alert: a, IndividualID: t.IndividualID, Channel: t.Channel, Endpoint: t.Endpoint, Locale: t.Locale, } body, err := json.Marshal(envelope) if err != nil { logger.Warn("marshal envelope", "err", err) continue } subject := broker.DeliveriesSubject(t.Channel, companyID) if _, err := js.PublishAsync(subject, body); err != nil { logger.Warn("publish delivery", "err", err, "subject", subject) continue } enqueued++ } logger.Info("routed", "alert_id", a.ID, "company", companyID, "source_id", a.SourceID, "severity", string(a.Severity), "recipients", len(targets), "enqueued", enqueued, ) _ = m.Ack() } // deliveryEnvelope is the wire shape published on // deliveries... M2 swaps FCMToken+Locale for // the channel-agnostic Channel+Endpoint so the same shape works // for telegram / sms / email / etc. in M3+. type deliveryEnvelope struct { Alert alert.Alert `json:"alert"` IndividualID string `json:"individual_id"` Channel string `json:"channel"` Endpoint string `json:"endpoint"` Locale string `json:"locale,omitempty"` } var _ = fmt.Sprintf // keep import