diff --git a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java index b48b08dd..e1f35b70 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 returnType); + + 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 returnType){ + return this.callSubOrchestrator(name, input, null, returnType); + } + 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..629d2400 100644 --- a/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java +++ b/sdk/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java @@ -226,6 +226,45 @@ 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)); + } + + //TODO:replace this with a deterministic GUID generation so that it's safe for replay, + // please find potentail bug here https://github.com/microsoft/durabletask-dotnet/issues/9 + 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 +471,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 sub-orchestration 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 SubOrchestrationInstanceCompleted 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: Sub-orchestrator '%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 +631,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/ErrorHandlingIntegrationTests.java b/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java index 5be325c5..993a5081 100644 --- a/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java +++ b/sdk/src/test/java/com/microsoft/durabletask/ErrorHandlingIntegrationTests.java @@ -96,4 +96,54 @@ void activityException(boolean handleException) { } } } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void subOrchestrationException(boolean handleException){ + final String orchestratorName = "OrchestrationWithBustedSubOrchestrator"; + final String subOrchestratorName = "BustedSubOrchestrator"; + final String errorMessage = "Kah-BOOOOOM!!!"; + + DurableTaskGrpcWorker worker = this.createWorkerBuilder() + .addOrchestrator(orchestratorName, ctx -> { + try { + String result = ctx.callSubOrchestrator(subOrchestratorName, "", String.class).get(); + ctx.complete(result); + } catch (TaskFailedException ex) { + if (handleException) { + ctx.complete("handled"); + } else { + throw ex; + } + } + }) + .addOrchestrator(subOrchestratorName, ctx -> { + throw new RuntimeException(errorMessage); + }) + .buildAndStart(); + DurableTaskClient client = DurableTaskGrpcClient.newBuilder().build(); + try(worker; client){ + String instanceId = client.scheduleNewOrchestrationInstance(orchestratorName, 1); + OrchestrationMetadata instance = client.waitForInstanceCompletion(instanceId, defaultTimeout, true); + assertNotNull(instance); + if (handleException){ + assertEquals(OrchestrationRuntimeStatus.COMPLETED, instance.getRuntimeStatus()); + 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", + subOrchestratorName, + 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/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";