diff --git a/CHANGELOG.md b/CHANGELOG.md index 035ab4328b..e6d5de697b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,17 @@ to docs, or any other relevant information. ## [Unreleased] +### :boom: Breaking Changes +- Experimental `NexusClientCallsInterceptor.GetNexusOperationResultInput` now takes a single nullable + `NexusSerializationContext` in place of separate endpoint, service and operation arguments, and + exposes it through `getSerializationContext()` in place of `getEndpoint()`, `getService()` and + `getOperation()`. + +### Fixed +- A standalone Nexus operation handle returned when an ID conflict policy of use-existing reuses a + running operation now uses that operation's endpoint, service and operation for its serialization + context. Previously it used the ones named by the start request, which may differ. + ### Changed - Release notes for all future releases are now in a single CHANGELOG.md file. `releases` directory with old release notes is kept for historical reference. diff --git a/temporal-sdk/src/main/java/io/temporal/client/UntypedNexusServiceClientImpl.java b/temporal-sdk/src/main/java/io/temporal/client/UntypedNexusServiceClientImpl.java index 73222e3ade..b79e8e490a 100644 --- a/temporal-sdk/src/main/java/io/temporal/client/UntypedNexusServiceClientImpl.java +++ b/temporal-sdk/src/main/java/io/temporal/client/UntypedNexusServiceClientImpl.java @@ -4,6 +4,7 @@ import io.temporal.common.Experimental; import io.temporal.common.converter.DataConverter; import io.temporal.common.interceptors.NexusClientCallsInterceptor; +import io.temporal.common.interceptors.NexusClientCallsInterceptor.DescribeNexusOperationExecutionInput; import io.temporal.common.interceptors.NexusClientCallsInterceptor.StartNexusOperationExecutionInput; import io.temporal.common.interceptors.NexusClientCallsInterceptor.StartNexusOperationExecutionOutput; import io.temporal.internal.client.NexusClientResolvedOptions; @@ -44,15 +45,29 @@ class UntypedNexusServiceClientImpl implements UntypedNexusServiceClient { @Override public UntypedNexusOperationHandle start( String operation, StartNexusOperationOptions options, @Nullable Object arg) { - Payload payload = serializeInput(arg, operation); + NexusSerializationContext serializationContext = + new NexusSerializationContext(endpoint, serviceName, operation); + Payload payload = serializeInput(arg, serializationContext); StartNexusOperationExecutionInput input = new StartNexusOperationExecutionInput( endpoint, serviceName, operation, payload, options, Collections.emptyMap()); StartNexusOperationExecutionOutput output = invoker.startNexusOperationExecution(input); - // The handle keeps what the start request was for, including when the server returned an - // operation that was already running, so the result is decoded the way it was encoded. + // An already-running operation reused by ID may have been started with a different endpoint, + // service or operation than this request names. The start response does not say which, so + // describe it. + if (!output.isStarted()) { + NexusOperationExecutionDescription description = + invoker + .describeNexusOperationExecution( + new DescribeNexusOperationExecutionInput( + output.getOperationId(), output.getRunId())) + .getDescription(); + serializationContext = + new NexusSerializationContext( + description.getEndpoint(), description.getService(), description.getOperation()); + } return new NexusOperationHandleImpl( - output.getOperationId(), output.getRunId(), invoker, endpoint, serviceName, operation); + output.getOperationId(), output.getRunId(), invoker, serializationContext); } @Override @@ -75,13 +90,14 @@ public R execute( return NexusOperationHandle.fromUntyped(handle, resultClass, resultType).getResult(); } - private @Nullable Payload serializeInput(@Nullable Object arg, String operation) { + private @Nullable Payload serializeInput( + @Nullable Object arg, NexusSerializationContext serializationContext) { if (arg == null) { return null; } Class argClass = arg.getClass(); return dataConverter - .withContext(new NexusSerializationContext(endpoint, serviceName, operation)) + .withContext(serializationContext) .toPayload(arg) .orElseThrow( () -> diff --git a/temporal-sdk/src/main/java/io/temporal/common/interceptors/NexusClientCallsInterceptor.java b/temporal-sdk/src/main/java/io/temporal/common/interceptors/NexusClientCallsInterceptor.java index 5f0bd59980..f4490efa8d 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/interceptors/NexusClientCallsInterceptor.java +++ b/temporal-sdk/src/main/java/io/temporal/common/interceptors/NexusClientCallsInterceptor.java @@ -251,14 +251,12 @@ final class GetNexusOperationResultInput { private final @Nonnull Deadline deadline; private final Class resultClass; private final @Nullable Type resultType; - private final @Nullable String endpoint; - private final @Nullable String service; - private final @Nullable String operation; + private final @Nullable NexusSerializationContext serializationContext; /** * Equivalent to {@link #GetNexusOperationResultInput(String, String, Deadline, Class, Type, - * String, String, String)} with no endpoint, service or operation, which is the case for a - * handle obtained by operation ID rather than by starting an operation. + * NexusSerializationContext)} with no serialization context, which is the case for a handle + * obtained by operation ID rather than by starting an operation. */ public GetNexusOperationResultInput( String operationId, @@ -266,21 +264,12 @@ public GetNexusOperationResultInput( @Nonnull Deadline deadline, Class resultClass, @Nullable Type resultType) { - this(operationId, runId, deadline, resultClass, resultType, null, null, null); + this(operationId, runId, deadline, resultClass, resultType, null); } /** - * The endpoint, service and operation identify the Nexus operation the result is being read - * for, and are used to build the {@link NexusSerializationContext} the result and failure are - * decoded with. They must all be set or all be {@code null}: a partially identified operation - * would silently decode without a context, which for a converter that varies by context means - * reading the payload the wrong way rather than failing. - * - * @param endpoint Nexus endpoint the operation was started on, or {@code null} if the operation - * was not started through this handle - * @param service Nexus service the operation was started on, or {@code null} - * @param operation Nexus operation that was started, or {@code null} - * @throws IllegalArgumentException if only some of endpoint, service and operation are set + * @param serializationContext serialization context of the start request, or {@code null} for a + * handle obtained by operation ID rather than by starting one */ public GetNexusOperationResultInput( String operationId, @@ -288,28 +277,13 @@ public GetNexusOperationResultInput( @Nonnull Deadline deadline, Class resultClass, @Nullable Type resultType, - @Nullable String endpoint, - @Nullable String service, - @Nullable String operation) { - boolean anySet = endpoint != null || service != null || operation != null; - boolean allSet = endpoint != null && service != null && operation != null; - if (anySet && !allSet) { - throw new IllegalArgumentException( - "endpoint, service and operation must all be set or all be null, got endpoint=" - + endpoint - + ", service=" - + service - + ", operation=" - + operation); - } + @Nullable NexusSerializationContext serializationContext) { this.operationId = operationId; this.runId = runId; this.deadline = deadline; this.resultClass = resultClass; this.resultType = resultType; - this.endpoint = endpoint; - this.service = service; - this.operation = operation; + this.serializationContext = serializationContext; } public String getOperationId() { @@ -335,31 +309,12 @@ public Type getResultType() { } /** - * Nexus endpoint the operation was started on. {@code null} when the operation was not started - * through this handle, in which case {@link #getService()} and {@link #getOperation()} are - * {@code null} too and the result is decoded without a Nexus serialization context. - */ - @Nullable - public String getEndpoint() { - return endpoint; - } - - /** - * Nexus service the operation was started on, or {@code null}. Set exactly when {@link - * #getEndpoint()} is set. + * Serialization context of the start request, or {@code null} for a handle obtained by + * operation ID rather than by starting one. */ @Nullable - public String getService() { - return service; - } - - /** - * Nexus operation that was started, or {@code null}. Set exactly when {@link #getEndpoint()} is - * set. - */ - @Nullable - public String getOperation() { - return operation; + public NexusSerializationContext getSerializationContext() { + return serializationContext; } } diff --git a/temporal-sdk/src/main/java/io/temporal/internal/client/NexusOperationHandleImpl.java b/temporal-sdk/src/main/java/io/temporal/internal/client/NexusOperationHandleImpl.java index ac6ca98f07..613b77c0a9 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/client/NexusOperationHandleImpl.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/client/NexusOperationHandleImpl.java @@ -10,6 +10,7 @@ import io.temporal.common.interceptors.NexusClientCallsInterceptor.GetNexusOperationResultOutput; import io.temporal.common.interceptors.NexusClientCallsInterceptor.RequestCancelNexusOperationExecutionInput; import io.temporal.common.interceptors.NexusClientCallsInterceptor.TerminateNexusOperationExecutionInput; +import io.temporal.payload.context.NexusSerializationContext; import java.lang.reflect.Type; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -25,25 +26,19 @@ public final class NexusOperationHandleImpl implements UntypedNexusOperationHand private final String operationId; private final @Nullable String runId; private final NexusClientCallsInterceptor interceptor; - // What the operation was started on, retained so its result and failure are decoded with the - // same serialization context that encoded them. All null for a handle obtained by operation ID, - // which never saw a start request. - private final @Nullable String endpoint; - private final @Nullable String service; - private final @Nullable String operation; + // Null for a handle obtained by operation ID rather than by starting one. + private final @Nullable NexusSerializationContext serializationContext; public NexusOperationHandleImpl( String operationId, @Nullable String runId, NexusClientCallsInterceptor interceptor) { - this(operationId, runId, interceptor, null, null, null); + this(operationId, runId, interceptor, null); } public NexusOperationHandleImpl( String operationId, @Nullable String runId, NexusClientCallsInterceptor interceptor, - @Nullable String endpoint, - @Nullable String service, - @Nullable String operation) { + @Nullable NexusSerializationContext serializationContext) { if (operationId == null) { throw new IllegalArgumentException("operationId is required"); } @@ -53,9 +48,7 @@ public NexusOperationHandleImpl( this.operationId = operationId; this.runId = runId; this.interceptor = interceptor; - this.endpoint = endpoint; - this.service = service; - this.operation = operation; + this.serializationContext = serializationContext; } @Override @@ -140,9 +133,7 @@ public R getResult( Deadline.after(timeout, unit), resultClass, resultType, - endpoint, - service, - operation); + serializationContext); return interceptor.getNexusOperationResult(input).getResult(); } @@ -162,9 +153,7 @@ public CompletableFuture getResultAsync( Deadline.after(timeout, unit), resultClass, resultType, - endpoint, - service, - operation); + serializationContext); return interceptor .getNexusOperationResultAsync(input) .thenApply(GetNexusOperationResultOutput::getResult); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/client/RootNexusClientInvoker.java b/temporal-sdk/src/main/java/io/temporal/internal/client/RootNexusClientInvoker.java index d92e2f3989..1c757208b1 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/client/RootNexusClientInvoker.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/client/RootNexusClientInvoker.java @@ -149,8 +149,8 @@ public DescribeNexusOperationExecutionOutput describeNexusOperationExecution( } catch (StatusRuntimeException e) { throw mapNotFound(input.getOperationId(), input.getRunId().orElse(null), e); } - // The response names the endpoint, service and operation, so the description decodes its - // payloads and failures with the same context the operation was started with. + // The response names the endpoint, service and operation, so its payloads and failures use that + // context. NexusOperationExecutionInfo info = response.getInfo(); DataConverter dataConverter = clientOptions @@ -163,23 +163,12 @@ public DescribeNexusOperationExecutionOutput describeNexusOperationExecution( response, dataConverter, clientOptions.getNamespace())); } - /** - * The client's data converter scoped to the Nexus operation the result is being read for, or left - * as-is when the operation is unknown, which is the case for a handle obtained by operation ID. - * - *

{@link GetNexusOperationResultInput} guarantees the endpoint, service and operation are set - * together or not at all, so one null means all three are null. Absence is tested with {@code - * null} rather than emptiness so an operation genuinely named with an empty string still gets a - * context. - */ private DataConverter dataConverterFor(GetNexusOperationResultInput input) { DataConverter dataConverter = clientOptions.getDataConverter(); - if (input.getEndpoint() == null) { - return dataConverter; - } - return dataConverter.withContext( - new NexusSerializationContext( - input.getEndpoint(), input.getService(), input.getOperation())); + NexusSerializationContext serializationContext = input.getSerializationContext(); + return serializationContext != null + ? dataConverter.withContext(serializationContext) + : dataConverter; } private DescribeNexusOperationExecutionRequest buildDescribeRequest( diff --git a/temporal-sdk/src/main/java/io/temporal/internal/nexus/NexusTaskHandlerImpl.java b/temporal-sdk/src/main/java/io/temporal/internal/nexus/NexusTaskHandlerImpl.java index 0da2ebf077..fbea85edb4 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/nexus/NexusTaskHandlerImpl.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/nexus/NexusTaskHandlerImpl.java @@ -160,13 +160,10 @@ public Result handle(NexusTask task, Scope metricsScope) throws TimeoutException } /** - * Records the serialization context for the operation this task is for, so that the data - * converter used for its input, result and failures is scoped to the endpoint, service and - * operation the request names. + * Records the serialization context for this task's endpoint, service and operation. * - *

Note that servers before 1.30.0 do not report the endpoint the task was addressed to, so the - * context is scoped by an empty endpoint there and does not agree with the caller's. Nexus - * serialization context on the handler side requires server 1.30.0 or later. + *

The endpoint is empty on servers before 1.30.0, which do not report the endpoint a task was + * addressed to. */ private void setSerializationContext(String service, String operation) { InternalNexusOperationContext nexusContext = CurrentNexusOperationContext.get(); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusTaskHandler.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusTaskHandler.java index f11eefc8a9..4d628eddf3 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusTaskHandler.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusTaskHandler.java @@ -26,8 +26,7 @@ class Result { @Nullable private final HandlerException handlerException; // Serialization context of the operation the task was for. Carried on the result because the // reply is encoded after the handler has returned, by which point the per-task context is no - // longer in scope. Null when the task named no operation, or when the server did not report - // the endpoint it was addressed to. + // longer in scope. Null when the task failed before its operation was known. @Nullable private final NexusSerializationContext serializationContext; public Result(@Nonnull Response response) { diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusWorker.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusWorker.java index 9cd2e14aa4..22a3047cc1 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusWorker.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/NexusWorker.java @@ -551,9 +551,8 @@ private void sendReply( .setTaskToken(taskToken) .setIdentity(options.getIdentity()) .setNamespace(namespace); - // The caller decodes this failure with the operation's context, so it has to be encoded - // with the same one. The context rides on the result because it is no longer in scope by - // the time the reply is built. + // The context rides on the result because it is no longer in scope by the time the reply + // is built. NexusSerializationContext serializationContext = response.getSerializationContext(); DataConverter dataConverterWithContext = serializationContext != null diff --git a/temporal-sdk/src/main/java/io/temporal/payload/context/NexusSerializationContext.java b/temporal-sdk/src/main/java/io/temporal/payload/context/NexusSerializationContext.java index 3e355c1a14..cdac598653 100644 --- a/temporal-sdk/src/main/java/io/temporal/payload/context/NexusSerializationContext.java +++ b/temporal-sdk/src/main/java/io/temporal/payload/context/NexusSerializationContext.java @@ -14,11 +14,7 @@ * operation results, and encoding failures produced while handling a Nexus task. * *

The context is not propagated to the eventual result of an asynchronous operation, because the - * operation is completed out of band rather than by the task the handler was invoked for. A - * standalone operation handle uses the context of its start request, including when the start - * request returns an already-running operation; a handle obtained by operation ID without starting - * an operation has no endpoint, service, or operation to build a context from and therefore - * serializes without one. + * operation is completed out of band rather than by the task the handler was invoked for. * *

Failure conversion is not symmetric: a failure is encoded by the handler and decoded by the * caller, so an implementation sees this context on only one side of a given failure, and for some diff --git a/temporal-sdk/src/test/java/io/temporal/client/nexus/StandaloneNexusSerializationContextTest.java b/temporal-sdk/src/test/java/io/temporal/client/nexus/StandaloneNexusSerializationContextTest.java index e8a100c2f6..366bdaabdc 100644 --- a/temporal-sdk/src/test/java/io/temporal/client/nexus/StandaloneNexusSerializationContextTest.java +++ b/temporal-sdk/src/test/java/io/temporal/client/nexus/StandaloneNexusSerializationContextTest.java @@ -4,6 +4,7 @@ import com.google.protobuf.ByteString; import io.temporal.api.common.v1.Payload; +import io.temporal.api.enums.v1.NexusOperationIdConflictPolicy; import io.temporal.api.nexus.v1.Endpoint; import io.temporal.client.NexusClient; import io.temporal.client.NexusClientOptions; @@ -17,6 +18,7 @@ import io.temporal.common.converter.DataConverter; import io.temporal.common.converter.DefaultDataConverter; import io.temporal.common.converter.FailureConverter; +import io.temporal.failure.ApplicationFailure; import io.temporal.failure.DefaultFailureConverter; import io.temporal.payload.codec.PayloadCodec; import io.temporal.payload.context.NexusSerializationContext; @@ -114,6 +116,44 @@ public void startedHandleUsesItsStartRequestContext() { CODEC.nexusContexts().contains(expectedContext())); } + @Test + public void startedHandleReusingARunningOperationUsesThatOperationsContext() { + NexusClient client = nexusClient(); + UntypedNexusServiceClient serviceClient = + client.newUntypedNexusServiceClient( + testWorkflowRule.getNexusEndpoint().getSpec().getName(), SERVICE); + String operationId = UUID.randomUUID().toString(); + UntypedNexusOperationHandle running = + serviceClient.start( + OPERATION, + StartNexusOperationOptions.newBuilder() + .setId(operationId) + .setScheduleToCloseTimeout(Duration.ofSeconds(30)) + .build(), + EchoNexusServiceImpl.ASYNC_PREFIX + UUID.randomUUID()); + + // Reusing the ID returns the running operation, not a new one for the operation named here. + UntypedNexusOperationHandle reused = + serviceClient.start( + "otherOperation", + StartNexusOperationOptions.newBuilder() + .setId(operationId) + .setScheduleToCloseTimeout(Duration.ofSeconds(30)) + .setIdConflictPolicy( + NexusOperationIdConflictPolicy.NEXUS_OPERATION_ID_CONFLICT_POLICY_USE_EXISTING) + .build(), + "ignored"); + Assert.assertEquals(running.getNexusOperationRunId(), reused.getNexusOperationRunId()); + reused.terminate("done"); + FAILURE_CONVERTER.reset(); + + Assert.assertThrows(NexusOperationFailedException.class, () -> reused.getResult(String.class)); + Assert.assertEquals( + "the reused handle should convert the result under the running operation's context", + Collections.singletonList(expectedContext()), + FAILURE_CONVERTER.nexusContexts()); + } + @Test public void failureUsesTheOperationsContext() { UntypedNexusOperationHandle handle = @@ -128,6 +168,25 @@ public void failureUsesTheOperationsContext() { FAILURE_CONVERTER.nexusContexts()); } + @Test + public void handlerFailureUsesTheOperationsContext() { + UntypedNexusOperationHandle handle = + startOperation(EchoNexusServiceImpl.HANDLER_FAIL_PREFIX + UUID.randomUUID()); + + // The worker encodes this failure after the handler has returned, outside the operation's + // scope. The strict codec fails the caller's decode if the worker used a different context. + NexusOperationFailedException exception = + Assert.assertThrows( + NexusOperationFailedException.class, () -> handle.getResult(String.class)); + Throwable cause = exception; + while (cause != null && !(cause instanceof ApplicationFailure)) { + cause = cause.getCause(); + } + Assert.assertNotNull("expected an ApplicationFailure in " + exception, cause); + Assert.assertEquals( + "failure-detail", ((ApplicationFailure) cause).getDetails().get(String.class)); + } + @Test public void describeUsesContextFromTheResponse() { String input = "ping-" + UUID.randomUUID(); diff --git a/temporal-sdk/src/test/java/io/temporal/internal/client/RootNexusClientInvokerTest.java b/temporal-sdk/src/test/java/io/temporal/internal/client/RootNexusClientInvokerTest.java index 4b4779a34f..ce72ddff7c 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/client/RootNexusClientInvokerTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/client/RootNexusClientInvokerTest.java @@ -38,47 +38,8 @@ public class RootNexusClientInvokerTest { NexusClientOptions.getDefaultInstance().getIdentity())); @Test - public void resultInputRejectsPartiallyIdentifiedOperation() { - // A partially identified operation would decode without a Nexus context, which for a converter - // that varies by context means reading the payload the wrong way rather than failing. - Assert.assertThrows( - IllegalArgumentException.class, - () -> - new GetNexusOperationResultInput<>( - "op-1", - null, - Deadline.after(10, TimeUnit.SECONDS), - String.class, - String.class, - "endpoint", - null, - "operation")); - } - - @Test - public void resultInputAcceptsFullyIdentifiedOperation() { - GetNexusOperationResultInput input = - new GetNexusOperationResultInput<>( - "op-1", - null, - Deadline.after(10, TimeUnit.SECONDS), - String.class, - String.class, - "endpoint", - "service", - "operation"); - Assert.assertEquals("endpoint", input.getEndpoint()); - Assert.assertEquals("service", input.getService()); - Assert.assertEquals("operation", input.getOperation()); - } - - @Test - public void resultInputAcceptsUnidentifiedOperation() { - // A handle obtained by operation ID never saw a start request. - GetNexusOperationResultInput input = input(); - Assert.assertNull(input.getEndpoint()); - Assert.assertNull(input.getService()); - Assert.assertNull(input.getOperation()); + public void resultInputForHandleObtainedByIdHasNoContext() { + Assert.assertNull(input().getSerializationContext()); } private static GetNexusOperationResultInput input() { diff --git a/temporal-sdk/src/test/java/io/temporal/internal/nexus/NexusTaskHandlerSerializationContextTest.java b/temporal-sdk/src/test/java/io/temporal/internal/nexus/NexusTaskHandlerSerializationContextTest.java index 35328b15f3..9f0f46ff8f 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/nexus/NexusTaskHandlerSerializationContextTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/nexus/NexusTaskHandlerSerializationContextTest.java @@ -182,8 +182,7 @@ public void serializerWithoutTaskInScopeUsesContextlessConverter() { @Test public void taskWithoutAnEndpointIsScopedByAnEmptyEndpoint() throws TimeoutException { // Servers before 1.30.0 do not report the endpoint a Nexus task was addressed to. The handler - // still scopes by service and operation, with an empty endpoint, which will not agree with the - // caller's context but is a Nexus context rather than an absent one. + // still uses a Nexus context there, with an empty endpoint, rather than none. NexusSerializationContext expected = new NexusSerializationContext("", SERVICE, OPERATION); DataConverter callerConverter = signingConverter(new SigningCodec()); SigningCodec handlerCodec = new SigningCodec(); diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/shared/EchoNexusServiceImpl.java b/temporal-sdk/src/test/java/io/temporal/workflow/shared/EchoNexusServiceImpl.java index 7b1961fa0e..12ac31557d 100644 --- a/temporal-sdk/src/test/java/io/temporal/workflow/shared/EchoNexusServiceImpl.java +++ b/temporal-sdk/src/test/java/io/temporal/workflow/shared/EchoNexusServiceImpl.java @@ -1,6 +1,7 @@ package io.temporal.workflow.shared; import io.nexusrpc.OperationException; +import io.nexusrpc.handler.HandlerException; import io.nexusrpc.handler.OperationCancelDetails; import io.nexusrpc.handler.OperationContext; import io.nexusrpc.handler.OperationHandler; @@ -8,6 +9,7 @@ import io.nexusrpc.handler.OperationStartDetails; import io.nexusrpc.handler.OperationStartResult; import io.nexusrpc.handler.ServiceImpl; +import io.temporal.failure.ApplicationFailure; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; @@ -18,6 +20,8 @@ *

    *
  • An input starting with {@link #FAIL_PREFIX} causes {@code start} to throw an {@link * OperationException#failed} so callers see a non-retryable handler failure. + *
  • An input starting with {@link #HANDLER_FAIL_PREFIX} causes {@code start} to throw a + * non-retryable {@link HandlerException} whose cause carries a detail payload. *
  • An input starting with {@link #ASYNC_PREFIX} causes {@code start} to return an * async-started result with a synthetic operation token; the operation stays in {@code * RUNNING} until something terminal (cancel that takes effect, terminate, schedule-to-close) @@ -35,6 +39,12 @@ public class EchoNexusServiceImpl { /** Inputs starting with this prefix make {@code start} throw, exercising the failure path. */ public static final String FAIL_PREFIX = "FAIL:"; + /** + * Inputs starting with this prefix make {@code start} throw a {@link HandlerException}, + * exercising the task failure path rather than the operation failure path. + */ + public static final String HANDLER_FAIL_PREFIX = "HANDLER_FAIL:"; + /** * Inputs starting with this prefix make {@code start} return an async-started result without ever * completing the operation. Used by cancel/terminate tests so the operation stays in {@code @@ -59,6 +69,13 @@ public OperationStartResult start( if (input != null && input.startsWith(FAIL_PREFIX)) { throw OperationException.failed("intentional failure: " + input); } + if (input != null && input.startsWith(HANDLER_FAIL_PREFIX)) { + throw new HandlerException( + HandlerException.ErrorType.BAD_REQUEST, + "intentional handler failure: " + input, + ApplicationFailure.newNonRetryableFailure( + "root cause", "ContextFailure", "failure-detail")); + } if (input != null && input.startsWith(ASYNC_PREFIX)) { return OperationStartResult.async("token-" + UUID.randomUUID()); }