From 8253a3c50bd2cdf7fc4c10371e0da43b26152dab Mon Sep 17 00:00:00 2001 From: Rob Holland Date: Fri, 25 Sep 2026 14:34:37 +0100 Subject: [PATCH] Emit client-observed benchmark runner metrics --- README.md | 17 ++++++++ cmd/runner/main.go | 60 +++++++++++++------------- cmd/runner/metrics.go | 64 +++++++++++++++++++++++++++- cmd/runner/metrics_test.go | 86 ++++++++++++++++++++++++++++++++++++++ 4 files changed, 195 insertions(+), 32 deletions(-) create mode 100644 cmd/runner/metrics_test.go diff --git a/README.md b/README.md index 8005e03..96e9ee0 100644 --- a/README.md +++ b/README.md @@ -140,6 +140,23 @@ kubectl run benchmark-runner --image ghcr.io/temporalio/benchmark-workers:main \ --command -- runner -t ExecuteActivity '{ "Count": 3, "Activity": "Echo", "Input": { "Message": "test" } }' ``` +When `PROMETHEUS_ENDPOINT` is set, the runner serves its own `benchmark_` +metrics alongside the Temporal SDK metrics on `/metrics`: + +| Metric | Meaning | +| --- | --- | +| `benchmark_runner_invocations_started_total` | Top-level workflow attempts, counted before the client submits them | +| `benchmark_runner_invocations_completed_total` | Workflows the client waited for and observed complete successfully | +| `benchmark_runner_invocations_failed_total` | Submission, signalling, or completion failures observed by the client | +| `benchmark_runner_invocation_duration_seconds{outcome="completed"\|"failed"}` | Client-observed time from before submission through completion or failure | + +The histogram has fine-grained buckets from 5 ms through 60 s, then tail +buckets up to 10 minutes. Use the `completed` outcome for latency percentiles; +the `failed` outcome captures time to error. With `-w=false`, a successful +submission increments only `started_total`, because the client has not +observed completion. These are runner-side measurements; the SDK's own metrics +remain available for diagnosing internal behavior. + ## Workflows The worker provides the following workflows for you to use during benchmarking: diff --git a/cmd/runner/main.go b/cmd/runner/main.go index b001513..7530055 100644 --- a/cmd/runner/main.go +++ b/cmd/runner/main.go @@ -22,14 +22,14 @@ import ( ) var ( - nWorkflows = flag.Int("c", 10, "concurrent workflows") - sWorkflow = flag.String("t", "", "workflow type") - sSignalType = flag.String("s", "", "signal type") - bWait = flag.Bool("w", true, "wait for workflows to complete") - sNamespace = flag.String("n", "default", "namespace") - sTaskQueue = flag.String("tq", "benchmark", "task queue") - nMaxInterval = flag.Int("max-interval", 60, "maximum interval (in seconds) for exponential backoff") - nFactor = flag.Int("backoff-factor", 2, "factor for exponential backoff") + nWorkflows = flag.Int("c", 10, "concurrent workflows") + sWorkflow = flag.String("t", "", "workflow type") + sSignalType = flag.String("s", "", "signal type") + bWait = flag.Bool("w", true, "wait for workflows to complete") + sNamespace = flag.String("n", "default", "namespace") + sTaskQueue = flag.String("tq", "benchmark", "task queue") + nMaxInterval = flag.Int("max-interval", 60, "maximum interval (in seconds) for exponential backoff") + nFactor = flag.Int("backoff-factor", 2, "factor for exponential backoff") bDisableBackoff = flag.Bool("disable-backoff", false, "disable exponential backoff on errors") ) @@ -156,11 +156,12 @@ func main() { clientOptions.Credentials = client.NewAPIKeyStaticCredentials(apiKey) } + invocations := newInvocationMetrics() if os.Getenv("PROMETHEUS_ENDPOINT") != "" { clientOptions.MetricsHandler = sdktally.NewMetricsHandler(newPrometheusScope(prometheus.Configuration{ ListenAddress: os.Getenv("PROMETHEUS_ENDPOINT"), TimerType: "histogram", - })) + }, invocations.registry)) } c, err := client.Dial(clientOptions) @@ -217,32 +218,31 @@ func main() { go (func() { currentInterval := 1 errChan := make(chan error, concurrentWorkflows) - + for { pool.Submit(func() { - wf, err := starter() - if err != nil { - fmt.Fprintf(os.Stderr, "Unable to start workflow: %v\n", err) - errChan <- err - return - } - - if waitForCompletion { - err = wf.Get(context.Background(), nil) + err := invocations.run(waitForCompletion, func() error { + wf, err := starter() if err != nil { - fmt.Fprintf(os.Stderr, "Workflow failed: %v\n", err) - errChan <- err - return + fmt.Fprintf(os.Stderr, "Unable to start workflow: %v\n", err) + return err } - } - - errChan <- nil + if waitForCompletion { + err = wf.Get(context.Background(), nil) + if err != nil { + fmt.Fprintf(os.Stderr, "Workflow failed: %v\n", err) + } + return err + } + return nil + }) + errChan <- err }) - + var lastErr error updated := false - - drainLoop: + + drainLoop: for { select { case err := <-errChan: @@ -252,11 +252,11 @@ func main() { break drainLoop } } - + if disableBackOff || !updated { continue } - + if lastErr != nil { fmt.Fprintf(os.Stderr, "Waiting for %d seconds before retrying to start workflow...\n", currentInterval) time.Sleep(time.Duration(currentInterval) * time.Second) diff --git a/cmd/runner/metrics.go b/cmd/runner/metrics.go index a72e41a..dfb9788 100644 --- a/cmd/runner/metrics.go +++ b/cmd/runner/metrics.go @@ -10,10 +10,70 @@ import ( sdktally "go.temporal.io/sdk/contrib/tally" ) -func newPrometheusScope(c prometheus.Configuration) tally.Scope { +// Dense buckets below a minute make p90/p99 useful for short invocations, +// while the tail keeps slower workloads visible instead of pinning them at +// 60 seconds. +var invocationDurationBuckets = []float64{ + 0.005, 0.0075, 0.01, 0.015, 0.02, 0.03, 0.04, 0.05, + 0.075, 0.1, 0.15, 0.2, 0.25, 0.3, 0.4, 0.5, 0.75, + 1, 1.25, 1.5, 2, 2.5, 3, 4, 5, 6, 7.5, 10, + 12.5, 15, 20, 25, 30, 40, 45, 50, 60, + 90, 120, 180, 300, 600, +} + +type invocationMetrics struct { + registry *prom.Registry + started prom.Counter + completed prom.Counter + failed prom.Counter + duration *prom.HistogramVec +} + +func newInvocationMetrics() *invocationMetrics { + m := &invocationMetrics{ + registry: prom.NewRegistry(), + started: prom.NewCounter(prom.CounterOpts{ + Name: "benchmark_runner_invocations_started_total", + Help: "Number of top-level invocation attempts made by the runner.", + }), + completed: prom.NewCounter(prom.CounterOpts{ + Name: "benchmark_runner_invocations_completed_total", + Help: "Number of top-level invocations the runner waited for and observed complete successfully.", + }), + failed: prom.NewCounter(prom.CounterOpts{ + Name: "benchmark_runner_invocations_failed_total", + Help: "Number of top-level invocations that failed to start or failed while waiting for completion.", + }), + duration: prom.NewHistogramVec(prom.HistogramOpts{ + Name: "benchmark_runner_invocation_duration_seconds", + Help: "Client-observed wall time from before submission to completion or failure.", + Buckets: invocationDurationBuckets, + }, []string{"outcome"}), + } + m.registry.MustRegister(m.started, m.completed, m.failed, m.duration) + return m +} + +// run measures one top-level invocation. A successful non-waiting submission +// has no known completion time, so it is neither a completion nor a duration. +func (m *invocationMetrics) run(wait bool, invoke func() error) error { + m.started.Inc() + start := time.Now() + err := invoke() + if err != nil { + m.failed.Inc() + m.duration.WithLabelValues("failed").Observe(time.Since(start).Seconds()) + } else if wait { + m.completed.Inc() + m.duration.WithLabelValues("completed").Observe(time.Since(start).Seconds()) + } + return err +} + +func newPrometheusScope(c prometheus.Configuration, registry *prom.Registry) tally.Scope { reporter, err := c.NewReporter( prometheus.ConfigurationOptions{ - Registry: prom.NewRegistry(), + Registry: registry, OnError: func(err error) { log.Println("error in prometheus reporter", err) }, diff --git a/cmd/runner/metrics_test.go b/cmd/runner/metrics_test.go new file mode 100644 index 0000000..d67be3d --- /dev/null +++ b/cmd/runner/metrics_test.go @@ -0,0 +1,86 @@ +package main + +import ( + "errors" + "testing" +) + +func TestInvocationMetricsContract(t *testing.T) { + m := newInvocationMetrics() + failure := errors.New("invocation failed") + + if err := m.run(true, func() error { return nil }); err != nil { + t.Fatalf("waited success: %v", err) + } + if err := m.run(true, func() error { return failure }); !errors.Is(err, failure) { + t.Fatalf("waited failure: got %v, want %v", err, failure) + } + if err := m.run(false, func() error { return nil }); err != nil { + t.Fatalf("unwaited submission: %v", err) + } + if err := m.run(false, func() error { return failure }); !errors.Is(err, failure) { + t.Fatalf("unwaited submission failure: got %v, want %v", err, failure) + } + + families, err := m.registry.Gather() + if err != nil { + t.Fatal(err) + } + counters := map[string]float64{ + "benchmark_runner_invocations_started_total": 4, + "benchmark_runner_invocations_completed_total": 1, + "benchmark_runner_invocations_failed_total": 2, + } + histogramSeen := false + for _, family := range families { + if want, ok := counters[family.GetName()]; ok { + if len(family.GetMetric()) != 1 || family.GetMetric()[0].GetCounter().GetValue() != want { + t.Errorf("%s: got %v, want %v", family.GetName(), family.GetMetric(), want) + } + delete(counters, family.GetName()) + } + if family.GetName() != "benchmark_runner_invocation_duration_seconds" { + continue + } + histogramSeen = true + counts := map[string]uint64{"completed": 1, "failed": 2} + for _, metric := range family.GetMetric() { + if len(metric.GetLabel()) != 1 || metric.GetLabel()[0].GetName() != "outcome" { + t.Errorf("unexpected duration labels: %v", metric.GetLabel()) + continue + } + outcome := metric.GetLabel()[0].GetValue() + want, ok := counts[outcome] + if !ok { + t.Errorf("unexpected outcome: %s", outcome) + continue + } + histogram := metric.GetHistogram() + if histogram.GetSampleCount() != want { + t.Errorf("%s sample count: got %d, want %d", outcome, histogram.GetSampleCount(), want) + } + for _, bound := range []float64{0.005, 0.1, 1, 60, 600} { + found := false + for _, bucket := range histogram.GetBucket() { + if bucket.GetUpperBound() == bound { + found = true + break + } + } + if !found { + t.Errorf("%s histogram missing bucket %g", outcome, bound) + } + } + delete(counts, outcome) + } + if len(counts) != 0 { + t.Errorf("missing duration outcomes: %v", counts) + } + } + if len(counters) != 0 { + t.Errorf("missing counters: %v", counters) + } + if !histogramSeen { + t.Error("missing invocation duration histogram") + } +}