Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
60 changes: 30 additions & 30 deletions cmd/runner/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
)

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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:
Expand All @@ -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)
Expand Down
64 changes: 62 additions & 2 deletions cmd/runner/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
},
Expand Down
86 changes: 86 additions & 0 deletions cmd/runner/metrics_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
Loading