| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233 |
- // 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/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.<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
- 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.<channel>.<company_id> 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.<channel>.<company_id>. 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
|