diff --git a/CHANGELOG.md b/CHANGELOG.md index af4bac6c19..c94b1e8046 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,6 +27,9 @@ to docs, or any other relevant information. Child workflow overrides and one-time routing require Temporal Server 1.32.0 or later. - `WorkerFactoryOptions.Builder.setLoggerTagPrefix` that can be used to customized structured logging tags (MDC keys) set by Temporal SDK in worker context. +- Workflow tasks taking longer than 5 seconds now log a `[TMPRL1104]` warning reporting the task duration and the + external storage downloads and uploads that contributed to it. The threshold is configurable with the + `TEMPORAL_WORKFLOW_TASK_DURATION_WARN_SECONDS` environment variable. ### Changed - Release notes for all future releases are now in a single CHANGELOG.md file. `releases` directory with old release diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java index f88789fe42..8064c70e36 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java @@ -53,6 +53,14 @@ CompletableFuture> store( List payloads, @Nullable StorageDriverTargetInfo target, CancellationToken cancellationToken) { + return store(payloads, target, cancellationToken, null); + } + + CompletableFuture> store( + List payloads, + @Nullable StorageDriverTargetInfo target, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { StorageDriverSelectContext selectContext = new StorageDriverSelectContextImpl(target, cancellationToken); Map> batches; @@ -64,7 +72,7 @@ CompletableFuture> store( if (batches.isEmpty()) { return CompletableFuture.completedFuture(payloads); } - return runStoreDrivers(batches, target, cancellationToken) + return runStoreDrivers(batches, target, cancellationToken, metrics) .thenApply(referencePayloads -> applyPayloadReplacements(payloads, referencePayloads)); } @@ -94,16 +102,24 @@ private Map> buildStoreBatches( private CompletableFuture>> runStoreDrivers( Map> batches, @Nullable StorageDriverTargetInfo target, - CancellationToken cancellationToken) { + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { return withDriverScope( cancellationToken, scope -> { StorageDriverStoreContext context = new StorageDriverStoreContextImpl(target, scope.token()); for (Batch batch : batches.values()) { + long startNanos = System.nanoTime(); scope .attach(batch.driver.store(context, batch.values())) - .map(claims -> createReferencePayloads(batch, claims)); + .map( + claims -> { + List> references = + createReferencePayloads(batch, claims); + recordBatch(metrics, batch, serializedSizeOf(batch.values()), startNanos); + return references; + }); } return scope.awaitAll(ListUtils::flatten); }); @@ -164,6 +180,13 @@ private static List> createReferencePayloads( CompletableFuture> retrieve( List payloads, CancellationToken cancellationToken) { + return retrieve(payloads, cancellationToken, null); + } + + CompletableFuture> retrieve( + List payloads, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { Map> batches; try { batches = buildRetrieveBatches(payloads); @@ -173,7 +196,7 @@ CompletableFuture> retrieve( if (batches.isEmpty()) { return CompletableFuture.completedFuture(payloads); } - return runRetrieveDrivers(batches, cancellationToken) + return runRetrieveDrivers(batches, cancellationToken, metrics) .thenApply(retrievedPayloads -> applyPayloadReplacements(payloads, retrievedPayloads)); } @@ -200,16 +223,26 @@ private Map> buildRetrieveBatches(List>> runRetrieveDrivers( Map> batches, - CancellationToken cancellationToken) { + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { return withDriverScope( cancellationToken, scope -> { StorageDriverRetrieveContext context = new StorageDriverRetrieveContextImpl(scope.token()); for (Batch batch : batches.values()) { + long startNanos = System.nanoTime(); scope .attach(batch.driver.retrieve(context, batch.values())) - .map(payloads -> mapPayloadsToOriginalPositions(batch, payloads)); + .map( + payloads -> { + List> mapped = + mapPayloadsToOriginalPositions(batch, payloads); + // The reference's recorded size is not part of the exchange contract, so + // measure what the driver actually returned. + recordBatch(metrics, batch, serializedSizeOf(payloads), startNanos); + return mapped; + }); } return scope.awaitAll(ListUtils::flatten); }); @@ -237,6 +270,22 @@ private static List> mapPayloadsToOriginalPositions( return replacements; } + private static long serializedSizeOf(List payloads) { + long total = 0; + for (Payload payload : payloads) { + total += payload.getSerializedSize(); + } + return total; + } + + private static void recordBatch( + @Nullable StorageOperationMetrics metrics, Batch batch, long sizeBytes, long startNanos) { + if (metrics != null) { + metrics.recordBatch( + batch.size(), sizeBytes, startNanos, System.nanoTime(), batch.driver.getName()); + } + } + private static CompletableFuture failedFuture(Throwable t) { CompletableFuture future = new CompletableFuture<>(); future.completeExceptionally(t); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java index 8fbe763cfa..295bb202b0 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java @@ -44,8 +44,17 @@ public void store( @Nullable StorageDriverTargetInfo target, @Nullable MessageVisitor targetVisitor, CancellationToken cancellationToken) { + store(builder, target, targetVisitor, cancellationToken, null); + } + + public void store( + Message.Builder builder, + @Nullable StorageDriverTargetInfo target, + @Nullable MessageVisitor targetVisitor, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { getOrThrowIfCancelled( - storeAsync(builder, target, targetVisitor, cancellationToken), cancellationToken); + storeAsync(builder, target, targetVisitor, cancellationToken, metrics), cancellationToken); } public CompletableFuture storeAsync( @@ -53,17 +62,37 @@ public CompletableFuture storeAsync( @Nullable StorageDriverTargetInfo target, @Nullable MessageVisitor targetVisitor, CancellationToken cancellationToken) { - return PayloadVisitors.visit(builder, storeOptions(target, targetVisitor, cancellationToken)); + return storeAsync(builder, target, targetVisitor, cancellationToken, null); + } + + public CompletableFuture storeAsync( + Message.Builder builder, + @Nullable StorageDriverTargetInfo target, + @Nullable MessageVisitor targetVisitor, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { + return PayloadVisitors.visit( + builder, storeOptions(target, targetVisitor, cancellationToken, metrics)); } public T retrieve( T message, CancellationToken cancellationToken) { - return getOrThrowIfCancelled(retrieveAsync(message, cancellationToken), cancellationToken); + return retrieve(message, cancellationToken, null); + } + + public T retrieve( + T message, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { + return getOrThrowIfCancelled( + retrieveAsync(message, cancellationToken, metrics), cancellationToken); } public CompletableFuture retrieveAsync( - T message, CancellationToken cancellationToken) { - return PayloadVisitors.visit(message, retrieveOptions(cancellationToken)); + T message, + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { + return PayloadVisitors.visit(message, retrieveOptions(cancellationToken, metrics)); } /** @@ -116,10 +145,11 @@ private static T getOrThrowIfCancelled( private PayloadVisitorOptions storeOptions( @Nullable StorageDriverTargetInfo target, @Nullable MessageVisitor targetVisitor, - CancellationToken cancellationToken) { + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { return PayloadVisitorOptions.newBuilder( (visitedTarget, payloads) -> - payloadTransformer.store(payloads, visitedTarget, cancellationToken)) + payloadTransformer.store(payloads, visitedTarget, cancellationToken, metrics)) .setInitialContext(target) .setMessageVisitor(targetVisitor) .setConcurrency(payloadVisitConcurrency) @@ -128,9 +158,11 @@ private PayloadVisitorOptions storeOptions( } private PayloadVisitorOptions retrieveOptions( - CancellationToken cancellationToken) { + CancellationToken cancellationToken, + @Nullable StorageOperationMetrics metrics) { return PayloadVisitorOptions.newBuilder( - (context, payloads) -> payloadTransformer.retrieve(payloads, cancellationToken)) + (context, payloads) -> + payloadTransformer.retrieve(payloads, cancellationToken, metrics)) .setConcurrency(payloadVisitConcurrency) .setSkipSearchAttributes(true) .build(); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/StorageOperationMetrics.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/StorageOperationMetrics.java new file mode 100644 index 0000000000..669ab27f4e --- /dev/null +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/StorageOperationMetrics.java @@ -0,0 +1,69 @@ +package io.temporal.internal.payload.storage; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.Set; +import java.util.TreeSet; + +/** + * Accumulates external storage activity for a single workflow task. Batches may complete + * concurrently on driver threads, so every method is synchronized. + */ +public final class StorageOperationMetrics { + private int payloadCount; + private long totalSizeBytes; + private final Set driverNames = new TreeSet<>(); + private final List spans = new ArrayList<>(); + + /** + * Records a completed batch. A batch reports the start and end of its span rather than a + * duration, leaving {@link #getTotalDuration()} to decide how overlapping spans combine. + */ + public synchronized void recordBatch( + int count, long sizeBytes, long startNanos, long endNanos, String driverName) { + payloadCount += count; + totalSizeBytes += sizeBytes; + driverNames.add(driverName); + spans.add(new long[] {startNanos, endNanos}); + } + + public synchronized int getPayloadCount() { + return payloadCount; + } + + public synchronized long getTotalSizeBytes() { + return totalSizeBytes; + } + + public synchronized List getDriverNames() { + return new ArrayList<>(driverNames); + } + + /** + * Wall-clock time storage was in flight. Batches may run concurrently, so overlapping spans are + * counted once rather than summed. + */ + public synchronized Duration getTotalDuration() { + if (spans.isEmpty()) { + return Duration.ZERO; + } + List sorted = new ArrayList<>(spans); + sorted.sort(Comparator.comparingLong(span -> span[0])); + long total = 0; + long currentStart = sorted.get(0)[0]; + long currentEnd = sorted.get(0)[1]; + for (int i = 1; i < sorted.size(); i++) { + long[] span = sorted.get(i); + if (span[0] > currentEnd) { + total += currentEnd - currentStart; + currentStart = span[0]; + currentEnd = span[1]; + } else if (span[1] > currentEnd) { + currentEnd = span[1]; + } + } + return Duration.ofNanos(total + currentEnd - currentStart); + } +} diff --git a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowTaskHandler.java b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowTaskHandler.java index a4800a7f7f..34da175096 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowTaskHandler.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowTaskHandler.java @@ -24,6 +24,7 @@ import io.temporal.internal.common.ProtobufTimeUtils; import io.temporal.internal.common.WorkflowExecutionUtils; import io.temporal.internal.payload.storage.ExternalStorageRunner; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.internal.worker.*; import io.temporal.payload.context.WorkflowSerializationContext; import io.temporal.serviceclient.MetricsTag; @@ -38,6 +39,7 @@ import java.util.concurrent.CancellationException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; +import javax.annotation.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -74,16 +76,20 @@ public ReplayWorkflowTaskHandler( } @Override - public WorkflowTaskHandler.Result handleWorkflowTask(PollWorkflowTaskQueueResponse workflowTask) + public WorkflowTaskHandler.Result handleWorkflowTask( + PollWorkflowTaskQueueResponse workflowTask, @Nullable StorageOperationMetrics downloadMetrics) throws Exception { String workflowType = workflowTask.getWorkflowType().getName(); Scope metricsScope = options.getMetricsScope().tagged(ImmutableMap.of(MetricsTag.WORKFLOW_TYPE, workflowType)); - return handleWorkflowTaskWithQuery(workflowTask.toBuilder(), metricsScope); + return handleWorkflowTaskWithQuery(workflowTask.toBuilder(), metricsScope, downloadMetrics); } private Result handleWorkflowTaskWithQuery( - PollWorkflowTaskQueueResponse.Builder workflowTask, Scope metricsScope) throws Exception { + PollWorkflowTaskQueueResponse.Builder workflowTask, + Scope metricsScope, + @Nullable StorageOperationMetrics downloadMetrics) + throws Exception { boolean directQuery = workflowTask.hasQuery(); AtomicBoolean createdNew = new AtomicBoolean(); WorkflowExecution execution = workflowTask.getWorkflowExecution(); @@ -91,9 +97,10 @@ private Result handleWorkflowTaskWithQuery( boolean useCache = stickyTaskQueue != null; try { - workflowTask = retrieveStoredPayloads(workflowTask); + workflowTask = retrieveStoredPayloads(workflowTask, downloadMetrics); workflowRunTaskHandler = - getOrCreateWorkflowExecutor(useCache, workflowTask, metricsScope, createdNew); + getOrCreateWorkflowExecutor( + useCache, workflowTask, metricsScope, createdNew, downloadMetrics); logWorkflowTaskToBeProcessed(workflowTask, createdNew); ServiceWorkflowHistoryIterator historyIterator = @@ -103,7 +110,8 @@ private Result handleWorkflowTaskWithQuery( workflowTask, metricsScope, options.getExternalStorageRunner(), - options.getStorageCancellation()); + options.getStorageCancellation(), + downloadMetrics); boolean finalCommand; Result result; @@ -180,14 +188,15 @@ private Result handleWorkflowTaskWithQuery( } private PollWorkflowTaskQueueResponse.Builder retrieveStoredPayloads( - PollWorkflowTaskQueueResponse.Builder workflowTask) { + PollWorkflowTaskQueueResponse.Builder workflowTask, + @Nullable StorageOperationMetrics downloadMetrics) { ExternalStorageRunner externalStorageRunner = options.getExternalStorageRunner(); if (externalStorageRunner == null) { ExternalStorageRunner.throwIfContainsReference(workflowTask.build()); return workflowTask; } return externalStorageRunner - .retrieve(workflowTask.build(), options.getStorageCancellation()) + .retrieve(workflowTask.build(), options.getStorageCancellation(), downloadMetrics) .toBuilder(); } @@ -383,7 +392,8 @@ private WorkflowRunTaskHandler getOrCreateWorkflowExecutor( boolean useCache, PollWorkflowTaskQueueResponse.Builder workflowTask, Scope metricsScope, - AtomicBoolean createdNew) + AtomicBoolean createdNew, + @Nullable StorageOperationMetrics downloadMetrics) throws Exception { if (useCache) { return cache.getOrCreate( @@ -391,17 +401,20 @@ private WorkflowRunTaskHandler getOrCreateWorkflowExecutor( metricsScope, () -> { createdNew.set(true); - return createStatefulHandler(workflowTask, metricsScope); + return createStatefulHandler(workflowTask, metricsScope, downloadMetrics); }); } else { createdNew.set(true); - return createStatefulHandler(workflowTask, metricsScope); + return createStatefulHandler(workflowTask, metricsScope, downloadMetrics); } } // TODO(maxim): Consider refactoring that avoids mutating workflow task. private WorkflowRunTaskHandler createStatefulHandler( - PollWorkflowTaskQueueResponse.Builder workflowTask, Scope metricsScope) throws Exception { + PollWorkflowTaskQueueResponse.Builder workflowTask, + Scope metricsScope, + @Nullable StorageOperationMetrics downloadMetrics) + throws Exception { WorkflowType workflowType = workflowTask.getWorkflowType(); WorkflowExecution workflowExecution = workflowTask.getWorkflowExecution(); List events = workflowTask.getHistory().getEventsList(); @@ -422,7 +435,8 @@ private WorkflowRunTaskHandler createStatefulHandler( ExternalStorageRunner.throwIfContainsReference(getHistoryResponse); } else { getHistoryResponse = - externalStorageRunner.retrieve(getHistoryResponse, options.getStorageCancellation()); + externalStorageRunner.retrieve( + getHistoryResponse, options.getStorageCancellation(), downloadMetrics); } workflowTask .setHistory(getHistoryResponse.getHistory()) diff --git a/temporal-sdk/src/main/java/io/temporal/internal/replay/ServiceWorkflowHistoryIterator.java b/temporal-sdk/src/main/java/io/temporal/internal/replay/ServiceWorkflowHistoryIterator.java index 8c9974a5ef..733724db9a 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/replay/ServiceWorkflowHistoryIterator.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/replay/ServiceWorkflowHistoryIterator.java @@ -14,6 +14,7 @@ import io.temporal.api.workflowservice.v1.PollWorkflowTaskQueueResponseOrBuilder; import io.temporal.common.CancellationToken; import io.temporal.internal.payload.storage.ExternalStorageRunner; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.internal.retryer.GrpcRetryer; import io.temporal.serviceclient.RpcRetryOptions; import io.temporal.serviceclient.WorkflowServiceStubs; @@ -35,6 +36,7 @@ class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator { private final GrpcRetryer grpcRetryer; private final @Nullable ExternalStorageRunner externalStorageRunner; private final CancellationToken storageCancellation; + private final @Nullable StorageOperationMetrics downloadMetrics; private Deadline deadline; private Iterator current; ByteString nextPageToken; @@ -44,7 +46,7 @@ class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator { String namespace, PollWorkflowTaskQueueResponseOrBuilder task, Scope metricsScope) { - this(service, namespace, task, metricsScope, null, CancellationToken.none()); + this(service, namespace, task, metricsScope, null, CancellationToken.none(), null); } ServiceWorkflowHistoryIterator( @@ -53,8 +55,10 @@ class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator { PollWorkflowTaskQueueResponseOrBuilder task, Scope metricsScope, @Nullable ExternalStorageRunner externalStorageRunner, - CancellationToken storageCancellation) { + CancellationToken storageCancellation, + @Nullable StorageOperationMetrics downloadMetrics) { this.storageCancellation = storageCancellation; + this.downloadMetrics = downloadMetrics; this.service = service; this.namespace = namespace; this.task = task; @@ -86,7 +90,7 @@ public boolean hasNext() { if (externalStorageRunner == null) { ExternalStorageRunner.throwIfContainsReference(history); } else { - history = externalStorageRunner.retrieve(history, storageCancellation); + history = externalStorageRunner.retrieve(history, storageCancellation, downloadMetrics); } current = history.getEventsList().iterator(); nextPageToken = response.getNextPageToken(); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/QueryReplayHelper.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/QueryReplayHelper.java index d438101b7e..858cb1c0c4 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/QueryReplayHelper.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/QueryReplayHelper.java @@ -73,7 +73,7 @@ private Optional queryWorkflowExecution( WorkflowType workflowType = started.getWorkflowType(); task.setWorkflowType(workflowType); task.setHistory(History.newBuilder().addAllEvents(events)); - WorkflowTaskHandler.Result result = handler.handleWorkflowTask(task.build()); + WorkflowTaskHandler.Result result = handler.handleWorkflowTask(task.build(), null); if (result.getQueryCompleted() != null) { RespondQueryTaskCompletedRequest r = result.getQueryCompleted(); if (!r.getErrorMessage().isEmpty()) { diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskHandler.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskHandler.java index 3f4a95959c..6b6014d883 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskHandler.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskHandler.java @@ -5,6 +5,7 @@ import io.temporal.api.workflowservice.v1.RespondQueryTaskCompletedRequest; import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest; import io.temporal.api.workflowservice.v1.RespondWorkflowTaskFailedRequest; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.serviceclient.RpcRetryOptions; import io.temporal.workflow.Functions; import javax.annotation.Nullable; @@ -137,12 +138,16 @@ public String getWorkflowType() { * Handles a single workflow task * * @param workflowTask The workflow task to handle. + * @param downloadMetrics Accumulates external storage retrievals performed for this task, or null + * to collect nothing. * @return One of the possible workflow task replies: RespondWorkflowTaskCompletedRequest, * RespondQueryTaskCompletedRequest, RespondWorkflowTaskFailedRequest * @throws Exception an original exception or error if the processing should be just abandoned * without replying to the server */ - Result handleWorkflowTask(PollWorkflowTaskQueueResponse workflowTask) throws Exception; + Result handleWorkflowTask( + PollWorkflowTaskQueueResponse workflowTask, @Nullable StorageOperationMetrics downloadMetrics) + throws Exception; /** True if this handler handles at least one workflow type. */ boolean isAnyTypeSupported(); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowWorker.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowWorker.java index cc9b5ef71b..605dc682d8 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowWorker.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowWorker.java @@ -21,9 +21,11 @@ import io.temporal.api.failure.v1.Failure; import io.temporal.api.workflowservice.v1.*; import io.temporal.failure.ApplicationFailure; +import io.temporal.internal.common.env.EnvironmentVariableUtils; import io.temporal.internal.logging.LoggerTag; import io.temporal.internal.logging.PrefixedMdc; import io.temporal.internal.payload.storage.ExternalStorageRunner; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.internal.payload.visitor.MessageVisitor; import io.temporal.internal.retryer.GrpcMessageTooLargeException; import io.temporal.internal.retryer.GrpcRetryer; @@ -36,6 +38,7 @@ import io.temporal.serviceclient.WorkflowServiceStubs; import io.temporal.worker.*; import io.temporal.worker.tuning.*; +import java.time.Duration; import java.util.*; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; @@ -399,21 +402,88 @@ public String toString() { options.getIdentity(), namespace, taskQueue); } + /** + * Warns when a workflow task took longer than the configured threshold, reporting the external + * storage work that contributed to it. The workflow ID, type and run ID are already on the MDC. + * + *

Storage fields are always present, zero when nothing was offloaded, so the message template + * stays constant for log aggregation. + */ + static void logWorkflowTaskDuration( + PollWorkflowTaskQueueResponse task, + String workflowType, + Duration taskDuration, + Duration warnThreshold, + StorageOperationMetrics downloadMetrics, + StorageOperationMetrics uploadMetrics) { + if (taskDuration.compareTo(warnThreshold) <= 0) { + return; + } + log.warn( + "[TMPRL1104] {}:{}:{} Workflow task duration exceeded {} seconds." + + " WorkflowType={}, WorkflowTaskDuration={}ms," + + " PayloadDownloadCount={}, PayloadDownloadSize={}," + + " PayloadDownloadDuration={}ms, PayloadDownloadDrivers={}," + + " PayloadUploadCount={}, PayloadUploadSize={}," + + " PayloadUploadDuration={}ms, PayloadUploadDrivers={}", + task.getWorkflowExecution().getRunId(), + task.getStartedEventId() + 1, + task.getAttempt(), + warnThreshold.getSeconds(), + workflowType, + taskDuration.toMillis(), + downloadMetrics.getPayloadCount(), + downloadMetrics.getTotalSizeBytes(), + downloadMetrics.getTotalDuration().toMillis(), + downloadMetrics.getDriverNames(), + uploadMetrics.getPayloadCount(), + uploadMetrics.getTotalSizeBytes(), + uploadMetrics.getTotalDuration().toMillis(), + uploadMetrics.getDriverNames()); + } + + private static final Duration DEFAULT_WFT_DURATION_WARN_THRESHOLD = Duration.ofSeconds(5); + + private static final Duration WFT_DURATION_WARN_THRESHOLD = + parseWftDurationWarnThreshold( + EnvironmentVariableUtils.readString("TEMPORAL_WORKFLOW_TASK_DURATION_WARN_SECONDS")); + + /** + * Separated from the environment read so it can be unit-tested without mutating the process + * environment. An unparsable value, including a negative one, falls back to the default rather + * than disabling the warning. + */ + static Duration parseWftDurationWarnThreshold(@Nullable String value) { + if (value != null) { + try { + long seconds = Long.parseLong(value.trim()); + if (seconds >= 0) { + return Duration.ofSeconds(seconds); + } + } catch (NumberFormatException e) { + // Fall through to the default. + } + } + return DEFAULT_WFT_DURATION_WARN_THRESHOLD; + } + private void storeOutboundPayloads( com.google.protobuf.Message.Builder builder, @Nullable StorageDriverTargetInfo target) { - storeOutboundPayloads(builder, target, null); + storeOutboundPayloads(builder, target, null, null); } private void storeOutboundPayloads( com.google.protobuf.Message.Builder builder, @Nullable StorageDriverTargetInfo target, - @Nullable MessageVisitor targetVisitor) { + @Nullable MessageVisitor targetVisitor, + @Nullable StorageOperationMetrics uploadMetrics) { ExternalStorageRunner externalStorageRunner = options.getExternalStorageRunner(); if (externalStorageRunner == null) { return; } try { - externalStorageRunner.store(builder, target, targetVisitor, options.getStorageCancellation()); + externalStorageRunner.store( + builder, target, targetVisitor, options.getStorageCancellation(), uploadMetrics); } catch (CancellationException e) { // if the worker is shutting down, extstore will throw a CancellationException and we need to // rethrow it here so the handle() method can decide what to do. @@ -573,8 +643,12 @@ public void handle(WorkflowTask task) throws Exception { PollWorkflowTaskQueueResponse currentTask = nextWFTResponse.get(); nextWFTResponse = Optional.empty(); boolean iterationFailed = false; + StorageOperationMetrics downloadMetrics = new StorageOperationMetrics(); + StorageOperationMetrics uploadMetrics = new StorageOperationMetrics(); + long iterationStartNanos = System.nanoTime(); try { - WorkflowTaskHandler.Result result = handleTask(currentTask, workflowTypeScope); + WorkflowTaskHandler.Result result = + handleTask(currentTask, workflowTypeScope, downloadMetrics); WorkflowTaskFailedCause taskFailedCause = null; try { RespondWorkflowTaskCompletedRequest taskCompleted = result.getTaskCompleted(); @@ -640,6 +714,7 @@ public void handle(WorkflowTask task) throws Exception { RespondWorkflowTaskCompletedRequest request = prepareTaskCompleted( currentTask.getTaskToken(), + uploadMetrics, requestBuilder, workflowStorageTarget(workflowExecution, workflowType), parentStorageTarget(result.getCompletionParentExecution())); @@ -806,6 +881,15 @@ public void handle(WorkflowTask task) throws Exception { if (iterationFailed) { taskCounter.recordFailed(); } + if (!options.getStorageCancellation().isCancellationRequested()) { + logWorkflowTaskDuration( + currentTask, + workflowType, + Duration.ofNanos(System.nanoTime() - iterationStartNanos), + WFT_DURATION_WARN_THRESHOLD, + downloadMetrics, + uploadMetrics); + } } } while (nextWFTResponse.isPresent()); } finally { @@ -837,11 +921,14 @@ public Throwable wrapFailure(WorkflowTask task, Throwable failure) { } private WorkflowTaskHandler.Result handleTask( - PollWorkflowTaskQueueResponse task, Scope workflowTypeMetricsScope) throws Exception { + PollWorkflowTaskQueueResponse task, + Scope workflowTypeMetricsScope, + StorageOperationMetrics downloadMetrics) + throws Exception { Stopwatch sw = workflowTypeMetricsScope.timer(MetricsType.WORKFLOW_TASK_EXECUTION_LATENCY).start(); try { - return handler.handleWorkflowTask(task); + return handler.handleWorkflowTask(task, downloadMetrics); } catch (Throwable e) { workflowTypeMetricsScope.counter(MetricsType.WORKFLOW_TASK_NO_COMPLETION_COUNTER).inc(1); // Make sure that the task failure metric has the correct type @@ -869,6 +956,7 @@ private WorkflowTaskHandler.Result handleTask( @SuppressWarnings("deprecation") private RespondWorkflowTaskCompletedRequest prepareTaskCompleted( ByteString taskToken, + StorageOperationMetrics uploadMetrics, RespondWorkflowTaskCompletedRequest.Builder taskCompleted, @Nullable StorageDriverTargetInfo storageTarget, @Nullable StorageDriverTargetInfo completionTarget) { @@ -892,7 +980,7 @@ private RespondWorkflowTaskCompletedRequest prepareTaskCompleted( MessageVisitor storageTargetVisitor = (current, message) -> deriveStorageTarget(namespace, current, message, completionTarget); - storeOutboundPayloads(taskCompleted, storageTarget, storageTargetVisitor); + storeOutboundPayloads(taskCompleted, storageTarget, storageTargetVisitor, uploadMetrics); return taskCompleted.build(); } diff --git a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/StorageOperationMetricsTest.java b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/StorageOperationMetricsTest.java new file mode 100644 index 0000000000..cdd2dd1f4d --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/StorageOperationMetricsTest.java @@ -0,0 +1,60 @@ +package io.temporal.internal.payload.storage; + +import static org.junit.Assert.assertEquals; + +import java.time.Duration; +import java.util.Arrays; +import java.util.concurrent.TimeUnit; +import org.junit.Test; + +public class StorageOperationMetricsTest { + + private static long nanos(long millis) { + return TimeUnit.MILLISECONDS.toNanos(millis); + } + + @Test + public void noBatchesReportsZero() { + StorageOperationMetrics metrics = new StorageOperationMetrics(); + assertEquals(Duration.ZERO, metrics.getTotalDuration()); + assertEquals(0, metrics.getPayloadCount()); + assertEquals(0, metrics.getTotalSizeBytes()); + } + + @Test + public void aggregatesCountSizeAndDriverNames() { + StorageOperationMetrics metrics = new StorageOperationMetrics(); + metrics.recordBatch(2, 1024, nanos(0), nanos(100), "s3"); + metrics.recordBatch(3, 2048, nanos(0), nanos(100), "gcs"); + metrics.recordBatch(1, 512, nanos(0), nanos(100), "s3"); + assertEquals(6, metrics.getPayloadCount()); + assertEquals(3584, metrics.getTotalSizeBytes()); + assertEquals(Arrays.asList("gcs", "s3"), metrics.getDriverNames()); + } + + @Test + public void concurrentBatchesCountedOnce() { + StorageOperationMetrics metrics = new StorageOperationMetrics(); + // Summing each batch would report 200ms; the real wall-clock span is 150ms. + metrics.recordBatch(1, 1, nanos(0), nanos(100), "s3"); + metrics.recordBatch(1, 1, nanos(50), nanos(150), "gcs"); + assertEquals(Duration.ofMillis(150), metrics.getTotalDuration()); + } + + @Test + public void disjointBatchesSummed() { + StorageOperationMetrics metrics = new StorageOperationMetrics(); + metrics.recordBatch(1, 1, nanos(0), nanos(100), "s3"); + metrics.recordBatch(1, 1, nanos(200), nanos(300), "s3"); + assertEquals(Duration.ofMillis(200), metrics.getTotalDuration()); + } + + @Test + public void adjacentAndNestedBatchesMerged() { + StorageOperationMetrics metrics = new StorageOperationMetrics(); + metrics.recordBatch(1, 1, nanos(0), nanos(100), "s3"); + metrics.recordBatch(1, 1, nanos(100), nanos(200), "s3"); + metrics.recordBatch(1, 1, nanos(120), nanos(180), "s3"); + assertEquals(Duration.ofMillis(200), metrics.getTotalDuration()); + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/internal/replay/GetVersionInterleavedUpdateReplayTaskHandlerTest.java b/temporal-sdk/src/test/java/io/temporal/internal/replay/GetVersionInterleavedUpdateReplayTaskHandlerTest.java index 8b5084e925..a2fca413d4 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/replay/GetVersionInterleavedUpdateReplayTaskHandlerTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/replay/GetVersionInterleavedUpdateReplayTaskHandlerTest.java @@ -9,6 +9,7 @@ import io.temporal.api.workflowservice.v1.PollWorkflowTaskQueueResponse; import io.temporal.client.WorkflowClient; import io.temporal.common.WorkflowExecutionHistory; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.internal.worker.QueryReplayHelper; import io.temporal.serviceclient.WorkflowServiceStubs; import io.temporal.testing.CloudTestExclusion.RequiresLocalServer; @@ -98,10 +99,11 @@ private static ReplayWorkflowRunTaskHandler createStatefulHandler( ReplayWorkflowTaskHandler.class.getDeclaredMethod( "createStatefulHandler", PollWorkflowTaskQueueResponse.Builder.class, - com.uber.m3.tally.Scope.class); + com.uber.m3.tally.Scope.class, + StorageOperationMetrics.class); method.setAccessible(true); return (ReplayWorkflowRunTaskHandler) - method.invoke(replayTaskHandler, replayTask, new NoopScope()); + method.invoke(replayTaskHandler, replayTask, new NoopScope(), null); } private static T getField(Object target, String fieldName, Class expectedType) diff --git a/temporal-sdk/src/test/java/io/temporal/internal/replay/ReplayWorkflowRunTaskHandlerTaskHandlerTests.java b/temporal-sdk/src/test/java/io/temporal/internal/replay/ReplayWorkflowRunTaskHandlerTaskHandlerTests.java index 9de63d27e3..f7964dec3a 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/replay/ReplayWorkflowRunTaskHandlerTaskHandlerTests.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/replay/ReplayWorkflowRunTaskHandlerTaskHandlerTests.java @@ -152,7 +152,7 @@ public void ifStickyExecutionAttributesAreNotSetThenWorkflowsAreNotCached() thro // Act WorkflowTaskHandler.Result result = - taskHandler.handleWorkflowTask(HistoryUtils.generateWorkflowTaskWithInitialHistory()); + taskHandler.handleWorkflowTask(HistoryUtils.generateWorkflowTaskWithInitialHistory(), null); // Assert assertEquals(0, cache.size()); assertNotNull(result.getTaskCompleted()); @@ -206,7 +206,8 @@ public void workflowTaskFailOnIncompleteHistory() throws Throwable { HistoryUtils.generateWorkflowTaskWithInitialHistory().toBuilder() .setHistory(History.newBuilder().build()) .setNextPageToken(ByteString.EMPTY) - .build()); + .build(), + null); // Assert assertEquals(0, cache.size()); @@ -264,7 +265,7 @@ public void resolvesExternalStorageReferencesInTheWorkflowTaskItself() throws Th client, null); - taskHandler.handleWorkflowTask(fullTask.toBuilder().setHistory(storedHistory).build()); + taskHandler.handleWorkflowTask(fullTask.toBuilder().setHistory(storedHistory).build(), null); ArgumentCaptor event = ArgumentCaptor.forClass(HistoryEvent.class); verify(workflow).start(event.capture(), any()); @@ -327,7 +328,8 @@ public void aCancelledDownloadIsNotReportedAsAWorkflowTaskFailure() throws Throw "stopping storage must not be turned into a workflow task failure", CancellationException.class, () -> - taskHandler.handleWorkflowTask(fullTask.toBuilder().setHistory(storedHistory).build())); + taskHandler.handleWorkflowTask( + fullTask.toBuilder().setHistory(storedHistory).build(), null)); } @Test @@ -366,7 +368,8 @@ public void aFailedDownloadIsReportedAsAWorkflowTaskFailure() throws Throwable { null); WorkflowTaskHandler.Result result = - taskHandler.handleWorkflowTask(fullTask.toBuilder().setHistory(storedHistory).build()); + taskHandler.handleWorkflowTask( + fullTask.toBuilder().setHistory(storedHistory).build(), null); assertNotNull( "a failed download must be reported rather than ending the task", result.getTaskFailed()); @@ -433,7 +436,7 @@ public void resolvesExternalStorageReferencesInFetchedFullHistory() throws Throw null); taskHandler.handleWorkflowTask( - fullTask.toBuilder().setHistory(History.getDefaultInstance()).build()); + fullTask.toBuilder().setHistory(History.getDefaultInstance()).build(), null); ArgumentCaptor event = ArgumentCaptor.forClass(HistoryEvent.class); verify(workflow).start(event.capture(), any()); @@ -502,7 +505,7 @@ public void ifStickyExecutionAttributesAreSetThenWorkflowsAreCached() throws Thr PollWorkflowTaskQueueResponse workflowTask = HistoryUtils.generateWorkflowTaskWithInitialHistory(); - WorkflowTaskHandler.Result result = taskHandler.handleWorkflowTask(workflowTask); + WorkflowTaskHandler.Result result = taskHandler.handleWorkflowTask(workflowTask, null); assertTrue(result.isCompletionCommand()); assertEquals(0, cache.size()); // do not cache if completion command @@ -534,7 +537,7 @@ public void setsSdkNameAndVersionIfNotSetInHistory() throws Throwable { PollWorkflowTaskQueueResponse workflowTask = HistoryUtils.generateWorkflowTaskWithInitialHistory(); - WorkflowTaskHandler.Result result = taskHandler.handleWorkflowTask(workflowTask); + WorkflowTaskHandler.Result result = taskHandler.handleWorkflowTask(workflowTask, null); assertTrue(result.isCompletionCommand()); assertEquals(Version.SDK_NAME, result.getTaskCompleted().getSdkMetadata().getSdkName()); diff --git a/temporal-sdk/src/test/java/io/temporal/internal/replay/ServiceWorkflowHistoryIteratorTest.java b/temporal-sdk/src/test/java/io/temporal/internal/replay/ServiceWorkflowHistoryIteratorTest.java index 672b0a7515..51049f7408 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/replay/ServiceWorkflowHistoryIteratorTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/replay/ServiceWorkflowHistoryIteratorTest.java @@ -153,7 +153,7 @@ private static ServiceWorkflowHistoryIterator fetchingIterator( PollWorkflowTaskQueueResponse workflowTask = PollWorkflowTaskQueueResponse.newBuilder().setNextPageToken(NEXT_PAGE_TOKEN).build(); return new ServiceWorkflowHistoryIterator( - null, "default", workflowTask, null, storage, storageCancellation) { + null, "default", workflowTask, null, storage, storageCancellation, null) { @Override GetWorkflowExecutionHistoryResponse queryWorkflowExecutionHistory() { return GetWorkflowExecutionHistoryResponse.newBuilder().setHistory(page).build(); diff --git a/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowWorkerTest.java b/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowWorkerTest.java index d85d913494..c2dea56846 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowWorkerTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowWorkerTest.java @@ -43,6 +43,7 @@ import io.temporal.internal.common.InternalUtils; import io.temporal.internal.concurrent.structured.CancelSource; import io.temporal.internal.payload.storage.ExternalStorageRunner; +import io.temporal.internal.payload.storage.StorageOperationMetrics; import io.temporal.internal.payload.storage.TestStorageDriver; import io.temporal.internal.replay.ReplayWorkflow; import io.temporal.internal.replay.ReplayWorkflowFactory; @@ -65,12 +66,16 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import javax.annotation.Nullable; +import junitparams.JUnitParamsRunner; +import junitparams.Parameters; import org.junit.Test; +import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.stubbing.Answer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +@RunWith(JUnitParamsRunner.class) public class WorkflowWorkerTest { private static final Logger log = LoggerFactory.getLogger(WorkflowWorkerTest.class); private final TestStatsReporter reporter = new TestStatsReporter(); @@ -168,7 +173,7 @@ public void concurrentPollRequestLockTest() throws Exception { }); CountDownLatch handleTaskLatch = new CountDownLatch(1); - when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class))) + when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class), any())) .thenAnswer( (Answer) invocation -> { @@ -247,7 +252,7 @@ public void concurrentPollRequestLockTest() throws Exception { // Cleanup worker.shutdown(new ShutdownManager(), false).get(); // Verify we only handled two tasks - verify(taskHandler, times(2)).handleWorkflowTask(any()); + verify(taskHandler, times(2)).handleWorkflowTask(any(), any()); } @Test @@ -330,7 +335,7 @@ public void respondWorkflowTaskFailureMetricTest() throws Exception { CountDownLatch handleTaskLatch = new CountDownLatch(1); - when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class))) + when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class), any())) .thenAnswer( (Answer) invocation -> { @@ -392,9 +397,10 @@ public void resetWorkflowIdFromWorkflowTaskTest() throws Throwable { WorkflowTaskHandler taskHandler = new WorkflowTaskHandler() { @Override - public WorkflowTaskHandler.Result handleWorkflowTask(PollWorkflowTaskQueueResponse task) + public WorkflowTaskHandler.Result handleWorkflowTask( + PollWorkflowTaskQueueResponse task, StorageOperationMetrics downloadMetrics) throws Exception { - WorkflowTaskHandler.Result result = rootTaskHandler.handleWorkflowTask(task); + WorkflowTaskHandler.Result result = rootTaskHandler.handleWorkflowTask(task, null); return new WorkflowTaskHandler.Result( result.getWorkflowType(), result.getTaskCompleted(), @@ -558,7 +564,7 @@ public void aTaskAbandonedWhileShuttingDownIsNotReported() throws Exception { }); // The task is abandoned part way through, which is what stopping storage looks like. - when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class))) + when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class), any())) .thenAnswer( (Answer) invocation -> { @@ -851,7 +857,7 @@ private void runOneTask( }); CountDownLatch handled = new CountDownLatch(1); - when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class))) + when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class), any())) .thenAnswer( (Answer) invocation -> { @@ -962,7 +968,7 @@ public void aStoreThatFailsWhileShuttingDownIsNotTreatedAsAProblem() throws Exce return null; }); - when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class))) + when(taskHandler.handleWorkflowTask(any(PollWorkflowTaskQueueResponse.class), any())) .thenAnswer( (Answer) invocation -> @@ -1137,4 +1143,104 @@ public void deriveStorageTargetKeepsTheCurrentTargetForOtherCommands() { assertSame(current, WorkflowWorker.deriveStorageTarget("ns", current, command)); } + + private static ListAppender captureWorkerLogs() { + LoggerContext loggerContext = (LoggerContext) LoggerFactory.getILoggerFactory(); + ListAppender logs = new ListAppender<>(); + logs.setContext(loggerContext); + logs.start(); + loggerContext.getLogger(WorkflowWorker.class.getName()).addAppender(logs); + return logs; + } + + private static void stopCapturingWorkerLogs(ListAppender logs) { + LoggerContext loggerContext = (LoggerContext) LoggerFactory.getILoggerFactory(); + loggerContext.getLogger(WorkflowWorker.class.getName()).detachAppender(logs); + logs.stop(); + } + + private static PollWorkflowTaskQueueResponse durationLogTask() { + return PollWorkflowTaskQueueResponse.newBuilder() + .setWorkflowExecution(WorkflowExecution.newBuilder().setRunId("run-1").build()) + .setStartedEventId(11) + .setAttempt(3) + .build(); + } + + @Test + public void aTaskWithinTheDurationThresholdIsNotWarnedAbout() { + ListAppender logs = captureWorkerLogs(); + try { + WorkflowWorker.logWorkflowTaskDuration( + durationLogTask(), + "MyWorkflow", + Duration.ofSeconds(5), + Duration.ofSeconds(5), + new StorageOperationMetrics(), + new StorageOperationMetrics()); + assertTrue(logs.list.isEmpty()); + } finally { + stopCapturingWorkerLogs(logs); + } + } + + @Test + public void aTaskOverTheDurationThresholdWarnsWithItsStorageMetrics() { + StorageOperationMetrics download = new StorageOperationMetrics(); + download.recordBatch(2, 1024, 0, TimeUnit.MILLISECONDS.toNanos(40), "s3"); + download.recordBatch(1, 512, 0, TimeUnit.MILLISECONDS.toNanos(40), "gcs"); + StorageOperationMetrics upload = new StorageOperationMetrics(); + upload.recordBatch(3, 2048, 0, TimeUnit.MILLISECONDS.toNanos(70), "s3"); + + ListAppender logs = captureWorkerLogs(); + try { + WorkflowWorker.logWorkflowTaskDuration( + durationLogTask(), + "MyWorkflow", + Duration.ofMillis(6500), + Duration.ofSeconds(5), + download, + upload); + + assertEquals(1, logs.list.size()); + ILoggingEvent event = logs.list.get(0); + assertEquals(Level.WARN, event.getLevel()); + String message = event.getFormattedMessage(); + // run id, the completion event id (startedEventId + 1) and the attempt. + assertTrue(message, message.contains("[TMPRL1104] run-1:12:3")); + assertTrue(message, message.contains("exceeded 5 seconds")); + assertTrue(message, message.contains("WorkflowType=MyWorkflow")); + assertTrue(message, message.contains("WorkflowTaskDuration=6500ms")); + assertTrue(message, message.contains("PayloadDownloadCount=3")); + assertTrue(message, message.contains("PayloadDownloadSize=1536")); + // Both download batches ran over the same 40ms, so the union is 40ms rather than 80ms. + assertTrue(message, message.contains("PayloadDownloadDuration=40ms")); + assertTrue(message, message.contains("PayloadDownloadDrivers=[gcs, s3]")); + assertTrue(message, message.contains("PayloadUploadCount=3")); + assertTrue(message, message.contains("PayloadUploadDuration=70ms")); + assertTrue(message, message.contains("PayloadUploadDrivers=[s3]")); + } finally { + stopCapturingWorkerLogs(logs); + } + } + + /** Cases are supplied as objects rather than CSV so null, blank and padded values survive. */ + private Object[] thresholdCases() { + return new Object[][] { + {null, Duration.ofSeconds(5)}, + {"10", Duration.ofSeconds(10)}, + {" 30 ", Duration.ofSeconds(30)}, + {"0", Duration.ZERO}, + {"-5", Duration.ofSeconds(5)}, + {"abc", Duration.ofSeconds(5)}, + {"1.5", Duration.ofSeconds(5)}, + {"", Duration.ofSeconds(5)}, + }; + } + + @Test + @Parameters(method = "thresholdCases") + public void parsesThresholdOrFallsBackToTheDefault(String value, Duration expected) { + assertEquals(expected, WorkflowWorker.parseWftDurationWarnThreshold(value)); + } }