|
@@ -0,0 +1,110 @@
|
|
|
|
|
+// 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(),
|
|
|
|
|
+ )
|
|
|
|
|
+}
|