From d11802d50225249e066676060f57275847b2dba3 Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Sat, 3 Oct 2026 17:16:00 +0000 Subject: [PATCH] Report activity schedule to start latency once per task without the activity type --- .../internal/worker/ActivityPollTask.java | 6 ------ .../internal/worker/ActivityWorker.java | 13 +++++++------ .../internal/worker/AsyncActivityPollTask.java | 6 ------ .../common/reporter/TestStatsReporter.java | 17 +++++++++++++++++ .../java/io/temporal/workflow/MetricsTest.java | 12 ++++++++++++ 5 files changed, 36 insertions(+), 18 deletions(-) diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityPollTask.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityPollTask.java index f0d3e649f0..4405bad8c8 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityPollTask.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityPollTask.java @@ -10,7 +10,6 @@ import io.temporal.api.workflowservice.v1.GetSystemInfoResponse; import io.temporal.api.workflowservice.v1.PollActivityTaskQueueRequest; import io.temporal.api.workflowservice.v1.PollActivityTaskQueueResponse; -import io.temporal.internal.common.ProtobufTimeUtils; import io.temporal.serviceclient.MetricsTag; import io.temporal.serviceclient.WorkflowServiceStubs; import io.temporal.worker.MetricsType; @@ -120,11 +119,6 @@ public ActivityTask poll() { metricsScope.counter(MetricsType.ACTIVITY_POLL_NO_TASK_COUNTER).inc(1); return null; } - metricsScope - .timer(MetricsType.ACTIVITY_SCHEDULE_TO_START_LATENCY) - .record( - ProtobufTimeUtils.toM3Duration( - response.getStartedTime(), response.getCurrentAttemptScheduledTime())); isSuccessful = true; pollerTracker.pollSucceeded(); return new ActivityTask( diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityWorker.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityWorker.java index 5aa5293175..b6f59a92a9 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityWorker.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/ActivityWorker.java @@ -333,6 +333,13 @@ public void handle(ActivityTask task) throws Exception { ActivityTaskHandler.Result result = null; boolean taskFailed = false; try { + // Schedule to start latency does not depend on the activity type, so it is reported with + // the worker scope. + workerMetricsScope + .timer(MetricsType.ACTIVITY_SCHEDULE_TO_START_LATENCY) + .record( + ProtobufTimeUtils.toM3Duration( + pollResponse.getStartedTime(), pollResponse.getCurrentAttemptScheduledTime())); result = handleActivity(task, metricsScope); if (result.getTaskFailed() != null && !io.temporal.internal.common.FailureUtils.isBenignApplicationFailure( @@ -382,12 +389,6 @@ private ActivityTaskHandler.Result handleActivity(ActivityTask task, Scope metri false); } PollActivityTaskQueueResponseOrBuilder pollResponse = task.getResponse(); - metricsScope - .timer(MetricsType.ACTIVITY_SCHEDULE_TO_START_LATENCY) - .record( - ProtobufTimeUtils.toM3Duration( - pollResponse.getStartedTime(), pollResponse.getCurrentAttemptScheduledTime())); - ActivityTaskHandler.Result result; Stopwatch sw = metricsScope.timer(MetricsType.ACTIVITY_EXEC_LATENCY).start(); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/AsyncActivityPollTask.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/AsyncActivityPollTask.java index 1e8791bd02..5d6b9cc5f2 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/AsyncActivityPollTask.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/AsyncActivityPollTask.java @@ -12,7 +12,6 @@ import io.temporal.api.workflowservice.v1.PollActivityTaskQueueRequest; import io.temporal.api.workflowservice.v1.PollActivityTaskQueueResponse; import io.temporal.internal.common.GrpcUtils; -import io.temporal.internal.common.ProtobufTimeUtils; import io.temporal.serviceclient.MetricsTag; import io.temporal.serviceclient.WorkflowServiceStubs; import io.temporal.worker.MetricsType; @@ -121,11 +120,6 @@ public CompletableFuture poll(SlotPermit permit) { return null; } pollerTracker.pollSucceeded(); - metricsScope - .timer(MetricsType.ACTIVITY_SCHEDULE_TO_START_LATENCY) - .record( - ProtobufTimeUtils.toM3Duration( - r.getStartedTime(), r.getCurrentAttemptScheduledTime())); return new ActivityTask( r, permit, diff --git a/temporal-sdk/src/test/java/io/temporal/common/reporter/TestStatsReporter.java b/temporal-sdk/src/test/java/io/temporal/common/reporter/TestStatsReporter.java index 6e03278164..d0c9735a2c 100644 --- a/temporal-sdk/src/test/java/io/temporal/common/reporter/TestStatsReporter.java +++ b/temporal-sdk/src/test/java/io/temporal/common/reporter/TestStatsReporter.java @@ -39,6 +39,9 @@ public synchronized void assertNoMetric(String name, Map tags) { + counters.get(metricName).get() + "'"); } + if (timers.containsKey(metricName)) { + fail("Timer '" + metricName + "' was reported"); + } } public synchronized void assertCounter(String name, Map tags, long expected) { @@ -89,6 +92,20 @@ public synchronized void assertTimer(String name, Map tags) { } } + public synchronized void assertTimer( + String name, Map tags, Predicate isExpected) { + String metricName = getMetricName(name, tags); + StatsAccumulator value = timers.get(metricName); + if (value == null) { + fail( + "No metric '" + + metricName + + "', reported metrics: \n " + + String.join("\n ", timers.keySet())); + } + assertTrue(metricName + " count: " + value.count(), isExpected.test(value)); + } + public synchronized void assertTimerMinDuration( String name, Map tags, Duration minDuration) { String metricName = getMetricName(name, tags); diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/MetricsTest.java b/temporal-sdk/src/test/java/io/temporal/workflow/MetricsTest.java index 1a4a3d2227..743838f6d0 100644 --- a/temporal-sdk/src/test/java/io/temporal/workflow/MetricsTest.java +++ b/temporal-sdk/src/test/java/io/temporal/workflow/MetricsTest.java @@ -307,6 +307,18 @@ public void testWorkflowMetrics() throws InterruptedException { reporter.assertCounter(TEMPORAL_REQUEST, workflowTags, 1); reporter.assertTimer(TEMPORAL_REQUEST_LATENCY, workflowTags); + // The workflow ran a single activity task. Schedule to start latency is independent of the + // activity type, so it must be reported exactly once and without the activity type tags. + reporter.assertTimer( + ACTIVITY_SCHEDULE_TO_START_LATENCY, TAGS_ACTIVITY_WORKER, stats -> stats.count() == 1); + reporter.assertNoMetric( + ACTIVITY_SCHEDULE_TO_START_LATENCY, + new ImmutableMap.Builder() + .putAll(TAGS_ACTIVITY_WORKER) + .put(MetricsTag.ACTIVITY_TYPE, "Execute") + .put(MetricsTag.WORKFLOW_TYPE, "NoArgsWorkflow") + .build()); + Map workflowTaskCompletionTags = new ImmutableMap.Builder() .putAll(TAGS_WORKFLOW_WORKER)