From 2dd507754fdf2cc37d6b777ac22def855a399dd7 Mon Sep 17 00:00:00 2001 From: Chris Gillum Date: Tue, 29 Mar 2022 17:58:38 -0700 Subject: [PATCH 1/2] - Updated for latest protobuf protocols - Renamed ErrorDetails to FailureDetails - Made FailureDetails a property of OrchestrationMetadata - Removed JSON serialization of errors - Split error testing into new JUnit class - Misc. cleanup --- .../durabletask/DurableTaskGrpcWorker.java | 336 +----------------- .../microsoft/durabletask/ErrorDetails.java | 53 --- .../microsoft/durabletask/FailureDetails.java | 62 ++++ .../durabletask/JacksonDataConverter.java | 4 +- .../durabletask/OrchestrationMetadata.java | 8 +- .../durabletask/OrchestrationRunner.java | 2 +- .../durabletask/TaskCanceledException.java | 2 +- .../durabletask/TaskFailedException.java | 21 +- .../durabletask/TaskOrchestrationContext.java | 2 +- .../TaskOrchestrationExecutor.java | 32 +- .../ErrorHandlingIntegrationTests.java | 99 ++++++ .../durabletask/IntegrationTestBase.java | 65 ++++ .../durabletask/IntegrationTests.java | 124 +------ submodules/durabletask-protobuf | 2 +- 14 files changed, 269 insertions(+), 543 deletions(-) delete mode 100644 sdk/src/main/java/com/microsoft/durabletask/ErrorDetails.java create mode 100644 sdk/src/main/java/com/microsoft/durabletask/FailureDetails.java create mode 100644 sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java create mode 100644 sdk/src/test/java/com/microsoft/durabletask/IntegrationTestBase.java diff --git a/sdk/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java b/sdk/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java index ae840dd4..cab99131 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java +++ b/sdk/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java @@ -11,16 +11,10 @@ import io.grpc.*; -import java.io.IOException; -import java.time.Duration; -import java.time.Instant; import java.util.*; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; import java.util.logging.Level; import java.util.logging.Logger; -import java.util.stream.Collectors; -import java.util.stream.IntStream; public class DurableTaskGrpcWorker implements AutoCloseable { private static final int DEFAULT_PORT = 4001; @@ -137,9 +131,9 @@ public void runAndBlock() throws InterruptedException { activityRequest.getTaskId()); } catch (Throwable e) { failureDetails = TaskFailureDetails.newBuilder() - .setErrorName(e.getClass().getName()) + .setErrorType(e.getClass().getName()) .setErrorMessage(e.getMessage()) - .setErrorDetails(ErrorDetails.getFullStackTrace(e)) + .setStackTrace(StringValue.of(FailureDetails.getFullStackTrace(e))) .build(); } @@ -179,280 +173,6 @@ public void stop() { this.close(); } - /** - * Main launches the worker from the command line. - */ - public static void main(String[] args) throws IOException, InterruptedException { - DurableTaskGrpcWorker.Builder builder = DurableTaskGrpcWorker.newBuilder(); - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "ActivityChaining"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - int initial = ctx.getInput(int.class); - - int x = ctx.callActivity("PlusOne", initial, int.class).get(); - int y = ctx.callActivity("PlusOne", x, int.class).get(); - int z = ctx.callActivity("PlusOne", y, int.class).get(); - - ctx.complete(z); - }; - } - }); - builder.addActivity(new TaskActivityFactory() { - @Override - public String getName() { return "PlusOne"; } - - @Override - public TaskActivity create() { - return ctx -> ctx.getInput(int.class) + 1; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { - return "Test"; - } - - @Override - public TaskOrchestration create() { - return ctx -> { - String output = String.format("Finished '%s', ID = %s", ctx.getName(), ctx.getInstanceId()); - ctx.complete(output); - }; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { - return "OrchestrationWithTimer"; - } - - @Override - public TaskOrchestration create() { - return ctx -> { - Task timer = ctx.createTimer(Duration.ofSeconds(3)); - timer.thenRun(() -> ctx.complete(ctx.getInput(Object.class))); - }; - } - }); - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { - return "OrchestrationWithTimer2"; - } - - @Override - public TaskOrchestration create() { - return ctx -> { - ctx.createTimer(Duration.ofSeconds(3)).get(); - ctx.complete(ctx.getInput(Object.class)); - }; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "TwoTimerReplayTester"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - ArrayList list = new ArrayList<>(); - list.add(ctx.getIsReplaying()); - ctx.createTimer(Duration.ofSeconds(0)) - .thenRun(() -> list.add(ctx.getIsReplaying())) - .thenCompose(Void -> ctx.createTimer(Duration.ofSeconds(0))) - .thenRun(() -> { - list.add(ctx.getIsReplaying()); - ctx.complete(list); - }); - }; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "OrchestrationWithActivity"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - String name = ctx.getInput(String.class); - Task task = ctx.callActivity("SayHello", name, String.class); - task.thenAccept(ctx::complete); - }; - } - }); - builder.addActivity(new TaskActivityFactory() { - @Override - public String getName() { return "SayHello"; } - - @Override - public TaskActivity create() { - return ctx -> { - String name = ctx.getInput(String.class); - return String.format("Hello, %s!", name); - }; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "CurrentDateTimeUtc"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - Instant instant1 = ctx.getCurrentInstant(); - Task t1 = ctx.callActivity("Echo", instant1, Instant.class); - t1.thenAccept(result1 -> { - if (!result1.equals(instant1)) { - ctx.complete(false); - return; - } - - Instant instant2 = ctx.getCurrentInstant(); - Task t2 = ctx.callActivity("Echo", instant2, Instant.class); - t2.thenAccept(result2 -> { - boolean success = result2.equals(instant2); - ctx.complete(success); - }); - }); - }; - } - }); - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "CurrentDateTimeUtc2"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - Instant instant1 = ctx.getCurrentInstant(); - Instant result1 = ctx.callActivity("Echo", instant1, Instant.class).get(); - if (!result1.equals(instant1)) { - ctx.complete(false); - return; - } - - Instant instant2 = ctx.getCurrentInstant(); - Instant result2 = ctx.callActivity("Echo", instant2, Instant.class).get(); - - boolean success = result2.equals(instant2); - ctx.complete(success); - }; - } - }); - builder.addActivity(new TaskActivityFactory() { - @Override - public String getName() { return "Echo"; } - - @Override - public TaskActivity create() { - return ctx -> ctx.getInput(Object.class); - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "OrchestrationsWithActivityChain"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - // Java requires us to wrap value in an AtomicReference in order - // for it to be mutated by a lambda function. - AtomicReference value = new AtomicReference<>(0); - - // Each iteration of the for loop appends a new callback stage to the sequence - Task task = ctx.completedTask(null); - for (int i = 0; i < 10; i++) { - task = task.thenCompose(Void -> ctx - .callActivity("PlusOne", value.get(), Integer.class) - .thenAccept(value::set)); - } - - task.thenRun(() -> ctx.complete(value.get())); - }; - } - }); - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "OrchestrationsWithActivityChain2"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - int value = 0; - for (int i = 0; i < 10; i++) { - value = ctx.callActivity("PlusOne", value, int.class).get(); - } - - ctx.complete(value); - }; - } - }); - - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "ActivityFanOut"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - // Schedule each task to run in parallel - List> parallelTasks = IntStream.range(0, 10) - .mapToObj(i -> ctx.callActivity("ToString", i, String.class)) - .collect(Collectors.toList()); - - // Wait for all tasks to complete, then sort and reverse the results - ctx.allOf(parallelTasks).thenAccept(results -> { - Collections.sort(results); - Collections.reverse(results); - ctx.complete(results); - }); - }; - } - }); - builder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return "ActivityFanOut2"; } - - @Override - public TaskOrchestration create() { - return ctx -> { - // Schedule each task to run in parallel - List> parallelTasks = IntStream.range(0, 10) - .mapToObj(i -> ctx.callActivity("ToString", i, String.class)) - .collect(Collectors.toList()); - - // Wait for all tasks to complete, then sort and reverse the results - List results = ctx.allOf(parallelTasks).get(); - Collections.sort(results); - Collections.reverse(results); - ctx.complete(results); - }; - } - }); - builder.addActivity(new TaskActivityFactory() { - @Override - public String getName() { return "ToString"; } - - @Override - public TaskActivity create() { - return ctx -> ctx.getInput(Object.class).toString(); - } - }); - - final DurableTaskGrpcWorker server = builder.build(); - server.runAndBlock(); - } - public static Builder newBuilder() { return new Builder(); } @@ -517,56 +237,4 @@ public DurableTaskGrpcWorker build() { return new DurableTaskGrpcWorker(this); } } - - // class TaskHubWorkerServiceImpl extends TaskHubWorkerServiceImplBase { - - // private final TaskOrchestrationExecutor taskOrchestrationExecutor; - // private final TaskActivityExecutor taskActivityExecutor; - // private final DataConverter dataConverter; - - // public TaskHubWorkerServiceImpl() { - // this.dataConverter = DurableTaskGrpcWorker.this.dataConverter; - // this.taskOrchestrationExecutor = new TaskOrchestrationExecutor( - // DurableTaskGrpcWorker.this.orchestrationFactories, - // DurableTaskGrpcWorker.this.dataConverter, - // DurableTaskGrpcWorker.logger); - // this.taskActivityExecutor = new TaskActivityExecutor( - // DurableTaskGrpcWorker.this.activityFactories, - // DurableTaskGrpcWorker.this.dataConverter, - // DurableTaskGrpcWorker.logger); - // } - - // @Override - // public void executeOrchestrator(OrchestratorRequest req, StreamObserver responses) { - // // TODO: Error handling for when the orchestrator isn't registered - // Collection actions = this.taskOrchestrationExecutor.execute( - // req.getPastEventsList(), - // req.getNewEventsList()); - // OrchestratorResponse response = OrchestratorResponse.newBuilder() - // .addAllActions(actions) - // .build(); - // responses.onNext(response); - // responses.onCompleted(); - // } - - // @Override - // public void executeActivity(ActivityRequest request, StreamObserver responses) { - // // TODO: Error handling for when the activity isn't registered - // String activityName = request.getName(); - // String activityInput = request.getInput().getValue(); - // try { - // ActivityResponse.Builder response = ActivityResponse.newBuilder(); - // String output = this.taskActivityExecutor.execute(activityName, activityInput); - // if (output != null) { - // response.setResult(StringValue.of(output)); - // } - // responses.onNext(response.build()); - // responses.onCompleted(); - // } catch (Exception ex) { - // String details = ErrorDetails.getFullStackTrace(ex); - // Status errorStatus = Status.UNKNOWN.withDescription(details); - // responses.onError(new StatusException(errorStatus)); - // } - // } - // } } \ No newline at end of file diff --git a/sdk/src/main/java/com/microsoft/durabletask/ErrorDetails.java b/sdk/src/main/java/com/microsoft/durabletask/ErrorDetails.java deleted file mode 100644 index 9b20bdf3..00000000 --- a/sdk/src/main/java/com/microsoft/durabletask/ErrorDetails.java +++ /dev/null @@ -1,53 +0,0 @@ -// Copyright (c) Microsoft Corporation. All rights reserved. -// Licensed under the MIT License. -package com.microsoft.durabletask; - -import java.io.PrintWriter; -import java.io.StringWriter; - -import com.fasterxml.jackson.annotation.*; - -// NOTE: This type is serializable and represents a strict error contract for Durable Tasks -// across multiple languages. -class ErrorDetails { - private final String errorName; - private final String errorMessage; - private final String errorDetails; - - @JsonCreator - public ErrorDetails( - @JsonProperty("errorName") String errorName, - @JsonProperty("errorMessage") String errorMessage, - @JsonProperty("errorDetails") String errorDetails) { - this.errorName = errorName; - this.errorMessage = errorMessage; - this.errorDetails = errorDetails; - } - - public ErrorDetails(Exception exception) { - this(exception.getClass().getName(), exception.getMessage(), getFullStackTrace(exception)); - } - - @JsonGetter("errorName") - public String getErrorName() { - return this.errorName; - } - - @JsonGetter("errorMessage") - public String getErrorMessage() { - return this.errorMessage; - } - - @JsonGetter("errorDetails") - public String getErrorDetails() { - return this.errorDetails; - } - - static String getFullStackTrace(Throwable e) { - StringWriter writer = new StringWriter(); - PrintWriter printWriter = new PrintWriter( writer ); - e.printStackTrace(printWriter); - printWriter.flush(); - return writer.toString(); - } -} \ No newline at end of file diff --git a/sdk/src/main/java/com/microsoft/durabletask/FailureDetails.java b/sdk/src/main/java/com/microsoft/durabletask/FailureDetails.java new file mode 100644 index 00000000..af67bad8 --- /dev/null +++ b/sdk/src/main/java/com/microsoft/durabletask/FailureDetails.java @@ -0,0 +1,62 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import com.google.protobuf.StringValue; +import com.microsoft.durabletask.protobuf.OrchestratorService.TaskFailureDetails; + +class FailureDetails { + private final String errorType; + private final String errorMessage; + private final String stackTrace; + + public FailureDetails( + String errorType, + String errorMessage, + String errorDetails) { + this.errorType = errorType; + this.errorMessage = errorMessage; + this.stackTrace = errorDetails; + } + + public FailureDetails(Exception exception) { + this(exception.getClass().getName(), exception.getMessage(), getFullStackTrace(exception)); + } + + public FailureDetails(TaskFailureDetails proto) { + this.errorType = proto.getErrorType(); + this.errorMessage = proto.getErrorMessage(); + this.stackTrace = proto.getStackTrace().getValue(); + } + + public String getErrorType() { + return this.errorType; + } + + public String getErrorMessage() { + return this.errorMessage; + } + + public String getStackTrace() { + return this.stackTrace; + } + + static String getFullStackTrace(Throwable e) { + StackTraceElement[] elements = e.getStackTrace(); + + // Plan for 256 characters per stack frame (which is likely on the high-end) + StringBuilder sb = new StringBuilder(elements.length * 256); + for (StackTraceElement element : elements) { + sb.append("\tat ").append(element.toString()).append(System.lineSeparator()); + } + return sb.toString(); + } + + TaskFailureDetails toProto() { + return TaskFailureDetails.newBuilder() + .setErrorType(this.getErrorType()) + .setErrorMessage(this.getErrorMessage()) + .setStackTrace(StringValue.of(this.getStackTrace())) + .build(); + } +} \ No newline at end of file diff --git a/sdk/src/main/java/com/microsoft/durabletask/JacksonDataConverter.java b/sdk/src/main/java/com/microsoft/durabletask/JacksonDataConverter.java index 78ec9796..a354982e 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/JacksonDataConverter.java +++ b/sdk/src/main/java/com/microsoft/durabletask/JacksonDataConverter.java @@ -21,7 +21,9 @@ public String serialize(Object value) { try { return jsonObjectMapper.writeValueAsString(value); } catch (JsonProcessingException e) { - throw this.wrapConverterException("Failed to serialize output argument.", e); + throw this.wrapConverterException( + String.format("Failed to serialize argument of type '%s'.", value.getClass().getName()), + e); } } diff --git a/sdk/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java b/sdk/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java index b76a67e9..2c6e0dba 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java +++ b/sdk/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java @@ -19,6 +19,7 @@ public class OrchestrationMetadata { private final String serializedInput; private final String serializedOutput; private final String serializedCustomStatus; + private final FailureDetails failureDetails; OrchestrationMetadata( OrchestratorService.GetInstanceResponse fetchResponse, @@ -36,6 +37,7 @@ public class OrchestrationMetadata { this.serializedInput = state.getInput().getValue(); this.serializedOutput = state.getOutput().getValue(); this.serializedCustomStatus = state.getCustomStatus().getValue(); + this.failureDetails = new FailureDetails(state.getFailureDetails()); } public String getName() { @@ -66,6 +68,10 @@ public String getSerializedOutput() { return this.serializedOutput; } + public FailureDetails getFailureDetails() { + return this.failureDetails; + } + public boolean isRunning() { return this.runtimeStatus == OrchestrationRuntimeStatus.RUNNING; } @@ -115,7 +121,7 @@ public String toString() { if (this.serializedInput != null) { sb.append(", Input: '").append(getTrimmedPayload(this.serializedInput)).append('\''); } - + if (this.serializedOutput != null) { sb.append(", Output: '").append(getTrimmedPayload(this.serializedOutput)).append('\''); } diff --git a/sdk/src/main/java/com/microsoft/durabletask/OrchestrationRunner.java b/sdk/src/main/java/com/microsoft/durabletask/OrchestrationRunner.java index c1e26334..cbec349a 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/OrchestrationRunner.java +++ b/sdk/src/main/java/com/microsoft/durabletask/OrchestrationRunner.java @@ -10,7 +10,6 @@ import java.util.HashMap; import java.util.logging.Logger; -// TODO: Move this into the SDK since it shouldn't have any dependencies on Functions anymore // TODO: JavaDoc public final class OrchestrationRunner { private static final Logger logger = Logger.getLogger(OrchestrationRunner.class.getPackage().getName()); @@ -18,6 +17,7 @@ public final class OrchestrationRunner { public static String loadAndRun( String triggerStateProtoBase64String, OrchestratorFunction orchestratorFunc) { + // Example string: CiBhOTMyYjdiYWM5MmI0MDM5YjRkMTYxMDIwNzlmYTM1YSIaCP///////////wESCwi254qRBhDk+rgocgAicgj///////////8BEgwIs+eKkQYQzMXjnQMaVwoLSGVsbG9DaXRpZXMSACJGCiBhOTMyYjdiYWM5MmI0MDM5YjRkMTYxMDIwNzlmYTM1YRIiCiA3ODEwOTA2N2Q4Y2Q0ODg1YWU4NjQ0OTNlMmRlMGQ3OA== byte[] decodedBytes = Base64.getDecoder().decode(triggerStateProtoBase64String); byte[] resultBytes = loadAndRun(decodedBytes, orchestratorFunc); return Base64.getEncoder().encodeToString(resultBytes); diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskCanceledException.java b/sdk/src/main/java/com/microsoft/durabletask/TaskCanceledException.java index fe3dbf1d..8f670ee5 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskCanceledException.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskCanceledException.java @@ -6,6 +6,6 @@ public class TaskCanceledException extends TaskFailedException { // Only intended to be created within this package TaskCanceledException(String message, String taskName, int taskId) { - super(message, taskName, taskId, new ErrorDetails(TaskCanceledException.class.getName(), message, "")); + super(message, taskName, taskId, new FailureDetails(TaskCanceledException.class.getName(), message, "")); } } diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskFailedException.java b/sdk/src/main/java/com/microsoft/durabletask/TaskFailedException.java index 7dca4f62..c397d0d6 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskFailedException.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskFailedException.java @@ -3,11 +3,15 @@ package com.microsoft.durabletask; public class TaskFailedException extends Exception { - private final ErrorDetails details; + private final FailureDetails details; private final String taskName; private final int taskId; - TaskFailedException(String message, String taskName, int taskId, ErrorDetails details) { + TaskFailedException(String taskName, int taskId, FailureDetails details) { + this(getExceptionMessage(taskName, taskId, details), taskName, taskId, details); + } + + protected TaskFailedException(String message, String taskName, int taskId, FailureDetails details) { super(message); this.taskName = taskName; this.taskId = taskId; @@ -23,7 +27,7 @@ public String getTaskName() { } public String getExceptionName() { - return this.details.getErrorName(); + return this.details.getErrorType(); } public String getExceptionMessage() { @@ -31,10 +35,17 @@ public String getExceptionMessage() { } public String getExceptionDetails() { - return this.details.getErrorDetails(); + return this.details.getStackTrace(); } - ErrorDetails getErrorDetails() { + FailureDetails getErrorDetails() { return this.details; } + + private static String getExceptionMessage(String taskName, int taskId, FailureDetails details) { + return String.format("Task '%s' (#%d) failed with an unhandled exception: %s", + taskName, + taskId, + details.getErrorMessage()); + } } diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java index 8e3c6690..b48b08dd 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java @@ -23,7 +23,7 @@ default Task> anyOf(Task... tasks) { Task createTimer(Duration delay); void complete(Object output); - void fail(Object errorOutput); + void fail(FailureDetails failureDetails); Task callActivity(String name, Object input, Class returnType); diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java index ab952487..bcf346c3 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java @@ -4,7 +4,6 @@ import com.google.protobuf.StringValue; import com.google.protobuf.Timestamp; -import com.microsoft.durabletask.DataConverter.DataConverterException; import com.microsoft.durabletask.protobuf.OrchestratorService.*; import com.microsoft.durabletask.protobuf.OrchestratorService.ScheduleTaskAction.Builder; @@ -47,7 +46,7 @@ public Collection execute(List pastEvents, Lis // The orchestrator threw an unhandled exception - fail it // TODO: What's the right way to log this? logger.warning("The orchestrator failed with an unhandled exception: " + e.toString()); - context.fail(new ErrorDetails(e)); + context.fail(new FailureDetails(e)); } catch (OrchestratorBlockedEvent orchestratorBlockedEvent) { logger.fine("The orchestrator has yielded and will await for new events."); } @@ -334,21 +333,7 @@ private void handleTaskFailed(HistoryEvent e) { return; } - // The taskFailed.details field is expected to contain a structured payload - // describing the failure details - String reason = failedEvent.getDetails().getValue(); - ErrorDetails details; - try { - details = this.dataConverter.deserialize(reason, ErrorDetails.class); - } catch (DataConverterException deserializeException) { - // Not expected - but we try to handle it gracefully. - details = new ErrorDetails( - "_UnknownException", - String.format( - "The exception details could not be deserialized. See the error details for the raw error payload: %s", - deserializeException.getMessage()), - reason); - } + FailureDetails details = new FailureDetails(failedEvent.getFailureDetails()); if (!this.isReplaying) { // TODO: Log task failure, including the number of bytes in the result @@ -356,7 +341,6 @@ private void handleTaskFailed(HistoryEvent e) { CompletableTask task = record.getTask(); TaskFailedException exception = new TaskFailedException( - String.format("Activity task '%s' with ID %d failed with an unhandled exception.", record.taskName, taskId), record.taskName, taskId, details); @@ -450,16 +434,16 @@ public void handleTimerFired(HistoryEvent e) { @Override public void complete(Object output) { - this.completeInternal(output, OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); + this.completeInternal(output, null, OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); } @Override - public void fail(Object errorOutput) { + public void fail(FailureDetails failureDetails) { // TODO: How does a parent orchestration use the output to construct an exception? - this.completeInternal(errorOutput, OrchestrationStatus.ORCHESTRATION_STATUS_FAILED); + this.completeInternal(null, failureDetails, OrchestrationStatus.ORCHESTRATION_STATUS_FAILED); } - private void completeInternal(Object output, OrchestrationStatus runtimeStatus) { + private void completeInternal(Object output, FailureDetails failureDetails, OrchestrationStatus runtimeStatus) { if (this.isComplete) { throw new IllegalStateException("The orchestrator was already completed."); } @@ -473,6 +457,10 @@ private void completeInternal(Object output, OrchestrationStatus runtimeStatus) builder.setResult(StringValue.of(resultAsJson)); } + if (failureDetails != null) { + builder.setFailureDetails(failureDetails.toProto()); + } + if (!this.isReplaying) { // TODO: Log completion, including the number of bytes in the output } diff --git a/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java b/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java new file mode 100644 index 00000000..5be325c5 --- /dev/null +++ b/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java @@ -0,0 +1,99 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.microsoft.durabletask; + +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * These integration tests are designed to exercise the core, high-level error-handling features of the Durable Task + * programming model. + *

+ * These tests currently require a sidecar process to be running on the local machine (the sidecar is what accepts the + * client operations and sends invocation instructions to the DurableTaskWorker). + */ +@Tag("integration") +public class ErrorHandlingIntegrationTests extends IntegrationTestBase { + @Test + void orchestratorException() { + final String orchestratorName = "OrchestratorWithException"; + final String errorMessage = "Kah-BOOOOOM!!!"; + + DurableTaskGrpcWorker worker = this.createWorkerBuilder() + .addOrchestrator(orchestratorName, ctx -> { + throw new RuntimeException(errorMessage); + }) + .buildAndStart(); + + DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); + try (worker; client) { + String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, 0); + OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); + assertNotNull(instance); + assertEquals(OrchestrationRuntimeStatus.FAILED, instance.getRuntimeStatus()); + + FailureDetails details = instance.getFailureDetails(); + assertNotNull(details); + assertEquals("java.lang.RuntimeException", details.getErrorType()); + assertTrue(details.getErrorMessage().contains(errorMessage)); + assertNotNull(details.getStackTrace()); + } + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void activityException(boolean handleException) { + final String orchestratorName = "OrchestratorWithActivityException"; + final String activityName = "Throw"; + final String errorMessage = "Kah-BOOOOOM!!!"; + + DurableTaskGrpcWorker worker = this.createWorkerBuilder() + .addOrchestrator(orchestratorName, ctx -> { + try { + ctx.callActivity(activityName).get(); + } catch (TaskFailedException ex) { + if (handleException) { + ctx.complete("handled"); + } else { + throw ex; + } + } + }) + .addActivity(activityName, ctx -> { + throw new RuntimeException(errorMessage); + }) + .buildAndStart(); + + DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); + try (worker; client) { + String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, ""); + OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); + assertNotNull(instance); + + if (handleException) { + String result = instance.readOutputAs(String.class); + assertNotNull(result); + assertEquals("handled", result); + } else { + assertEquals(OrchestrationRuntimeStatus.FAILED, instance.getRuntimeStatus()); + + FailureDetails details = instance.getFailureDetails(); + assertNotNull(details); + + String expectedMessage = String.format( + "Task '%s' (#0) failed with an unhandled exception: %s", + activityName, + errorMessage); + assertEquals(expectedMessage, details.getErrorMessage()); + assertEquals("com.microsoft.durabletask.TaskFailedException", details.getErrorType()); + assertNotNull(details.getStackTrace()); + // CONSIDER: Additional validation of getErrorDetails? + } + } + } +} diff --git a/sdk/src/test/java/com/microsoft/durabletask/IntegrationTestBase.java b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTestBase.java new file mode 100644 index 00000000..74a4c371 --- /dev/null +++ b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTestBase.java @@ -0,0 +1,65 @@ +package com.microsoft.durabletask; + +import org.junit.jupiter.api.AfterEach; + +import java.time.Duration; + +public class IntegrationTestBase { + protected static final Duration defaultTimeout = Duration.ofSeconds(10); + + // All tests that create a server should save it to this variable for proper shutdown + private DurableTaskGrpcWorker server; + + @AfterEach + private void shutdown() throws InterruptedException { + if (this.server != null) { + this.server.stop(); + } + } + + protected TestDurableTaskWorkerBuilder createWorkerBuilder() { + return new TestDurableTaskWorkerBuilder(); + } + + public class TestDurableTaskWorkerBuilder { + final DurableTaskGrpcWorker.Builder innerBuilder; + + private TestDurableTaskWorkerBuilder() { + this.innerBuilder = DurableTaskGrpcWorker.newBuilder(); + } + + public DurableTaskGrpcWorker buildAndStart() { + DurableTaskGrpcWorker server = this.innerBuilder.build(); + IntegrationTestBase.this.server = server; + server.start(); + return server; + } + + public TestDurableTaskWorkerBuilder addOrchestrator( + String name, + TaskOrchestration implementation) { + this.innerBuilder.addOrchestration(new TaskOrchestrationFactory() { + @Override + public String getName() { return name; } + + @Override + public TaskOrchestration create() { return implementation; } + }); + return this; + } + + public IntegrationTests.TestDurableTaskWorkerBuilder addActivity( + String name, + TaskActivity implementation) + { + this.innerBuilder.addActivity(new TaskActivityFactory() { + @Override + public String getName() { return name; } + + @Override + public TaskActivity create() { return implementation; } + }); + return this; + } + } +} diff --git a/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java index b15398e0..1668e828 100644 --- a/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java +++ b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java @@ -31,7 +31,7 @@ * sends invocation instructions to the DurableTaskWorker). */ @Tag("integration") -public class IntegrationTests { +public class IntegrationTests extends IntegrationTestBase { static final Duration defaultTimeout = Duration.ofSeconds(100); // All tests that create a server should save it to this variable for proper shutdown @@ -224,82 +224,6 @@ void activityChain() throws IOException { } } - @Test - void orchestratorException() throws IOException { - final String orchestratorName = "OrchestratorWithException"; - final String errorMessage = "Kah-BOOOOOM!!!"; - - DurableTaskGrpcWorker worker = this.createWorkerBuilder() - .addOrchestrator(orchestratorName, ctx -> { - throw new RuntimeException(errorMessage); - }) - .buildAndStart(); - - DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); - try (worker; client) { - String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, 0); - OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); - assertNotNull(instance); - assertEquals(OrchestrationRuntimeStatus.FAILED, instance.getRuntimeStatus()); - - ErrorDetails details = instance.readOutputAs(ErrorDetails.class); - assertNotNull(details); - assertTrue(details.getErrorDetails().contains(errorMessage)); - } - } - - @ParameterizedTest - @ValueSource(booleans = {true, false}) - void activityException(boolean handleException) throws IOException { - final String orchestratorName = "OrchestratorWithActivityException"; - final String activityName = "Throw"; - final String errorMessage = "Kah-BOOOOOM!!!"; - - DurableTaskGrpcWorker worker = this.createWorkerBuilder() - .addOrchestrator(orchestratorName, ctx -> { - try { - ctx.callActivity(activityName).get(); - } catch (TaskFailedException ex) { - if (handleException) { - ctx.complete("handled"); - } else { - throw ex; - } - } - }) - .addActivity(activityName, ctx -> { - throw new RuntimeException(errorMessage); - }) - .buildAndStart(); - - DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); - try (worker; client) { - String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, ""); - OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); - assertNotNull(instance); - - if (handleException) { - String result = instance.readOutputAs(String.class); - assertNotNull(result); - assertEquals("handled", result); - } else { - assertEquals(OrchestrationRuntimeStatus.FAILED, instance.getRuntimeStatus()); - - ErrorDetails details = instance.readOutputAs(ErrorDetails.class); - assertNotNull(details); - - String expectedMessage = String.format( - "Activity task '%s' with ID 0 failed with an unhandled exception.", - activityName, - errorMessage); - assertEquals(expectedMessage, details.getErrorMessage()); - assertEquals("com.microsoft.durabletask.TaskFailedException", details.getErrorName()); - assertNotNull(details.getErrorDetails()); - // CONSIDER: Additional validation of getErrorDetails? - } - } - } - @Test void activityFanOut() throws IOException { final String orchestratorName = "ActivityFanOut"; @@ -419,50 +343,4 @@ void externalEventsWithTimeouts(boolean raiseEvent) throws IOException { } } } - - private TestDurableTaskWorkerBuilder createWorkerBuilder() { - return new TestDurableTaskWorkerBuilder(); - } - - public class TestDurableTaskWorkerBuilder { - final DurableTaskGrpcWorker.Builder innerBuilder; - - private TestDurableTaskWorkerBuilder() { - this.innerBuilder = DurableTaskGrpcWorker.newBuilder(); - } - - public DurableTaskGrpcWorker buildAndStart() { - DurableTaskGrpcWorker server = this.innerBuilder.build(); - IntegrationTests.this.server = server; - server.start(); - return server; - } - - public TestDurableTaskWorkerBuilder addOrchestrator( - String name, - TaskOrchestration implementation) { - this.innerBuilder.addOrchestration(new TaskOrchestrationFactory() { - @Override - public String getName() { return name; } - - @Override - public TaskOrchestration create() { return implementation; } - }); - return this; - } - - public TestDurableTaskWorkerBuilder addActivity( - String name, - TaskActivity implementation) - { - this.innerBuilder.addActivity(new TaskActivityFactory() { - @Override - public String getName() { return name; } - - @Override - public TaskActivity create() { return implementation; } - }); - return this; - } - } } diff --git a/submodules/durabletask-protobuf b/submodules/durabletask-protobuf index f2df040b..bf27c7da 160000 --- a/submodules/durabletask-protobuf +++ b/submodules/durabletask-protobuf @@ -1 +1 @@ -Subproject commit f2df040b29ef28e8f2ad03bae420217c52dc1181 +Subproject commit bf27c7da57dbcbb6a85dd5398801b79dc2524814 From 7f1311d1b9e73e560cc08af5831bc08b057719b0 Mon Sep 17 00:00:00 2001 From: kaibocai <89094811+kaibocai@users.noreply.github.com> Date: Wed, 30 Mar 2022 12:57:10 -0500 Subject: [PATCH 2/2] add logics to support suborchestration--add one intergration tests --- .../durabletask/TaskOrchestrationContext.java | 14 +++ .../TaskOrchestrationExecutor.java | 115 +++++++++++++++++- .../durabletask/IntegrationTests.java | 21 ++++ 3 files changed, 144 insertions(+), 6 deletions(-) diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java index b48b08dd..b5e23d2e 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java @@ -35,6 +35,20 @@ default Task callActivity(String name, Object input) { return this.callActivity(name, input, Void.class); } + Task callSubOrchestrator(String name, Object input, String instanceId, Class voidClass); + + default Task callSubOrchestrator(String name){ + return this.callSubOrchestrator(name, null); + } + + default Task callSubOrchestrator(String name, Object input){ + return this.callSubOrchestrator(name, input, null); + } + + default Task callSubOrchestrator(String name, Object input, Class voidClass){ + return this.callSubOrchestrator(name, input, null, voidClass); + } + Task waitForExternalEvent(String name, Duration timeout, Class dataType) throws TaskCanceledException; default Task waitForExternalEvent(String name, Duration timeout) throws TaskCanceledException { diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java index bcf346c3..423b51f4 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java @@ -226,6 +226,43 @@ public Task callActivity(String name, Object input, Class returnType) return task; } + @Override + public Task callSubOrchestrator(String name, Object input, String instanceId, Class returnType){ + Helpers.throwIfArgumentNull(name, "name"); + Helpers.throwIfArgumentNull(returnType, "returnType"); + int id = this.sequenceNumber++; + + String serializedInput = this.dataConverter.serialize(input); + CreateSubOrchestrationAction.Builder createSubOrchestrationActionBuilder = CreateSubOrchestrationAction.newBuilder().setName(name); + if (serializedInput != null) { + createSubOrchestrationActionBuilder.setInput(StringValue.of(serializedInput)); + } + + if (instanceId == null) { + instanceId = UUID.randomUUID().toString(); + } + createSubOrchestrationActionBuilder.setInstanceId(instanceId); + + this.pendingActions.put(id, OrchestratorAction.newBuilder() + .setId(id) + .setCreateSubOrchestration(createSubOrchestrationActionBuilder) + .build()); + + if (!this.isReplaying) { + this.logger.fine(() -> String.format( + "%s: calling sub-orchestration '%s' (#%d) with serialized input: %s", + this.instanceId, + name, + id, + serializedInput != null ? serializedInput : "(null)")); + } + + CompletableTask task = new CompletableTask<>(); + TaskRecord record = new TaskRecord<>(task, name, returnType); + this.openTasks.put(id, record); + return task; + } + public Task waitForExternalEvent(String name, Duration timeout, Class dataType) { Helpers.throwIfArgumentNull(name, "name"); Helpers.throwIfArgumentNull(dataType, "dataType"); @@ -432,6 +469,69 @@ public void handleTimerFired(HistoryEvent e) { task.complete(null); } + private void handleSubOrchestrationCreated(HistoryEvent e) { + int taskId = e.getEventId(); + SubOrchestrationInstanceCreatedEvent subOrchestrationInstanceCreated = e.getSubOrchestrationInstanceCreated(); + OrchestratorAction taskAction = this.pendingActions.remove(taskId); + if (taskAction == null) { + String message = String.format( + "Non-deterministic orchestrator detected: a history event scheduling an activity task with sequence ID %d and name '%s' was replayed but the current orchestrator implementation didn't actually schedule this task. Was a change made to the orchestrator code after this instance had already started running?", + taskId, + subOrchestrationInstanceCreated.getName()); + throw new NonDeterministicOrchestratorException(message); + } + } + + private void handleSubOrchestrationCompleted(HistoryEvent e) { + SubOrchestrationInstanceCompletedEvent subOrchestrationInstanceCompletedEvent = e.getSubOrchestrationInstanceCompleted(); + int taskId = subOrchestrationInstanceCompletedEvent.getTaskScheduledId(); + TaskRecord record = this.openTasks.remove(taskId); + if (record == null) { + this.logger.warning("Discarding a potentially duplicate TaskCompleted event with ID = " + taskId); + return; + } + String rawResult = subOrchestrationInstanceCompletedEvent.getResult().getValue(); + + if (!this.isReplaying) { + // TODO: Structured logging + // TODO: Would it make more sense to put this log in the activity executor? + this.logger.fine(() -> String.format( + "%s: Activity '%s' (#%d) completed with serialized output: %s", + this.instanceId, + record.getTaskName(), + taskId, + rawResult != null ? rawResult : "(null)")); + + } + + Object result = this.dataConverter.deserialize(rawResult, record.getDataType()); + CompletableTask task = record.getTask(); + task.complete(result); + } + + private void handleSubOrchestrationFailed(HistoryEvent e){ + SubOrchestrationInstanceFailedEvent subOrchestrationInstanceFailedEvent = e.getSubOrchestrationInstanceFailed(); + int taskId = subOrchestrationInstanceFailedEvent.getTaskScheduledId(); + TaskRecord record = this.openTasks.remove(taskId); + if (record == null) { + // TODO: Log a warning about a potential duplicate task completion event + return; + } + + FailureDetails details = new FailureDetails(subOrchestrationInstanceFailedEvent.getFailureDetails()); + + if (!this.isReplaying) { + // TODO: Log task failure, including the number of bytes in the result + } + + CompletableTask task = record.getTask(); + TaskFailedException exception = new TaskFailedException( + record.taskName, + taskId, + details); + task.completeExceptionally(exception); + } + @Override public void complete(Object output) { this.completeInternal(output, null, OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); @@ -529,12 +629,15 @@ private void processEvent(HistoryEvent e) throws TaskFailedException, Orchestrat case TIMERFIRED: this.handleTimerFired(e); break; -// case SUBORCHESTRATIONINSTANCECREATED: -// break; -// case SUBORCHESTRATIONINSTANCECOMPLETED: -// break; -// case SUBORCHESTRATIONINSTANCEFAILED: -// break; + case SUBORCHESTRATIONINSTANCECREATED: + this.handleSubOrchestrationCreated(e); + break; + case SUBORCHESTRATIONINSTANCECOMPLETED: + this.handleSubOrchestrationCompleted(e); + break; + case SUBORCHESTRATIONINSTANCEFAILED: + this.handleSubOrchestrationFailed(e); + break; // case EVENTSENT: // break; case EVENTRAISED: diff --git a/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java index 1668e828..f5feffd9 100644 --- a/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java +++ b/sdk/src/test/java/com/microsoft/durabletask/IntegrationTests.java @@ -224,6 +224,27 @@ void activityChain() throws IOException { } } + @Test + void subOrchestration(){ + final String orchestratorName = "SubOrchestration"; + DurableTaskGrpcWorker worker = this.createWorkerBuilder().addOrchestrator(orchestratorName, ctx -> { + int result = 5; + int input = ctx.getInput(int.class); + if (input < 3){ + result += ctx.callSubOrchestrator(orchestratorName, input + 1, int.class).get(); + } + ctx.complete(result); + }).buildAndStart(); + DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); + try(worker; client){ + String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, 1); + OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); + assertNotNull(instance); + assertEquals(OrchestrationRuntimeStatus.COMPLETED, instance.getRuntimeStatus()); + assertEquals(15, instance.readOutputAs(int.class)); + } + } + @Test void activityFanOut() throws IOException { final String orchestratorName = "ActivityFanOut";