Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 22 additions & 70 deletions core/src/main/java/tech/ydb/core/grpc/GrpcTransportBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -148,12 +148,25 @@ public Supplier<ScheduledExecutorService> getSchedulerFactory() {
return schedulerFactory;
}

/**
* use {@link GrpcTransportBuilder#getBalancingSettings()} instead
* @deprecated
* @return current local dc
*/
@Deprecated
Comment thread
alex268 marked this conversation as resolved.
public String getLocalDc() {
return localDc;
}

public BalancingSettings getBalancingSettings() {
return balancingSettings;
if (balancingSettings != null) {
return balancingSettings;
}
if (localDc != null) {
return BalancingSettings.fromLocation(localDc);
}

return BalancingSettings.defaultInstance();
}

public Executor getCallExecutor() {
Expand All @@ -164,15 +177,16 @@ public AuthRpcProvider<? super GrpcAuthRpc> getAuthProvider() {
return authProvider;
}

/**
* use {@link tech.ydb.core.settings.BaseRequestSettings#getRequestTimeout()} instead
* @return default request timeout
* @deprecated
*/
@Deprecated
public long getReadTimeoutMillis() {
return readTimeoutMillis;
}

@Deprecated
public long getConnectTimeoutMillis() {
return 10_000;
}

public long getDiscoveryTimeoutMillis() {
return discoveryTimeoutMillis;
}
Expand Down Expand Up @@ -236,24 +250,6 @@ public GrpcTransportBuilder addChannelInitializer(Consumer<? super ManagedChanne
return this;
}

/**
* Set a custom initialization of {@link io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder} <br>
* This method is deprecated. Use
* {@link GrpcTransportBuilder#withChannelFactoryBuilder(tech.ydb.core.impl.pool.ManagedChannelFactory.Builder)}
* instead
*
* @param ci custom NettyChannelBuilder initializer
* @return this
* @deprecated
*/
@Deprecated
public GrpcTransportBuilder withChannelInitializer(
Consumer<io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder> ci
) {
this.channelFactoryBuilder = tech.ydb.core.impl.pool.ShadedNettyChannelFactory.withInterceptor(ci);
return this;
}

/**
* use {@link GrpcTransportBuilder#withBalancingSettings(tech.ydb.core.grpc.BalancingSettings) } instead
* @param dc preferable location
Expand Down Expand Up @@ -321,7 +317,7 @@ public GrpcTransportBuilder withGrpcCompression(@Nonnull GrpcCompression compres
}

/**
* use tech.ydb.table.settings.RequestSettings#setTimeout(java.time.Duration) instead
* use {@link tech.ydb.core.settings.BaseRequestSettings.BaseBuilder#withRequestTimeout(java.time.Duration)} instead
* @param timeout global timeout for grpc calls
* @return this
* @deprecated
Expand All @@ -334,7 +330,7 @@ public GrpcTransportBuilder withReadTimeout(Duration timeout) {
}

/**
* use tech.ydb.table.settings.RequestSettings#setTimeout(long, java.time.TimeUnit) instead
* use {@link tech.ydb.core.settings.BaseRequestSettings.BaseBuilder#withRequestTimeout(java.time.Duration)} instead
* @param timeout size of global timeout for grpc calls
* @param unit time unit of global timeout for grpc calls
* @return this
Expand All @@ -347,16 +343,6 @@ public GrpcTransportBuilder withReadTimeout(long timeout, TimeUnit unit) {
return this;
}

@Deprecated
public GrpcTransportBuilder withConnectTimeout(Duration timeout) {
return this;
}

@Deprecated
public GrpcTransportBuilder withConnectTimeout(long timeout, TimeUnit unit) {
return this;
}

public GrpcTransportBuilder withDiscoveryTimeout(Duration timeout) {
this.discoveryTimeoutMillis = timeout.toMillis();
Preconditions.checkArgument(discoveryTimeoutMillis > 0, "discoveryTimeoutMillis must be greater than 0");
Expand Down Expand Up @@ -444,28 +430,6 @@ public GrpcTransportBuilder withExtraBuildInfo(String extraBuildInfo) {
return this;
}

/**
* use {@link GrpcTransportBuilder#withGrpcRetry(boolean) } instead
* @return this
* @deprecated
*/
@Deprecated
public GrpcTransportBuilder enableRetry() {
this.grpcRetry = true;
return this;
}

/**
* use {@link GrpcTransportBuilder#withGrpcRetry(boolean) } instead
* @return this
* @deprecated
*/
@Deprecated
public GrpcTransportBuilder disableRetry() {
this.grpcRetry = false;
return this;
}

public GrpcTransport build() {
YdbTransportImpl impl = new YdbTransportImpl(this);
try {
Expand All @@ -476,16 +440,4 @@ public GrpcTransport build() {
throw ex;
}
}

@Deprecated
public GrpcTransport buildAsync(Runnable ready) {
YdbTransportImpl impl = new YdbTransportImpl(this);
try {
impl.startAsync(ready);
return impl;
} catch (RuntimeException ex) {
impl.close();
throw ex;
}
}
}
28 changes: 1 addition & 27 deletions core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public class YdbTransportImpl extends BaseGrpcTransport {

public YdbTransportImpl(GrpcTransportBuilder builder) {
super(builder);
BalancingSettings balancingSettings = getBalancingSettings(builder);
BalancingSettings balancingSettings = builder.getBalancingSettings();
Duration discoveryTimeout = Duration.ofMillis(builder.getDiscoveryTimeoutMillis());

this.database = Strings.nullToEmpty(builder.getDatabase());
Expand Down Expand Up @@ -77,18 +77,6 @@ public String toString() {
return "YdbTransport{endpoint=" + serverEndpoint + ", database=" + database + "}";
}

@Deprecated
public void startAsync(Runnable readyWatcher) {
endpointPool.setNewState(null, Collections.singletonList(serverEndpoint));
discovery.start();
if (readyWatcher != null) {
scheduler.execute(() -> {
discovery.waitReady(-1);
readyWatcher.run();
});
}
}

@Override
protected void shutdown() {
discovery.stop();
Expand All @@ -98,20 +86,6 @@ protected void shutdown() {
YdbSchedulerFactory.shutdownScheduler(scheduler);
}

private static BalancingSettings getBalancingSettings(GrpcTransportBuilder builder) {
BalancingSettings balancingSettings = builder.getBalancingSettings();
if (balancingSettings != null) {
return balancingSettings;
}

String localDc = builder.getLocalDc();
if (localDc != null) {
return BalancingSettings.fromLocation(builder.getLocalDc());
}

return BalancingSettings.defaultInstance();
}

@Override
public ScheduledExecutorService getScheduler() {
return scheduler;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ public AuthCallOptions() {
this.readTimeoutMillis = 0;
}

@SuppressWarnings("deprecation")
public AuthCallOptions(
ScheduledExecutorService scheduler,
List<EndpointRecord> endpoints,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package tech.ydb.core.grpc;

import org.junit.Assert;
import org.junit.Test;

/**
*
* @author Aleksandr Gorshenin
*/
public class GrpcTransportBuilderTest {
private final static String YDB = "grpcs://some-host:3456/database";

@Test
public void defaultValuesTest() {
GrpcTransportBuilder builder = GrpcTransport.forConnectionString(YDB);

Assert.assertEquals("/database", builder.getDatabase());
Assert.assertEquals("some-host:3456", builder.getEndpoint());
Assert.assertNull(builder.getCert());
Assert.assertEquals(GrpcTransportBuilder.InitMode.SYNC, builder.getInitMode());
Assert.assertNotNull(builder.getManagedChannelFactory());
}

@Test
public void balancingSettingsTest() {
BalancingSettings def = GrpcTransport.forConnectionString(YDB).getBalancingSettings();
Assert.assertNotNull(def);
Assert.assertEquals(BalancingSettings.Policy.USE_ALL_NODES, def.getPolicy());
Assert.assertNull(def.getPreferableLocation());

@SuppressWarnings("deprecation")
BalancingSettings localDc = GrpcTransport.forConnectionString(YDB)
.withLocalDataCenter("ca").getBalancingSettings();
Assert.assertNotNull(localDc);
Assert.assertEquals(BalancingSettings.Policy.USE_PREFERABLE_LOCATION, localDc.getPolicy());
Assert.assertEquals("ca", localDc.getPreferableLocation());

BalancingSettings detect = GrpcTransport.forConnectionString(YDB)
.withBalancingSettings(BalancingSettings.detectLocalDs()).getBalancingSettings();
Assert.assertNotNull(detect);
Assert.assertEquals(BalancingSettings.Policy.USE_DETECT_LOCAL_DC, detect.getPolicy());
Assert.assertNull(detect.getPreferableLocation());

@SuppressWarnings("deprecation")
BalancingSettings mixed = GrpcTransport.forConnectionString(YDB)
.withBalancingSettings(BalancingSettings.detectLocalDs())
.withLocalDataCenter("ca")
.getBalancingSettings();

Assert.assertNotNull(mixed);
Assert.assertEquals(BalancingSettings.Policy.USE_DETECT_LOCAL_DC, mixed.getPolicy());
Assert.assertNull(mixed.getPreferableLocation());
}
}
47 changes: 0 additions & 47 deletions core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@


import java.time.Duration;
import java.time.ZoneId;
import java.util.HashMap;
import java.util.Map;
import java.util.Queue;
Expand Down Expand Up @@ -140,52 +139,6 @@ public void defaultBuildNotReadyTest() {
Assert.assertNull(ex.getCause());
}

@Test
public void asyncBuildGoodTest() {
Ticker tickerRequests = new Ticker();
MockedScheduler scheduler = new MockedScheduler(MockedClock.create(ZoneId.of("UTC")));

Mockito.when(discoveryChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getListEndpointsMethod()), Mockito.any()))
.thenReturn(MockedCall.discovery("self", new EndpointRecord("node", 2136)));

Mockito.when(discoveryChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getWhoAmIMethod()), Mockito.any()))
.thenReturn(MockedCall.whoAmICall(tickerRequests, "i am discovery"));
Mockito.when(transportChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getWhoAmIMethod()), Mockito.any()))
.thenReturn(MockedCall.whoAmICall(tickerRequests, "i am node"));

CompletableFuture<Void> isReady = new CompletableFuture<>();

@SuppressWarnings("deprecation")
GrpcTransport transport = GrpcTransport.forConnectionString("grpc://mocked:2136/local")
.withSchedulerFactory(() -> scheduler)
.withChannelFactoryBuilder(builder -> channelFactory)
.buildAsync(() -> isReady.complete(null));

Assert.assertNotNull(transport);
Assert.assertFalse(isReady.isDone());

// before discovery completed we send requests to discovery endpoint
tickerRequests.noTasks();
CompletableFuture<Result<DiscoveryProtos.WhoAmIResult>> f1 = whoAmI(transport);
Assert.assertFalse(f1.isDone());

tickerRequests.runNextTask().noTasks();
Assert.assertTrue(f1.isDone());
Assert.assertEquals("i am discovery", f1.join().getValue().getUser());

// Complete discovery
Assert.assertFalse(isReady.isDone());
scheduler.hasTasksCount(2).runNextTask().runNextTask();

CompletableFuture<Result<DiscoveryProtos.WhoAmIResult>> f2 = whoAmI(transport);
Assert.assertFalse(f2.isDone());
tickerRequests.runNextTask().noTasks();
Assert.assertTrue(f2.isDone());
Assert.assertEquals("i am node", f2.join().getValue().getUser());

Assert.assertTrue(isReady.isDone());
}

@Test
@SuppressWarnings("SleepWhileInLoop")
public void asyncWaitingForReadyTest() throws Exception {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@

import java.io.IOException;
import java.net.Socket;
import java.net.SocketAddress;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
Expand Down Expand Up @@ -43,7 +42,6 @@ public class EndpointPoolTest {
public void setUp() throws IOException {
mocks = MockitoAnnotations.openMocks(this);
threadLocalStaticMock.when(ThreadLocalRandom::current).thenReturn(random);
Mockito.doNothing().when(socket).connect(Mockito.any(SocketAddress.class));
Mockito.when(socketFactory.createSocket()).thenReturn(socket);
}

Expand Down
Loading
Loading