From 0213af656ce51bebc35e37a408d12d8507aaf2db Mon Sep 17 00:00:00 2001 From: jmaeagle99 <44687433+jmaeagle99@users.noreply.github.com> Date: Fri, 28 Aug 2026 15:26:52 -0700 Subject: [PATCH] Add workflow task completion pagination --- .../worker/NamespaceCapabilities.java | 31 ++++ .../WorkflowTaskCompletionPaginator.java | 113 ++++++++++++ .../internal/worker/WorkflowWorker.java | 173 ++++++++++++++++-- .../io/temporal/worker/WorkerFactory.java | 1 + .../WorkflowTaskCompletionPaginatorTest.java | 120 ++++++++++++ ...skCompletionPaginationIntegrationTest.java | 107 +++++++++++ .../io/temporal/serviceclient/MetricsTag.java | 1 + 7 files changed, 534 insertions(+), 12 deletions(-) create mode 100644 temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginator.java create mode 100644 temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginatorTest.java create mode 100644 temporal-sdk/src/test/java/io/temporal/worker/WorkflowTaskCompletionPaginationIntegrationTest.java diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/NamespaceCapabilities.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/NamespaceCapabilities.java index 4bddd45d9e..e9866089dc 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/NamespaceCapabilities.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/NamespaceCapabilities.java @@ -1,7 +1,9 @@ package io.temporal.internal.worker; import io.temporal.api.namespace.v1.NamespaceInfo.Capabilities; +import io.temporal.api.namespace.v1.NamespaceInfo.Limits; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; /** * Holds namespace-level capabilities discovered from the server's DescribeNamespace response. A @@ -14,6 +16,8 @@ public final class NamespaceCapabilities { private final AtomicBoolean gracefulPollShutdown = new AtomicBoolean(false); private final AtomicBoolean workerHeartbeats = new AtomicBoolean(false); private final AtomicBoolean workerCommands = new AtomicBoolean(false); + private final AtomicBoolean workflowTaskCompletionPagination = new AtomicBoolean(false); + private final AtomicLong workflowTaskCompletionSizeLimit = new AtomicLong(0); public void setFromCapabilities(Capabilities capabilities) { if (capabilities.getPollerAutoscalingAutoEnroll()) { @@ -31,6 +35,13 @@ public void setFromCapabilities(Capabilities capabilities) { if (capabilities.getWorkerCommands()) { workerCommands.set(true); } + if (capabilities.getWorkflowTaskCompletionPagination()) { + workflowTaskCompletionPagination.set(true); + } + } + + public void setFromLimits(Limits limits) { + workflowTaskCompletionSizeLimit.set(limits.getWorkflowTaskCompletionSizeLimitError()); } public boolean isPollerAutoscaling() { @@ -64,4 +75,24 @@ public boolean isWorkerCommands() { public void setWorkerCommands(boolean value) { workerCommands.set(value); } + + public boolean isWorkflowTaskCompletionPagination() { + return workflowTaskCompletionPagination.get(); + } + + public void setWorkflowTaskCompletionPagination(boolean value) { + workflowTaskCompletionPagination.set(value); + } + + /** + * The namespace's limit on the recombined size in bytes of a single workflow task completion, or + * 0 when the namespace advertises no explicit limit. + */ + public long getWorkflowTaskCompletionSizeLimit() { + return workflowTaskCompletionSizeLimit.get(); + } + + public void setWorkflowTaskCompletionSizeLimit(long value) { + workflowTaskCompletionSizeLimit.set(value); + } } diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginator.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginator.java new file mode 100644 index 0000000000..d72fef25d1 --- /dev/null +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginator.java @@ -0,0 +1,113 @@ +package io.temporal.internal.worker; + +import io.temporal.api.command.v1.Command; +import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest; +import java.util.ArrayList; +import java.util.List; + +/** + * Splits an oversized {@link RespondWorkflowTaskCompletedRequest} into pages that each stay under + * the gRPC request size limit, so a completion carrying more command bytes than a single request + * can hold is delivered across multiple requests sharing one task token. The server buffers the + * commands of the intermediate pages and merges them with the final page when it arrives. + */ +final class WorkflowTaskCompletionPaginator { + + /** + * Maximum encoded size of a single completion page, kept below the ~4 MiB gRPC frame limit. This + * per-page cap is distinct from the namespace's limit on the recombined completion size. + * + *

Pages are packed by summing command body sizes only; the 512 KiB of headroom below 4 MiB + * absorbs everything that sum omits: the per-request overhead (task token, identity, namespace) + * and the per-command wire framing (a field tag plus a length varint, up to 6 bytes each). At the + * server's default per-workflow history-count limit (~51,200 events), worst-case framing is ~300 + * KiB, so this headroom covers even a page of many tiny commands and lets us skip per-command + * accounting. + */ + static final int MAX_PAGE_BYTES = 4 * 1024 * 1024 - 512 * 1024; + + /** The result of splitting a completion: zero or more intermediate pages plus the final page. */ + static final class Pages { + final List intermediatePages; + final RespondWorkflowTaskCompletedRequest finalPage; + + Pages( + List intermediatePages, + RespondWorkflowTaskCompletedRequest finalPage) { + this.intermediatePages = intermediatePages; + this.finalPage = finalPage; + } + + /** True when the completion was split; false when the final page should be sent as-is. */ + boolean isPaginated() { + return !intermediatePages.isEmpty(); + } + } + + /** + * Splits {@code request} into intermediate pages that each stay under {@code maxPageBytes} by + * distributing its commands across them in order. The final page carries the remaining metadata + * and messages, and its page number is the count of intermediate pages. + * + *

Returns a {@link Pages} with no intermediate pages (send {@code request} as-is) when the + * request already fits, has no commands to distribute, or has a single command that alone exceeds + * a page (which the server then rejects). + */ + static Pages paginate(RespondWorkflowTaskCompletedRequest request, int maxPageBytes) { + if (request.getSerializedSize() <= maxPageBytes) { + return new Pages(new ArrayList<>(), request); + } + + List commands = request.getCommandsList(); + // Only commands can be split across pages, so pagination cannot help when there are none, or + // when + // a single command alone exceeds a page. + if (commands.isEmpty()) { + return new Pages(new ArrayList<>(), request); + } + for (Command command : commands) { + if (command.getSerializedSize() > maxPageBytes) { + return new Pages(new ArrayList<>(), request); + } + } + + List intermediatePages = new ArrayList<>(); + List current = new ArrayList<>(); + int currentLen = 0; + for (Command command : commands) { + int commandLen = command.getSerializedSize(); + if (!current.isEmpty() && currentLen + commandLen > maxPageBytes) { + intermediatePages.add(newIntermediatePage(request, current, intermediatePages.size())); + current = new ArrayList<>(); + currentLen = 0; + } + currentLen += commandLen; + current.add(command); + } + if (!current.isEmpty()) { + intermediatePages.add(newIntermediatePage(request, current, intermediatePages.size())); + } + + RespondWorkflowTaskCompletedRequest finalPage = + request.toBuilder() + .clearCommands() + .setPageNumber(intermediatePages.size()) + .setIntermediatePage(false) + .build(); + return new Pages(intermediatePages, finalPage); + } + + private static RespondWorkflowTaskCompletedRequest newIntermediatePage( + RespondWorkflowTaskCompletedRequest request, List commands, int pageNumber) { + return RespondWorkflowTaskCompletedRequest.newBuilder() + .setTaskToken(request.getTaskToken()) + .setIdentity(request.getIdentity()) + .setNamespace(request.getNamespace()) + .setIntermediatePage(true) + .setPageNumber(pageNumber) + .addAllCommands(commands) + .build(); + } + + private WorkflowTaskCompletionPaginator() {} +} 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 98660034d4..1da77d85eb 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 @@ -9,11 +9,14 @@ import com.uber.m3.tally.Scope; import com.uber.m3.tally.Stopwatch; import com.uber.m3.util.ImmutableMap; +import io.grpc.Status; import io.grpc.StatusRuntimeException; +import io.temporal.api.command.v1.Command; import io.temporal.api.common.v1.WorkflowExecution; import io.temporal.api.enums.v1.QueryResultType; import io.temporal.api.enums.v1.TaskQueueKind; import io.temporal.api.enums.v1.WorkflowTaskFailedCause; +import io.temporal.api.errordetails.v1.WorkflowTaskCompletionBufferLostFailure; import io.temporal.api.failure.v1.Failure; import io.temporal.api.workflowservice.v1.*; import io.temporal.failure.ApplicationFailure; @@ -23,6 +26,7 @@ import io.temporal.payload.context.WorkflowSerializationContext; import io.temporal.serviceclient.MetricsTag; import io.temporal.serviceclient.RpcRetryOptions; +import io.temporal.serviceclient.StatusUtils; import io.temporal.serviceclient.WorkflowServiceStubs; import io.temporal.worker.*; import io.temporal.worker.tuning.*; @@ -39,6 +43,13 @@ final class WorkflowWorker implements SuspendableWorker { private static final Logger log = LoggerFactory.getLogger(WorkflowWorker.class); + // Backoff between resends of a paginated completion after the server reports its buffered pages + // were lost. Buffer loss is transient; the loop is bounded by the server eventually timing the + // task + // out (after which the stale token fails with a different error) or by worker shutdown. + private static final long WFT_COMPLETION_PAGE_RESEND_INITIAL_BACKOFF_MS = 100; + private static final long WFT_COMPLETION_PAGE_RESEND_MAX_BACKOFF_MS = 5000; + private final WorkflowRunLockManager runLocks; private final WorkflowServiceStubs service; @@ -477,7 +488,28 @@ public void handle(WorkflowTask task) throws Exception { } } else { try { - if (taskCompleted != null) { + WorkflowTaskFailedCause requestTooLargeCause = + taskCompleted == null + ? null + : completionExceedingSizeLimitCause(taskCompleted); + if (requestTooLargeCause != null) { + // A completion whose recombined command bytes exceed the namespace limit would + // be + // rejected and the workflow terminated by the server, so fail it proactively + // rather than sending doomed pages. + taskFailedCause = requestTooLargeCause; + RespondWorkflowTaskFailedRequest.Builder taskFailedBuilder = + RespondWorkflowTaskFailedRequest.newBuilder() + .setFailure( + requestTooLargeFailure( + workflowExecution.getWorkflowId(), taskCompleted)) + .setCause(requestTooLargeCause); + sendTaskFailed( + currentTask.getTaskToken(), + taskFailedBuilder, + result.getRequestRetryOptions(), + workflowTypeScope); + } else if (taskCompleted != null) { RespondWorkflowTaskCompletedRequest.Builder requestBuilder = taskCompleted.toBuilder(); try (EagerActivitySlotsReservation activitySlotsReservation = @@ -566,6 +598,9 @@ public void handle(WorkflowTask task) throws Exception { case WORKFLOW_TASK_FAILED_CAUSE_GRPC_MESSAGE_TOO_LARGE: taskFailureType = MetricsTag.TASK_FAILURE_VALUE_GRPC_MESSAGE_TOO_LARGE; break; + case WORKFLOW_TASK_FAILED_CAUSE_REQUEST_TOO_LARGE: + taskFailureType = MetricsTag.TASK_FAILURE_VALUE_REQUEST_TOO_LARGE; + break; default: taskFailureType = MetricsTag.TASK_FAILURE_VALUE_WORKFLOW_ERROR; } @@ -654,10 +689,6 @@ private RespondWorkflowTaskCompletedResponse sendTaskCompleted( RespondWorkflowTaskCompletedRequest.Builder taskCompleted, RpcRetryOptions retryOptions, Scope workflowTypeMetricsScope) { - GrpcRetryer.GrpcRetryerOptions grpcRetryOptions = - new GrpcRetryer.GrpcRetryerOptions( - RpcRetryOptions.newBuilder().buildWithDefaultsFrom(retryOptions), null); - taskCompleted .setIdentity(options.getIdentity()) .setNamespace(namespace) @@ -676,13 +707,81 @@ private RespondWorkflowTaskCompletedResponse sendTaskCompleted( taskCompleted.setBinaryChecksum(options.getBuildId()); } - return grpcRetryer.retryWithResult( - () -> - service - .blockingStub() - .withOption(METRICS_TAGS_CALL_OPTIONS_KEY, workflowTypeMetricsScope) - .respondWorkflowTaskCompleted(taskCompleted.build()), - grpcRetryOptions); + RespondWorkflowTaskCompletedRequest request = taskCompleted.build(); + GrpcRetryer.GrpcRetryerOptions grpcRetryOptions = + new GrpcRetryer.GrpcRetryerOptions( + RpcRetryOptions.newBuilder().buildWithDefaultsFrom(retryOptions), null); + + if (!namespaceCapabilities.isWorkflowTaskCompletionPagination()) { + return grpcRetryer.retryWithResult( + () -> respondWorkflowTaskCompleted(request, workflowTypeMetricsScope), + grpcRetryOptions); + } + + WorkflowTaskCompletionPaginator.Pages pages = + WorkflowTaskCompletionPaginator.paginate( + request, WorkflowTaskCompletionPaginator.MAX_PAGE_BYTES); + if (!pages.isPaginated()) { + return grpcRetryer.retryWithResult( + () -> respondWorkflowTaskCompleted(pages.finalPage, workflowTypeMetricsScope), + grpcRetryOptions); + } + return sendPaginatedTaskCompleted(pages, retryOptions, workflowTypeMetricsScope); + } + + /** + * Sends a paginated completion, resending every page from page 0 on buffer loss. Buffer loss — + * the server dropping the pages it had buffered for this token — is transient, so this backs + * off and retries. The gRPC retry layer does not retry buffer loss (it is excluded via a + * DoNotRetryItem below), so this loop is its sole handler; it bails on worker shutdown, and the + * server bounds it by eventually timing the task out. + */ + private RespondWorkflowTaskCompletedResponse sendPaginatedTaskCompleted( + WorkflowTaskCompletionPaginator.Pages pages, + RpcRetryOptions retryOptions, + Scope workflowTypeMetricsScope) { + // Buffer loss requires resending every page, which a single-page gRPC retry cannot do, so it + // is handled by this loop instead of the retryer. + GrpcRetryer.GrpcRetryerOptions pageRetryOptions = + new GrpcRetryer.GrpcRetryerOptions( + RpcRetryOptions.newBuilder( + RpcRetryOptions.newBuilder().buildWithDefaultsFrom(retryOptions)) + .addDoNotRetry(Status.Code.ABORTED, WorkflowTaskCompletionBufferLostFailure.class) + .validateBuildWithDefaults(), + null); + long backoffMs = WFT_COMPLETION_PAGE_RESEND_INITIAL_BACKOFF_MS; + while (true) { + try { + for (RespondWorkflowTaskCompletedRequest page : pages.intermediatePages) { + grpcRetryer.retryWithResult( + () -> respondWorkflowTaskCompleted(page, workflowTypeMetricsScope), + pageRetryOptions); + } + return grpcRetryer.retryWithResult( + () -> respondWorkflowTaskCompleted(pages.finalPage, workflowTypeMetricsScope), + pageRetryOptions); + } catch (StatusRuntimeException e) { + if (!StatusUtils.hasFailure(e, WorkflowTaskCompletionBufferLostFailure.class) + || isShutdown()) { + throw e; + } + try { + Thread.sleep(backoffMs); + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw e; + } + backoffMs = Math.min(backoffMs * 2, WFT_COMPLETION_PAGE_RESEND_MAX_BACKOFF_MS); + } + } + } + + private RespondWorkflowTaskCompletedResponse respondWorkflowTaskCompleted( + RespondWorkflowTaskCompletedRequest request, Scope workflowTypeMetricsScope) { + return service + .blockingStub() + .withOption(METRICS_TAGS_CALL_OPTIONS_KEY, workflowTypeMetricsScope) + .respondWorkflowTaskCompleted(request); } @SuppressWarnings("deprecation") @@ -760,6 +859,56 @@ private void handleReportingFailure( workflowExecution, workflowTypeScope, "Failed result reporting to the server", e); } + /** + * Returns the fail cause when {@code taskCompleted}'s recombined command bytes exceed the + * namespace's completion size limit, or null otherwise. The limit governs the server's + * recombined page buffer, so it only applies when pagination is enabled and the completion is + * large enough to be paginated; a completion that fits in a single request is never buffered + * and is left for the server to accept. Only command bytes count toward the limit, not messages + * or metadata. + */ + private WorkflowTaskFailedCause completionExceedingSizeLimitCause( + RespondWorkflowTaskCompletedRequest taskCompleted) { + if (!namespaceCapabilities.isWorkflowTaskCompletionPagination() + || taskCompleted.getSerializedSize() <= WorkflowTaskCompletionPaginator.MAX_PAGE_BYTES) { + return null; + } + long sizeLimit = namespaceCapabilities.getWorkflowTaskCompletionSizeLimit(); + if (sizeLimit <= 0) { + return null; + } + long commandBytes = 0; + for (Command command : taskCompleted.getCommandsList()) { + commandBytes += command.getSerializedSize(); + } + if (commandBytes <= sizeLimit) { + return null; + } + return WorkflowTaskFailedCause.WORKFLOW_TASK_FAILED_CAUSE_REQUEST_TOO_LARGE; + } + + private Failure requestTooLargeFailure( + String workflowId, RespondWorkflowTaskCompletedRequest taskCompleted) { + long commandBytes = 0; + for (Command command : taskCompleted.getCommandsList()) { + commandBytes += command.getSerializedSize(); + } + String message = + String.format( + "Workflow task completion command size %d exceeds the namespace limit of %d bytes", + commandBytes, namespaceCapabilities.getWorkflowTaskCompletionSizeLimit()); + ApplicationFailure applicationFailure = + ApplicationFailure.newBuilder() + .setMessage(message) + .setType("WorkflowTaskCompletionRequestTooLarge") + .build(); + applicationFailure.setStackTrace(new StackTraceElement[0]); // don't serialize stack trace + return options + .getDataConverter() + .withContext(new WorkflowSerializationContext(namespace, workflowId)) + .exceptionToFailure(applicationFailure); + } + private Failure grpcMessageTooLargeFailure( String workflowId, GrpcMessageTooLargeException e, String messagePrefix) { ApplicationFailure applicationFailure = diff --git a/temporal-sdk/src/main/java/io/temporal/worker/WorkerFactory.java b/temporal-sdk/src/main/java/io/temporal/worker/WorkerFactory.java index 70bcf28c76..b6fa12650f 100644 --- a/temporal-sdk/src/main/java/io/temporal/worker/WorkerFactory.java +++ b/temporal-sdk/src/main/java/io/temporal/worker/WorkerFactory.java @@ -275,6 +275,7 @@ public synchronized void start() { .build()); namespaceCapabilities.setFromCapabilities( describeNamespaceResponse.getNamespaceInfo().getCapabilities()); + namespaceCapabilities.setFromLimits(describeNamespaceResponse.getNamespaceInfo().getLimits()); // Build plugin execution chain (reverse order for proper nesting) Consumer startChain = WorkerFactory::doStart; diff --git a/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginatorTest.java b/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginatorTest.java new file mode 100644 index 0000000000..e77c93d713 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowTaskCompletionPaginatorTest.java @@ -0,0 +1,120 @@ +package io.temporal.internal.worker; + +import static org.junit.Assert.*; + +import com.google.protobuf.ByteString; +import io.temporal.api.command.v1.Command; +import io.temporal.api.command.v1.RecordMarkerCommandAttributes; +import io.temporal.api.common.v1.Payload; +import io.temporal.api.common.v1.Payloads; +import io.temporal.api.enums.v1.CommandType; +import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest; +import java.util.ArrayList; +import java.util.List; +import org.junit.Test; + +public class WorkflowTaskCompletionPaginatorTest { + + private static Command commandWithPayload(int dataSize) { + return Command.newBuilder() + .setCommandType(CommandType.COMMAND_TYPE_RECORD_MARKER) + .setRecordMarkerCommandAttributes( + RecordMarkerCommandAttributes.newBuilder() + .setMarkerName("marker") + .putDetails( + "data", + Payloads.newBuilder() + .addPayloads( + Payload.newBuilder().setData(ByteString.copyFrom(new byte[dataSize]))) + .build())) + .build(); + } + + private static RespondWorkflowTaskCompletedRequest requestWith(List commands) { + return RespondWorkflowTaskCompletedRequest.newBuilder() + .setTaskToken(ByteString.copyFromUtf8("task-token")) + .setIdentity("identity") + .setNamespace("namespace") + .addAllCommands(commands) + .build(); + } + + @Test + public void completionWithinLimitIsASingleFinalPage() { + RespondWorkflowTaskCompletedRequest request = + requestWith(java.util.Collections.singletonList(commandWithPayload(16))); + + WorkflowTaskCompletionPaginator.Pages pages = + WorkflowTaskCompletionPaginator.paginate(request, 4096); + + assertFalse(pages.isPaginated()); + assertEquals(0, pages.finalPage.getPageNumber()); + assertFalse(pages.finalPage.getIntermediatePage()); + assertEquals(1, pages.finalPage.getCommandsCount()); + } + + @Test + public void largeCompletionSplitsCommandsAcrossPages() { + int maxPageBytes = 1024; + int commandCount = 6; + List commands = new ArrayList<>(); + for (int i = 0; i < commandCount; i++) { + commands.add(commandWithPayload(400)); + } + RespondWorkflowTaskCompletedRequest request = requestWith(commands); + assertTrue(request.getSerializedSize() > maxPageBytes); + + WorkflowTaskCompletionPaginator.Pages pages = + WorkflowTaskCompletionPaginator.paginate(request, maxPageBytes); + + assertTrue(pages.isPaginated()); + assertFalse(pages.finalPage.getIntermediatePage()); + assertEquals(0, pages.finalPage.getCommandsCount()); + assertEquals(pages.intermediatePages.size(), pages.finalPage.getPageNumber()); + assertTrue(pages.finalPage.getSerializedSize() <= maxPageBytes); + assertEquals(ByteString.copyFromUtf8("task-token"), pages.finalPage.getTaskToken()); + + int totalCommands = 0; + for (int i = 0; i < pages.intermediatePages.size(); i++) { + RespondWorkflowTaskCompletedRequest page = pages.intermediatePages.get(i); + assertTrue(page.getIntermediatePage()); + assertEquals(i, page.getPageNumber()); + assertEquals(ByteString.copyFromUtf8("task-token"), page.getTaskToken()); + assertTrue( + "intermediate page " + i + " over limit", page.getSerializedSize() <= maxPageBytes); + totalCommands += page.getCommandsCount(); + } + // Every command is preserved exactly once across the intermediate pages. + assertEquals(commandCount, totalCommands); + } + + @Test + public void singleCommandLargerThanAPageIsNotSplit() { + int maxPageBytes = 1024; + RespondWorkflowTaskCompletedRequest request = + requestWith(java.util.Collections.singletonList(commandWithPayload(4096))); + + WorkflowTaskCompletionPaginator.Pages pages = + WorkflowTaskCompletionPaginator.paginate(request, maxPageBytes); + + assertFalse(pages.isPaginated()); + assertEquals(1, pages.finalPage.getCommandsCount()); + assertFalse(pages.finalPage.getIntermediatePage()); + } + + @Test + public void noCommandsIsNotSplit() { + RespondWorkflowTaskCompletedRequest request = + RespondWorkflowTaskCompletedRequest.newBuilder() + .setTaskToken(ByteString.copyFromUtf8("task-token")) + .setIdentity("identity") + .setNamespace("namespace") + .build(); + + WorkflowTaskCompletionPaginator.Pages pages = + WorkflowTaskCompletionPaginator.paginate(request, 1); + + assertFalse(pages.isPaginated()); + assertEquals(0, pages.finalPage.getCommandsCount()); + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/worker/WorkflowTaskCompletionPaginationIntegrationTest.java b/temporal-sdk/src/test/java/io/temporal/worker/WorkflowTaskCompletionPaginationIntegrationTest.java new file mode 100644 index 0000000000..52c3140bbc --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/worker/WorkflowTaskCompletionPaginationIntegrationTest.java @@ -0,0 +1,107 @@ +package io.temporal.worker; + +import static org.junit.Assume.assumeTrue; + +import io.temporal.activity.ActivityInterface; +import io.temporal.activity.ActivityMethod; +import io.temporal.activity.ActivityOptions; +import io.temporal.api.namespace.v1.NamespaceInfo.Capabilities; +import io.temporal.api.workflowservice.v1.DescribeNamespaceRequest; +import io.temporal.api.workflowservice.v1.DescribeNamespaceResponse; +import io.temporal.client.WorkflowOptions; +import io.temporal.testing.internal.SDKTestWorkflowRule; +import io.temporal.workflow.Async; +import io.temporal.workflow.Promise; +import io.temporal.workflow.Workflow; +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import org.junit.Rule; +import org.junit.Test; + +public class WorkflowTaskCompletionPaginationIntegrationTest { + + // Six 1 MiB activity inputs scheduled in a single workflow task produce a ~6 MiB completion, well + // over the ~4 MiB gRPC request limit, so the workflow completes only if the completion is + // paginated. + private static final int ACTIVITY_COUNT = 6; + private static final int ACTIVITY_INPUT_BYTES = 1024 * 1024; + + @Rule + public SDKTestWorkflowRule testWorkflowRule = + SDKTestWorkflowRule.newBuilder() + .setWorkflowTypes(LargeCompletionWorkflowImpl.class) + .setActivityImplementations(new NoopActivityImpl()) + .build(); + + @Test + public void largeCompletionIsPaginated() { + assumeTrue( + "Requires a real server with workflow task completion pagination support", + SDKTestWorkflowRule.useExternalService); + assumeTrue( + "Server does not support workflow task completion pagination", + getNamespaceCapabilities().getWorkflowTaskCompletionPagination()); + + LargeCompletionWorkflow workflow = + testWorkflowRule + .getWorkflowClient() + .newWorkflowStub( + LargeCompletionWorkflow.class, + WorkflowOptions.newBuilder() + .setTaskQueue(testWorkflowRule.getTaskQueue()) + .setWorkflowExecutionTimeout(Duration.ofMinutes(1)) + .build()); + // Completes without error only when the oversized completion is delivered across pages. + workflow.run(); + } + + private Capabilities getNamespaceCapabilities() { + DescribeNamespaceResponse response = + testWorkflowRule + .getWorkflowClient() + .getWorkflowServiceStubs() + .blockingStub() + .describeNamespace( + DescribeNamespaceRequest.newBuilder() + .setNamespace(testWorkflowRule.getWorkflowClient().getOptions().getNamespace()) + .build()); + return response.getNamespaceInfo().getCapabilities(); + } + + @WorkflowInterface + public interface LargeCompletionWorkflow { + @WorkflowMethod + void run(); + } + + public static class LargeCompletionWorkflowImpl implements LargeCompletionWorkflow { + private final NoopActivity activity = + Workflow.newActivityStub( + NoopActivity.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(10)).build()); + + @Override + public void run() { + byte[] input = new byte[ACTIVITY_INPUT_BYTES]; + List> promises = new ArrayList<>(ACTIVITY_COUNT); + for (int i = 0; i < ACTIVITY_COUNT; i++) { + promises.add(Async.procedure(activity::process, input)); + } + Promise.allOf(promises).get(); + } + } + + @ActivityInterface + public interface NoopActivity { + @ActivityMethod + void process(byte[] input); + } + + public static class NoopActivityImpl implements NoopActivity { + @Override + public void process(byte[] input) {} + } +} diff --git a/temporal-serviceclient/src/main/java/io/temporal/serviceclient/MetricsTag.java b/temporal-serviceclient/src/main/java/io/temporal/serviceclient/MetricsTag.java index d894709d74..69b64ecc6b 100644 --- a/temporal-serviceclient/src/main/java/io/temporal/serviceclient/MetricsTag.java +++ b/temporal-serviceclient/src/main/java/io/temporal/serviceclient/MetricsTag.java @@ -26,6 +26,7 @@ public class MetricsTag { public static final String TASK_FAILURE_TYPE = "failure_reason"; public static final String TASK_FAILURE_VALUE_NON_DETERMINISM_ERROR = "NonDeterminismError"; public static final String TASK_FAILURE_VALUE_GRPC_MESSAGE_TOO_LARGE = "GrpcMessageTooLarge"; + public static final String TASK_FAILURE_VALUE_REQUEST_TOO_LARGE = "RequestTooLarge"; public static final String TASK_FAILURE_VALUE_WORKFLOW_ERROR = "WorkflowError"; public static final String TASK_FAILURE_VALUE_ACTIVITY_ERROR = "ActivityError"; public static final String TASK_FAILURE_VALUE_OPERATION_FAILED = "operation_failed";