main.go 2.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  1. // cmd/deliverd-bench is a no-op deliverd for M10-bench only.
  2. // It consumes from NATS (ba.<company>.deliveries) and ACKs every message
  3. // immediately — no Postgres, no FCM, no Telegram, no retry, no DLQ.
  4. // Purpose: prove the broker+router can sustain 50k/s without the
  5. // delivery tier becoming a bottleneck.
  6. //
  7. // Usage (bench profile):
  8. // docker compose --profile bench up -d
  9. package main
  10. import (
  11. "context"
  12. "fmt"
  13. "log/slog"
  14. "net/http"
  15. "os"
  16. "os/signal"
  17. "sync/atomic"
  18. "syscall"
  19. "time"
  20. "github.com/nats-io/nats.go"
  21. )
  22. func main() {
  23. natsURL := os.Getenv("BA_NATS_URL")
  24. if natsURL == "" {
  25. natsURL = "nats://localhost:4222"
  26. }
  27. logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))
  28. logger.Info("starting", "nats_url", natsURL)
  29. ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
  30. defer stop()
  31. nc, err := nats.Connect(natsURL,
  32. nats.Name("deliverd-bench"),
  33. nats.MaxReconnects(-1),
  34. )
  35. if err != nil {
  36. logger.Error("nats connect", "err", err)
  37. os.Exit(1)
  38. }
  39. defer nc.Close()
  40. logger.Info("nats connected", "url", natsURL)
  41. js, err := nc.JetStream()
  42. if err != nil {
  43. logger.Error("jetstream", "err", err)
  44. os.Exit(1)
  45. }
  46. var (
  47. consumed atomic.Uint64
  48. failed atomic.Uint64
  49. )
  50. sub, err := js.Subscribe("ba.*.deliveries", func(msg *nats.Msg) {
  51. consumed.Add(1)
  52. if err := msg.Ack(); err != nil {
  53. failed.Add(1)
  54. logger.Debug("ack failed", "err", err)
  55. }
  56. }, nats.Durable("deliverd-bench"), nats.AckExplicit())
  57. if err != nil {
  58. logger.Error("subscribe", "err", err)
  59. os.Exit(1)
  60. }
  61. defer sub.Unsubscribe()
  62. logger.Info("deliverd-bench listening on ba.*.deliveries")
  63. // Prometheus metrics endpoint.
  64. if addr := os.Getenv("BA_METRICS_ADDR"); addr != "" {
  65. mux := http.NewServeMux()
  66. mux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) {
  67. fmt.Fprintf(w, "deliverd_bench_consumed_total %d\n", consumed.Load())
  68. fmt.Fprintf(w, "deliverd_bench_failed_total %d\n", failed.Load())
  69. })
  70. srv := &http.Server{Addr: addr, Handler: mux, ReadHeaderTimeout: 5 * time.Second}
  71. go srv.ListenAndServe()
  72. logger.Info("metrics", "addr", addr)
  73. }
  74. // Status ticker every 10s.
  75. go func() {
  76. ticker := time.NewTicker(10 * time.Second)
  77. defer ticker.Stop()
  78. for {
  79. select {
  80. case <-ctx.Done():
  81. return
  82. case <-ticker.C:
  83. cons := consumed.Load()
  84. fail := failed.Load()
  85. logger.Info("bench stats", "consumed", cons, "failed", fail)
  86. if cons > 100 && fail > cons/10 {
  87. logger.Warn("high ack failure rate", "consumed", cons, "failed", fail)
  88. }
  89. }
  90. }
  91. }()
  92. <-ctx.Done()
  93. logger.Info("deliverd-bench shutting down",
  94. "consumed", consumed.Load(),
  95. "failed", failed.Load(),
  96. )
  97. }