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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand All @@ -75,13 +90,14 @@ public <R> 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(
() ->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -251,65 +251,39 @@ final class GetNexusOperationResultInput<R> {
private final @Nonnull Deadline deadline;
private final Class<R> 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,
@Nullable String runId,
@Nonnull Deadline deadline,
Class<R> 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,
@Nullable String runId,
@Nonnull Deadline deadline,
Class<R> 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() {
Expand All @@ -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;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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");
}
Expand All @@ -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
Expand Down Expand Up @@ -140,9 +133,7 @@ public <R> R getResult(
Deadline.after(timeout, unit),
resultClass,
resultType,
endpoint,
service,
operation);
serializationContext);
return interceptor.getNexusOperationResult(input).getResult();
}

Expand All @@ -162,9 +153,7 @@ public <R> CompletableFuture<R> getResultAsync(
Deadline.after(timeout, unit),
resultClass,
resultType,
endpoint,
service,
operation);
serializationContext);
return interceptor
.getNexusOperationResultAsync(input)
.thenApply(GetNexusOperationResultOutput::getResult);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
*
* <p>{@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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
* <p>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.
* <p>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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,11 +14,7 @@
* operation results, and encoding failures produced while handling a Nexus task.
*
* <p>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.
*
* <p>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
Expand Down
Loading
Loading