From 86b842cb0be7c0bb77f1da9657c8dcdef18dae1c Mon Sep 17 00:00:00 2001 From: Yann Massard Date: Tue, 6 Oct 2026 16:12:14 +0000 Subject: [PATCH] [Session] Throw IoTDBConnectionException instead of NPE when retries are exhausted callWithRetryAndReconnect returned a null result after all attempts failed with a TException. Callers dereferenced it and threw a NullPointerException, losing the real cause. SessionPool treated the NPE as a RuntimeException and put the broken session back into the pool, so every later call failed the same way until the client JVM was restarted. Throw IoTDBConnectionException with the last TException as the cause, so the error is diagnosable and SessionPool evicts the broken session. --- .../iotdb/session/SessionConnection.java | 19 +- .../session/BrokenSessionConnectionTest.java | 222 ++++++++++++++++++ 2 files changed, 237 insertions(+), 4 deletions(-) create mode 100644 iotdb-client/session/src/test/java/org/apache/iotdb/session/BrokenSessionConnectionTest.java diff --git a/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java b/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java index 6bcc41859a8c..6f26246845f1 100644 --- a/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java +++ b/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java @@ -969,7 +969,8 @@ private RetryResult callWithRetry(TFunction rpc) { return new RetryResult<>(status, lastTException, i); } - private RetryResult callWithRetryAndReconnect(TFunction rpc) { + private RetryResult callWithRetryAndReconnect(TFunction rpc) + throws IoTDBConnectionException { return callWithRetryAndReconnect( rpc, status -> status.isSetNeedRetry() && status.isNeedRetry(), @@ -977,7 +978,7 @@ private RetryResult callWithRetryAndReconnect(TFunction rpc) } private RetryResult callWithRetryAndReconnect( - TFunction rpc, Function statusGetter) { + TFunction rpc, Function statusGetter) throws IoTDBConnectionException { return callWithRetryAndReconnect( rpc, t -> { @@ -989,9 +990,15 @@ private RetryResult callWithRetryAndReconnect( == TSStatusCode.PLAN_FAILED_NETWORK_PARTITION.getStatusCode()); } - /** reconnect if the remote datanode is unreachable retry if the status is set to needRetry */ + /** + * reconnect if the remote datanode is unreachable retry if the status is set to needRetry + * + * @throws IoTDBConnectionException if no attempt produced a result, i.e. the last attempt failed + * with a TException. The TException is the cause. + */ private RetryResult callWithRetryAndReconnect( - TFunction rpc, Predicate shouldRetry, Predicate forceReconnect) { + TFunction rpc, Predicate shouldRetry, Predicate forceReconnect) + throws IoTDBConnectionException { TException lastTException = null; T result = null; int retryAttempt; @@ -1039,6 +1046,10 @@ private RetryResult callWithRetryAndReconnect( } } + if (result == null) { + // all attempts failed with a TException, callers must not see a null result + throw new IoTDBConnectionException(lastTException); + } return new RetryResult<>(result, lastTException, retryAttempt); } diff --git a/iotdb-client/session/src/test/java/org/apache/iotdb/session/BrokenSessionConnectionTest.java b/iotdb-client/session/src/test/java/org/apache/iotdb/session/BrokenSessionConnectionTest.java new file mode 100644 index 000000000000..952edb49f646 --- /dev/null +++ b/iotdb-client/session/src/test/java/org/apache/iotdb/session/BrokenSessionConnectionTest.java @@ -0,0 +1,222 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.session; + +import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.isession.ISession; +import org.apache.iotdb.isession.SessionConfig; +import org.apache.iotdb.rpc.DeepCopyRpcTransportFactory; +import org.apache.iotdb.rpc.IoTDBConnectionException; +import org.apache.iotdb.rpc.RpcUtils; +import org.apache.iotdb.service.rpc.thrift.IClientRPCService; +import org.apache.iotdb.service.rpc.thrift.TSExecuteStatementReq; +import org.apache.iotdb.session.pool.SessionPool; + +import org.apache.thrift.protocol.TBinaryProtocol; +import org.apache.thrift.transport.TTransport; +import org.apache.thrift.transport.TTransportException; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.powermock.reflect.Whitebox; + +import java.net.ServerSocket; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.ConcurrentLinkedDeque; +import java.util.function.Supplier; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +/** + * A broken connection must fail with {@link IoTDBConnectionException}, not a NullPointerException, + * so that SessionPool evicts the session instead of handing it out again. + * + *

Setup: a {@link SessionConnection} whose Thrift transport was closed (as {@code reconnect()} + * does) and whose reconnect targets are unreachable. Every RPC fails with "Cannot write to null + * outputStream", every reconnect fails, and {@code callWithRetryAndReconnect} runs out of retries. + */ +public class BrokenSessionConnectionTest { + + private static final String SQL = "select count(*) from root.sg.d1"; + + /** Nothing listens here. Connecting is refused right away. */ + private TEndPoint deadEndPoint; + + private TTransport closedTransport; + + @Before + public void setUp() throws Exception { + // Open a real transport against a short-lived server, then close both. The transport ends up + // in the same state as after SessionConnection.reconnect() closed it: its outputStream is null. + try (ServerSocket server = new ServerSocket(0)) { + deadEndPoint = new TEndPoint("127.0.0.1", server.getLocalPort()); + closedTransport = + DeepCopyRpcTransportFactory.getInstance( + SessionConfig.DEFAULT_INITIAL_BUFFER_CAPACITY, + SessionConfig.DEFAULT_MAX_FRAME_SIZE) + .getTransport(deadEndPoint.getIp(), deadEndPoint.getPort(), 1000); + closedTransport.open(); + } + closedTransport.close(); + } + + @After + public void tearDown() { + if (closedTransport != null) { + closedTransport.close(); + } + } + + /** + * This just checks if the test setup works as expected, meaning we have a TTransport with a null + * outputStream. + */ + @Test + public void closedTransportFailsWithNullOutputStream() throws Exception { + // Sanity check: the client fails with the same error as in the bug report. + IClientRPCService.Iface client = newClient(closedTransport); + try { + client.executeQueryStatementV2(new TSExecuteStatementReq(0, SQL, 0)); + fail("expected TTransportException"); + } catch (TTransportException e) { + assertEquals("Cannot write to null outputStream", e.getMessage()); + } + } + + /** + * Retries run out and {@code executeQueryStatement} should throw {@link IoTDBConnectionException} + * with the last TException as the cause. On unfixed code it throws a NullPointerException from + * {@code execResp.getStatus()}. + */ + @Test + public void exhaustedRetriesThrowConnectionExceptionInsteadOfNpe() throws Exception { + SessionConnection connection = newBrokenConnection(newSession()); + + try { + connection.executeQueryStatement(SQL, 1000); + fail("expected IoTDBConnectionException"); + } catch (IoTDBConnectionException e) { + assertTrue( + "cause should be the retained TException, was " + e.getCause(), + e.getCause() instanceof TTransportException); + } catch (NullPointerException e) { + throw new AssertionError("exhausted retries surfaced as NPE, root cause was lost", e); + } + } + + /** The same bug on a TSStatus-returning path: the null status NPEs in RpcUtils.verifySuccess. */ + @Test + public void exhaustedRetriesOnStatusPathThrowConnectionException() throws Exception { + SessionConnection connection = newBrokenConnection(newSession()); + + try { + connection.setStorageGroup("root.sg"); + fail("expected IoTDBConnectionException"); + } catch (IoTDBConnectionException e) { + // expected + } catch (NullPointerException e) { + throw new AssertionError("exhausted retries surfaced as NPE, root cause was lost", e); + } + } + + /** + * A session whose connection fails after all retries now surfaces {@link + * IoTDBConnectionException}, so SessionPool's existing eviction path removes it. Before, the NPE + * took the {@code RuntimeException -> putBack} branch and the session was queued again. + */ + @Test + public void brokenConnectionLeadsToSessionEviction() throws Exception { + Session brokenSession = newSession(); + Whitebox.setInternalState( + brokenSession, "defaultSessionConnection", newBrokenConnection(brokenSession)); + + SessionPool pool = + new SessionPool.Builder() + .nodeUrls( + Collections.singletonList(deadEndPoint.getIp() + ":" + deadEndPoint.getPort())) + .user("root") + .password("root") + .maxSize(1) + .waitToGetSessionTimeoutInMs(1000) + .enableAutoFetch(false) + .build(); + try { + ConcurrentLinkedDeque queue = Whitebox.getInternalState(pool, "queue"); + queue.add(brokenSession); + Whitebox.setInternalState(pool, "size", 1); + + // Each call fails, but with a connection error, and the broken session must not be reused. + for (int call = 1; call <= 3; call++) { + try { + pool.executeQueryStatement(SQL); + fail("call " + call + ": expected IoTDBConnectionException"); + } catch (IoTDBConnectionException e) { + // expected: the server is unreachable + } catch (NullPointerException e) { + throw new AssertionError( + "call " + + call + + ": the pool reused the broken session" + + (queue.contains(brokenSession) ? " (still queued)" : ""), + e); + } + assertFalse( + "call " + call + ": broken session was put back into the pool", + queue.contains(brokenSession)); + } + } finally { + pool.close(); + } + } + + private Session newSession() { + return new Session.Builder() + .nodeUrls(Collections.singletonList(deadEndPoint.getIp() + ":" + deadEndPoint.getPort())) + .username("root") + .password("root") + .enableAutoFetch(false) + .build(); + } + + /** + * A connection whose transport was closed and whose reconnect targets are unreachable. Retry + * interval is 1 ms to keep the test fast (the default is 500 ms with 11 attempts). + */ + private SessionConnection newBrokenConnection(Session session) { + SessionConnection connection = new SessionConnection(Session.TREE); + Whitebox.setInternalState(connection, "session", session); + Whitebox.setInternalState(connection, "transport", closedTransport); + Whitebox.setInternalState(connection, "client", newClient(closedTransport)); + Whitebox.setInternalState(connection, "endPoint", deadEndPoint); + Supplier> availableNodes = () -> Collections.singletonList(deadEndPoint); + Whitebox.setInternalState(connection, "availableNodes", availableNodes); + Whitebox.setInternalState(connection, "retryIntervalInMs", 1L); + return connection; + } + + private static IClientRPCService.Iface newClient(TTransport transport) { + return RpcUtils.newSynchronizedClient( + new IClientRPCService.Client(new TBinaryProtocol(transport))); + } +}