Skip to content

Commit 86b842c

Browse files
committed
[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.
1 parent 40e7ce6 commit 86b842c

2 files changed

Lines changed: 237 additions & 4 deletions

File tree

‎iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java‎

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -969,15 +969,16 @@ private RetryResult<TSStatus> callWithRetry(TFunction<TSStatus> rpc) {
969969
return new RetryResult<>(status, lastTException, i);
970970
}
971971

972-
private RetryResult<TSStatus> callWithRetryAndReconnect(TFunction<TSStatus> rpc) {
972+
private RetryResult<TSStatus> callWithRetryAndReconnect(TFunction<TSStatus> rpc)
973+
throws IoTDBConnectionException {
973974
return callWithRetryAndReconnect(
974975
rpc,
975976
status -> status.isSetNeedRetry() && status.isNeedRetry(),
976977
status -> status.getCode() == TSStatusCode.PLAN_FAILED_NETWORK_PARTITION.getStatusCode());
977978
}
978979

979980
private <T> RetryResult<T> callWithRetryAndReconnect(
980-
TFunction<T> rpc, Function<T, TSStatus> statusGetter) {
981+
TFunction<T> rpc, Function<T, TSStatus> statusGetter) throws IoTDBConnectionException {
981982
return callWithRetryAndReconnect(
982983
rpc,
983984
t -> {
@@ -989,9 +990,15 @@ private <T> RetryResult<T> callWithRetryAndReconnect(
989990
== TSStatusCode.PLAN_FAILED_NETWORK_PARTITION.getStatusCode());
990991
}
991992

992-
/** reconnect if the remote datanode is unreachable retry if the status is set to needRetry */
993+
/**
994+
* reconnect if the remote datanode is unreachable retry if the status is set to needRetry
995+
*
996+
* @throws IoTDBConnectionException if no attempt produced a result, i.e. the last attempt failed
997+
* with a TException. The TException is the cause.
998+
*/
993999
private <T> RetryResult<T> callWithRetryAndReconnect(
994-
TFunction<T> rpc, Predicate<T> shouldRetry, Predicate<T> forceReconnect) {
1000+
TFunction<T> rpc, Predicate<T> shouldRetry, Predicate<T> forceReconnect)
1001+
throws IoTDBConnectionException {
9951002
TException lastTException = null;
9961003
T result = null;
9971004
int retryAttempt;
@@ -1039,6 +1046,10 @@ private <T> RetryResult<T> callWithRetryAndReconnect(
10391046
}
10401047
}
10411048

1049+
if (result == null) {
1050+
// all attempts failed with a TException, callers must not see a null result
1051+
throw new IoTDBConnectionException(lastTException);
1052+
}
10421053
return new RetryResult<>(result, lastTException, retryAttempt);
10431054
}
10441055

Lines changed: 222 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,222 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.session;
21+
22+
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
23+
import org.apache.iotdb.isession.ISession;
24+
import org.apache.iotdb.isession.SessionConfig;
25+
import org.apache.iotdb.rpc.DeepCopyRpcTransportFactory;
26+
import org.apache.iotdb.rpc.IoTDBConnectionException;
27+
import org.apache.iotdb.rpc.RpcUtils;
28+
import org.apache.iotdb.service.rpc.thrift.IClientRPCService;
29+
import org.apache.iotdb.service.rpc.thrift.TSExecuteStatementReq;
30+
import org.apache.iotdb.session.pool.SessionPool;
31+
32+
import org.apache.thrift.protocol.TBinaryProtocol;
33+
import org.apache.thrift.transport.TTransport;
34+
import org.apache.thrift.transport.TTransportException;
35+
import org.junit.After;
36+
import org.junit.Before;
37+
import org.junit.Test;
38+
import org.powermock.reflect.Whitebox;
39+
40+
import java.net.ServerSocket;
41+
import java.util.Collections;
42+
import java.util.List;
43+
import java.util.concurrent.ConcurrentLinkedDeque;
44+
import java.util.function.Supplier;
45+
46+
import static org.junit.Assert.assertEquals;
47+
import static org.junit.Assert.assertFalse;
48+
import static org.junit.Assert.assertTrue;
49+
import static org.junit.Assert.fail;
50+
51+
/**
52+
* A broken connection must fail with {@link IoTDBConnectionException}, not a NullPointerException,
53+
* so that SessionPool evicts the session instead of handing it out again.
54+
*
55+
* <p>Setup: a {@link SessionConnection} whose Thrift transport was closed (as {@code reconnect()}
56+
* does) and whose reconnect targets are unreachable. Every RPC fails with "Cannot write to null
57+
* outputStream", every reconnect fails, and {@code callWithRetryAndReconnect} runs out of retries.
58+
*/
59+
public class BrokenSessionConnectionTest {
60+
61+
private static final String SQL = "select count(*) from root.sg.d1";
62+
63+
/** Nothing listens here. Connecting is refused right away. */
64+
private TEndPoint deadEndPoint;
65+
66+
private TTransport closedTransport;
67+
68+
@Before
69+
public void setUp() throws Exception {
70+
// Open a real transport against a short-lived server, then close both. The transport ends up
71+
// in the same state as after SessionConnection.reconnect() closed it: its outputStream is null.
72+
try (ServerSocket server = new ServerSocket(0)) {
73+
deadEndPoint = new TEndPoint("127.0.0.1", server.getLocalPort());
74+
closedTransport =
75+
DeepCopyRpcTransportFactory.getInstance(
76+
SessionConfig.DEFAULT_INITIAL_BUFFER_CAPACITY,
77+
SessionConfig.DEFAULT_MAX_FRAME_SIZE)
78+
.getTransport(deadEndPoint.getIp(), deadEndPoint.getPort(), 1000);
79+
closedTransport.open();
80+
}
81+
closedTransport.close();
82+
}
83+
84+
@After
85+
public void tearDown() {
86+
if (closedTransport != null) {
87+
closedTransport.close();
88+
}
89+
}
90+
91+
/**
92+
* This just checks if the test setup works as expected, meaning we have a TTransport with a null
93+
* outputStream.
94+
*/
95+
@Test
96+
public void closedTransportFailsWithNullOutputStream() throws Exception {
97+
// Sanity check: the client fails with the same error as in the bug report.
98+
IClientRPCService.Iface client = newClient(closedTransport);
99+
try {
100+
client.executeQueryStatementV2(new TSExecuteStatementReq(0, SQL, 0));
101+
fail("expected TTransportException");
102+
} catch (TTransportException e) {
103+
assertEquals("Cannot write to null outputStream", e.getMessage());
104+
}
105+
}
106+
107+
/**
108+
* Retries run out and {@code executeQueryStatement} should throw {@link IoTDBConnectionException}
109+
* with the last TException as the cause. On unfixed code it throws a NullPointerException from
110+
* {@code execResp.getStatus()}.
111+
*/
112+
@Test
113+
public void exhaustedRetriesThrowConnectionExceptionInsteadOfNpe() throws Exception {
114+
SessionConnection connection = newBrokenConnection(newSession());
115+
116+
try {
117+
connection.executeQueryStatement(SQL, 1000);
118+
fail("expected IoTDBConnectionException");
119+
} catch (IoTDBConnectionException e) {
120+
assertTrue(
121+
"cause should be the retained TException, was " + e.getCause(),
122+
e.getCause() instanceof TTransportException);
123+
} catch (NullPointerException e) {
124+
throw new AssertionError("exhausted retries surfaced as NPE, root cause was lost", e);
125+
}
126+
}
127+
128+
/** The same bug on a TSStatus-returning path: the null status NPEs in RpcUtils.verifySuccess. */
129+
@Test
130+
public void exhaustedRetriesOnStatusPathThrowConnectionException() throws Exception {
131+
SessionConnection connection = newBrokenConnection(newSession());
132+
133+
try {
134+
connection.setStorageGroup("root.sg");
135+
fail("expected IoTDBConnectionException");
136+
} catch (IoTDBConnectionException e) {
137+
// expected
138+
} catch (NullPointerException e) {
139+
throw new AssertionError("exhausted retries surfaced as NPE, root cause was lost", e);
140+
}
141+
}
142+
143+
/**
144+
* A session whose connection fails after all retries now surfaces {@link
145+
* IoTDBConnectionException}, so SessionPool's existing eviction path removes it. Before, the NPE
146+
* took the {@code RuntimeException -> putBack} branch and the session was queued again.
147+
*/
148+
@Test
149+
public void brokenConnectionLeadsToSessionEviction() throws Exception {
150+
Session brokenSession = newSession();
151+
Whitebox.setInternalState(
152+
brokenSession, "defaultSessionConnection", newBrokenConnection(brokenSession));
153+
154+
SessionPool pool =
155+
new SessionPool.Builder()
156+
.nodeUrls(
157+
Collections.singletonList(deadEndPoint.getIp() + ":" + deadEndPoint.getPort()))
158+
.user("root")
159+
.password("root")
160+
.maxSize(1)
161+
.waitToGetSessionTimeoutInMs(1000)
162+
.enableAutoFetch(false)
163+
.build();
164+
try {
165+
ConcurrentLinkedDeque<ISession> queue = Whitebox.getInternalState(pool, "queue");
166+
queue.add(brokenSession);
167+
Whitebox.setInternalState(pool, "size", 1);
168+
169+
// Each call fails, but with a connection error, and the broken session must not be reused.
170+
for (int call = 1; call <= 3; call++) {
171+
try {
172+
pool.executeQueryStatement(SQL);
173+
fail("call " + call + ": expected IoTDBConnectionException");
174+
} catch (IoTDBConnectionException e) {
175+
// expected: the server is unreachable
176+
} catch (NullPointerException e) {
177+
throw new AssertionError(
178+
"call "
179+
+ call
180+
+ ": the pool reused the broken session"
181+
+ (queue.contains(brokenSession) ? " (still queued)" : ""),
182+
e);
183+
}
184+
assertFalse(
185+
"call " + call + ": broken session was put back into the pool",
186+
queue.contains(brokenSession));
187+
}
188+
} finally {
189+
pool.close();
190+
}
191+
}
192+
193+
private Session newSession() {
194+
return new Session.Builder()
195+
.nodeUrls(Collections.singletonList(deadEndPoint.getIp() + ":" + deadEndPoint.getPort()))
196+
.username("root")
197+
.password("root")
198+
.enableAutoFetch(false)
199+
.build();
200+
}
201+
202+
/**
203+
* A connection whose transport was closed and whose reconnect targets are unreachable. Retry
204+
* interval is 1 ms to keep the test fast (the default is 500 ms with 11 attempts).
205+
*/
206+
private SessionConnection newBrokenConnection(Session session) {
207+
SessionConnection connection = new SessionConnection(Session.TREE);
208+
Whitebox.setInternalState(connection, "session", session);
209+
Whitebox.setInternalState(connection, "transport", closedTransport);
210+
Whitebox.setInternalState(connection, "client", newClient(closedTransport));
211+
Whitebox.setInternalState(connection, "endPoint", deadEndPoint);
212+
Supplier<List<TEndPoint>> availableNodes = () -> Collections.singletonList(deadEndPoint);
213+
Whitebox.setInternalState(connection, "availableNodes", availableNodes);
214+
Whitebox.setInternalState(connection, "retryIntervalInMs", 1L);
215+
return connection;
216+
}
217+
218+
private static IClientRPCService.Iface newClient(TTransport transport) {
219+
return RpcUtils.newSynchronizedClient(
220+
new IClientRPCService.Client(new TBinaryProtocol(transport)));
221+
}
222+
}

0 commit comments

Comments
 (0)