From 4c92222b941a6b42329bd0ba59c975524f742a21 Mon Sep 17 00:00:00 2001 From: Igor Melnichenko Date: Tue, 4 Aug 2026 21:25:14 +0300 Subject: [PATCH 1/3] Fix nanosecond/millisecond mix-up in transport discovery YdbTransportImpl.getChannel computed the time left until the request deadline in nanoseconds (deadlineAfter is a System.nanoTime() value) and passed it to YdbDiscovery.waitReady(long millis), which awaits on a Condition with MILLISECONDS. A 5 second deadline therefore became a wait of roughly 57 days, so any request issued before the endpoint pool was populated blocked indefinitely instead of failing at its own deadline. --- .../tech/ydb/core/impl/YdbTransportImpl.java | 10 ++++-- .../ydb/core/impl/YdbTransportImplTest.java | 36 +++++++++++++++++++ 2 files changed, 43 insertions(+), 3 deletions(-) diff --git a/core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java b/core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java index a8bc6255a..d0d8e1e4d 100644 --- a/core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java +++ b/core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java @@ -6,6 +6,7 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import com.google.common.base.Strings; import org.slf4j.Logger; @@ -135,11 +136,14 @@ public AuthCallOptions getAuthCallOptions() { protected GrpcChannel getChannel(GrpcRequestSettings settings) { EndpointRecord endpoint = endpointPool.getEndpoint(channelPool.getReadyEndpoints(), settings); if (endpoint == null) { - long timeout = -1; + // negative value tells waitReady to use the default discovery timeout + long timeoutMs = -1; if (settings.getDeadlineAfter() != 0) { - timeout = settings.getDeadlineAfter() - System.nanoTime(); + long leftNanos = settings.getDeadlineAfter() - System.nanoTime(); + // an already expired deadline must not fall back to the default timeout + timeoutMs = Math.max(TimeUnit.NANOSECONDS.toMillis(leftNanos), 1); } - discovery.waitReady(timeout); + discovery.waitReady(timeoutMs); endpoint = endpointPool.getEndpoint(Collections.emptySet(), settings); } return channelPool.getChannel(endpoint); diff --git a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java index d708ab576..c1a61dcf4 100644 --- a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java +++ b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java @@ -27,6 +27,7 @@ import tech.ydb.core.UnexpectedResultException; import tech.ydb.core.grpc.GrpcRequestSettings; import tech.ydb.core.grpc.GrpcTransport; +import tech.ydb.core.grpc.GrpcTransportBuilder; import tech.ydb.core.grpc.YdbHeaders; import tech.ydb.core.impl.pool.EndpointRecord; import tech.ydb.core.impl.pool.ManagedChannelFactory; @@ -183,6 +184,41 @@ public void asyncBuildGoodTest() { Assert.assertTrue(isReady.isDone()); } + @Test(timeout = 30_000) + public void requestDeadlineLimitsWaitingForDiscovery() { + Mockito.when(discoveryChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getListEndpointsMethod()), Mockito.any())) + .thenReturn(MockedCall.neverAnswer(testScheduler)); + + // deliberately much longer than the deadline of the request below + Duration discoveryTimeout = Duration.ofMinutes(10); + + try (GrpcTransport transport = GrpcTransport.forConnectionString("grpc://mocked:2136/local") + .withInitMode(GrpcTransportBuilder.InitMode.ASYNC) + .withDiscoveryTimeout(discoveryTimeout) + .withChannelFactoryBuilder(builder -> channelFactory) + .build() + ) { + + GrpcRequestSettings settings = GrpcRequestSettings.newBuilder() + .withDeadline(Duration.ofMillis(200)) + .build(); + + long startedAt = System.nanoTime(); + Result result = transport + .unaryCall( + DiscoveryServiceGrpc.getWhoAmIMethod(), + settings, + DiscoveryProtos.WhoAmIRequest.newBuilder().build() + ) + .join(); + Duration elapsed = Duration.ofNanos(System.nanoTime() - startedAt); + + Assert.assertFalse(result.isSuccess()); + // the call must give up near its own deadline instead of waiting for the discovery timeout + Assert.assertTrue("discovery waiting took " + elapsed, elapsed.compareTo(discoveryTimeout) < 0); + } + } + @Test public void failFastOnMissingPort() { String endpoint = "127.1.2.3"; From 1cc23c7a8b8298bff4182bb96a02c0a71bce4a6c Mon Sep 17 00:00:00 2001 From: Alexandr Gorshenin Date: Wed, 12 Aug 2026 15:35:04 +0100 Subject: [PATCH 2/3] Fixed timeouts in test --- .../test/java/tech/ydb/core/impl/YdbTransportImplTest.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java index c1a61dcf4..f28c37be7 100644 --- a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java +++ b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java @@ -189,8 +189,7 @@ public void requestDeadlineLimitsWaitingForDiscovery() { Mockito.when(discoveryChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getListEndpointsMethod()), Mockito.any())) .thenReturn(MockedCall.neverAnswer(testScheduler)); - // deliberately much longer than the deadline of the request below - Duration discoveryTimeout = Duration.ofMinutes(10); + Duration discoveryTimeout = Duration.ofSeconds(5); try (GrpcTransport transport = GrpcTransport.forConnectionString("grpc://mocked:2136/local") .withInitMode(GrpcTransportBuilder.InitMode.ASYNC) @@ -200,7 +199,7 @@ public void requestDeadlineLimitsWaitingForDiscovery() { ) { GrpcRequestSettings settings = GrpcRequestSettings.newBuilder() - .withDeadline(Duration.ofMillis(200)) + .withDeadline(Duration.ofMillis(50)) .build(); long startedAt = System.nanoTime(); From 5612a66a7b0346c1ecd77c1adf0ac330c9c9f2e4 Mon Sep 17 00:00:00 2001 From: Alexandr Gorshenin Date: Thu, 13 Aug 2026 11:24:29 +0100 Subject: [PATCH 3/3] Updated unit test to coverage all lines --- .../ydb/core/impl/YdbTransportImplTest.java | 63 ++++++++++++++----- 1 file changed, 47 insertions(+), 16 deletions(-) diff --git a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java index f28c37be7..08bbcbef9 100644 --- a/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java +++ b/core/src/test/java/tech/ydb/core/impl/YdbTransportImplTest.java @@ -8,9 +8,11 @@ import java.util.Queue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import javax.annotation.Nonnull; @@ -184,37 +186,66 @@ public void asyncBuildGoodTest() { Assert.assertTrue(isReady.isDone()); } - @Test(timeout = 30_000) - public void requestDeadlineLimitsWaitingForDiscovery() { + @Test + public void asyncWaitingForReadyTest() throws Exception { + CountDownLatch discoveryLatch = new CountDownLatch(1); + Queue lazyTasks = new ConcurrentLinkedQueue<>(); + Executor lazyExecutor = (Runnable r) -> { + lazyTasks.add(r); + discoveryLatch.countDown(); + }; + Mockito.when(discoveryChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getListEndpointsMethod()), Mockito.any())) - .thenReturn(MockedCall.neverAnswer(testScheduler)); + .thenReturn(MockedCall.discovery(lazyExecutor, "self", new EndpointRecord("node", 2136))); + Mockito.when(transportChannel.newCall(Mockito.eq(DiscoveryServiceGrpc.getWhoAmIMethod()), Mockito.any())) + .thenReturn(MockedCall.whoAmICall("i am node")); Duration discoveryTimeout = Duration.ofSeconds(5); + Duration callTimeout = Duration.ofMillis(50); + Assert.assertEquals(0, lazyTasks.size()); try (GrpcTransport transport = GrpcTransport.forConnectionString("grpc://mocked:2136/local") .withInitMode(GrpcTransportBuilder.InitMode.ASYNC) .withDiscoveryTimeout(discoveryTimeout) .withChannelFactoryBuilder(builder -> channelFactory) .build() ) { - - GrpcRequestSettings settings = GrpcRequestSettings.newBuilder() - .withDeadline(Duration.ofMillis(50)) - .build(); + Assert.assertTrue(discoveryLatch.await(1, TimeUnit.SECONDS)); + Assert.assertEquals(1, lazyTasks.size()); // a discovery call long startedAt = System.nanoTime(); - Result result = transport - .unaryCall( + CompletableFuture> call1 = CompletableFuture.supplyAsync( + () ->transport.unaryCall( DiscoveryServiceGrpc.getWhoAmIMethod(), - settings, + GrpcRequestSettings.newBuilder().withDeadline(callTimeout).build(), DiscoveryProtos.WhoAmIRequest.newBuilder().build() - ) - .join(); - Duration elapsed = Duration.ofNanos(System.nanoTime() - startedAt); + ).join(), + testScheduler + ); - Assert.assertFalse(result.isSuccess()); - // the call must give up near its own deadline instead of waiting for the discovery timeout - Assert.assertTrue("discovery waiting took " + elapsed, elapsed.compareTo(discoveryTimeout) < 0); + Assert.assertFalse(call1.isDone()); + Result res1 = call1.get(5, TimeUnit.SECONDS); + Assert.assertEquals(StatusCode.CLIENT_INTERNAL_ERROR, res1.getStatus().getCode()); + Assert.assertEquals(1, lazyTasks.size()); // a new call wasn't executed + + long nanos = System.nanoTime() - startedAt; + Assert.assertTrue(nanos >= callTimeout.toNanos()); + Assert.assertTrue(nanos < discoveryTimeout.toNanos()); + + CompletableFuture> call2 = CompletableFuture.supplyAsync( + () ->transport.unaryCall( + DiscoveryServiceGrpc.getWhoAmIMethod(), + GrpcRequestSettings.newBuilder().build(), + DiscoveryProtos.WhoAmIRequest.newBuilder().build() + ).join(), + testScheduler + ); + Assert.assertFalse(call2.isDone()); + Assert.assertEquals(1, lazyTasks.size()); // a new call wasn't executed + + lazyTasks.poll().run(); // complete discovery + Result res2 = call2.join(); + Assert.assertTrue(res2.isSuccess()); } }