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
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();
}

Expand Down Expand Up @@ -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<Void> 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<Boolean> 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<String> 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<Instant> t1 = ctx.callActivity("Echo", instant1, Instant.class);
t1.thenAccept(result1 -> {
if (!result1.equals(instant1)) {
ctx.complete(false);
return;
}

Instant instant2 = ctx.getCurrentInstant();
Task<Instant> 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<T> in order
// for it to be mutated by a lambda function.
AtomicReference<Integer> value = new AtomicReference<>(0);

// Each iteration of the for loop appends a new callback stage to the sequence
Task<Void> 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<Task<String>> 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<Task<String>> 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<String> 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();
}
Expand Down Expand Up @@ -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<OrchestratorResponse> responses) {
// // TODO: Error handling for when the orchestrator isn't registered
// Collection<OrchestratorAction> 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<ActivityResponse> 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));
// }
// }
// }
}
Loading