// cmd/deliverd-bench is a no-op deliverd for M10-bench only. // It consumes from NATS (ba..deliveries) and ACKs every message // immediately — no Postgres, no FCM, no Telegram, no retry, no DLQ. // Purpose: prove the broker+router can sustain 50k/s without the // delivery tier becoming a bottleneck. // // Usage (bench profile): // docker compose --profile bench up -d package main import ( "context" "fmt" "log/slog" "net/http" "os" "os/signal" "sync/atomic" "syscall" "time" "github.com/nats-io/nats.go" ) func main() { natsURL := os.Getenv("BA_NATS_URL") if natsURL == "" { natsURL = "nats://localhost:4222" } logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo})) logger.Info("starting", "nats_url", natsURL) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() nc, err := nats.Connect(natsURL, nats.Name("deliverd-bench"), nats.MaxReconnects(-1), ) if err != nil { logger.Error("nats connect", "err", err) os.Exit(1) } defer nc.Close() logger.Info("nats connected", "url", natsURL) js, err := nc.JetStream() if err != nil { logger.Error("jetstream", "err", err) os.Exit(1) } var ( consumed atomic.Uint64 failed atomic.Uint64 ) sub, err := js.Subscribe("ba.*.deliveries", func(msg *nats.Msg) { consumed.Add(1) if err := msg.Ack(); err != nil { failed.Add(1) logger.Debug("ack failed", "err", err) } }, nats.Durable("deliverd-bench"), nats.AckExplicit()) if err != nil { logger.Error("subscribe", "err", err) os.Exit(1) } defer sub.Unsubscribe() logger.Info("deliverd-bench listening on ba.*.deliveries") // Prometheus metrics endpoint. if addr := os.Getenv("BA_METRICS_ADDR"); addr != "" { mux := http.NewServeMux() mux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { fmt.Fprintf(w, "deliverd_bench_consumed_total %d\n", consumed.Load()) fmt.Fprintf(w, "deliverd_bench_failed_total %d\n", failed.Load()) }) srv := &http.Server{Addr: addr, Handler: mux, ReadHeaderTimeout: 5 * time.Second} go srv.ListenAndServe() logger.Info("metrics", "addr", addr) } // Status ticker every 10s. go func() { ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: cons := consumed.Load() fail := failed.Load() logger.Info("bench stats", "consumed", cons, "failed", fail) if cons > 100 && fail > cons/10 { logger.Warn("high ack failure rate", "consumed", cons, "failed", fail) } } } }() <-ctx.Done() logger.Info("deliverd-bench shutting down", "consumed", consumed.Load(), "failed", failed.Load(), ) }