From f4c9b02908df4bb88060ec27b2922bea132ae9b2 Mon Sep 17 00:00:00 2001 From: Qing Wang Date: Wed, 21 Oct 2020 21:38:06 +0800 Subject: [PATCH] FIRST INITIAL --- .../org/restrpc/client/AsyncRpcFunction.java | 2 +- .../restrpc/client/AsyncRpcFunctionImpl.java | 8 ++--- .../main/java/org/restrpc/client/Codec.java | 25 ++++++------- .../org/restrpc/client/NativeRpcClient.java | 36 ++++++++++--------- .../java/org/restrpc/client/RestFuture.java | 2 ++ .../java/org/restrpc/client/RpcClient.java | 2 +- .../org/restrpc/test/BasicClientTest.java | 10 +++--- 7 files changed, 45 insertions(+), 40 deletions(-) diff --git a/java/src/main/java/org/restrpc/client/AsyncRpcFunction.java b/java/src/main/java/org/restrpc/client/AsyncRpcFunction.java index 4c3f314..45baeaf 100644 --- a/java/src/main/java/org/restrpc/client/AsyncRpcFunction.java +++ b/java/src/main/java/org/restrpc/client/AsyncRpcFunction.java @@ -10,7 +10,7 @@ public interface AsyncRpcFunction { // CompletableFuture invoke(Arg1Type arg1); - RestFuture invoke(Arg1Type arg1, Arg2Type arg2); + CompletableFuture invoke(Class returnClz, Arg1Type arg1, Arg2Type arg2); // // CompletableFuture invoke(Arg1Type arg1, Arg2Type arg2, Arg3Type arg3); diff --git a/java/src/main/java/org/restrpc/client/AsyncRpcFunctionImpl.java b/java/src/main/java/org/restrpc/client/AsyncRpcFunctionImpl.java index b5e1f17..e71924a 100644 --- a/java/src/main/java/org/restrpc/client/AsyncRpcFunctionImpl.java +++ b/java/src/main/java/org/restrpc/client/AsyncRpcFunctionImpl.java @@ -23,9 +23,9 @@ public class AsyncRpcFunctionImpl implements AsyncRpcFunction { // } public - RestFuture invoke(Arg1Type arg1, Arg2Type arg2) { + CompletableFuture invoke(Class returnClz, Arg1Type arg1, Arg2Type arg2) { Object[] args = new Object[] {arg1, arg2}; - return internalInvoke(args); + return internalInvoke(returnClz, args); } // public @@ -53,7 +53,7 @@ public class AsyncRpcFunctionImpl implements AsyncRpcFunction { // return internalInvoke(args); // } - private RestFuture internalInvoke(Object[] args) { - return rpcClient.invoke(funcName, args); + private CompletableFuture internalInvoke(Class returnClz, Object[] args) { + return rpcClient.invoke(returnClz, funcName, args); } } diff --git a/java/src/main/java/org/restrpc/client/Codec.java b/java/src/main/java/org/restrpc/client/Codec.java index 4213e94..f19ead3 100644 --- a/java/src/main/java/org/restrpc/client/Codec.java +++ b/java/src/main/java/org/restrpc/client/Codec.java @@ -46,8 +46,8 @@ public class Codec { return messagePacker.toByteArray(); } - public Object decodeSingleValue(String targetTypeName, byte[] encodedBytes) throws IOException { - if (targetTypeName == null) { + public Object decodeReturnValue(Class returnClz, byte[] encodedBytes) throws IOException { + if (returnClz == null) { throw new RuntimeException("Internal bug."); } @@ -57,17 +57,18 @@ public class Codec { // TODO(qwang): unpack nil. MessageUnpacker messageUnpacker = MessagePack.newDefaultUnpacker(encodedBytes); - switch (targetTypeName) { - case INT_TYPE_NAME: - return messageUnpacker.unpackInt(); - case LONG_TYPE_NAME: - return messageUnpacker.unpackLong(); - case STRING_TYPE_NAME: - return messageUnpacker.unpackString(); - default: - throw new RuntimeException("Unknown type: " + targetTypeName); - } + // Unpack unnecessary fields. + messageUnpacker.unpackArrayHeader(); + messageUnpacker.unpackInt(); + if (Integer.class.equals(returnClz)) { + return messageUnpacker.unpackInt(); + } else if (Long.class.equals(returnClz)) { + return messageUnpacker.unpackLong(); + } else if (String.class.equals(returnClz)) { + return messageUnpacker.unpackString(); + } + throw new RuntimeException("Unknown type: " + returnClz); } public int myDecodeInt(byte[] encodedBytes) throws IOException { diff --git a/java/src/main/java/org/restrpc/client/NativeRpcClient.java b/java/src/main/java/org/restrpc/client/NativeRpcClient.java index 69836cb..bda6bab 100644 --- a/java/src/main/java/org/restrpc/client/NativeRpcClient.java +++ b/java/src/main/java/org/restrpc/client/NativeRpcClient.java @@ -1,7 +1,7 @@ package org.restrpc.client; + import java.io.IOException; -import java.util.Objects; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; @@ -16,9 +16,9 @@ public class NativeRpcClient implements RpcClient { private Codec codec; // The map to cache return type. - private ConcurrentHashMap localFutureReturnTypenameCache = new ConcurrentHashMap<>(); + private ConcurrentHashMap> localFutureReturnTypenameCache = new ConcurrentHashMap<>(); - private ConcurrentHashMap> localFutureCache = new ConcurrentHashMap<>(); + private ConcurrentHashMap> localFutureCache = new ConcurrentHashMap<>(); public NativeRpcClient() { rpcClientPointer = nativeNewRpcClient(); @@ -39,7 +39,7 @@ public class NativeRpcClient implements RpcClient { return new AsyncRpcFunctionImpl(this, funcName); } - public RestFuture invoke(String funcName, Object[] args) { + public CompletableFuture invoke(Class returnClz, String funcName, Object[] args) { if (rpcClientPointer == -1) { throw new RuntimeException("no init"); } @@ -55,10 +55,14 @@ public class NativeRpcClient implements RpcClient { return null; } - final long requestId = nativeInvoke(rpcClientPointer, encodedBytes); - RestFuture futureToReturn = new RestFuture() {}; - localFutureCache.put(requestId, futureToReturn); - return futureToReturn; + synchronized (this) { + final long requestId = nativeInvoke(rpcClientPointer, encodedBytes); + CompletableFuture futureToReturn = new CompletableFuture(); + localFutureReturnTypenameCache.put(requestId, returnClz); + + localFutureCache.put(requestId, futureToReturn); + return futureToReturn; + } } public void close() { @@ -76,15 +80,13 @@ public class NativeRpcClient implements RpcClient { */ private void onResultReceived(long requestId, byte[] encodedReturnValueBytes) throws IOException { // if (requestId not is local_cache) {//error} -// codec.decodeReturnValue(encodedReturnValueBytes); - System.out.println("-----------oo java-----------"); - System.out.println("result is " + codec.myDecodeInt(encodedReturnValueBytes)); - System.out.println("-----------end java-----------"); - - RestFuture future = localFutureCache.get(requestId); - future.getClass().getGenericSuperclass().getTypeName(); - Object o = codec.decodeSingleValue( - future.getClass().getGenericSuperclass().getTypeName(), encodedReturnValueBytes); +// System.out.println("result is " + codec.myDecodeInt(encodedReturnValueBytes)); + synchronized (this) { + final Class returnClz = localFutureReturnTypenameCache.get(requestId); + CompletableFuture future = localFutureCache.get(requestId); + Object o = codec.decodeReturnValue(returnClz, encodedReturnValueBytes); + future.complete(o); + } } private native long nativeNewRpcClient(); diff --git a/java/src/main/java/org/restrpc/client/RestFuture.java b/java/src/main/java/org/restrpc/client/RestFuture.java index c8e0517..f0ecbb6 100644 --- a/java/src/main/java/org/restrpc/client/RestFuture.java +++ b/java/src/main/java/org/restrpc/client/RestFuture.java @@ -2,6 +2,7 @@ package org.restrpc.client; import java.io.IOException; +import java.util.concurrent.CompletableFuture; public class RestFuture { @@ -12,6 +13,7 @@ public class RestFuture { public RestFuture(Class metaType) { this.metaType = metaType; } + private byte[] encodedData; private static Codec codec = new Codec(); diff --git a/java/src/main/java/org/restrpc/client/RpcClient.java b/java/src/main/java/org/restrpc/client/RpcClient.java index 36903ad..d8235b8 100644 --- a/java/src/main/java/org/restrpc/client/RpcClient.java +++ b/java/src/main/java/org/restrpc/client/RpcClient.java @@ -8,7 +8,7 @@ public interface RpcClient { AsyncRpcFunction asyncFunc(String funcName); - RestFuture invoke(String funcName, Object[] args); + CompletableFuture invoke(Class returnClz, String funcName, Object[] args); void close(); } diff --git a/java/src/test/java/org/restrpc/test/BasicClientTest.java b/java/src/test/java/org/restrpc/test/BasicClientTest.java index 38f9335..60f42ff 100644 --- a/java/src/test/java/org/restrpc/test/BasicClientTest.java +++ b/java/src/test/java/org/restrpc/test/BasicClientTest.java @@ -12,17 +12,17 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.OutputStream; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; public class BasicClientTest { @Test - public void testBasic() throws IOException, InterruptedException { + public void testBasic() throws IOException, InterruptedException, ExecutionException { RpcClient rpcClient = new NativeRpcClient(); rpcClient.connect("127.0.0.1:9000"); - RestFuture obj = rpcClient.asyncFunc("add").invoke(2, 3); - TimeUnit.SECONDS.sleep(10000); - System.out.println("The result of add(2, 3) is " + obj.get()); + CompletableFuture future = rpcClient.asyncFunc("echo").invoke(String.class, 2, 3); + System.out.println("The result of add(2, 3) is " + future.get()); } private String getClassStr(Object o) { @@ -50,7 +50,7 @@ public class BasicClientTest { @Test public void testDecode() throws IOException { - CompletableFuture i = new CompletableFuture() {}; +// CompletableFuture i = new CompletableFuture() {}; // System.out.println(i.getClass().getDeclare); // // OutputStream os = new ByteArrayOutputStream();