From 085bdf77f8787b6fa57682f7d7555da3a94a0c78 Mon Sep 17 00:00:00 2001 From: LM Date: Wed, 11 Dec 2019 20:58:53 -0800 Subject: [PATCH 1/8] ActorRuntime --- sdk/src/main/java/io/dapr/actors/ActorId.java | 14 ++ .../java/io/dapr/actors/DaprAsyncClient.java | 2 +- .../io/dapr/actors/DaprClientBuilder.java | 2 +- .../io/dapr/actors/DaprHttpAsyncClient.java | 4 +- .../io/dapr/actors/runtime/ActorManager.java | 11 ++ .../io/dapr/actors/runtime/ActorRuntime.java | 164 ++++++++++++++++++ .../io/dapr/actors/runtime/ActorService.java | 8 + 7 files changed, 201 insertions(+), 4 deletions(-) create mode 100644 sdk/src/main/java/io/dapr/actors/ActorId.java create mode 100644 sdk/src/main/java/io/dapr/actors/runtime/ActorManager.java create mode 100644 sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java create mode 100644 sdk/src/main/java/io/dapr/actors/runtime/ActorService.java diff --git a/sdk/src/main/java/io/dapr/actors/ActorId.java b/sdk/src/main/java/io/dapr/actors/ActorId.java new file mode 100644 index 0000000000..412f729452 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/ActorId.java @@ -0,0 +1,14 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors; + +/** + * Stub + */ +public class ActorId { + public ActorId(String id) { + } +} diff --git a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java b/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java index 9b0345f429..47f87caeef 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java @@ -10,7 +10,7 @@ /** * Interface for interacting with Dapr runtime. */ -interface DaprAsyncClient { +public interface DaprAsyncClient { /** * Invokes an Actor method on Dapr. diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java index 4ac91f30a1..418221db06 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java +++ b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java @@ -10,7 +10,7 @@ /** * Builds an instance of DaprAsyncClient or DaprClient. */ -class DaprClientBuilder { +public class DaprClientBuilder { /** * Default port for Dapr after checking environment variable. diff --git a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java index b4e153e5d8..b44a3f237c 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java @@ -16,7 +16,7 @@ /** * Http client to call Dapr's API for actors. */ -class DaprHttpAsyncClient implements DaprAsyncClient { +public class DaprHttpAsyncClient implements DaprAsyncClient { /** * Defines the standard application/json type for HTTP calls in Dapr. @@ -47,7 +47,7 @@ class DaprHttpAsyncClient implements DaprAsyncClient { * @param port Port for calling Dapr. (e.g. 3500) * @param httpClient RestClient used for all API calls in this new instance. */ - DaprHttpAsyncClient(int port, OkHttpClient httpClient) + public DaprHttpAsyncClient(int port, OkHttpClient httpClient) { this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; this.httpClient = httpClient; diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorManager.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorManager.java new file mode 100644 index 0000000000..3974dd43a6 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorManager.java @@ -0,0 +1,11 @@ +package io.dapr.actors.runtime; + +// stub +public class ActorManager { + + public ActorManager(ActorService actorService) { + + } +} + + diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java new file mode 100644 index 0000000000..dd1c216567 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -0,0 +1,164 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.runtime; + +import java.util.HashMap; +import java.util.Set; +import java.util.function.Function; + +import io.dapr.actors.*; + +/** + * Contains methods to register actor types. Registering the types allows the runtime to create instances of the actor. + */ +public class ActorRuntime { + private static volatile ActorRuntime instance; + private static DaprAsyncClient daprAsyncClient; + private final String TraceType = "ActorRuntime"; + private final HashMap actorManagers; + + private ActorRuntime() throws IllegalStateException{ + if (instance != null) { + throw new IllegalStateException("ActorRuntime should only be constructed once"); + } + + this.actorManagers = new HashMap(); + daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); + } + + /** TODO: maybe this should not be a singleton + * Returns an ActorRuntime object. + * @return An ActorRuntime object. + */ + public static ActorRuntime getInstance() { + if (instance == null) { + synchronized (ActorRuntime.class) { + if (instance == null) { + instance = new ActorRuntime(); + } + } + } + + return instance; + } + + /** + * + * @return Actor type names registered with the runtime. + */ + public Set getRegisteredActorTypes() { + return this.actorManagers.keySet(); + } + + /** + * Registers an actor with the runtime. + * @param clazz The type of actor. + */ + void RegisterActor(Class clazz) { + RegisterActor(clazz, null); + } + + /** + * Registers an actor with the runtime. + * @param clazz The type of actor. + * @param actorServiceFactory An optional delegate to create actor service. This can be used for dependency injection into actors. + */ + void RegisterActor(Class clazz, Function actorServiceFactory) + { + ActorTypeInformation actorTypeInfo = ActorTypeInformation.create(clazz); + + ActorService actorService; + if (actorServiceFactory != null) + { + actorService = actorServiceFactory.apply(actorTypeInfo); + } + else + { + actorService = new ActorService(actorTypeInfo); + } + + // Create ActorManagers, override existing entry if registered again. + this.actorManagers.put(actorTypeInfo.getName(), new ActorManager(actorService)); + } + + /** + * Activates an actor for an actor type with given actor id. + * @param actorTypeName Actor type name to activate the actor for. + * @param actorId Actor id for the actor to be activated. + */ + static void Activate(String actorTypeName, String actorId) + { + // uncomment when ActorManager implemented + // return instance.GetActorManager(actorTypeName).ActivateActor(new ActorId(actorId)); + } + + /** + * Deactivates an actor for an actor type with given actor id. + * @param actorTypeName Actor type name to deactivate the actor for. + * @param actorId Actor id for the actor to be deactivated. + */ + static void Deactivate(String actorTypeName, String actorId) + { + // uncomment when ActorManager implemented + // return instance.GetActorManager(actorTypeName).DeactivateActor(new ActorId(actorId)); + } + + /** + * Invokes the specified method for the actor, this is mainly used for cross language invocation. + * @param actorTypeName Actor type name to invoke the method for. + * @param actorId Actor id for the actor for which method will be invoked. + * @param actorMethodName Method name on actor type which will be invoked. + * @param requestBodyStream Payload for the actor method. + * @param responseBodyStream Response for the actor method. + * @return + */ + static void Dispatch(String actorTypeName, String actorId, String actorMethodName, byte[] requestBodyStream, byte[] responseBodyStream) + { + // uncomment when ActorManager implemented + // return instance.GetActorManager(actorTypeName).Dispatch(new ActorId(actorId), actorMethodName, requestBodyStream, responseBodyStream); + } + + /** + * Fires a reminder for the Actor. + * @param actorTypeName Actor type name to invoke the method for. + * @param actorId Actor id for the actor for which method will be invoked. + * @param reminderName The name of reminder provided during registration. + * @param requestBodyStream Payload for the actor method + */ + static void FireReminder(String actorTypeName, String actorId, String reminderName, byte[] requestBodyStream) + { + // uncomment when ActorManager implemented + // return instance.GetActorManager(actorTypeName).FireReminder(new ActorId(actorId), reminderName, requestBodyStream); + } + + /** + * Fires a timer for the Actor. + * @param actorTypeName Actor type name to invoke the method for. + * @param actorId Actor id for the actor for which method will be invoked. + * @param timerName The name of timer provided during registration. + */ + static void FireTimer(String actorTypeName, String actorId, String timerName) + { + // uncomment when ActorManager implemented + // return instance.GetActorManager(actorTypeName).FireTimerAsync(new ActorId(actorId), timerName); + } + + private ActorManager GetActorManager(String actorTypeName) throws IllegalStateException + { + ActorManager actorManager = this.actorManagers.get(actorTypeName); + + if (actorManager == null) + { + String errorMsg = String.format("Actor type %s is not registered with Actor runtime.", actorTypeName); + + // TODO - figure out logging. For now just print it straight out + System.out.println(errorMsg); + throw new IllegalStateException(errorMsg); + } + + return actorManager; + } +} \ No newline at end of file diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorService.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorService.java new file mode 100644 index 0000000000..110f6fc248 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorService.java @@ -0,0 +1,8 @@ +package io.dapr.actors.runtime; + +// stub +public class ActorService { + public ActorService(ActorTypeInformation actorTypeInformation) { + + } +} From 13ba95d80bc4230410cbd726466e7236d780d1a3 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 14:49:52 -0800 Subject: [PATCH 2/8] code review --- .../main/java/io/dapr/actors/ActorTrace.java | 23 +++++++++++++++++++ .../io/dapr/actors/runtime/AbstractActor.java | 2 +- .../java/io/dapr/actors/runtime/Actor.java | 2 +- .../io/dapr/actors/runtime/ActorRuntime.java | 9 ++++---- 4 files changed, 29 insertions(+), 7 deletions(-) create mode 100644 sdk/src/main/java/io/dapr/actors/ActorTrace.java diff --git a/sdk/src/main/java/io/dapr/actors/ActorTrace.java b/sdk/src/main/java/io/dapr/actors/ActorTrace.java new file mode 100644 index 0000000000..f3c7d5fea1 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/ActorTrace.java @@ -0,0 +1,23 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors; + +/** + * Stub + */ +public class ActorTrace { + public static void WriteInfo(String text) { + System.out.println(text); + } + + public static void WriteWarning(String text) { + System.out.println("Warning: " + text); + } + + public static void WriteError(String text) { + System.err.println(text); + } +} diff --git a/sdk/src/main/java/io/dapr/actors/runtime/AbstractActor.java b/sdk/src/main/java/io/dapr/actors/runtime/AbstractActor.java index 81eeaac341..1162c254f2 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/AbstractActor.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AbstractActor.java @@ -6,7 +6,7 @@ package io.dapr.actors.runtime; /** - * TODO + * TODO - this is the base class Actor implementations (user code) will extend. */ public abstract class AbstractActor { } diff --git a/sdk/src/main/java/io/dapr/actors/runtime/Actor.java b/sdk/src/main/java/io/dapr/actors/runtime/Actor.java index 25c8193252..abc4af567e 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/Actor.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/Actor.java @@ -6,7 +6,7 @@ package io.dapr.actors.runtime; /** - * TODO + * TODO - this is the interface user Actor methods should implement to receive calls. */ public interface Actor { } diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index dd1c216567..4a52600ff9 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -29,7 +29,7 @@ private ActorRuntime() throws IllegalStateException{ daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); } - /** TODO: maybe this should not be a singleton + /** * Returns an ActorRuntime object. * @return An ActorRuntime object. */ @@ -57,7 +57,7 @@ public Set getRegisteredActorTypes() { * Registers an actor with the runtime. * @param clazz The type of actor. */ - void RegisterActor(Class clazz) { + void RegisterActor(Class clazz) { RegisterActor(clazz, null); } @@ -66,7 +66,7 @@ void RegisterActor(Class clazz) { * @param clazz The type of actor. * @param actorServiceFactory An optional delegate to create actor service. This can be used for dependency injection into actors. */ - void RegisterActor(Class clazz, Function actorServiceFactory) + void RegisterActor(Class clazz, Function actorServiceFactory) { ActorTypeInformation actorTypeInfo = ActorTypeInformation.create(clazz); @@ -154,8 +154,7 @@ private ActorManager GetActorManager(String actorTypeName) throws IllegalStateEx { String errorMsg = String.format("Actor type %s is not registered with Actor runtime.", actorTypeName); - // TODO - figure out logging. For now just print it straight out - System.out.println(errorMsg); + ActorTrace.WriteError(errorMsg); throw new IllegalStateException(errorMsg); } From 4995aa35967726f8a985964236ff4480f33ff107 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 15:37:25 -0800 Subject: [PATCH 3/8] temp --- .../io/dapr/actors/ActorProxyToAppClient.java | 21 ++++ .../ActorProxyToAppHttpAsyncClient.java | 64 ++++++++++++ .../java/io/dapr/actors/DaprAsyncClient.java | 12 +-- .../java/io/dapr/actors/DaprClientBase.java | 99 +++++++++++++++++++ .../io/dapr/actors/DaprClientBuilder.java | 4 +- .../io/dapr/actors/DaprHttpAsyncClient.java | 87 +--------------- .../io/dapr/actors/runtime/ActorRuntime.java | 4 +- 7 files changed, 194 insertions(+), 97 deletions(-) create mode 100644 sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java create mode 100644 sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java create mode 100644 sdk/src/main/java/io/dapr/actors/DaprClientBase.java diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java new file mode 100644 index 0000000000..68a2215387 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java @@ -0,0 +1,21 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors; + +import reactor.core.publisher.Mono; + +public interface ActorProxyToAppClient { + + /** + * Invokes an Actor method on Dapr. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param methodName Method name to invoke. + * @param jsonPayload Serialized body. + * @return Asynchronous result with the Actor's response. + */ + Mono invokeActorMethod(String actorType, String actorId, String methodName, String jsonPayload); +} diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java new file mode 100644 index 0000000000..0e1fbce267 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java @@ -0,0 +1,64 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors; + +import com.fasterxml.jackson.databind.ObjectMapper; +import okhttp3.*; +import reactor.core.publisher.Mono; + +import java.io.IOException; +import java.net.URL; +import java.util.UUID; + +/** + * Http client to call Dapr's API for actors. + */ +public class ActorProxyToAppHttpAsyncClient extends DaprClientBase implements ActorProxyToAppClient { + + /** + * Defines the standard application/json type for HTTP calls in Dapr. + */ + private static final MediaType MEDIA_TYPE_APPLICATION_JSON = MediaType.get("application/json; charset=utf-8"); + + /** + * Shared object representing an empty request body in JSON. + */ + private static final RequestBody REQUEST_BODY_EMPTY_JSON = RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, ""); + + /** + * JSON Object Mapper. + */ + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + /** + * Base Url for calling Dapr. (e.g. http://localhost:3500/) + */ + private final String baseUrl; + + /** + * Http client used for all API calls. + */ + private final OkHttpClient httpClient; + + /** + * Creates a new instance of {@link ActorProxyToAppHttpAsyncClient}. + * @param port Port for calling Dapr. (e.g. 3500) + * @param httpClient RestClient used for all API calls in this new instance. + */ + public ActorProxyToAppHttpAsyncClient(int port, OkHttpClient httpClient) + { + this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; + this.httpClient = httpClient; + } + + /** + * {@inheritDoc} + */ + @Override + public Mono invokeActorMethod(String actorType, String actorId, String methodName, String jsonPayload) { + String url = String.format(Constants.ACTOR_METHOD_RELATIVE_URL_FORMAT, actorType, actorId, methodName); + return invokeAPI("PUT", url, jsonPayload); + } +} diff --git a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java b/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java index 47f87caeef..523550e54a 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java @@ -10,17 +10,7 @@ /** * Interface for interacting with Dapr runtime. */ -public interface DaprAsyncClient { - - /** - * Invokes an Actor method on Dapr. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param methodName Method name to invoke. - * @param jsonPayload Serialized body. - * @return Asynchronous result with the Actor's response. - */ - Mono invokeActorMethod(String actorType, String actorId, String methodName, String jsonPayload); +public interface AppToDaprAsyncClient { /** * Gets a state from Dapr's Actor. diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBase.java b/sdk/src/main/java/io/dapr/actors/DaprClientBase.java new file mode 100644 index 0000000000..f11c2b3252 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/DaprClientBase.java @@ -0,0 +1,99 @@ +package io.dapr.actors; + + + +//import com.fasterxml.jackson.databind.ObjectMapper; +//import okhttp3.*; + +import okhttp3.Request; +import okhttp3.RequestBody; +import okhttp3.Response; +import reactor.core.publisher.Mono; + +import java.io.IOException; +import java.net.URL; +import java.util.UUID; + +// base class of hierarchy +public class DaprClientBase { + + // common methods + /** + * Invokes an API asynchronously that returns Void. + * @param method HTTP method. + * @param urlString url as String. + * @param json JSON payload or null. + * @return Asynchronous Void + */ + protected final Mono invokeAPIVoid(String method, String urlString, String json) { + return this.invokeAPI(method, urlString, json).then(); + } + + /** + * Invokes an API asynchronously that returns a text payload. + * @param method HTTP method. + * @param urlString url as String. + * @param json JSON payload or null. + * @return Asynchronous text + */ + protected final Mono invokeAPI(String method, String urlString, String json) { + return Mono.fromSupplier(() -> { + try { + return tryInvokeAPI(method, urlString, json); + } catch (IOException e) { + e.printStackTrace(); + throw new RuntimeException(e); + } + }); + } + + /** + * Invokes an API synchronously and returns a text payload. + * @param method HTTP method. + * @param urlString url as String. + * @param json JSON payload or null. + * @return text + */ + protected final String tryInvokeAPI(String method, String urlString, String json) throws IOException { + String requestId = UUID.randomUUID().toString(); + RequestBody body = json != null ? RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, json) : REQUEST_BODY_EMPTY_JSON; + + Request request = new Request.Builder() + .url(new URL(this.baseUrl + urlString)) + .method(method, body) + .addHeader(Constants.HEADER_DAPR_REQUEST_ID, requestId) + .build(); + + // TODO: make this call async as well. + Response response = this.httpClient.newCall(request).execute(); + if (!response.isSuccessful()) + { + DaprError error = parseDaprError(response.body().string()); + if ((error != null) && (error.getErrorCode() != null) && (error.getMessage() != null)) { + throw new DaprException(error); + } + + throw new DaprException("UNKNOWN", String.format("Dapr's Actor API %s failed with return code %d %s", urlString, response.code())); + } + + return response.body().string(); + } + + /** + * Tries to parse an error from Dapr response body. + * @param json Response body from Dapr. + * @return DaprError or null if could not parse. + */ + protected static DaprError parseDaprError(String json) { + if (json == null) { + return null; + } + + try { + return OBJECT_MAPPER.readValue(json, DaprError.class); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } +} diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java index 418221db06..0a00917522 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java +++ b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java @@ -21,10 +21,10 @@ public class DaprClientBuilder { * Builds an async client. * @return Builds an async client. */ - public DaprAsyncClient buildAsyncClient() { + public AppToDaprAsyncClient buildAsyncClient() { OkHttpClient.Builder builder = new OkHttpClient.Builder(); // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. - return new DaprHttpAsyncClient(this.port, builder.build()); + return new AppToDaprHttpAsyncClient(this.port, builder.build()); } /** diff --git a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java index b44a3f237c..1c4c40c493 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java @@ -9,6 +9,7 @@ import okhttp3.*; import reactor.core.publisher.Mono; + import java.io.IOException; import java.net.URL; import java.util.UUID; @@ -16,7 +17,8 @@ /** * Http client to call Dapr's API for actors. */ -public class DaprHttpAsyncClient implements DaprAsyncClient { +//public class DaprHttpAsyncClient implements DaprAsyncClient { +public class AppToDaprHttpAsyncClient extends DaprClientBase implements AppToDaprAsyncClient { /** * Defines the standard application/json type for HTTP calls in Dapr. @@ -43,11 +45,11 @@ public class DaprHttpAsyncClient implements DaprAsyncClient { private final OkHttpClient httpClient; /** - * Creates a new instance of {@link DaprHttpAsyncClient}. + * Creates a new instance of {@link AppToDaprHttpAsyncClient}. * @param port Port for calling Dapr. (e.g. 3500) * @param httpClient RestClient used for all API calls in this new instance. */ - public DaprHttpAsyncClient(int port, OkHttpClient httpClient) + public AppToDaprHttpAsyncClient(int port, OkHttpClient httpClient) { this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; this.httpClient = httpClient; @@ -124,83 +126,4 @@ public Mono unregisterTimerAsync(String actorType, String actorId, String String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); return invokeAPIVoid("DELETE", url, null); } - - /** - * Invokes an API asynchronously that returns Void. - * @param method HTTP method. - * @param urlString url as String. - * @param json JSON payload or null. - * @return Asynchronous Void - */ - private final Mono invokeAPIVoid(String method, String urlString, String json) { - return this.invokeAPI(method, urlString, json).then(); - } - - /** - * Invokes an API asynchronously that returns a text payload. - * @param method HTTP method. - * @param urlString url as String. - * @param json JSON payload or null. - * @return Asynchronous text - */ - private final Mono invokeAPI(String method, String urlString, String json) { - return Mono.fromSupplier(() -> { - try { - return tryInvokeAPI(method, urlString, json); - } catch (IOException e) { - e.printStackTrace(); - throw new RuntimeException(e); - } - }); - } - - /** - * Invokes an API synchronously and returns a text payload. - * @param method HTTP method. - * @param urlString url as String. - * @param json JSON payload or null. - * @return text - */ - private final String tryInvokeAPI(String method, String urlString, String json) throws IOException { - String requestId = UUID.randomUUID().toString(); - RequestBody body = json != null ? RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, json) : REQUEST_BODY_EMPTY_JSON; - - Request request = new Request.Builder() - .url(new URL(this.baseUrl + urlString)) - .method(method, body) - .addHeader(Constants.HEADER_DAPR_REQUEST_ID, requestId) - .build(); - - // TODO: make this call async as well. - Response response = this.httpClient.newCall(request).execute(); - if (!response.isSuccessful()) - { - DaprError error = parseDaprError(response.body().string()); - if ((error != null) && (error.getErrorCode() != null) && (error.getMessage() != null)) { - throw new DaprException(error); - } - - throw new DaprException("UNKNOWN", String.format("Dapr's Actor API %s failed with return code %d %s", urlString, response.code())); - } - - return response.body().string(); - } - - /** - * Tries to parse an error from Dapr response body. - * @param json Response body from Dapr. - * @return DaprError or null if could not parse. - */ - private static DaprError parseDaprError(String json) { - if (json == null) { - return null; - } - - try { - return OBJECT_MAPPER.readValue(json, DaprError.class); - } catch (IOException e) { - e.printStackTrace(); - return null; - } - } } diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index 4a52600ff9..c7d1615ddf 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -16,7 +16,7 @@ */ public class ActorRuntime { private static volatile ActorRuntime instance; - private static DaprAsyncClient daprAsyncClient; + private static AppToDaprHttpAsyncClient appToDaprHttpAsyncClient; private final String TraceType = "ActorRuntime"; private final HashMap actorManagers; @@ -26,7 +26,7 @@ private ActorRuntime() throws IllegalStateException{ } this.actorManagers = new HashMap(); - daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); + appToDaprHttpAsyncClient = new DaprClientBuilder().buildAsyncClient(); } /** From 27d6a48b88c61cb0f9c2f28e7cfc11570a66b8d8 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 16:57:47 -0800 Subject: [PATCH 4/8] split DaprAsyncClient hierarchy into 2 hierarchies for the different directions of communication --- ...t.java => ActorProxyToAppAsyncClient.java} | 5 +- .../ActorProxyToAppHttpAsyncClient.java | 30 +-- .../main/java/io/dapr/actors/Constants.java | 2 +- .../java/io/dapr/actors/DaprClientBase.java | 39 +++- .../io/dapr/actors/DaprClientBuilder.java | 5 +- .../io/dapr/actors/runtime/ActorRuntime.java | 2 +- .../AppToDaprAsyncClient.java} | 159 ++++++------- .../AppToDaprHttpAsyncClient.java} | 221 ++++++++---------- .../actors/runtime/DaprClientBuilder.java | 59 +++++ .../io/dapr/actors/DaprHttpAsyncClientIT.java | 2 +- 10 files changed, 278 insertions(+), 246 deletions(-) rename sdk/src/main/java/io/dapr/actors/{ActorProxyToAppClient.java => ActorProxyToAppAsyncClient.java} (86%) rename sdk/src/main/java/io/dapr/actors/{DaprAsyncClient.java => runtime/AppToDaprAsyncClient.java} (92%) rename sdk/src/main/java/io/dapr/actors/{DaprHttpAsyncClient.java => runtime/AppToDaprHttpAsyncClient.java} (63%) create mode 100644 sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java similarity index 86% rename from sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java rename to sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java index 68a2215387..d79a3a8b4c 100644 --- a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppClient.java +++ b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java @@ -7,7 +7,10 @@ import reactor.core.publisher.Mono; -public interface ActorProxyToAppClient { +/** + * Interface to invoke actor methods. + */ +interface ActorProxyToAppAsyncClient { /** * Invokes an Actor method on Dapr. diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java index 0e1fbce267..07513a08e4 100644 --- a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java @@ -14,33 +14,10 @@ import java.util.UUID; /** - * Http client to call Dapr's API for actors. + * Http client to call actors methods. */ -public class ActorProxyToAppHttpAsyncClient extends DaprClientBase implements ActorProxyToAppClient { +class ActorProxyToAppHttpAsyncClient extends DaprClientBase implements ActorProxyToAppAsyncClient { - /** - * Defines the standard application/json type for HTTP calls in Dapr. - */ - private static final MediaType MEDIA_TYPE_APPLICATION_JSON = MediaType.get("application/json; charset=utf-8"); - - /** - * Shared object representing an empty request body in JSON. - */ - private static final RequestBody REQUEST_BODY_EMPTY_JSON = RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, ""); - - /** - * JSON Object Mapper. - */ - private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); - /** - * Base Url for calling Dapr. (e.g. http://localhost:3500/) - */ - private final String baseUrl; - - /** - * Http client used for all API calls. - */ - private final OkHttpClient httpClient; /** * Creates a new instance of {@link ActorProxyToAppHttpAsyncClient}. @@ -49,8 +26,7 @@ public class ActorProxyToAppHttpAsyncClient extends DaprClientBase implements Ac */ public ActorProxyToAppHttpAsyncClient(int port, OkHttpClient httpClient) { - this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; - this.httpClient = httpClient; + super(port, httpClient); } /** diff --git a/sdk/src/main/java/io/dapr/actors/Constants.java b/sdk/src/main/java/io/dapr/actors/Constants.java index 66f4dbfc3f..86ea80fac8 100644 --- a/sdk/src/main/java/io/dapr/actors/Constants.java +++ b/sdk/src/main/java/io/dapr/actors/Constants.java @@ -8,7 +8,7 @@ /** * Useful constants for the Dapr's Actor SDK. */ -final class Constants { +public final class Constants { /** * Dapr API used in this client. diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBase.java b/sdk/src/main/java/io/dapr/actors/DaprClientBase.java index f11c2b3252..375584298f 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBase.java +++ b/sdk/src/main/java/io/dapr/actors/DaprClientBase.java @@ -2,12 +2,9 @@ -//import com.fasterxml.jackson.databind.ObjectMapper; -//import okhttp3.*; +import com.fasterxml.jackson.databind.ObjectMapper; -import okhttp3.Request; -import okhttp3.RequestBody; -import okhttp3.Response; +import okhttp3.*; import reactor.core.publisher.Mono; import java.io.IOException; @@ -16,6 +13,38 @@ // base class of hierarchy public class DaprClientBase { + /** + * Defines the standard application/json type for HTTP calls in Dapr. + */ + private static final MediaType MEDIA_TYPE_APPLICATION_JSON = MediaType.get("application/json; charset=utf-8"); + + /** + * Shared object representing an empty request body in JSON. + */ + private static final RequestBody REQUEST_BODY_EMPTY_JSON = RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, ""); + + /** + * JSON Object Mapper. + */ + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + + private final String baseUrl; + + /** + * Http client used for all API calls. + */ + private final OkHttpClient httpClient; + + /** + * Creates a new instance of {@link DaprClientBase}. + * @param port Port for calling Dapr. (e.g. 3500) + * @param httpClient RestClient used for all API calls in this new instance. + */ + public DaprClientBase(int port, OkHttpClient httpClient) + { + this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; + this.httpClient = httpClient; + } // common methods /** diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java index 0a00917522..14affd3110 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java +++ b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java @@ -6,6 +6,7 @@ package io.dapr.actors; import okhttp3.OkHttpClient; +import io.dapr.actors.*; /** * Builds an instance of DaprAsyncClient or DaprClient. @@ -21,10 +22,10 @@ public class DaprClientBuilder { * Builds an async client. * @return Builds an async client. */ - public AppToDaprAsyncClient buildAsyncClient() { + public ActorProxyToAppAsyncClient buildAsyncClient() { OkHttpClient.Builder builder = new OkHttpClient.Builder(); // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. - return new AppToDaprHttpAsyncClient(this.port, builder.build()); + return new ActorProxyToAppHttpAsyncClient(this.port, builder.build()); } /** diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index c7d1615ddf..979abc3aca 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -16,7 +16,7 @@ */ public class ActorRuntime { private static volatile ActorRuntime instance; - private static AppToDaprHttpAsyncClient appToDaprHttpAsyncClient; + private static AppToDaprAsyncClient appToDaprHttpAsyncClient; private final String TraceType = "ActorRuntime"; private final HashMap actorManagers; diff --git a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java similarity index 92% rename from sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java rename to sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java index 523550e54a..f50156f6df 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java @@ -1,79 +1,80 @@ -/* - * Copyright (c) Microsoft Corporation. - * Licensed under the MIT License. - */ - -package io.dapr.actors; - -import reactor.core.publisher.Mono; - -/** - * Interface for interacting with Dapr runtime. - */ -public interface AppToDaprAsyncClient { - - /** - * Gets a state from Dapr's Actor. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param keyName State name. - * @return Asynchronous result with current state value. - */ - Mono getState(String actorType, String actorId, String keyName); - - /** - * Removes Actor state in Dapr. This is temporary until the Dapr runtime implements the Batch state update. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param keyName State name. - * @return Asynchronous void result. - */ - Mono removeState(String actorType, String actorId, String keyName); - - /** - * Saves state batch to Dapr. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param data State to be saved. - * @return Asynchronous void result. - */ - Mono saveStateTransactionally(String actorType, String actorId, String data); - - /** - * Register a reminder. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param reminderName Name of reminder to be registered. - * @param data JSON reminder data as per Dapr's spec. - * @return Asynchronous void result. - */ - Mono registerReminder(String actorType, String actorId, String reminderName, String data); - - /** - * Unregisters a reminder. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param reminderName Name of reminder to be unregistered. - * @return Asynchronous void result. - */ - Mono unregisterReminder(String actorType, String actorId, String reminderName); - - /** - * Registers a timer. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param timerName Name of timer to be registered. - * @param data JSON reminder data as per Dapr's spec. - * @return Asynchronous void result. - */ - Mono registerTimer(String actorType, String actorId, String timerName, String data); - - /** - * Unregisters a timer. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param timerName Name of timer to be unregistered. - * @return Asynchronous void result. - */ - Mono unregisterTimerAsync(String actorType, String actorId, String timerName); -} +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.runtime; + +import reactor.core.publisher.Mono; + + +/** + * Interface for interacting from the actor app to the Dapr runtime. + */ +interface AppToDaprAsyncClient { + + /** + * Gets a state from Dapr's Actor. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param keyName State name. + * @return Asynchronous result with current state value. + */ + Mono getState(String actorType, String actorId, String keyName); + + /** + * Removes Actor state in Dapr. This is temporary until the Dapr runtime implements the Batch state update. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param keyName State name. + * @return Asynchronous void result. + */ + Mono removeState(String actorType, String actorId, String keyName); + + /** + * Saves state batch to Dapr. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param data State to be saved. + * @return Asynchronous void result. + */ + Mono saveStateTransactionally(String actorType, String actorId, String data); + + /** + * Register a reminder. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param reminderName Name of reminder to be registered. + * @param data JSON reminder data as per Dapr's spec. + * @return Asynchronous void result. + */ + Mono registerReminder(String actorType, String actorId, String reminderName, String data); + + /** + * Unregisters a reminder. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param reminderName Name of reminder to be unregistered. + * @return Asynchronous void result. + */ + Mono unregisterReminder(String actorType, String actorId, String reminderName); + + /** + * Registers a timer. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param timerName Name of timer to be registered. + * @param data JSON reminder data as per Dapr's spec. + * @return Asynchronous void result. + */ + Mono registerTimer(String actorType, String actorId, String timerName, String data); + + /** + * Unregisters a timer. + * @param actorType Type of actor. + * @param actorId Actor Identifier. + * @param timerName Name of timer to be unregistered. + * @return Asynchronous void result. + */ + Mono unregisterTimerAsync(String actorType, String actorId, String timerName); +} diff --git a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java similarity index 63% rename from sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java rename to sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java index 1c4c40c493..849b1381a8 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java @@ -1,129 +1,92 @@ -/* - * Copyright (c) Microsoft Corporation. - * Licensed under the MIT License. - */ - -package io.dapr.actors; - -import com.fasterxml.jackson.databind.ObjectMapper; -import okhttp3.*; -import reactor.core.publisher.Mono; - - -import java.io.IOException; -import java.net.URL; -import java.util.UUID; - -/** - * Http client to call Dapr's API for actors. - */ -//public class DaprHttpAsyncClient implements DaprAsyncClient { -public class AppToDaprHttpAsyncClient extends DaprClientBase implements AppToDaprAsyncClient { - - /** - * Defines the standard application/json type for HTTP calls in Dapr. - */ - private static final MediaType MEDIA_TYPE_APPLICATION_JSON = MediaType.get("application/json; charset=utf-8"); - - /** - * Shared object representing an empty request body in JSON. - */ - private static final RequestBody REQUEST_BODY_EMPTY_JSON = RequestBody.create(MEDIA_TYPE_APPLICATION_JSON, ""); - - /** - * JSON Object Mapper. - */ - private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); - /** - * Base Url for calling Dapr. (e.g. http://localhost:3500/) - */ - private final String baseUrl; - - /** - * Http client used for all API calls. - */ - private final OkHttpClient httpClient; - - /** - * Creates a new instance of {@link AppToDaprHttpAsyncClient}. - * @param port Port for calling Dapr. (e.g. 3500) - * @param httpClient RestClient used for all API calls in this new instance. - */ - public AppToDaprHttpAsyncClient(int port, OkHttpClient httpClient) - { - this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; - this.httpClient = httpClient; - } - - /** - * {@inheritDoc} - */ - @Override - public Mono invokeActorMethod(String actorType, String actorId, String methodName, String jsonPayload) { - String url = String.format(Constants.ACTOR_METHOD_RELATIVE_URL_FORMAT, actorType, actorId, methodName); - return invokeAPI("PUT", url, jsonPayload); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono getState(String actorType, String actorId, String keyName) { - String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); - return invokeAPI("GET", url, null); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono removeState(String actorType, String actorId, String keyName) { - String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); - return invokeAPIVoid("DELETE", url, null); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono saveStateTransactionally(String actorType, String actorId, String data) { - String url = String.format(Constants.ACTOR_STATE_RELATIVE_URL_FORMAT, actorType, actorId); - return invokeAPIVoid("PUT", url, data); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono registerReminder(String actorType, String actorId, String reminderName, String data) { - String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); - return invokeAPIVoid("PUT", url, data); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono unregisterReminder(String actorType, String actorId, String reminderName) { - String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); - return invokeAPIVoid("DELETE", url, null); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono registerTimer(String actorType, String actorId, String timerName, String data) { - String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); - return invokeAPIVoid("PUT", url, data); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono unregisterTimerAsync(String actorType, String actorId, String timerName) { - String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); - return invokeAPIVoid("DELETE", url, null); - } -} +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.runtime; + +import io.dapr.actors.Constants; +import io.dapr.actors.DaprClientBase; +import okhttp3.*; +import reactor.core.publisher.Mono; + + +/** + * Http client to call Dapr's API for actors. + */ +//public class DaprHttpAsyncClient implements DaprAsyncClient { +class AppToDaprHttpAsyncClient extends DaprClientBase implements AppToDaprAsyncClient { + + /** + * Creates a new instance of {@link AppToDaprHttpAsyncClient}. + * @param port Port for calling Dapr. (e.g. 3500) + * @param httpClient RestClient used for all API calls in this new instance. + */ + public AppToDaprHttpAsyncClient(int port, OkHttpClient httpClient) + { + super(port, httpClient); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono getState(String actorType, String actorId, String keyName) { + String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); + return invokeAPI("GET", url, null); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono removeState(String actorType, String actorId, String keyName) { + String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); + return invokeAPIVoid("DELETE", url, null); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono saveStateTransactionally(String actorType, String actorId, String data) { + String url = String.format(Constants.ACTOR_STATE_RELATIVE_URL_FORMAT, actorType, actorId); + return invokeAPIVoid("PUT", url, data); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono registerReminder(String actorType, String actorId, String reminderName, String data) { + String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); + return invokeAPIVoid("PUT", url, data); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono unregisterReminder(String actorType, String actorId, String reminderName) { + String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); + return invokeAPIVoid("DELETE", url, null); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono registerTimer(String actorType, String actorId, String timerName, String data) { + String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); + return invokeAPIVoid("PUT", url, data); + } + + /** + * {@inheritDoc} + */ + @Override + public Mono unregisterTimerAsync(String actorType, String actorId, String timerName) { + String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); + return invokeAPIVoid("DELETE", url, null); + } +} diff --git a/sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java new file mode 100644 index 0000000000..d76b050c92 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java @@ -0,0 +1,59 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.runtime; + +import okhttp3.OkHttpClient; +import io.dapr.actors.*; + +/** + * Builds an instance of DaprAsyncClient or DaprClient. + */ +public class DaprClientBuilder { + + /** + * Default port for Dapr after checking environment variable. + */ + private int port = DaprClientBuilder.GetEnvPortOrDefault(); + + /** + * Builds an async client. + * @return Builds an async client. + */ + public AppToDaprAsyncClient buildAsyncClient() { + OkHttpClient.Builder builder = new OkHttpClient.Builder(); + // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. + return new AppToDaprHttpAsyncClient(this.port, builder.build()); + } + + /** + * Overrides the port. + * @param port New port. + * @return This instance. + */ + public DaprClientBuilder withPort(int port) { + this.port = port; + return this; + } + + /** + * Tries to get a valid port from environment variable or returns default. + * @return Port defined in env variable or default. + */ + private static int GetEnvPortOrDefault() { + String envPort = System.getenv(Constants.ENV_DAPR_HTTP_PORT); + if (envPort == null) { + return Constants.DEFAULT_PORT; + } + + try { + return Integer.parseInt(envPort.trim()); + } catch (NumberFormatException e) { + e.printStackTrace(); + } + + return Constants.DEFAULT_PORT; + } +} diff --git a/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java b/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java index 763f8e38b9..f9fa0e0966 100644 --- a/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java +++ b/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java @@ -20,7 +20,7 @@ public class DaprHttpAsyncClientIT { */ @Test(expected = RuntimeException.class) public void invokeUnknownActor() { - DaprAsyncClient daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); + ActorProxyToAppAsyncClient daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); daprAsyncClient .invokeActorMethod("ActorThatDoesNotExist", "100", "GetData", null) .doOnError(x -> { From a73fcf07e8b4d6cf45fc59fb512d8e764e06d502 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 17:10:42 -0800 Subject: [PATCH 5/8] Rename DaprClientBase to AbstractDaprClient --- .../actors/{DaprClientBase.java => AbstractDaprClient.java} | 6 +++--- .../java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java | 2 +- .../io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java | 4 ++-- 3 files changed, 6 insertions(+), 6 deletions(-) rename sdk/src/main/java/io/dapr/actors/{DaprClientBase.java => AbstractDaprClient.java} (95%) diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBase.java b/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java similarity index 95% rename from sdk/src/main/java/io/dapr/actors/DaprClientBase.java rename to sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java index 375584298f..f51983b0ae 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBase.java +++ b/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java @@ -12,7 +12,7 @@ import java.util.UUID; // base class of hierarchy -public class DaprClientBase { +public abstract class AbstractDaprClient { /** * Defines the standard application/json type for HTTP calls in Dapr. */ @@ -36,11 +36,11 @@ public class DaprClientBase { private final OkHttpClient httpClient; /** - * Creates a new instance of {@link DaprClientBase}. + * Creates a new instance of {@link AbstractDaprClient}. * @param port Port for calling Dapr. (e.g. 3500) * @param httpClient RestClient used for all API calls in this new instance. */ - public DaprClientBase(int port, OkHttpClient httpClient) + public AbstractDaprClient(int port, OkHttpClient httpClient) { this.baseUrl = String.format("http://%s:%d/", Constants.DEFAULT_HOSTNAME, port);; this.httpClient = httpClient; diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java index 07513a08e4..a7eb50c2e9 100644 --- a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java @@ -16,7 +16,7 @@ /** * Http client to call actors methods. */ -class ActorProxyToAppHttpAsyncClient extends DaprClientBase implements ActorProxyToAppAsyncClient { +class ActorProxyToAppHttpAsyncClient extends AbstractDaprClient implements ActorProxyToAppAsyncClient { /** diff --git a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java index 849b1381a8..f3028ef257 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java @@ -6,7 +6,7 @@ package io.dapr.actors.runtime; import io.dapr.actors.Constants; -import io.dapr.actors.DaprClientBase; +import io.dapr.actors.AbstractDaprClient; import okhttp3.*; import reactor.core.publisher.Mono; @@ -15,7 +15,7 @@ * Http client to call Dapr's API for actors. */ //public class DaprHttpAsyncClient implements DaprAsyncClient { -class AppToDaprHttpAsyncClient extends DaprClientBase implements AppToDaprAsyncClient { +class AppToDaprHttpAsyncClient extends AbstractDaprClient implements AppToDaprAsyncClient { /** * Creates a new instance of {@link AppToDaprHttpAsyncClient}. From cf5f25dfefaeed546b64bf379a8002f06aa0cd9d Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 18:17:55 -0800 Subject: [PATCH 6/8] code review --- ...uilder.java => AbstractClientBuilder.java} | 22 ++--- .../io/dapr/actors/DaprClientBuilder.java | 59 ------------- .../java/io/dapr/actors/DaprException.java | 4 +- .../ActorProxyAsyncClient.java} | 4 +- .../client/ActorProxyClientBuilder.java | 30 +++++++ .../ActorProxyHttpAsyncClient.java} | 17 ++-- .../io/dapr/actors/runtime/ActorRuntime.java | 6 +- .../actors/runtime/AppToDaprAsyncClient.java | 9 -- .../runtime/AppToDaprClientBuilder.java | 30 +++++++ .../runtime/AppToDaprHttpAsyncClient.java | 21 ++--- .../{ => client}/DaprHttpAsyncClientIT.java | 85 ++++++++++--------- 11 files changed, 129 insertions(+), 158 deletions(-) rename sdk/src/main/java/io/dapr/actors/{runtime/DaprClientBuilder.java => AbstractClientBuilder.java} (57%) delete mode 100644 sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java rename sdk/src/main/java/io/dapr/actors/{ActorProxyToAppAsyncClient.java => client/ActorProxyAsyncClient.java} (89%) create mode 100644 sdk/src/main/java/io/dapr/actors/client/ActorProxyClientBuilder.java rename sdk/src/main/java/io/dapr/actors/{ActorProxyToAppHttpAsyncClient.java => client/ActorProxyHttpAsyncClient.java} (59%) create mode 100644 sdk/src/main/java/io/dapr/actors/runtime/AppToDaprClientBuilder.java rename sdk/src/test/java/io/dapr/actors/{ => client}/DaprHttpAsyncClientIT.java (89%) diff --git a/sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/AbstractClientBuilder.java similarity index 57% rename from sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java rename to sdk/src/main/java/io/dapr/actors/AbstractClientBuilder.java index d76b050c92..5af3cc78ac 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/DaprClientBuilder.java +++ b/sdk/src/main/java/io/dapr/actors/AbstractClientBuilder.java @@ -3,37 +3,27 @@ * Licensed under the MIT License. */ -package io.dapr.actors.runtime; +package io.dapr.actors; import okhttp3.OkHttpClient; import io.dapr.actors.*; /** - * Builds an instance of DaprAsyncClient or DaprClient. + * Base class for client builders */ -public class DaprClientBuilder { +public abstract class AbstractClientBuilder { /** * Default port for Dapr after checking environment variable. */ - private int port = DaprClientBuilder.GetEnvPortOrDefault(); - - /** - * Builds an async client. - * @return Builds an async client. - */ - public AppToDaprAsyncClient buildAsyncClient() { - OkHttpClient.Builder builder = new OkHttpClient.Builder(); - // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. - return new AppToDaprHttpAsyncClient(this.port, builder.build()); - } + private int port = AbstractClientBuilder.GetEnvPortOrDefault(); /** * Overrides the port. * @param port New port. * @return This instance. */ - public DaprClientBuilder withPort(int port) { + public AbstractClientBuilder withPort(int port) { this.port = port; return this; } @@ -42,7 +32,7 @@ public DaprClientBuilder withPort(int port) { * Tries to get a valid port from environment variable or returns default. * @return Port defined in env variable or default. */ - private static int GetEnvPortOrDefault() { + protected static int GetEnvPortOrDefault() { String envPort = System.getenv(Constants.ENV_DAPR_HTTP_PORT); if (envPort == null) { return Constants.DEFAULT_PORT; diff --git a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java deleted file mode 100644 index 14affd3110..0000000000 --- a/sdk/src/main/java/io/dapr/actors/DaprClientBuilder.java +++ /dev/null @@ -1,59 +0,0 @@ -/* - * Copyright (c) Microsoft Corporation. - * Licensed under the MIT License. - */ - -package io.dapr.actors; - -import okhttp3.OkHttpClient; -import io.dapr.actors.*; - -/** - * Builds an instance of DaprAsyncClient or DaprClient. - */ -public class DaprClientBuilder { - - /** - * Default port for Dapr after checking environment variable. - */ - private int port = DaprClientBuilder.GetEnvPortOrDefault(); - - /** - * Builds an async client. - * @return Builds an async client. - */ - public ActorProxyToAppAsyncClient buildAsyncClient() { - OkHttpClient.Builder builder = new OkHttpClient.Builder(); - // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. - return new ActorProxyToAppHttpAsyncClient(this.port, builder.build()); - } - - /** - * Overrides the port. - * @param port New port. - * @return This instance. - */ - public DaprClientBuilder withPort(int port) { - this.port = port; - return this; - } - - /** - * Tries to get a valid port from environment variable or returns default. - * @return Port defined in env variable or default. - */ - private static int GetEnvPortOrDefault() { - String envPort = System.getenv(Constants.ENV_DAPR_HTTP_PORT); - if (envPort == null) { - return Constants.DEFAULT_PORT; - } - - try { - return Integer.parseInt(envPort.trim()); - } catch (NumberFormatException e) { - e.printStackTrace(); - } - - return Constants.DEFAULT_PORT; - } -} diff --git a/sdk/src/main/java/io/dapr/actors/DaprException.java b/sdk/src/main/java/io/dapr/actors/DaprException.java index 2c7b8348f7..50ad94a9b9 100644 --- a/sdk/src/main/java/io/dapr/actors/DaprException.java +++ b/sdk/src/main/java/io/dapr/actors/DaprException.java @@ -10,7 +10,7 @@ /** * A Dapr's specific exception. */ -class DaprException extends IOException { +public class DaprException extends IOException { /** * Dapr's error code for this exception. @@ -39,7 +39,7 @@ class DaprException extends IOException { * Returns the exception's error code. * @return Error code. */ - String getErrorCode() { + public String getErrorCode() { return this.errorCode; } } diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java b/sdk/src/main/java/io/dapr/actors/client/ActorProxyAsyncClient.java similarity index 89% rename from sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java rename to sdk/src/main/java/io/dapr/actors/client/ActorProxyAsyncClient.java index d79a3a8b4c..b4f6b0bfda 100644 --- a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/client/ActorProxyAsyncClient.java @@ -3,14 +3,14 @@ * Licensed under the MIT License. */ -package io.dapr.actors; +package io.dapr.actors.client; import reactor.core.publisher.Mono; /** * Interface to invoke actor methods. */ -interface ActorProxyToAppAsyncClient { +interface ActorProxyAsyncClient { /** * Invokes an Actor method on Dapr. diff --git a/sdk/src/main/java/io/dapr/actors/client/ActorProxyClientBuilder.java b/sdk/src/main/java/io/dapr/actors/client/ActorProxyClientBuilder.java new file mode 100644 index 0000000000..4f7d88fb9f --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/client/ActorProxyClientBuilder.java @@ -0,0 +1,30 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.client; + +import okhttp3.OkHttpClient; +import io.dapr.actors.*; + +/** + * Builds an instance of ActorProxyAsyncClient. + */ +class ActorProxyClientBuilder extends AbstractClientBuilder { + + /** + * Default port for Dapr after checking environment variable. + */ + private int port = ActorProxyClientBuilder.GetEnvPortOrDefault(); + + /** + * Builds an async client. + * @return Builds an async client. + */ + public ActorProxyAsyncClient buildAsyncClient() { + OkHttpClient.Builder builder = new OkHttpClient.Builder(); + // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. + return new ActorProxyHttpAsyncClient(this.port, builder.build()); + } +} diff --git a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java similarity index 59% rename from sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java rename to sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java index a7eb50c2e9..1b1923e920 100644 --- a/sdk/src/main/java/io/dapr/actors/ActorProxyToAppHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java @@ -3,28 +3,25 @@ * Licensed under the MIT License. */ -package io.dapr.actors; +package io.dapr.actors.client; -import com.fasterxml.jackson.databind.ObjectMapper; +import io.dapr.actors.*; +import io.dapr.actors.AbstractDaprClient; import okhttp3.*; import reactor.core.publisher.Mono; -import java.io.IOException; -import java.net.URL; -import java.util.UUID; - /** * Http client to call actors methods. */ -class ActorProxyToAppHttpAsyncClient extends AbstractDaprClient implements ActorProxyToAppAsyncClient { +class ActorProxyHttpAsyncClient extends AbstractDaprClient implements ActorProxyAsyncClient { /** - * Creates a new instance of {@link ActorProxyToAppHttpAsyncClient}. + * Creates a new instance of {@link ActorProxyHttpAsyncClient}. * @param port Port for calling Dapr. (e.g. 3500) * @param httpClient RestClient used for all API calls in this new instance. */ - public ActorProxyToAppHttpAsyncClient(int port, OkHttpClient httpClient) + public ActorProxyHttpAsyncClient(int port, OkHttpClient httpClient) { super(port, httpClient); } @@ -35,6 +32,6 @@ public ActorProxyToAppHttpAsyncClient(int port, OkHttpClient httpClient) @Override public Mono invokeActorMethod(String actorType, String actorId, String methodName, String jsonPayload) { String url = String.format(Constants.ACTOR_METHOD_RELATIVE_URL_FORMAT, actorType, actorId, methodName); - return invokeAPI("PUT", url, jsonPayload); + return super.invokeAPI("PUT", url, jsonPayload); } } diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index 979abc3aca..af758b7cd8 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -26,7 +26,7 @@ private ActorRuntime() throws IllegalStateException{ } this.actorManagers = new HashMap(); - appToDaprHttpAsyncClient = new DaprClientBuilder().buildAsyncClient(); + appToDaprHttpAsyncClient = new AppToDaprClientBuilder().buildAsyncClient(); } /** @@ -57,7 +57,7 @@ public Set getRegisteredActorTypes() { * Registers an actor with the runtime. * @param clazz The type of actor. */ - void RegisterActor(Class clazz) { + public void RegisterActor(Class clazz) { RegisterActor(clazz, null); } @@ -66,7 +66,7 @@ void RegisterActor(Class clazz) { * @param clazz The type of actor. * @param actorServiceFactory An optional delegate to create actor service. This can be used for dependency injection into actors. */ - void RegisterActor(Class clazz, Function actorServiceFactory) + public void RegisterActor(Class clazz, Function actorServiceFactory) { ActorTypeInformation actorTypeInfo = ActorTypeInformation.create(clazz); diff --git a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java index f50156f6df..66fe3d1757 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprAsyncClient.java @@ -22,15 +22,6 @@ interface AppToDaprAsyncClient { */ Mono getState(String actorType, String actorId, String keyName); - /** - * Removes Actor state in Dapr. This is temporary until the Dapr runtime implements the Batch state update. - * @param actorType Type of actor. - * @param actorId Actor Identifier. - * @param keyName State name. - * @return Asynchronous void result. - */ - Mono removeState(String actorType, String actorId, String keyName); - /** * Saves state batch to Dapr. * @param actorType Type of actor. diff --git a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprClientBuilder.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprClientBuilder.java new file mode 100644 index 0000000000..1476057ca3 --- /dev/null +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprClientBuilder.java @@ -0,0 +1,30 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.runtime; + +import okhttp3.OkHttpClient; +import io.dapr.actors.*; + +/** + * Builds an instance of AppToDaprAsyncClient. + */ +class AppToDaprClientBuilder extends AbstractClientBuilder { + + /** + * Default port for Dapr after checking environment variable. + */ + private int port = AppToDaprClientBuilder.GetEnvPortOrDefault(); + + /** + * Builds an async client. + * @return Builds an async client. + */ + public AppToDaprAsyncClient buildAsyncClient() { + OkHttpClient.Builder builder = new OkHttpClient.Builder(); + // TODO: Expose configurations for OkHttpClient or com.microsoft.rest.RestClient. + return new AppToDaprHttpAsyncClient(this.port, builder.build()); + } +} diff --git a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java index f3028ef257..6c4041042b 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/AppToDaprHttpAsyncClient.java @@ -33,16 +33,7 @@ public AppToDaprHttpAsyncClient(int port, OkHttpClient httpClient) @Override public Mono getState(String actorType, String actorId, String keyName) { String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); - return invokeAPI("GET", url, null); - } - - /** - * {@inheritDoc} - */ - @Override - public Mono removeState(String actorType, String actorId, String keyName) { - String url = String.format(Constants.ACTOR_STATE_KEY_RELATIVE_URL_FORMAT, actorType, actorId, keyName); - return invokeAPIVoid("DELETE", url, null); + return super.invokeAPI("GET", url, null); } /** @@ -51,7 +42,7 @@ public Mono removeState(String actorType, String actorId, String keyName) @Override public Mono saveStateTransactionally(String actorType, String actorId, String data) { String url = String.format(Constants.ACTOR_STATE_RELATIVE_URL_FORMAT, actorType, actorId); - return invokeAPIVoid("PUT", url, data); + return super.invokeAPIVoid("PUT", url, data); } /** @@ -60,7 +51,7 @@ public Mono saveStateTransactionally(String actorType, String actorId, Str @Override public Mono registerReminder(String actorType, String actorId, String reminderName, String data) { String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); - return invokeAPIVoid("PUT", url, data); + return super.invokeAPIVoid("PUT", url, data); } /** @@ -69,7 +60,7 @@ public Mono registerReminder(String actorType, String actorId, String remi @Override public Mono unregisterReminder(String actorType, String actorId, String reminderName) { String url = String.format(Constants.ACTOR_REMINDER_RELATIVE_URL_FORMAT, actorType, actorId, reminderName); - return invokeAPIVoid("DELETE", url, null); + return super.invokeAPIVoid("DELETE", url, null); } /** @@ -78,7 +69,7 @@ public Mono unregisterReminder(String actorType, String actorId, String re @Override public Mono registerTimer(String actorType, String actorId, String timerName, String data) { String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); - return invokeAPIVoid("PUT", url, data); + return super.invokeAPIVoid("PUT", url, data); } /** @@ -87,6 +78,6 @@ public Mono registerTimer(String actorType, String actorId, String timerNa @Override public Mono unregisterTimerAsync(String actorType, String actorId, String timerName) { String url = String.format(Constants.ACTOR_TIMER_RELATIVE_URL_FORMAT, actorType, actorId, timerName); - return invokeAPIVoid("DELETE", url, null); + return super.invokeAPIVoid("DELETE", url, null); } } diff --git a/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java b/sdk/src/test/java/io/dapr/actors/client/DaprHttpAsyncClientIT.java similarity index 89% rename from sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java rename to sdk/src/test/java/io/dapr/actors/client/DaprHttpAsyncClientIT.java index f9fa0e0966..3096f13a89 100644 --- a/sdk/src/test/java/io/dapr/actors/DaprHttpAsyncClientIT.java +++ b/sdk/src/test/java/io/dapr/actors/client/DaprHttpAsyncClientIT.java @@ -1,42 +1,43 @@ -/* - * Copyright (c) Microsoft Corporation. - * Licensed under the MIT License. - */ - -package io.dapr.actors; - -import org.junit.Assert; -import org.junit.Test; - -/** - * Integration test for the HTTP Async Client. - * - * Requires Dapr running. - */ -public class DaprHttpAsyncClientIT { - - /** - * Checks if the error is correctly parsed when trying to invoke a function on an unknown actor type. - */ - @Test(expected = RuntimeException.class) - public void invokeUnknownActor() { - ActorProxyToAppAsyncClient daprAsyncClient = new DaprClientBuilder().buildAsyncClient(); - daprAsyncClient - .invokeActorMethod("ActorThatDoesNotExist", "100", "GetData", null) - .doOnError(x -> { - Assert.assertTrue(x instanceof RuntimeException); - RuntimeException runtimeException = (RuntimeException)x; - - Throwable cause = runtimeException.getCause(); - Assert.assertTrue(cause instanceof DaprException); - DaprException daprException = (DaprException)cause; - - Assert.assertNotNull(daprException); - Assert.assertEquals("ERR_INVOKE_ACTOR", daprException.getErrorCode()); - Assert.assertNotNull(daprException.getMessage()); - Assert.assertFalse(daprException.getMessage().isEmpty()); - }) - .doOnSuccess(x -> Assert.fail("This call should fail.")) - .block(); - } -} +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT License. + */ + +package io.dapr.actors.client; + +import io.dapr.actors.*; +import org.junit.Assert; +import org.junit.Test; + +/** + * Integration test for the HTTP Async Client. + * + * Requires Dapr running. + */ +public class DaprHttpAsyncClientIT { + + /** + * Checks if the error is correctly parsed when trying to invoke a function on an unknown actor type. + */ + @Test(expected = RuntimeException.class) + public void invokeUnknownActor() { + ActorProxyAsyncClient daprAsyncClient = new ActorProxyClientBuilder().buildAsyncClient(); + daprAsyncClient + .invokeActorMethod("ActorThatDoesNotExist", "100", "GetData", null) + .doOnError(x -> { + Assert.assertTrue(x instanceof RuntimeException); + RuntimeException runtimeException = (RuntimeException)x; + + Throwable cause = runtimeException.getCause(); + Assert.assertTrue(cause instanceof DaprException); + DaprException daprException = (DaprException)cause; + + Assert.assertNotNull(daprException); + Assert.assertEquals("ERR_INVOKE_ACTOR", daprException.getErrorCode()); + Assert.assertNotNull(daprException.getMessage()); + Assert.assertFalse(daprException.getMessage().isEmpty()); + }) + .doOnSuccess(x -> Assert.fail("This call should fail.")) + .block(); + } +} From d7b23c41399988317ab5b97c7ae2344a864e1f40 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 19:16:57 -0800 Subject: [PATCH 7/8] more code review --- .../client/ActorProxyHttpAsyncClient.java | 2 -- .../io/dapr/actors/runtime/ActorRuntime.java | 35 +++++++++++++++---- 2 files changed, 29 insertions(+), 8 deletions(-) diff --git a/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java b/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java index 1b1923e920..2b44b35eab 100644 --- a/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java +++ b/sdk/src/main/java/io/dapr/actors/client/ActorProxyHttpAsyncClient.java @@ -14,8 +14,6 @@ * Http client to call actors methods. */ class ActorProxyHttpAsyncClient extends AbstractDaprClient implements ActorProxyAsyncClient { - - /** * Creates a new instance of {@link ActorProxyHttpAsyncClient}. * @param port Port for calling Dapr. (e.g. 3500) diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index af758b7cd8..84fb2c5f0e 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -5,6 +5,8 @@ package io.dapr.actors.runtime; +import java.util.Collection; +import java.util.Collections; import java.util.HashMap; import java.util.Set; import java.util.function.Function; @@ -15,18 +17,37 @@ * Contains methods to register actor types. Registering the types allows the runtime to create instances of the actor. */ public class ActorRuntime { + /** + * Gets an instance to the ActorRuntime. There is only 1. + */ private static volatile ActorRuntime instance; - private static AppToDaprAsyncClient appToDaprHttpAsyncClient; - private final String TraceType = "ActorRuntime"; + + /** + * A client used to communicate from the actor to the Dapr runtime. + */ + private static AppToDaprAsyncClient appToDaprAsyncClient; + + /** + * A trace type used when logging. + */ + private static final String TraceType = "ActorRuntime"; + + /** + * Map of ActorType --> ActorManager. + */ private final HashMap actorManagers; + /** + * The default constructor. This should not be called directly. + * @throws IllegalStateException + */ private ActorRuntime() throws IllegalStateException{ if (instance != null) { throw new IllegalStateException("ActorRuntime should only be constructed once"); } this.actorManagers = new HashMap(); - appToDaprHttpAsyncClient = new AppToDaprClientBuilder().buildAsyncClient(); + appToDaprAsyncClient = new AppToDaprClientBuilder().buildAsyncClient(); } /** @@ -49,8 +70,8 @@ public static ActorRuntime getInstance() { * * @return Actor type names registered with the runtime. */ - public Set getRegisteredActorTypes() { - return this.actorManagers.keySet(); + public Collection getRegisteredActorTypes() { + return Collections.unmodifiableCollection(this.actorManagers.keySet()); } /** @@ -81,7 +102,9 @@ public void RegisterActor(Class clazz, Fu } // Create ActorManagers, override existing entry if registered again. - this.actorManagers.put(actorTypeInfo.getName(), new ActorManager(actorService)); + synchronized (this.actorManagers) { + this.actorManagers.put(actorTypeInfo.getName(), new ActorManager(actorService)); + } } /** From 8a08e7feb2bab2344ff066e8cec8ec954f74cca0 Mon Sep 17 00:00:00 2001 From: LM Date: Thu, 12 Dec 2019 19:25:15 -0800 Subject: [PATCH 8/8] more cr --- sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java | 3 +++ sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java | 4 ++-- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java b/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java index f51983b0ae..2be14579b2 100644 --- a/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java +++ b/sdk/src/main/java/io/dapr/actors/AbstractDaprClient.java @@ -28,6 +28,9 @@ public abstract class AbstractDaprClient { */ private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + /** + * The base url used for form urls. This is typically "http://localhost:3500". + */ private final String baseUrl; /** diff --git a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java index 84fb2c5f0e..6574cb8658 100644 --- a/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java +++ b/sdk/src/main/java/io/dapr/actors/runtime/ActorRuntime.java @@ -78,7 +78,7 @@ public Collection getRegisteredActorTypes() { * Registers an actor with the runtime. * @param clazz The type of actor. */ - public void RegisterActor(Class clazz) { + public void RegisterActor(Class clazz) { RegisterActor(clazz, null); } @@ -87,7 +87,7 @@ public void RegisterActor(Class clazz) { * @param clazz The type of actor. * @param actorServiceFactory An optional delegate to create actor service. This can be used for dependency injection into actors. */ - public void RegisterActor(Class clazz, Function actorServiceFactory) + public void RegisterActor(Class clazz, Function actorServiceFactory) { ActorTypeInformation actorTypeInfo = ActorTypeInformation.create(clazz);