From 293068dc7fe59c7adc4d98fdbf6f37fa20f7f2f2 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 13 Aug 2026 17:55:00 +0300 Subject: [PATCH] YDBAPPTEAM-1650 Add session closed metric --- .../tech/ydb/core/impl/Observability.java | 2 +- .../tech/ydb/core/impl/ObservabilityTest.java | 8 +- .../java/tech/ydb/query/impl/SessionImpl.java | 11 +- .../java/tech/ydb/query/impl/SessionPool.java | 62 +++- .../tech/ydb/query/impl/PoolMetricsTest.java | 346 +++++++++++++++++- .../query/impl/SessionImplAttachHintTest.java | 19 +- .../tech/ydb/table/impl/pool/PoolMetrics.java | 14 +- .../tech/ydb/table/impl/pool/SessionPool.java | 29 +- .../ydb/table/impl/pool/StatefulSession.java | 6 + .../ydb/table/impl/pool/WaitingQueue.java | 32 +- .../ydb/table/impl/pool/MockedTableRpc.java | 9 + .../ydb/table/impl/pool/PoolMetricsTest.java | 67 +++- .../ydb/table/impl/pool/WaitingQueueTest.java | 2 +- 13 files changed, 533 insertions(+), 74 deletions(-) diff --git a/core/src/main/java/tech/ydb/core/impl/Observability.java b/core/src/main/java/tech/ydb/core/impl/Observability.java index 9a1a6f2b7..1b3f647df 100644 --- a/core/src/main/java/tech/ydb/core/impl/Observability.java +++ b/core/src/main/java/tech/ydb/core/impl/Observability.java @@ -12,7 +12,7 @@ */ public final class Observability { public static final String TRACING_CHAIN = ";ydb-sdk-tracing/0.1.0"; - public static final String METRICS_CHAIN = ";ydb-sdk-metrics/0.1.0"; + public static final String METRICS_CHAIN = ";ydb-sdk-metrics/0.2.0"; private static volatile boolean isTracingEnabled = false; private static volatile boolean isMetricsEnabled = false; diff --git a/core/src/test/java/tech/ydb/core/impl/ObservabilityTest.java b/core/src/test/java/tech/ydb/core/impl/ObservabilityTest.java index 0274dc8fa..61d16aa99 100644 --- a/core/src/test/java/tech/ydb/core/impl/ObservabilityTest.java +++ b/core/src/test/java/tech/ydb/core/impl/ObservabilityTest.java @@ -31,19 +31,19 @@ public void baseTest() { Assert.assertEquals(BASE, Observability.getDiscoveryBuildInfo(BASE)); Observability.reportMetricsUsage(new Meter() { }); - Assert.assertEquals(BASE + ";ydb-sdk-metrics/0.1.0", Observability.getDiscoveryBuildInfo(BASE)); + Assert.assertEquals(BASE + ";ydb-sdk-metrics/0.2.0", Observability.getDiscoveryBuildInfo(BASE)); Observability.reportTracingUsage((String spanName, SpanKind spanKind) -> Span.NOOP); Assert.assertEquals( - BASE + ";ydb-sdk-tracing/0.1.0;ydb-sdk-metrics/0.1.0", + BASE + ";ydb-sdk-tracing/0.1.0;ydb-sdk-metrics/0.2.0", Observability.getDiscoveryBuildInfo(BASE) ); Observability.reportMetricsUsage(Meter.NOOP); Observability.reportTracingUsage(NoopTracer.getInstance()); Assert.assertEquals( - BASE + ";ydb-sdk-tracing/0.1.0;ydb-sdk-metrics/0.1.0", + BASE + ";ydb-sdk-tracing/0.1.0;ydb-sdk-metrics/0.2.0", Observability.getDiscoveryBuildInfo(BASE) ); } -} \ No newline at end of file +} diff --git a/query/src/main/java/tech/ydb/query/impl/SessionImpl.java b/query/src/main/java/tech/ydb/query/impl/SessionImpl.java index 1f3cbaa54..c130af5fe 100644 --- a/query/src/main/java/tech/ydb/query/impl/SessionImpl.java +++ b/query/src/main/java/tech/ydb/query/impl/SessionImpl.java @@ -120,6 +120,8 @@ public QueryTransaction createNewTransaction(TxMode txMode) { public abstract void updateSessionState(Status status); + abstract void closeSession(String reason); + @Override public CompletableFuture> beginTransaction(TxMode tx, BeginTransactionSettings settings) { YdbQuery.BeginTransactionRequest request = YdbQuery.BeginTransactionRequest.newBuilder() @@ -170,21 +172,20 @@ public CompletableFuture start(GrpcReadStream.Observer observer) } StatusCode code = StatusCode.fromProto(message.getStatus()); Status status = Status.of(code, Issue.fromPb(message.getIssuesList())); - updateSessionState(status); + observer.onNext(status); // The hint is sent by the server with a success status. switch (message.getSessionHintCase()) { case NODE_SHUTDOWN: pessimizationHook.set(nodeID != 0); - updateSessionState(Status.of(StatusCode.BAD_SESSION)); + closeSession("node_shutdown"); break; case SESSION_SHUTDOWN: - updateSessionState(Status.of(StatusCode.BAD_SESSION)); + closeSession("session_shutdown"); break; default: + updateSessionState(status); break; } - - observer.onNext(status); }); } diff --git a/query/src/main/java/tech/ydb/query/impl/SessionPool.java b/query/src/main/java/tech/ydb/query/impl/SessionPool.java index 6c3a7718a..289f840ad 100644 --- a/query/src/main/java/tech/ydb/query/impl/SessionPool.java +++ b/query/src/main/java/tech/ydb/query/impl/SessionPool.java @@ -9,6 +9,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.LongAdder; import java.util.function.BiConsumer; @@ -163,11 +164,10 @@ private boolean tryComplete(CompletableFuture> future, Pool private class PooledQuerySession extends SessionImpl { private final GrpcReadStream attachStream; + private final AtomicBoolean isBroken = new AtomicBoolean(); private volatile Instant lastActive; private volatile boolean isStarted = false; - private volatile boolean isBroken = false; - private volatile boolean isStopped = false; PooledQuerySession(QueryServiceRpc rpc, YdbQuery.CreateSessionResponse response) { super(rpc, response); @@ -179,18 +179,41 @@ private class PooledQuerySession extends SessionImpl { @Override public void updateSessionState(Status status) { this.lastActive = clock.instant(); - boolean isStatusBroken = - status.getCode() == StatusCode.BAD_SESSION || - status.getCode() == StatusCode.SESSION_BUSY || - status.getCode() == StatusCode.INTERNAL_ERROR || - status.getCode() == StatusCode.CLIENT_DEADLINE_EXCEEDED || - status.getCode() == StatusCode.CLIENT_DEADLINE_EXPIRED || - status.getCode() == StatusCode.CLIENT_CANCELLED || - status.getCode() == StatusCode.TRANSPORT_UNAVAILABLE; - if (isStatusBroken) { - logger.warn("QuerySession[{}] broken with status {}", getId(), status); + String reason; + switch (status.getCode()) { + case BAD_SESSION: + case SESSION_EXPIRED: + reason = "bad_session"; + break; + case SESSION_BUSY: + reason = "session_busy"; + break; + case INTERNAL_ERROR: + reason = "internal_error"; + break; + case CLIENT_DEADLINE_EXCEEDED: + case CLIENT_DEADLINE_EXPIRED: + reason = "client_timeout"; + break; + case CLIENT_CANCELLED: + reason = "client_cancelled"; + break; + case TRANSPORT_UNAVAILABLE: + reason = "transport_error"; + break; + default: + return; + } + + logger.warn("QuerySession[{}] broken with status {}", getId(), status); + closeSession(reason); + } + + @Override + void closeSession(String reason) { + if (isStarted && isBroken.compareAndSet(false, true)) { + metrics.onSessionClosed(reason); } - isBroken = isBroken || isStatusBroken; } public Instant getLastActive() { @@ -209,14 +232,16 @@ public CompletableFuture> start() { return; } + isStarted = true; if (future.complete(ok)) { logger.debug("QuerySession[{}] attach message {}", getId(), status); - isStarted = true; return; } logger.trace("QuerySession[{}] attach message {}", getId(), status); }).whenComplete((status, th) -> { + closeSession(status != null && status.isSuccess() ? "attach_closed" : "transport_error"); + if (th != null) { logger.debug("QuerySession[{}] finished with exception", getId(), th); } @@ -235,7 +260,6 @@ public CompletableFuture> start() { private void clean() { logger.debug("QuerySession[{}] attach stream is stopped", getId()); - isStopped = true; if (!isStarted) { destroy(); } @@ -260,11 +284,11 @@ public void destroy() { @Override public void close() { - logger.trace("QuerySession[{}] closed with broken status {}", getId(), isBroken); + logger.trace("QuerySession[{}] closed with broken status {}", getId(), isBroken.get()); stats.released.increment(); metrics.onSessionReleased(); - if (isBroken || isStopped) { + if (isBroken.get()) { queue.delete(this); } else { queue.release(this); @@ -314,9 +338,9 @@ public CompletableFuture create() { } @Override - public void destroy(PooledQuerySession session) { + public void destroy(PooledQuerySession session, WaitingQueue.RemovalReason reason) { + session.closeSession(reason.toString()); stats.deleted.increment(); - metrics.onSessionDeleted(); // Execute deleteSession call outside current context to avoid cancellation and deadline propogation Context ctx = Context.ROOT.fork(); diff --git a/query/src/test/java/tech/ydb/query/impl/PoolMetricsTest.java b/query/src/test/java/tech/ydb/query/impl/PoolMetricsTest.java index a0ecd6e1f..727b0c665 100644 --- a/query/src/test/java/tech/ydb/query/impl/PoolMetricsTest.java +++ b/query/src/test/java/tech/ydb/query/impl/PoolMetricsTest.java @@ -16,8 +16,10 @@ import org.junit.Test; import org.mockito.ArgumentMatcher; +import tech.ydb.common.transaction.TxMode; import tech.ydb.core.Result; import tech.ydb.core.Status; +import tech.ydb.core.StatusCode; import tech.ydb.core.grpc.GrpcReadStream; import tech.ydb.core.grpc.GrpcRequestSettings; import tech.ydb.core.grpc.GrpcTransport; @@ -29,7 +31,11 @@ import tech.ydb.core.tracing.NoopTracer; import tech.ydb.proto.StatusCodesProtos.StatusIds; import tech.ydb.proto.query.YdbQuery; +import tech.ydb.query.QueryClient; import tech.ydb.query.QuerySession; +import tech.ydb.query.QueryStream; +import tech.ydb.query.result.QueryInfo; +import tech.ydb.table.TableClient; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyDouble; @@ -40,7 +46,9 @@ import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; public class PoolMetricsTest { @@ -57,6 +65,7 @@ public class PoolMetricsTest { private final DoubleHistogram createTime = mock(DoubleHistogram.class); private final Map counters = new HashMap<>(); private final Map> gauges = new HashMap<>(); + private Runnable cleaner; private final ArgumentMatcher poolName = a -> attr(a, "pool.name", POOL); private final ArgumentMatcher stateIdle = a -> attr(a, "state", "idle"); @@ -65,8 +74,11 @@ public class PoolMetricsTest { @Before public void setup() { - when(scheduler.scheduleAtFixedRate(any(), anyLong(), anyLong(), any())) - .thenAnswer(inv -> mock(ScheduledFuture.class)); + when(scheduler.scheduleAtFixedRate(any(Runnable.class), anyLong(), anyLong(), any())) + .thenAnswer(inv -> { + cleaner = inv.getArgument(0); + return mock(ScheduledFuture.class); + }); when(scheduler.schedule(any(Runnable.class), anyLong(), any())) .thenAnswer(inv -> mock(ScheduledFuture.class)); @@ -83,11 +95,11 @@ public void setup() { public void allInstrumentsAreCreated() { try (SessionPool pool = createPool(0, 2)) { verify(meter).createCounter(eq(PREFIX + "created"), eq("{session}"), anyString()); - verify(meter).createCounter(eq(PREFIX + "deleted"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "acquired"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "released"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "requested"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "failed"), eq("{session}"), anyString()); + verify(meter).createCounter(eq(PREFIX + "closed"), eq("{session}"), anyString()); verify(meter).createHistogram(eq(PREFIX + "create_time"), eq("s"), anyString()); verify(meter).createLongGauge(eq(PREFIX + "max"), eq("{session}"), anyString(), any()); verify(meter).createLongGauge(eq(PREFIX + "min"), eq("{session}"), anyString(), any()); @@ -96,6 +108,18 @@ public void allInstrumentsAreCreated() { } } + @Test + public void queryAndTableClientsUseSameClosedCounter() { + GrpcTransport transport = mock(GrpcTransport.class); + when(transport.getScheduler()).thenReturn(scheduler); + when(transport.getTracer()).thenReturn(NoopTracer.getInstance()); + + try (QueryClient queryClient = QueryClient.newClient(transport).withMeter(meter, POOL).build(); + TableClient tableClient = QueryClient.newTableClient(transport).withMeter(meter, POOL).build()) { + verify(meter, times(2)).createCounter(eq(PREFIX + "closed"), eq("{session}"), anyString()); + } + } + @Test public void sessionLifecycleRecordsCounters() { try (SessionPool pool = createPool(0, 2)) { @@ -109,7 +133,6 @@ public void sessionLifecycleRecordsCounters() { verify(counter("released")).add(eq(1L), argThat(poolName)); } - verify(counter("deleted")).add(eq(1L), argThat(poolName)); verify(counter("failed"), never()).add(anyLong(), any()); } @@ -162,6 +185,224 @@ public void gaugesObserveStats() { } } + @Test + public void poolGracefulShutdownRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + acquireReady(pool).close(); + } + + verifyClosed("pool_graceful_shutdown"); + } + + @Test + public void poolIdleTimeoutRecordsClosedCounter() { + SessionPool pool = new SessionPool(clock, rpc, scheduler, 0, 2, Duration.ZERO, meter, POOL); + acquireReady(pool).close(); + cleaner.run(); + pool.close(); + + verifyClosed("pool_idle_timeout"); + } + + @Test + public void poolResizeRecordsClosedCounter() { + SessionPool pool = createPool(0, 2); + QuerySession first = acquireReady(pool); + QuerySession second = acquireReady(pool); + + pool.updateMaxSize(1); + first.close(); + + verifyClosed("pool_resize"); + second.close(); + pool.close(); + } + + @Test + public void nodeShutdownRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + rpc.sendAttachMessage(YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .setNodeShutdown(YdbQuery.NodeShutdownHint.getDefaultInstance()) + .build()); + rpc.completeAttach(Status.SUCCESS); + session.close(); + } + + verifyClosed("node_shutdown"); + } + + @Test + public void sessionShutdownRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + rpc.sendAttachMessage(YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .setSessionShutdown(YdbQuery.SessionShutdownHint.getDefaultInstance()) + .build()); + session.close(); + } + + verifyClosed("session_shutdown"); + } + + @Test + public void firstAttachShutdownHintRecordsClosedCounter() { + rpc.initialAttachMessage = YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .setNodeShutdown(YdbQuery.NodeShutdownHint.getDefaultInstance()) + .build(); + + try (SessionPool pool = createPool(0, 2)) { + acquireReady(pool).close(); + } + + verifyClosed("node_shutdown"); + } + + @Test + public void attachStreamEofRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + rpc.completeAttach(Status.SUCCESS); + session.close(); + } + + verifyClosed("attach_closed"); + } + + @Test + public void attachStreamFailureRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + rpc.failAttach(new RuntimeException("transport failure")); + session.close(); + } + + verifyClosed("transport_error"); + } + + @Test + public void clientQueryTimeoutRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.completeQuery(Status.of(StatusCode.CLIENT_DEADLINE_EXCEEDED)); + Assert.assertFalse(query.join().isSuccess()); + + rpc.sendAttachMessage(YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .setNodeShutdown(YdbQuery.NodeShutdownHint.getDefaultInstance()) + .build()); + rpc.completeAttach(Status.SUCCESS); + session.close(); + } + + verifyClosed("client_timeout"); + } + + @Test + public void queryStreamCancellationRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + QueryStream query = session.createQuery("SELECT 1", TxMode.NONE); + CompletableFuture> result = query.execute(); + query.cancel(); + Assert.assertFalse(result.join().isSuccess()); + + rpc.sendAttachMessage(YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .setSessionShutdown(YdbQuery.SessionShutdownHint.getDefaultInstance()) + .build()); + rpc.completeAttach(Status.SUCCESS); + session.close(); + } + + verifyClosed("client_cancelled"); + } + + @Test + public void badSessionRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.sendQueryMessage(StatusIds.StatusCode.BAD_SESSION); + rpc.completeQuery(Status.SUCCESS); + Assert.assertFalse(query.join().isSuccess()); + session.close(); + } + + verifyClosed("bad_session"); + } + + @Test + public void sessionExpiredRecordsBadSessionCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.sendQueryMessage(StatusIds.StatusCode.SESSION_EXPIRED); + rpc.completeQuery(Status.SUCCESS); + Assert.assertFalse(query.join().isSuccess()); + session.close(); + } + + verifyClosed("bad_session"); + } + + @Test + public void sessionBusyRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.sendQueryMessage(StatusIds.StatusCode.SESSION_BUSY); + rpc.completeQuery(Status.SUCCESS); + Assert.assertFalse(query.join().isSuccess()); + session.close(); + } + + verifyClosed("session_busy"); + } + + @Test + public void transportUnavailableRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.completeQuery(Status.of(StatusCode.TRANSPORT_UNAVAILABLE)); + Assert.assertFalse(query.join().isSuccess()); + session.close(); + } + + verifyClosed("transport_error"); + } + + @Test + public void internalErrorRecordsClosedCounter() { + try (SessionPool pool = createPool(0, 2)) { + QuerySession session = acquireReady(pool); + CompletableFuture> query = session.createQuery("SELECT 1", TxMode.NONE).execute(); + rpc.sendQueryMessage(StatusIds.StatusCode.INTERNAL_ERROR); + rpc.completeQuery(Status.SUCCESS); + Assert.assertFalse(query.join().isSuccess()); + session.close(); + } + + verifyClosed("internal_error"); + } + + @Test + public void primaryAttachFailureDoesNotRecordClosedCounter() { + rpc.initialAttachMessage = YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.BAD_SESSION) + .build(); + try (SessionPool pool = createPool(0, 2)) { + Assert.assertFalse(pool.acquire(TIMEOUT).join().isSuccess()); + } + + verify(counter("closed"), never()).add(anyLong(), any()); + } + private SessionPool createPool(int minSize, int maxSize) { return new SessionPool(clock, rpc, scheduler, minSize, maxSize, IDLE, meter, POOL); } @@ -176,6 +417,15 @@ private LongCounter counter(String shortName) { return counters.get(PREFIX + shortName); } + private void verifyClosed(String reason) { + verify(counter("closed")).add( + eq(1L), + argThat(poolName), + argThat(a -> a.getKey().equals("reason") && a.getValue().equals(reason)) + ); + verifyNoMoreInteractions(counter("closed")); + } + private static boolean attr(Attr attr, String shortKey, String value) { return attr.getKey().equals(PREFIX + shortKey) && attr.getValue().equals(value); } @@ -189,6 +439,11 @@ private static boolean attr(Attr attr, String shortKey, String value) { private static final class TestRpc extends QueryServiceRpc { private final AtomicInteger ids = new AtomicInteger(); private volatile boolean overloaded = false; + private volatile YdbQuery.SessionState initialAttachMessage = YdbQuery.SessionState.newBuilder() + .setStatus(StatusIds.StatusCode.SUCCESS) + .build(); + private TestAttachStream attachStream; + private TestQueryStream queryStream; TestRpc() { super(DUMMY_TRANSPORT); @@ -201,6 +456,7 @@ public CompletableFuture> createSession( YdbQuery.CreateSessionResponse response = YdbQuery.CreateSessionResponse.newBuilder() .setStatus(code) .setSessionId("session-" + ids.incrementAndGet()) + .setNodeId(42) .build(); return CompletableFuture.completedFuture(Result.success(response)); } @@ -208,20 +464,37 @@ public CompletableFuture> createSession( @Override public GrpcReadStream attachSession( YdbQuery.AttachSessionRequest request, GrpcRequestSettings settings) { - YdbQuery.SessionState message = YdbQuery.SessionState.newBuilder() - .setStatus(StatusIds.StatusCode.SUCCESS) - .build(); - return new GrpcReadStream() { - @Override - public CompletableFuture start(Observer observer) { - observer.onNext(message); - return new CompletableFuture<>(); - } + attachStream = new TestAttachStream(initialAttachMessage); + return attachStream; + } + + void sendAttachMessage(YdbQuery.SessionState message) { + attachStream.observer.onNext(message); + } + + void completeAttach(Status status) { + attachStream.completion.complete(status); + } + + void failAttach(Throwable th) { + attachStream.completion.completeExceptionally(th); + } + + @Override + public GrpcReadStream executeQuery( + YdbQuery.ExecuteQueryRequest request, GrpcRequestSettings settings) { + queryStream = new TestQueryStream(); + return queryStream; + } - @Override - public void cancel() { - } - }; + void sendQueryMessage(StatusIds.StatusCode status) { + queryStream.observer.onNext(YdbQuery.ExecuteQueryResponsePart.newBuilder() + .setStatus(status) + .build()); + } + + void completeQuery(Status status) { + queryStream.completion.complete(status); } @Override @@ -233,4 +506,43 @@ public CompletableFuture> deleteSession( .build())); } } + + private static final class TestAttachStream implements GrpcReadStream { + private final YdbQuery.SessionState initialMessage; + private final CompletableFuture completion = new CompletableFuture<>(); + private Observer observer; + + TestAttachStream(YdbQuery.SessionState initialMessage) { + this.initialMessage = initialMessage; + } + + @Override + public CompletableFuture start(Observer observer) { + this.observer = observer; + if (initialMessage != null) { + observer.onNext(initialMessage); + } + return completion; + } + + @Override + public void cancel() { + } + } + + private static final class TestQueryStream implements GrpcReadStream { + private final CompletableFuture completion = new CompletableFuture<>(); + private Observer observer; + + @Override + public CompletableFuture start(Observer observer) { + this.observer = observer; + return completion; + } + + @Override + public void cancel() { + completion.complete(Status.of(StatusCode.CLIENT_CANCELLED)); + } + } } diff --git a/query/src/test/java/tech/ydb/query/impl/SessionImplAttachHintTest.java b/query/src/test/java/tech/ydb/query/impl/SessionImplAttachHintTest.java index bd04633ee..a1a7eccfa 100644 --- a/query/src/test/java/tech/ydb/query/impl/SessionImplAttachHintTest.java +++ b/query/src/test/java/tech/ydb/query/impl/SessionImplAttachHintTest.java @@ -7,7 +7,6 @@ import org.junit.Test; import tech.ydb.core.Status; -import tech.ydb.core.StatusCode; import tech.ydb.core.grpc.GrpcReadStream; import tech.ydb.core.grpc.GrpcRequestSettings; import tech.ydb.core.grpc.GrpcTransport; @@ -58,7 +57,7 @@ public void nodeShutdownHint_setsPessimizationHookWhenNodeIdKnown() { Assert.assertNotNull(rpc.capturedSettings); Assert.assertNotNull(rpc.capturedSettings.getPessimizationHook()); Assert.assertTrue(rpc.capturedSettings.getPessimizationHook().getAsBoolean()); - Assert.assertEquals(StatusCode.BAD_SESSION, session.getLastSessionState().getCode()); + Assert.assertEquals("node_shutdown", session.getCloseReason()); } } @@ -82,7 +81,7 @@ public void nodeShutdownHint_doesNotPessimizeWhenNodeIdIsZero() { Assert.assertNotNull(rpc.capturedSettings); Assert.assertNotNull(rpc.capturedSettings.getPessimizationHook()); Assert.assertFalse(rpc.capturedSettings.getPessimizationHook().getAsBoolean()); - Assert.assertEquals(StatusCode.BAD_SESSION, session.getLastSessionState().getCode()); + Assert.assertEquals("node_shutdown", session.getCloseReason()); } } @@ -106,7 +105,7 @@ public void sessionShutdownHint_doesNotPessimizeEndpoint() { Assert.assertNotNull(rpc.capturedSettings); Assert.assertNotNull(rpc.capturedSettings.getPessimizationHook()); Assert.assertFalse(rpc.capturedSettings.getPessimizationHook().getAsBoolean()); - Assert.assertEquals(StatusCode.BAD_SESSION, session.getLastSessionState().getCode()); + Assert.assertEquals("session_shutdown", session.getCloseReason()); } } @@ -143,7 +142,7 @@ public GrpcReadStream attachSession( } private static final class TestSession extends SessionImpl { - private final AtomicReference lastState = new AtomicReference<>(); + private final AtomicReference closeReason = new AtomicReference<>(); TestSession(QueryServiceRpc rpc, YdbQuery.CreateSessionResponse response) { super(rpc, response); @@ -151,11 +150,15 @@ private static final class TestSession extends SessionImpl { @Override public void updateSessionState(Status status) { - lastState.set(status); } - Status getLastSessionState() { - return lastState.get(); + @Override + void closeSession(String reason) { + closeReason.set(reason); + } + + String getCloseReason() { + return closeReason.get(); } @Override diff --git a/table/src/main/java/tech/ydb/table/impl/pool/PoolMetrics.java b/table/src/main/java/tech/ydb/table/impl/pool/PoolMetrics.java index 7005b1a4a..4a6eb3b18 100644 --- a/table/src/main/java/tech/ydb/table/impl/pool/PoolMetrics.java +++ b/table/src/main/java/tech/ydb/table/impl/pool/PoolMetrics.java @@ -17,11 +17,11 @@ public final class PoolMetrics { private final String statusKey; private final LongCounter created; - private final LongCounter deleted; private final LongCounter acquired; private final LongCounter released; private final LongCounter requested; private final LongCounter failed; + private final LongCounter closed; private final DoubleHistogram createTime; public PoolMetrics(Meter meter, String name, String poolName, WaitingQueue queue, int minSize) { @@ -34,11 +34,11 @@ public PoolMetrics(Meter meter, String name, String poolName, WaitingQueue qu this.inUseAttrs = new Attr[]{poolNameAttr, Attr.of(prefix + "state", "in_use")}; this.created = meter.createCounter(prefix + "created", UNIT, "Total successful session creations."); - this.deleted = meter.createCounter(prefix + "deleted", UNIT, "Total session deletions."); this.acquired = meter.createCounter(prefix + "acquired", UNIT, "Total session acquires from the pool."); this.released = meter.createCounter(prefix + "released", UNIT, "Total session releases back to the pool."); this.requested = meter.createCounter(prefix + "requested", UNIT, "Total CreateSession calls."); this.failed = meter.createCounter(prefix + "failed", UNIT, "Total failed session creations."); + this.closed = meter.createCounter(prefix + "closed", UNIT, "Total closed sessions."); this.createTime = meter.createHistogram(prefix + "create_time", "s", "Session creation cost."); meter.createLongGauge(prefix + "max", UNIT, "Configured MaxPoolSize", @@ -67,10 +67,6 @@ public void onSessionCreated() { created.add(1L, poolAttrs); } - public void onSessionDeleted() { - deleted.add(1L, poolAttrs); - } - public void onSessionAcquired() { acquired.add(1L, poolAttrs); } @@ -80,6 +76,10 @@ public void onSessionReleased() { } public void onSessionFailed(Status status) { - failed.add(1L, new Attr[]{poolNameAttr, Attr.of(statusKey, status.getCode().name())}); + failed.add(1L, poolNameAttr, Attr.of(statusKey, status.getCode().name())); + } + + public void onSessionClosed(String reason) { + closed.add(1L, poolNameAttr, Attr.of("reason", reason)); } } diff --git a/table/src/main/java/tech/ydb/table/impl/pool/SessionPool.java b/table/src/main/java/tech/ydb/table/impl/pool/SessionPool.java index a825f8c49..a40a8100c 100644 --- a/table/src/main/java/tech/ydb/table/impl/pool/SessionPool.java +++ b/table/src/main/java/tech/ydb/table/impl/pool/SessionPool.java @@ -9,6 +9,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.LongAdder; import java.util.function.BiConsumer; @@ -160,12 +161,35 @@ private boolean validateSession(ClosableSession session, CompletableFuture create() { } @Override - public void destroy(ClosableSession session) { + public void destroy(ClosableSession session, WaitingQueue.RemovalReason reason) { + session.reportSessionClosed(reason.toString()); stats.deleted.increment(); - metrics.onSessionDeleted(); // Execute deleteSession call outside current context to avoid cancellation and deadline propogation Context ctx = Context.ROOT.fork(); Context previous = ctx.attach(); @@ -436,4 +460,3 @@ public void run() { } } } - diff --git a/table/src/main/java/tech/ydb/table/impl/pool/StatefulSession.java b/table/src/main/java/tech/ydb/table/impl/pool/StatefulSession.java index 745f9032f..9805ad84e 100644 --- a/table/src/main/java/tech/ydb/table/impl/pool/StatefulSession.java +++ b/table/src/main/java/tech/ydb/table/impl/pool/StatefulSession.java @@ -36,11 +36,16 @@ protected void updateSessionState(Throwable th, StatusCode code, boolean gracefu while (!state.compareAndSet(current, current.updated(clock.instant(), th, code, gracefulShutdown))) { current = state.get(); } + if (state.get().needShutdown()) { + onSessionClosed(th, code, gracefulShutdown); + } if (logger.isTraceEnabled()) { logger.trace("{} updated => {}, {}", this, state.get().status, state.get().lastUpdate); } } + protected abstract void onSessionClosed(Throwable th, StatusCode code, boolean gracefulShutdown); + public State state() { return state.get(); } @@ -77,6 +82,7 @@ private State(Status status, Instant lastActive, Instant lastUpdate) { public Instant lastActive() { return this.lastActive; } + public Instant lastUpdate() { return this.lastUpdate; } diff --git a/table/src/main/java/tech/ydb/table/impl/pool/WaitingQueue.java b/table/src/main/java/tech/ydb/table/impl/pool/WaitingQueue.java index c4a23d462..321116457 100644 --- a/table/src/main/java/tech/ydb/table/impl/pool/WaitingQueue.java +++ b/table/src/main/java/tech/ydb/table/impl/pool/WaitingQueue.java @@ -29,9 +29,27 @@ public class WaitingQueue implements AutoCloseable { private static final Logger logger = LoggerFactory.getLogger(WaitingQueue.class); + public enum RemovalReason { + EXPLICIT("explicit"), + IDLE_TIMEOUT("pool_idle_timeout"), + POOL_RESIZE("pool_resize"), + POOL_CLOSE("pool_graceful_shutdown"); + + private final String value; + + RemovalReason(String value) { + this.value = value; + } + + @Override + public String toString() { + return value; + } + } + public interface Handler { CompletableFuture create(); - void destroy(T object); + void destroy(T object, RemovalReason reason); } /** Limit of waiting requests = maxSize * constant */ @@ -112,7 +130,7 @@ public void release(T object) { // if queue is overflowed if (queueSize.get() > limits.maxSize) { queueSize.decrementAndGet(); - handler.destroy(object); + handler.destroy(object, RemovalReason.POOL_RESIZE); return; } @@ -135,7 +153,7 @@ public void delete(T object) { return; } queueSize.decrementAndGet(); - handler.destroy(object); + handler.destroy(object, RemovalReason.EXPLICIT); // After deleting one object we can try to create new pending if it needed checkNextWaitingAcquire(); @@ -185,7 +203,7 @@ private boolean safeAcquireObject(CompletableFuture acquire, T object) { acquire.completeExceptionally(new CancellationException("Queue is already closed")); if (used.remove(object, object)) { queueSize.decrementAndGet(); - handler.destroy(object); + handler.destroy(object, RemovalReason.POOL_CLOSE); } return true; } @@ -299,7 +317,7 @@ private void clear() { T nextIdle = idle.poll(); while (nextIdle != null) { queueSize.decrementAndGet(); - handler.destroy(nextIdle); + handler.destroy(nextIdle, RemovalReason.POOL_CLOSE); nextIdle = idle.poll(); } } @@ -329,7 +347,7 @@ public void accept(T object, Throwable th) { if (!pendingRequests.remove(pending, pending)) { acquire.completeExceptionally(new CancellationException("Queue is already closed")); if (object != null) { - handler.destroy(object); + handler.destroy(object, RemovalReason.POOL_CLOSE); } return; } @@ -378,7 +396,7 @@ public void remove() { return; } if (idle.removeLastOccurrence(lastRet)) { - handler.destroy(lastRet); + handler.destroy(lastRet, RemovalReason.IDLE_TIMEOUT); lastRet = null; queueSize.decrementAndGet(); checkNextWaitingAcquire(); diff --git a/table/src/test/java/tech/ydb/table/impl/pool/MockedTableRpc.java b/table/src/test/java/tech/ydb/table/impl/pool/MockedTableRpc.java index 33b0ba978..2d8ef4cf2 100644 --- a/table/src/test/java/tech/ydb/table/impl/pool/MockedTableRpc.java +++ b/table/src/test/java/tech/ydb/table/impl/pool/MockedTableRpc.java @@ -28,6 +28,7 @@ public class MockedTableRpc extends TableRpcStub { private final static Status BAD_SESSION = Status.of(StatusCode.BAD_SESSION); private final static Status OVERLOADED = Status.of(StatusCode.OVERLOADED); + private final static Status INTERNAL_ERROR = Status.of(StatusCode.INTERNAL_ERROR); private final static Status TRANSPORT_UNAVAILABLE = Status.of(StatusCode.TRANSPORT_UNAVAILABLE); private final Clock clock; @@ -230,6 +231,14 @@ public void completeOverloaded() { future.complete(Result.fail(OVERLOADED)); } + public void completeBadSession() { + future.complete(Result.fail(BAD_SESSION)); + } + + public void completeInternalError() { + future.complete(Result.fail(INTERNAL_ERROR)); + } + public void completeTransportUnavailable() { if (!activeSessions.contains(request.getSessionId())) { future.complete(Result.fail(BAD_SESSION)); diff --git a/table/src/test/java/tech/ydb/table/impl/pool/PoolMetricsTest.java b/table/src/test/java/tech/ydb/table/impl/pool/PoolMetricsTest.java index bdfa0861f..3159e06a3 100644 --- a/table/src/test/java/tech/ydb/table/impl/pool/PoolMetricsTest.java +++ b/table/src/test/java/tech/ydb/table/impl/pool/PoolMetricsTest.java @@ -20,6 +20,8 @@ import tech.ydb.core.metrics.LongMeasurement; import tech.ydb.core.metrics.Meter; import tech.ydb.table.Session; +import tech.ydb.table.query.DataQueryResult; +import tech.ydb.table.transaction.TxControl; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyDouble; @@ -31,6 +33,7 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; public class PoolMetricsTest extends FutureHelper { @@ -76,11 +79,11 @@ public void allInstrumentsAreCreated() { SessionPool pool = createPool(0, 2); verify(meter).createCounter(eq(PREFIX + "created"), eq("{session}"), anyString()); - verify(meter).createCounter(eq(PREFIX + "deleted"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "acquired"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "released"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "requested"), eq("{session}"), anyString()); verify(meter).createCounter(eq(PREFIX + "failed"), eq("{session}"), anyString()); + verify(meter).createCounter(eq(PREFIX + "closed"), eq("{session}"), anyString()); verify(meter).createHistogram(eq(PREFIX + "create_time"), eq("s"), anyString()); verify(meter).createLongGauge(eq(PREFIX + "max"), eq("{session}"), anyString(), any()); verify(meter).createLongGauge(eq(PREFIX + "min"), eq("{session}"), anyString(), any()); @@ -107,12 +110,63 @@ public void sessionLifecycleRecordsCounters() { verify(counter("released")).add(eq(1L), argThat(poolName)); pool.close(); - verify(counter("deleted")).add(eq(1L), argThat(poolName)); + verifyClosed("pool_graceful_shutdown"); tableRpc.completeSessionDeleteRequests(); verify(counter("failed"), never()).add(anyLong(), any()); } + @Test + public void gracefulShutdownRecordsClosedReason() { + SessionPool pool = createPool(0, 2); + CompletableFuture> acquire = pendingFuture(pool.acquire(TIMEOUT)); + tableRpc.nextCreateSession().completeSuccess(); + Session session = futureIsReady(acquire).getValue(); + + CompletableFuture> query = session.executeDataQuery("SELECT 1", TxControl.onlineRo()); + tableRpc.nextExecuteDataQuery().completeSuccessWithShutdownHook(); + futureIsReady(query); + session.close(); + + verifyClosed("session_shutdown"); + tableRpc.completeSessionDeleteRequests(); + pool.close(); + } + + @Test + public void badSessionRecordsClosedReason() { + SessionPool pool = createPool(0, 2); + CompletableFuture> acquire = pendingFuture(pool.acquire(TIMEOUT)); + tableRpc.nextCreateSession().completeSuccess(); + Session session = futureIsReady(acquire).getValue(); + + CompletableFuture> query = session.executeDataQuery("SELECT 1", TxControl.onlineRo()); + tableRpc.nextExecuteDataQuery().completeBadSession(); + futureIsReady(query); + session.close(); + + verifyClosed("bad_session"); + tableRpc.completeSessionDeleteRequests(); + pool.close(); + } + + @Test + public void internalErrorRecordsClosedReason() { + SessionPool pool = createPool(0, 2); + CompletableFuture> acquire = pendingFuture(pool.acquire(TIMEOUT)); + tableRpc.nextCreateSession().completeSuccess(); + Session session = futureIsReady(acquire).getValue(); + + CompletableFuture> query = session.executeDataQuery("SELECT 1", TxControl.onlineRo()); + tableRpc.nextExecuteDataQuery().completeInternalError(); + futureIsReady(query); + session.close(); + + verifyClosed("internal_error"); + tableRpc.completeSessionDeleteRequests(); + pool.close(); + } + @Test public void failedCreateRecordsFailedCounter() { SessionPool pool = createPool(0, 2); @@ -180,6 +234,15 @@ private LongCounter counter(String shortName) { return counters.get(PREFIX + shortName); } + private void verifyClosed(String reason) { + verify(counter("closed")).add( + eq(1L), + argThat(poolName), + argThat(a -> a.getKey().equals("reason") && a.getValue().equals(reason)) + ); + verifyNoMoreInteractions(counter("closed")); + } + private static boolean attr(Attr attr, String shortKey, String value) { return attr.getKey().equals(PREFIX + shortKey) && attr.getValue().equals(value); } diff --git a/table/src/test/java/tech/ydb/table/impl/pool/WaitingQueueTest.java b/table/src/test/java/tech/ydb/table/impl/pool/WaitingQueueTest.java index 76233a173..2e9e29d36 100644 --- a/table/src/test/java/tech/ydb/table/impl/pool/WaitingQueueTest.java +++ b/table/src/test/java/tech/ydb/table/impl/pool/WaitingQueueTest.java @@ -43,7 +43,7 @@ public CompletableFuture create() { } @Override - public void destroy(Resource object) { + public void destroy(Resource object, WaitingQueue.RemovalReason reason) { active.remove(object); }