Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import io.temporal.common.converter.DataConverter;
import io.temporal.internal.common.ProtobufTimeUtils;
import io.temporal.internal.common.WorkflowExecutionUtils;
import io.temporal.internal.payload.storage.ExternalStorageRunner;
import io.temporal.internal.worker.*;
import io.temporal.payload.context.WorkflowSerializationContext;
import io.temporal.serviceclient.MetricsTag;
Expand Down Expand Up @@ -77,6 +78,12 @@ public WorkflowTaskHandler.Result handleWorkflowTask(PollWorkflowTaskQueueRespon
String workflowType = workflowTask.getWorkflowType().getName();
Scope metricsScope =
options.getMetricsScope().tagged(ImmutableMap.of(MetricsTag.WORKFLOW_TYPE, workflowType));
ExternalStorageRunner externalStorage = options.getExternalStorage();
if (externalStorage == null) {
ExternalStorageRunner.throwIfContainsReference(workflowTask);
} else {
workflowTask = externalStorage.retrieve(workflowTask);
}
return handleWorkflowTaskWithQuery(workflowTask.toBuilder(), metricsScope);
}

Expand All @@ -94,7 +101,8 @@ private Result handleWorkflowTaskWithQuery(
logWorkflowTaskToBeProcessed(workflowTask, createdNew);

ServiceWorkflowHistoryIterator historyIterator =
new ServiceWorkflowHistoryIterator(service, namespace, workflowTask, metricsScope);
new ServiceWorkflowHistoryIterator(
service, namespace, workflowTask, metricsScope, options.getExternalStorage());
boolean finalCommand;
Result result;

Expand Down Expand Up @@ -395,6 +403,12 @@ private WorkflowRunTaskHandler createStatefulHandler(
.blockingStub()
.withOption(METRICS_TAGS_CALL_OPTIONS_KEY, metricsScope)
.getWorkflowExecutionHistory(getHistoryRequest);
ExternalStorageRunner externalStorage = options.getExternalStorage();
if (externalStorage == null) {
ExternalStorageRunner.throwIfContainsReference(getHistoryResponse);
} else {
getHistoryResponse = externalStorage.retrieve(getHistoryResponse);
}
workflowTask
.setHistory(getHistoryResponse.getHistory())
.setNextPageToken(getHistoryResponse.getNextPageToken());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,12 +12,14 @@
import io.temporal.api.workflowservice.v1.GetWorkflowExecutionHistoryRequest;
import io.temporal.api.workflowservice.v1.GetWorkflowExecutionHistoryResponse;
import io.temporal.api.workflowservice.v1.PollWorkflowTaskQueueResponseOrBuilder;
import io.temporal.internal.payload.storage.ExternalStorageRunner;
import io.temporal.internal.retryer.GrpcRetryer;
import io.temporal.serviceclient.RpcRetryOptions;
import io.temporal.serviceclient.WorkflowServiceStubs;
import java.time.Duration;
import java.util.Iterator;
import java.util.NoSuchElementException;
import javax.annotation.Nullable;

/** Supports iteration over history while loading new pages through calls to the service. */
class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator {
Expand All @@ -29,6 +31,7 @@ class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator {
private final Scope metricsScope;
private final PollWorkflowTaskQueueResponseOrBuilder task;
private final GrpcRetryer grpcRetryer;
private final @Nullable ExternalStorageRunner externalStorage;
private Deadline deadline;
private Iterator<HistoryEvent> current;
ByteString nextPageToken;
Expand All @@ -38,10 +41,20 @@ class ServiceWorkflowHistoryIterator implements WorkflowHistoryIterator {
String namespace,
PollWorkflowTaskQueueResponseOrBuilder task,
Scope metricsScope) {
this(service, namespace, task, metricsScope, null);
}

ServiceWorkflowHistoryIterator(
WorkflowServiceStubs service,
String namespace,
PollWorkflowTaskQueueResponseOrBuilder task,
Scope metricsScope,
@Nullable ExternalStorageRunner externalStorage) {
this.service = service;
this.namespace = namespace;
this.task = task;
this.metricsScope = metricsScope;
this.externalStorage = externalStorage;
// TODO Refactor WorkflowHistoryIteratorTest or WorkflowHistoryIterator to remove this check.
// `service == null` shouldn't be allowed as it's needed for a normal functioning of this
// class.
Expand All @@ -64,7 +77,13 @@ public boolean hasNext() {
// true.
GetWorkflowExecutionHistoryResponse response = queryWorkflowExecutionHistory();

current = response.getHistory().getEventsList().iterator();
History history = response.getHistory();
if (externalStorage == null) {
ExternalStorageRunner.throwIfContainsReference(history);
} else {
history = externalStorage.retrieve(history);
}
current = history.getEventsList().iterator();
nextPageToken = response.getNextPageToken();
// Server can return an empty page, but a valid nextPageToken that contains
// more events.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,29 @@
import com.google.common.base.Preconditions;
import com.google.common.base.Strings;
import com.google.protobuf.ByteString;
import com.google.protobuf.MessageOrBuilder;
import com.uber.m3.tally.Scope;
import com.uber.m3.tally.Stopwatch;
import com.uber.m3.util.ImmutableMap;
import io.grpc.StatusRuntimeException;
import io.temporal.api.command.v1.*;
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.failure.v1.Failure;
import io.temporal.api.workflowservice.v1.*;
import io.temporal.common.CancellationToken;
import io.temporal.failure.ApplicationFailure;
import io.temporal.internal.logging.LoggerTag;
import io.temporal.internal.payload.storage.ExternalStorageRunner;
import io.temporal.internal.payload.visitor.MessageVisitor;
import io.temporal.internal.retryer.GrpcMessageTooLargeException;
import io.temporal.internal.retryer.GrpcRetryer;
import io.temporal.payload.context.WorkflowSerializationContext;
import io.temporal.payload.storage.StorageDriverActivityInfo;
import io.temporal.payload.storage.StorageDriverTargetInfo;
import io.temporal.payload.storage.StorageDriverWorkflowInfo;
import io.temporal.serviceclient.MetricsTag;
import io.temporal.serviceclient.RpcRetryOptions;
import io.temporal.serviceclient.WorkflowServiceStubs;
Expand Down Expand Up @@ -381,6 +389,89 @@ public String toString() {
options.getIdentity(), namespace, taskQueue);
}

private void storeOutboundPayloads(
com.google.protobuf.Message.Builder builder, @Nullable StorageDriverTargetInfo target) {
ExternalStorageRunner externalStorage = options.getExternalStorage();
if (externalStorage != null) {
externalStorage.store(builder, target);
}
}

private void storeOutboundPayloads(
com.google.protobuf.Message.Builder builder,
@Nullable StorageDriverTargetInfo target,
MessageVisitor<StorageDriverTargetInfo> targetVisitor) {
ExternalStorageRunner externalStorage = options.getExternalStorage();
if (externalStorage != null) {
externalStorage.store(builder, target, targetVisitor, CancellationToken.none());
}
}

@Nullable
private StorageDriverTargetInfo workflowStorageTarget(
WorkflowExecution execution, String workflowType) {
if (options.getExternalStorage() == null) {
return null;
}
return new StorageDriverWorkflowInfo(
namespace, execution.getWorkflowId(), execution.getRunId(), workflowType);
}

static StorageDriverTargetInfo deriveStorageTarget(
String namespace, StorageDriverTargetInfo current, MessageOrBuilder message) {
if (!(message instanceof CommandOrBuilder)) {
return current;
}
CommandOrBuilder command = (CommandOrBuilder) message;
// Keep this exhaustive so new command attributes require an explicit target decision.
switch (command.getAttributesCase()) {
case SCHEDULE_ACTIVITY_TASK_COMMAND_ATTRIBUTES:
ScheduleActivityTaskCommandAttributesOrBuilder activity =
command.getScheduleActivityTaskCommandAttributesOrBuilder();
return new StorageDriverActivityInfo(
namespace, activity.getActivityId(), null, activity.getActivityType().getName());
case START_CHILD_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
StartChildWorkflowExecutionCommandAttributesOrBuilder child =
command.getStartChildWorkflowExecutionCommandAttributesOrBuilder();
return new StorageDriverWorkflowInfo(
namespace, child.getWorkflowId(), null, child.getWorkflowType().getName());
case SIGNAL_EXTERNAL_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
WorkflowExecution execution =
command.getSignalExternalWorkflowExecutionCommandAttributes().getExecution();
return new StorageDriverWorkflowInfo(
namespace, execution.getWorkflowId(), execution.getRunId(), null);
case CONTINUE_AS_NEW_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
if (current instanceof StorageDriverWorkflowInfo) {
ContinueAsNewWorkflowExecutionCommandAttributesOrBuilder continueAsNew =
command.getContinueAsNewWorkflowExecutionCommandAttributesOrBuilder();
StorageDriverWorkflowInfo currentWorkflow = (StorageDriverWorkflowInfo) current;
String workflowType = continueAsNew.getWorkflowType().getName();
return new StorageDriverWorkflowInfo(
namespace,
currentWorkflow.getId(),
null,
Strings.isNullOrEmpty(workflowType) ? currentWorkflow.getType() : workflowType);
}
return current;
case ATTRIBUTES_NOT_SET:
case START_TIMER_COMMAND_ATTRIBUTES:
case COMPLETE_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
case FAIL_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
case REQUEST_CANCEL_ACTIVITY_TASK_COMMAND_ATTRIBUTES:
case CANCEL_TIMER_COMMAND_ATTRIBUTES:
case CANCEL_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
case REQUEST_CANCEL_EXTERNAL_WORKFLOW_EXECUTION_COMMAND_ATTRIBUTES:
case RECORD_MARKER_COMMAND_ATTRIBUTES:
case UPSERT_WORKFLOW_SEARCH_ATTRIBUTES_COMMAND_ATTRIBUTES:
case PROTOCOL_MESSAGE_COMMAND_ATTRIBUTES:
case MODIFY_WORKFLOW_PROPERTIES_COMMAND_ATTRIBUTES:
case SCHEDULE_NEXUS_OPERATION_COMMAND_ATTRIBUTES:
case REQUEST_CANCEL_NEXUS_OPERATION_COMMAND_ATTRIBUTES:
return current;
}
throw new IllegalStateException("Unhandled command attributes: " + command.getAttributesCase());
}

private class TaskHandlerImpl implements PollTaskExecutor.TaskHandler<WorkflowTask> {

final WorkflowTaskHandler handler;
Expand Down Expand Up @@ -453,7 +544,10 @@ public void handle(WorkflowTask task) throws Exception {
if (queryCompleted != null) {
try {
sendDirectQueryCompletedResponse(
currentTask.getTaskToken(), queryCompleted.toBuilder(), workflowTypeScope);
currentTask.getTaskToken(),
queryCompleted.toBuilder(),
workflowTypeScope,
workflowStorageTarget(workflowExecution, workflowType));
} catch (StatusRuntimeException e) {
GrpcMessageTooLargeException tooLargeException =
GrpcMessageTooLargeException.tryWrap(e);
Expand All @@ -473,7 +567,10 @@ public void handle(WorkflowTask task) throws Exception {
.setErrorMessage(failure.getMessage())
.setFailure(failure);
sendDirectQueryCompletedResponse(
currentTask.getTaskToken(), queryFailedBuilder, workflowTypeScope);
currentTask.getTaskToken(),
queryFailedBuilder,
workflowTypeScope,
workflowStorageTarget(workflowExecution, workflowType));
}
} else {
try {
Expand All @@ -489,7 +586,8 @@ public void handle(WorkflowTask task) throws Exception {
currentTask.getTaskToken(),
requestBuilder,
result.getRequestRetryOptions(),
workflowTypeScope);
workflowTypeScope,
workflowStorageTarget(workflowExecution, workflowType));
// If we were processing a speculative WFT the server may instruct us that the
// task was dropped by resting out event ID.
long resetEventId = response.getResetHistoryEventId();
Expand All @@ -509,7 +607,8 @@ public void handle(WorkflowTask task) throws Exception {
currentTask.getTaskToken(),
taskFailed.toBuilder(),
result.getRequestRetryOptions(),
workflowTypeScope);
workflowTypeScope,
workflowStorageTarget(workflowExecution, workflowType));
}

// Apply post-completion metrics only if runnable present and the above succeeded
Expand Down Expand Up @@ -546,7 +645,8 @@ public void handle(WorkflowTask task) throws Exception {
currentTask.getTaskToken(),
taskFailedBuilder,
result.getRequestRetryOptions(),
workflowTypeScope);
workflowTypeScope,
workflowStorageTarget(workflowExecution, workflowType));
}
}
} catch (Exception e) {
Expand Down Expand Up @@ -651,7 +751,8 @@ private RespondWorkflowTaskCompletedResponse sendTaskCompleted(
ByteString taskToken,
RespondWorkflowTaskCompletedRequest.Builder taskCompleted,
RpcRetryOptions retryOptions,
Scope workflowTypeMetricsScope) {
Scope workflowTypeMetricsScope,
@Nullable StorageDriverTargetInfo storageTarget) {
GrpcRetryer.GrpcRetryerOptions grpcRetryOptions =
new GrpcRetryer.GrpcRetryerOptions(
RpcRetryOptions.newBuilder().buildWithDefaultsFrom(retryOptions), null);
Expand All @@ -674,12 +775,16 @@ private RespondWorkflowTaskCompletedResponse sendTaskCompleted(
taskCompleted.setBinaryChecksum(options.getBuildId());
}

MessageVisitor<StorageDriverTargetInfo> storageTargetVisitor =
(current, message) -> deriveStorageTarget(namespace, current, message);
storeOutboundPayloads(taskCompleted, storageTarget, storageTargetVisitor);
RespondWorkflowTaskCompletedRequest request = taskCompleted.build();
return grpcRetryer.retryWithResult(
() ->
service
.blockingStub()
.withOption(METRICS_TAGS_CALL_OPTIONS_KEY, workflowTypeMetricsScope)
.respondWorkflowTaskCompleted(taskCompleted.build()),
.respondWorkflowTaskCompleted(request),
grpcRetryOptions);
}

Expand All @@ -688,7 +793,8 @@ private void sendTaskFailed(
ByteString taskToken,
RespondWorkflowTaskFailedRequest.Builder taskFailed,
RpcRetryOptions retryOptions,
Scope workflowTypeMetricsScope) {
Scope workflowTypeMetricsScope,
@Nullable StorageDriverTargetInfo storageTarget) {
GrpcRetryer.GrpcRetryerOptions grpcRetryOptions =
new GrpcRetryer.GrpcRetryerOptions(
RpcRetryOptions.newBuilder().buildWithDefaultsFrom(retryOptions), null);
Expand All @@ -702,25 +808,30 @@ private void sendTaskFailed(
taskFailed.setWorkerVersion(options.workerVersionStamp());
}

storeOutboundPayloads(taskFailed, storageTarget);
RespondWorkflowTaskFailedRequest request = taskFailed.build();
grpcRetryer.retry(
() ->
service
.blockingStub()
.withOption(METRICS_TAGS_CALL_OPTIONS_KEY, workflowTypeMetricsScope)
.respondWorkflowTaskFailed(taskFailed.build()),
.respondWorkflowTaskFailed(request),
grpcRetryOptions);
}

private void sendDirectQueryCompletedResponse(
ByteString taskToken,
RespondQueryTaskCompletedRequest.Builder queryCompleted,
Scope workflowTypeMetricsScope) {
Scope workflowTypeMetricsScope,
@Nullable StorageDriverTargetInfo storageTarget) {
queryCompleted.setTaskToken(taskToken).setNamespace(namespace);
storeOutboundPayloads(queryCompleted, storageTarget);
RespondQueryTaskCompletedRequest request = queryCompleted.build();
// Do not retry query response
service
.blockingStub()
.withOption(METRICS_TAGS_CALL_OPTIONS_KEY, workflowTypeMetricsScope)
.respondQueryTaskCompleted(queryCompleted.build());
.respondQueryTaskCompleted(request);
}

private void logExceptionDuringResultReporting(
Expand Down
Loading
Loading