From cae35780ea0a4f0a88f0a44b559448a6aecc1275 Mon Sep 17 00:00:00 2001 From: Rupert Griffin Date: Thu, 25 Jun 2026 12:55:26 +0000 Subject: [PATCH] reorganized gRPC channel logic --- .../emissary/grpc/GrpcConnectionPlace.java | 4 +- .../java/emissary/grpc/GrpcRoutingPlace.java | 74 ++++++++---------- .../emissary/grpc/channel/ChannelManager.java | 46 ++++++++++++ .../ChannelPoolFactory.java} | 43 +++++++++-- .../grpc/exceptions/GrpcExceptionUtils.java | 2 +- .../emissary/grpc/invoker/GrpcInvoker.java | 45 ++++++----- .../grpc/pool/LoadBalancingPolicy.java | 11 --- .../emissary/grpc/pool/PoolException.java | 17 ----- .../grpc/pool/PoolRetrievalOrdering.java | 5 -- .../emissary/grpc/sample/GrpcSamplePlace.java | 4 +- .../ChannelPoolFactoryTest.java} | 75 ++++++++++--------- 11 files changed, 184 insertions(+), 142 deletions(-) create mode 100644 src/main/java/emissary/grpc/channel/ChannelManager.java rename src/main/java/emissary/grpc/{pool/ConnectionFactory.java => channel/ChannelPoolFactory.java} (92%) delete mode 100644 src/main/java/emissary/grpc/pool/LoadBalancingPolicy.java delete mode 100644 src/main/java/emissary/grpc/pool/PoolException.java delete mode 100644 src/main/java/emissary/grpc/pool/PoolRetrievalOrdering.java rename src/test/java/emissary/grpc/{pool/ConnectionFactoryTest.java => channel/ChannelPoolFactoryTest.java} (63%) diff --git a/src/main/java/emissary/grpc/GrpcConnectionPlace.java b/src/main/java/emissary/grpc/GrpcConnectionPlace.java index 6e09d5ae28..15582f28f1 100644 --- a/src/main/java/emissary/grpc/GrpcConnectionPlace.java +++ b/src/main/java/emissary/grpc/GrpcConnectionPlace.java @@ -1,7 +1,7 @@ package emissary.grpc; import emissary.config.Configurator; -import emissary.grpc.pool.ConnectionFactory; +import emissary.grpc.channel.ChannelManager; import emissary.grpc.retry.RetryHandler; import com.google.common.util.concurrent.ListenableFuture; @@ -27,7 +27,7 @@ * */ diff --git a/src/main/java/emissary/grpc/GrpcRoutingPlace.java b/src/main/java/emissary/grpc/GrpcRoutingPlace.java index 6ebdda8ce0..9993e70804 100644 --- a/src/main/java/emissary/grpc/GrpcRoutingPlace.java +++ b/src/main/java/emissary/grpc/GrpcRoutingPlace.java @@ -1,9 +1,9 @@ package emissary.grpc; import emissary.config.Configurator; +import emissary.grpc.channel.ChannelManager; +import emissary.grpc.channel.ChannelPoolFactory.PoolException; import emissary.grpc.invoker.GrpcInvoker; -import emissary.grpc.pool.ConnectionFactory; -import emissary.grpc.pool.PoolException; import emissary.grpc.retry.RetryHandler; import emissary.place.ServiceProviderPlace; @@ -15,7 +15,6 @@ import io.grpc.stub.AbstractBlockingStub; import io.grpc.stub.AbstractFutureStub; import jakarta.annotation.Nullable; -import org.apache.commons.pool2.ObjectPool; import java.io.IOException; import java.io.InputStream; @@ -38,7 +37,7 @@ * identifier for the given host:port *
  • {@code GRPC_PORT_{Target-ID}} - gRPC service port, where {@code Target-ID} is the unique identifier for the given * host:port
  • - *
  • See {@link ConnectionFactory} for supported pooling and gRPC channel configuration keys and defaults.
  • + *
  • See {@link ChannelManager} for supported gRPC channel configuration keys and defaults.
  • *
  • See {@link RetryHandler} for supported retry configuration keys and defaults.
  • * */ @@ -51,11 +50,7 @@ public abstract class GrpcRoutingPlace extends ServiceProviderPlace implements I public static final String GRPC_HOST = "GRPC_HOST_"; public static final String GRPC_PORT = "GRPC_PORT_"; - protected GrpcInvoker grpcInvoker; - - protected final Map hostnameTable = new HashMap<>(); - protected final Map portNumberTable = new HashMap<>(); - protected final Map> channelPoolTable = new HashMap<>(); + protected final Map invokerTable = new HashMap<>(); protected GrpcRoutingPlace() throws IOException { super(); @@ -98,28 +93,27 @@ protected GrpcRoutingPlace(@Nullable Configurator configs) throws IOException { } private void configureGrpc() { - if (configG == null) { - throw new IllegalStateException("gRPC configurations not found for " + this.getPlaceName()); - } + Objects.requireNonNull(configG); - hostnameTable.putAll(getHostnameConfigs()); - portNumberTable.putAll(getPortNumberConfigs()); + Map hosts = getHostnameConfigs(); + Map ports = getPortNumberConfigs(); - if (!hostnameTable.keySet().equals(portNumberTable.keySet())) { + if (!hosts.keySet().equals(ports.keySet())) { throw new IllegalArgumentException("gRPC hostname target-IDs do not match gRPC port number target-IDs"); } - if (hostnameTable.isEmpty()) { + Set targetIds = hosts.keySet(); + if (targetIds.isEmpty()) { throw new NullPointerException(String.format( "Missing required arguments: %s${Target-ID} and %s${Target-ID}", GRPC_HOST, GRPC_PORT)); } - Set targetIds = hostnameTable.keySet(); + RetryHandler retryHandler = new RetryHandler(configG, this.getPlaceName(), this::retryOnException); for (String id : targetIds) { - channelPoolTable.put(id, newConnectionPool(id)); + ChannelManager channelManager = new ChannelManager(hosts.get(id), ports.get(id), configG); + GrpcInvoker grpcInvoker = new GrpcInvoker(channelManager, retryHandler); + invokerTable.put(id, grpcInvoker); } - - grpcInvoker = new GrpcInvoker(new RetryHandler(configG, this.getPlaceName(), this::retryOnException)); } protected Map getHostnameConfigs() { @@ -147,17 +141,9 @@ protected boolean retryOnException(Throwable t) { return t instanceof PoolException; } - private ObjectPool newConnectionPool(String id) { - return newConnectionFactory(id).newConnectionPool(); - } - - private ConnectionFactory newConnectionFactory(String id) { - return new ConnectionFactory(hostnameTable.get(id), portNumberTable.get(id), Objects.requireNonNull(configG)); - } - /** - * Wrapper method for {@link GrpcInvoker#invoke(ObjectPool, Function, BiFunction, Message)} that executes a unary gRPC - * call to a given endpoint. + * Wrapper method for {@link GrpcInvoker#invoke(Function, BiFunction, Message)} that executes a unary gRPC call to a + * given endpoint. * * @param targetId the identifier used in the configs for the given gRPC endpoint * @param stubFactory function that creates the appropriate gRPC stub from a {@link ManagedChannel} @@ -170,12 +156,12 @@ private ConnectionFactory newConnectionFactory(String id) { */ protected > R invokeGrpc( String targetId, Function stubFactory, BiFunction callLogic, Q request) { - return grpcInvoker.invoke(channelPoolLookup(targetId), stubFactory, callLogic, request); + return getInvoker(targetId).invoke(stubFactory, callLogic, request); } /** - * Wrapper method for {@link GrpcInvoker#invokeAsync(ObjectPool, Function, BiFunction, Message)} that executes a unary - * gRPC call to a given endpoint and returns a {@link CompletableFuture future}. + * Wrapper method for {@link GrpcInvoker#invokeAsync(Function, BiFunction, Message)} that executes a unary gRPC call to + * a given endpoint and returns a {@link CompletableFuture future}. * * @param targetId the identifier used in the configs for the given gRPC endpoint * @param stubFactory function that creates the appropriate gRPC stub from a {@link ManagedChannel} @@ -188,25 +174,27 @@ protected > CompletableFuture invokeGrpcAsync( String targetId, Function stubFactory, BiFunction> callLogic, Q request) { - return grpcInvoker.invokeAsync(channelPoolLookup(targetId), stubFactory, callLogic, request); - } - - private ObjectPool channelPoolLookup(String targetId) { - return tableLookup(channelPoolTable, targetId); + return getInvoker(targetId).invokeAsync(stubFactory, callLogic, request); } public String getHostname(String targetId) { - return tableLookup(hostnameTable, targetId); + return getInvoker(targetId).getHost(); } public int getPortNumber(String targetId) { - return tableLookup(portNumberTable, targetId); + return getInvoker(targetId).getPort(); } - protected T tableLookup(Map table, String targetId) { - if (table.containsKey(targetId)) { - return table.get(targetId); + private GrpcInvoker getInvoker(String targetId) { + if (invokerTable.containsKey(targetId)) { + return invokerTable.get(targetId); } throw new IllegalArgumentException(String.format("Target-ID %s was never configured", targetId)); } + + @Override + public void shutDown() { + super.shutDown(); + invokerTable.values().forEach(GrpcInvoker::close); + } } diff --git a/src/main/java/emissary/grpc/channel/ChannelManager.java b/src/main/java/emissary/grpc/channel/ChannelManager.java new file mode 100644 index 0000000000..4668a07a86 --- /dev/null +++ b/src/main/java/emissary/grpc/channel/ChannelManager.java @@ -0,0 +1,46 @@ +package emissary.grpc.channel; + +import emissary.config.Configurator; + +import io.grpc.ManagedChannel; +import org.apache.commons.pool2.ObjectPool; + +/** + * Wrapper for an {@link ObjectPool} created by a {@link ChannelPoolFactory}. + */ +public class ChannelManager implements AutoCloseable { + private final ObjectPool channelPool; + private final String host; + private final int port; + + public ChannelManager(String host, int port, Configurator configG) { + this.channelPool = new ChannelPoolFactory(host, port, configG).newConnectionPool(); + this.host = host; + this.port = port; + } + + public ManagedChannel acquire() { + return ChannelPoolFactory.acquireChannel(channelPool); + } + + public void release(ManagedChannel channel) { + ChannelPoolFactory.returnChannel(channel, channelPool); + } + + public void shutdown(ManagedChannel channel) { + ChannelPoolFactory.invalidateChannel(channel, channelPool); + } + + @Override + public void close() { + channelPool.close(); + } + + public String getHost() { + return host; + } + + public int getPort() { + return port; + } +} diff --git a/src/main/java/emissary/grpc/pool/ConnectionFactory.java b/src/main/java/emissary/grpc/channel/ChannelPoolFactory.java similarity index 92% rename from src/main/java/emissary/grpc/pool/ConnectionFactory.java rename to src/main/java/emissary/grpc/channel/ChannelPoolFactory.java index ab17497857..236881e55f 100644 --- a/src/main/java/emissary/grpc/pool/ConnectionFactory.java +++ b/src/main/java/emissary/grpc/channel/ChannelPoolFactory.java @@ -1,4 +1,4 @@ -package emissary.grpc.pool; +package emissary.grpc.channel; import emissary.config.Configurator; @@ -45,7 +45,7 @@ * Source for default * gRPC configurations. */ -public class ConnectionFactory extends BasePooledObjectFactory { +public class ChannelPoolFactory extends BasePooledObjectFactory { public static final String GRPC_KEEP_ALIVE_MILLIS = "GRPC_KEEP_ALIVE_MILLIS"; public static final String GRPC_KEEP_ALIVE_TIMEOUT_MILLIS = "GRPC_KEEP_ALIVE_TIMEOUT_MILLIS"; public static final String GRPC_KEEP_ALIVE_WITHOUT_CALLS = "GRPC_KEEP_ALIVE_WITHOUT_CALLS"; @@ -62,7 +62,7 @@ public class ConnectionFactory extends BasePooledObjectFactory { private static final int MAX_PORT_NUMBER = 0xFFFF; - protected static final Logger logger = LoggerFactory.getLogger(ConnectionFactory.class); + protected static final Logger logger = LoggerFactory.getLogger(ChannelPoolFactory.class); private final GenericObjectPoolConfig poolConfig = new GenericObjectPoolConfig<>(); @@ -81,13 +81,13 @@ public class ConnectionFactory extends BasePooledObjectFactory { * Constructs a new gRPC connection factory using the provided host, port, and configuration. Initializes pool settings * and gRPC channel properties from the given configuration source. *

    - * See {@link ConnectionFactory} for supported configuration keys and defaults. + * See {@link ChannelPoolFactory} for supported configuration keys and defaults. * * @param host gRPC service hostname or DNS target * @param port gRPC service port * @param configG configuration provider for channel and pool parameters */ - public ConnectionFactory(String host, int port, Configurator configG) { + public ChannelPoolFactory(String host, int port, Configurator configG) { this.host = host; this.port = port; this.target = createTarget(); @@ -105,7 +105,7 @@ public ConnectionFactory(String host, int port, Configurator configG) { // Specifies how the client chooses between multiple backend addresses // e.g. "pick_first" uses the first address only, "round_robin" cycles through all of them for client-side balancing this.loadBalancingPolicy = configG.findObjectEntry( - GRPC_LOAD_BALANCING_POLICY, LoadBalancingPolicy::valueOf, LoadBalancingPolicy.ROUND_ROBIN).formattedName(); + GRPC_LOAD_BALANCING_POLICY, LoadBalancingPolicy::valueOf, LoadBalancingPolicy.ROUND_ROBIN).toString(); // Max size (in bytes) for incoming messages and message metadata from the server this.maxInboundMessageByteSize = configG.findIntEntry(GRPC_MAX_INBOUND_MESSAGE_BYTE_SIZE, 4 << 20); // 4 MiB @@ -331,4 +331,35 @@ public boolean getPoolIsLifo() { public boolean getPoolIsFifo() { return !this.poolConfig.getLifo(); } + + /** + * Exception type for failures with handling the gRPC connection pool, such as failed borrows. + */ + public static class PoolException extends RuntimeException { + + private static final long serialVersionUID = 1495483102825486040L; + + public PoolException(String errorMessage, Throwable err) { + super(errorMessage, err); + } + } + + public enum PoolRetrievalOrdering { + LIFO, FIFO; + } + + public enum LoadBalancingPolicy { + ROUND_ROBIN("round_robin"), PICK_FIRST("pick_first"); + + private final String policy; + + LoadBalancingPolicy(String policy) { + this.policy = policy; + } + + @Override + public String toString() { + return policy; + } + } } diff --git a/src/main/java/emissary/grpc/exceptions/GrpcExceptionUtils.java b/src/main/java/emissary/grpc/exceptions/GrpcExceptionUtils.java index 8d15c95738..1e51135f61 100644 --- a/src/main/java/emissary/grpc/exceptions/GrpcExceptionUtils.java +++ b/src/main/java/emissary/grpc/exceptions/GrpcExceptionUtils.java @@ -58,7 +58,7 @@ private static String getStatusCodeDescription(Status status) { private static String getStatusCodeDescription(Status status, String message) { String description = status.getDescription(); - if (description != null) { + if (description != null && !description.equals(message)) { return String.format("%s (%s)", description, message); } return message; diff --git a/src/main/java/emissary/grpc/invoker/GrpcInvoker.java b/src/main/java/emissary/grpc/invoker/GrpcInvoker.java index 80e9954018..56bf51175a 100644 --- a/src/main/java/emissary/grpc/invoker/GrpcInvoker.java +++ b/src/main/java/emissary/grpc/invoker/GrpcInvoker.java @@ -1,8 +1,8 @@ package emissary.grpc.invoker; +import emissary.grpc.channel.ChannelManager; import emissary.grpc.exceptions.GrpcExceptionUtils; import emissary.grpc.future.CompletableFutureAdaptors; -import emissary.grpc.pool.ConnectionFactory; import emissary.grpc.retry.RetryHandler; import com.google.common.util.concurrent.ListenableFuture; @@ -10,7 +10,6 @@ import io.grpc.ManagedChannel; import io.grpc.stub.AbstractBlockingStub; import io.grpc.stub.AbstractFutureStub; -import org.apache.commons.pool2.ObjectPool; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; @@ -19,10 +18,12 @@ import java.util.function.BiFunction; import java.util.function.Function; -public class GrpcInvoker { +public class GrpcInvoker implements AutoCloseable { + private final ChannelManager channelManager; private final RetryHandler retryHandler; - public GrpcInvoker(RetryHandler retryHandler) { + public GrpcInvoker(ChannelManager channelManager, RetryHandler retryHandler) { + this.channelManager = channelManager; this.retryHandler = retryHandler; } @@ -31,7 +32,6 @@ public GrpcInvoker(RetryHandler retryHandler) { * due to an allowed Exception, the call will be tried again per the configurations set using {@link RetryHandler}. All * other Exceptions are thrown on the spot. Will also throw an Exception once max attempts have been reached. * - * @param channelPool object pool of gRPC connections for a given endpoint * @param stubFactory function that creates the appropriate gRPC stub from a {@link ManagedChannel} * @param callLogic function that performs the actual gRPC call using the stub and request * @param request the protobuf request message to send @@ -41,17 +41,16 @@ public GrpcInvoker(RetryHandler retryHandler) { * @param the gRPC stub type */ public > R invoke( - ObjectPool channelPool, Function stubFactory, - BiFunction callLogic, Q request) { + Function stubFactory, BiFunction callLogic, Q request) { return retryHandler.execute(() -> { - ManagedChannel channel = ConnectionFactory.acquireChannel(channelPool); + ManagedChannel channel = channelManager.acquire(); try { S stub = stubFactory.apply(channel); return callLogic.apply(stub, request); } catch (RuntimeException e) { throw GrpcExceptionUtils.toContextualRuntimeException(e); } finally { - ConnectionFactory.returnChannel(channel, channelPool); + channelManager.release(channel); } }); } @@ -61,7 +60,6 @@ public > * to an allowed Exception, the call will be tried again per the configurations set using {@link RetryHandler}. All * other Exceptions are thrown on the spot. Will also throw an Exception once max attempts have been reached. * - * @param channelPool object pool of gRPC connections for a given endpoint * @param stubFactory function that creates the appropriate gRPC stub from a {@link ManagedChannel} * @param callLogic function that performs the actual gRPC call using the stub and request * @param request the protobuf request message to send @@ -71,16 +69,15 @@ public > * @param the gRPC stub type */ public > CompletableFuture invokeAsync( - ObjectPool channelPool, Function stubFactory, - BiFunction> callLogic, Q request) { + Function stubFactory, BiFunction> callLogic, Q request) { AtomicReference> listenableRef = new AtomicReference<>(); CompletionStage stage = retryHandler.executeAsync(() -> { - ManagedChannel channel = ConnectionFactory.acquireChannel(channelPool); + ManagedChannel channel = channelManager.acquire(); try { S stub = stubFactory.apply(channel); listenableRef.set(callLogic.apply(stub, request)); CompletableFuture completable = CompletableFutureAdaptors.fromListenableFuture(listenableRef.get()); - return attachHandlingHook(completable, channel, channelPool); + return attachHandlingHook(completable, channel); } catch (RuntimeException e) { throw GrpcExceptionUtils.toContextualRuntimeException(e); } @@ -95,14 +92,13 @@ public > C * * @param future the future to attach the hook to * @param channel the borrowed channel - * @param channelPool the pool the channel came from * @param response type * @return a future with explicit handling logic */ - private static CompletableFuture attachHandlingHook( - CompletableFuture future, ManagedChannel channel, ObjectPool channelPool) { + private CompletableFuture attachHandlingHook( + CompletableFuture future, ManagedChannel channel) { return future.handle((response, throwable) -> { - ConnectionFactory.returnChannel(channel, channelPool); + channelManager.release(channel); if (throwable == null) { return response; } @@ -144,4 +140,17 @@ private static CompletableFuture attachCancellationHook( }); return future; } + + public String getHost() { + return channelManager.getHost(); + } + + public int getPort() { + return channelManager.getPort(); + } + + @Override + public void close() { + channelManager.close(); + } } diff --git a/src/main/java/emissary/grpc/pool/LoadBalancingPolicy.java b/src/main/java/emissary/grpc/pool/LoadBalancingPolicy.java deleted file mode 100644 index 11119fb99e..0000000000 --- a/src/main/java/emissary/grpc/pool/LoadBalancingPolicy.java +++ /dev/null @@ -1,11 +0,0 @@ -package emissary.grpc.pool; - -import java.util.Locale; - -public enum LoadBalancingPolicy { - ROUND_ROBIN, PICK_FIRST; - - public String formattedName() { - return name().toLowerCase(Locale.ROOT); - } -} diff --git a/src/main/java/emissary/grpc/pool/PoolException.java b/src/main/java/emissary/grpc/pool/PoolException.java deleted file mode 100644 index 94291b86f6..0000000000 --- a/src/main/java/emissary/grpc/pool/PoolException.java +++ /dev/null @@ -1,17 +0,0 @@ -package emissary.grpc.pool; - -/** - * Exception type for failures with handling the gRPC connection pool, such as failed borrows. - */ -public class PoolException extends RuntimeException { - - private static final long serialVersionUID = 1495483102825486040L; - - public PoolException(String errorMessage) { - super(errorMessage); - } - - public PoolException(String errorMessage, Throwable err) { - super(errorMessage, err); - } -} diff --git a/src/main/java/emissary/grpc/pool/PoolRetrievalOrdering.java b/src/main/java/emissary/grpc/pool/PoolRetrievalOrdering.java deleted file mode 100644 index 0e6ed903c2..0000000000 --- a/src/main/java/emissary/grpc/pool/PoolRetrievalOrdering.java +++ /dev/null @@ -1,5 +0,0 @@ -package emissary.grpc.pool; - -public enum PoolRetrievalOrdering { - LIFO, FIFO; -} diff --git a/src/main/java/emissary/grpc/sample/GrpcSamplePlace.java b/src/main/java/emissary/grpc/sample/GrpcSamplePlace.java index bade4ec99d..90586ac1e8 100644 --- a/src/main/java/emissary/grpc/sample/GrpcSamplePlace.java +++ b/src/main/java/emissary/grpc/sample/GrpcSamplePlace.java @@ -48,13 +48,13 @@ public void processEndpoint(IBaseDataObject o, String endpoint) { } public void processEndpointsSequentially(IBaseDataObject o) { - hostnameTable.keySet().stream() + invokerTable.keySet().stream() .sorted(Comparator.naturalOrder()) .forEach(endpoint -> processEndpoint(o, endpoint)); } public void processEndpointsInParallel(IBaseDataObject o, @Nullable Function exceptionally) { - Map> futureMap = hostnameTable.keySet().stream() + Map> futureMap = invokerTable.keySet().stream() .collect(Collectors.toMap(k -> k, k -> invokeGrpcAsync( k, SampleServiceGrpc::newFutureStub, SampleServiceFutureStub::callSampleService, generateRequest(o)))); Map responseMap = CompletableFutureFinalizers.awaitAllAndGet(futureMap, HashMap::new, exceptionally); diff --git a/src/test/java/emissary/grpc/pool/ConnectionFactoryTest.java b/src/test/java/emissary/grpc/channel/ChannelPoolFactoryTest.java similarity index 63% rename from src/test/java/emissary/grpc/pool/ConnectionFactoryTest.java rename to src/test/java/emissary/grpc/channel/ChannelPoolFactoryTest.java index 209ac52121..114377b7d1 100644 --- a/src/test/java/emissary/grpc/pool/ConnectionFactoryTest.java +++ b/src/test/java/emissary/grpc/channel/ChannelPoolFactoryTest.java @@ -1,4 +1,4 @@ -package emissary.grpc.pool; +package emissary.grpc.channel; import emissary.config.Configurator; import emissary.config.ServiceConfigGuide; @@ -24,30 +24,30 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; -class ConnectionFactoryTest extends UnitTest { - private static ConnectionFactory buildConnectionFactory(Configurator configG) { - return new ConnectionFactory("localhost", 2222, configG); +class ChannelPoolFactoryTest extends UnitTest { + private static ChannelPoolFactory buildConnectionFactory(Configurator configG) { + return new ChannelPoolFactory("localhost", 2222, configG); } private static Configurator getDefaultConfigs() { Configurator configT = new ServiceConfigGuide(); - configT.addEntry(ConnectionFactory.GRPC_POOL_MIN_IDLE_CONNECTIONS, "1"); - configT.addEntry(ConnectionFactory.GRPC_POOL_MAX_IDLE_CONNECTIONS, "2"); - configT.addEntry(ConnectionFactory.GRPC_POOL_MAX_SIZE, "2"); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_MIN_IDLE_CONNECTIONS, "1"); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_MAX_IDLE_CONNECTIONS, "2"); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_MAX_SIZE, "2"); return configT; } @ParameterizedTest @ValueSource(strings = {"localhost", "dns:///foo.bar", "foo.bar"}) void testGoodHostName(String host) { - Runnable invocation = () -> new ConnectionFactory(host, 1, new ServiceConfigGuide()); + Runnable invocation = () -> new ChannelPoolFactory(host, 1, new ServiceConfigGuide()); assertDoesNotThrow(invocation::run); } @ParameterizedTest @NullAndEmptySource void testBadHostName(String host) { - Runnable invocation = () -> new ConnectionFactory(host, 1, new ServiceConfigGuide()); + Runnable invocation = () -> new ChannelPoolFactory(host, 1, new ServiceConfigGuide()); IllegalArgumentException e = assertThrows(IllegalArgumentException.class, invocation::run); assertEquals("Missing required gRPC host configuration", e.getMessage()); } @@ -55,14 +55,14 @@ void testBadHostName(String host) { @ParameterizedTest @ValueSource(ints = {1, 8000, 8001, 8080, 65535}) void testGoodPortNumber(int port) { - Runnable invocation = () -> new ConnectionFactory("dns:///foo.bar", port, new ServiceConfigGuide()); + Runnable invocation = () -> new ChannelPoolFactory("dns:///foo.bar", port, new ServiceConfigGuide()); assertDoesNotThrow(invocation::run); } @ParameterizedTest @ValueSource(ints = {-1, 0, 65536}) void testBadPortNumber(int port) { - Runnable invocation = () -> new ConnectionFactory("dns:///foo.bar", port, new ServiceConfigGuide()); + Runnable invocation = () -> new ChannelPoolFactory("dns:///foo.bar", port, new ServiceConfigGuide()); IllegalArgumentException e = assertThrows(IllegalArgumentException.class, invocation::run); assertEquals(String.format("Port \"%d\" is outside valid range [1, 65535]", port), e.getMessage()); } @@ -70,17 +70,17 @@ void testBadPortNumber(int port) { @Test void testBadPoolRetrievalOrderConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_POOL_RETRIEVAL_ORDER, "ZIFO"); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_RETRIEVAL_ORDER, "ZIFO"); IllegalArgumentException e = assertThrows(IllegalArgumentException.class, () -> buildConnectionFactory(configT)); - assertEquals("No enum constant emissary.grpc.pool.PoolRetrievalOrdering.ZIFO", e.getMessage()); + assertEquals("No enum constant emissary.grpc.channel.ChannelPoolFactory.PoolRetrievalOrdering.ZIFO", e.getMessage()); } @Test void testLifoPoolRetrievalOrderConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_POOL_RETRIEVAL_ORDER, PoolRetrievalOrdering.LIFO.name()); - ConnectionFactory factory = buildConnectionFactory(configT); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_RETRIEVAL_ORDER, ChannelPoolFactory.PoolRetrievalOrdering.LIFO.name()); + ChannelPoolFactory factory = buildConnectionFactory(configT); assertTrue(factory.getPoolIsLifo()); assertFalse(factory.getPoolIsFifo()); } @@ -88,8 +88,8 @@ void testLifoPoolRetrievalOrderConfig() { @Test void testFifoPoolRetrievalOrderConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_POOL_RETRIEVAL_ORDER, PoolRetrievalOrdering.FIFO.name()); - ConnectionFactory factory = buildConnectionFactory(configT); + configT.addEntry(ChannelPoolFactory.GRPC_POOL_RETRIEVAL_ORDER, ChannelPoolFactory.PoolRetrievalOrdering.FIFO.name()); + ChannelPoolFactory factory = buildConnectionFactory(configT); assertFalse(factory.getPoolIsLifo()); assertTrue(factory.getPoolIsFifo()); } @@ -97,31 +97,31 @@ void testFifoPoolRetrievalOrderConfig() { @Test void testBadLoadBalancingConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_LOAD_BALANCING_POLICY, "BAD_SCHEDULER"); + configT.addEntry(ChannelPoolFactory.GRPC_LOAD_BALANCING_POLICY, "BAD_SCHEDULER"); IllegalArgumentException e = assertThrows(IllegalArgumentException.class, () -> buildConnectionFactory(configT)); - assertEquals("No enum constant emissary.grpc.pool.LoadBalancingPolicy.BAD_SCHEDULER", e.getMessage()); + assertEquals("No enum constant emissary.grpc.channel.ChannelPoolFactory.LoadBalancingPolicy.BAD_SCHEDULER", e.getMessage()); } @Test void testRoundRobinLoadBalancingConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_LOAD_BALANCING_POLICY, LoadBalancingPolicy.ROUND_ROBIN.name()); - ConnectionFactory factory = buildConnectionFactory(configT); + configT.addEntry(ChannelPoolFactory.GRPC_LOAD_BALANCING_POLICY, ChannelPoolFactory.LoadBalancingPolicy.ROUND_ROBIN.name()); + ChannelPoolFactory factory = buildConnectionFactory(configT); assertEquals("round_robin", factory.getLoadBalancingPolicy()); } @Test void testPickFirstLoadBalancingConfig() { Configurator configT = getDefaultConfigs(); - configT.addEntry(ConnectionFactory.GRPC_LOAD_BALANCING_POLICY, LoadBalancingPolicy.PICK_FIRST.name()); - ConnectionFactory factory = buildConnectionFactory(configT); + configT.addEntry(ChannelPoolFactory.GRPC_LOAD_BALANCING_POLICY, ChannelPoolFactory.LoadBalancingPolicy.PICK_FIRST.name()); + ChannelPoolFactory factory = buildConnectionFactory(configT); assertEquals("pick_first", factory.getLoadBalancingPolicy()); } @Nested class PooledChannelTests { - private ConnectionFactory factory; + private ChannelPoolFactory factory; private ObjectPool pool; @BeforeEach @@ -133,38 +133,39 @@ void setUp() { @Test void testAcquireChannel() throws Exception { - ManagedChannel channel = ConnectionFactory.acquireChannel(pool); + ManagedChannel channel = ChannelPoolFactory.acquireChannel(pool); assertNotNull(channel); pool.returnObject(channel); } @Test void testInvalidateChannel() throws Exception { - ManagedChannel c1 = ConnectionFactory.acquireChannel(pool); - ConnectionFactory.invalidateChannel(c1, pool); - ManagedChannel c2 = ConnectionFactory.acquireChannel(pool); + ManagedChannel c1 = ChannelPoolFactory.acquireChannel(pool); + ChannelPoolFactory.invalidateChannel(c1, pool); + ManagedChannel c2 = ChannelPoolFactory.acquireChannel(pool); assertNotSame(c1, c2); pool.returnObject(c2); } @Test void testReturnChannel() { - ManagedChannel c1 = ConnectionFactory.acquireChannel(pool); - ConnectionFactory.returnChannel(c1, pool); - ManagedChannel c2 = ConnectionFactory.acquireChannel(pool); + ManagedChannel c1 = ChannelPoolFactory.acquireChannel(pool); + ChannelPoolFactory.returnChannel(c1, pool); + ManagedChannel c2 = ChannelPoolFactory.acquireChannel(pool); assertSame(c1, c2); - ConnectionFactory.returnChannel(c2, pool); + ChannelPoolFactory.returnChannel(c2, pool); } @Test void testMaxPoolSizeBlocks() { - ManagedChannel c1 = ConnectionFactory.acquireChannel(pool); - ManagedChannel c2 = ConnectionFactory.acquireChannel(pool); - PoolException exception = assertThrows(PoolException.class, () -> ConnectionFactory.acquireChannel(pool)); + ManagedChannel c1 = ChannelPoolFactory.acquireChannel(pool); + ManagedChannel c2 = ChannelPoolFactory.acquireChannel(pool); + ChannelPoolFactory.PoolException exception = + assertThrows(ChannelPoolFactory.PoolException.class, () -> ChannelPoolFactory.acquireChannel(pool)); assertTrue(Strings.CS.startsWith(exception.getMessage(), "Unable to borrow channel from pool: Timeout waiting for idle object")); - ConnectionFactory.returnChannel(c1, pool); - ConnectionFactory.returnChannel(c2, pool); + ChannelPoolFactory.returnChannel(c1, pool); + ChannelPoolFactory.returnChannel(c2, pool); } @Test