Skip to content
Merged
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 @@ -35,6 +35,20 @@ default Task<Void> callActivity(String name, Object input) {
return this.callActivity(name, input, Void.class);
}

<V> Task<V> callSubOrchestrator(String name, Object input, String instanceId, Class<V> returnType);

default Task<Void> callSubOrchestrator(String name){
return this.callSubOrchestrator(name, null);
}

default Task<Void> callSubOrchestrator(String name, Object input){
return this.callSubOrchestrator(name, input, null);
}

default <V>Task<V> callSubOrchestrator(String name, Object input, Class<V> returnType){
return this.callSubOrchestrator(name, input, null, returnType);
}

<V> Task<V> waitForExternalEvent(String name, Duration timeout, Class<V> dataType) throws TaskCanceledException;

default Task<Void> waitForExternalEvent(String name, Duration timeout) throws TaskCanceledException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,45 @@ public <V> Task<V> callActivity(String name, Object input, Class<V> returnType)
return task;
}

@Override
public <V> Task<V> callSubOrchestrator(String name, Object input, String instanceId, Class<V> 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();
Comment thread
kaibocai marked this conversation as resolved.
Outdated
}
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<V> task = new CompletableTask<>();
TaskRecord<V> record = new TaskRecord<>(task, name, returnType);
this.openTasks.put(id, record);
return task;
}

public <V> Task<V> waitForExternalEvent(String name, Duration timeout, Class<V> dataType) {
Helpers.throwIfArgumentNull(name, "name");
Helpers.throwIfArgumentNull(dataType, "dataType");
Expand Down Expand Up @@ -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();
Comment thread
kaibocai marked this conversation as resolved.
Outdated
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);
Expand Down Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Comment thread
kaibocai marked this conversation as resolved.
Outdated
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down