| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110 |
- // cmd/deliverd-bench is a no-op deliverd for M10-bench only.
- // It consumes from NATS (ba.<company>.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(),
- )
- }
|