Skip to content

Commit 4f65a99

Browse files
committed
Support IPv6 endpoint URLs
1 parent f3e1612 commit 4f65a99

31 files changed

Lines changed: 553 additions & 132 deletions

File tree

‎integration-test/src/main/java/org/apache/iotdb/it/env/cluster/env/AbstractEnv.java‎

Lines changed: 12 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@
6464
import org.apache.iotdb.jdbc.IoTDBConnection;
6565
import org.apache.iotdb.rpc.IoTDBConnectionException;
6666
import org.apache.iotdb.rpc.TSStatusCode;
67+
import org.apache.iotdb.rpc.UrlUtils;
6768
import org.apache.iotdb.session.Session;
6869
import org.apache.iotdb.session.TableSessionBuilder;
6970
import org.apache.iotdb.session.pool.SessionPool;
@@ -150,9 +151,8 @@ public List<String> getMetricPrometheusReporterContents(String authHeader) {
150151
final String configNodeMetricContent =
151152
getUrlContent(
152153
Config.IOTDB_HTTP_URL_PREFIX
153-
+ configNode.getIp()
154-
+ ":"
155-
+ configNode.getMetricPort()
154+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(
155+
configNode.getIp(), configNode.getMetricPort())
156156
+ "/metrics",
157157
authHeader);
158158
result.add(configNodeMetricContent);
@@ -162,9 +162,8 @@ public List<String> getMetricPrometheusReporterContents(String authHeader) {
162162
final String dataNodeMetricContent =
163163
getUrlContent(
164164
Config.IOTDB_HTTP_URL_PREFIX
165-
+ dataNode.getIp()
166-
+ ":"
167-
+ dataNode.getMetricPort()
165+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(
166+
dataNode.getIp(), dataNode.getMetricPort())
168167
+ "/metrics",
169168
authHeader);
170169
result.add(dataNodeMetricContent);
@@ -1038,7 +1037,8 @@ protected NodeConnection getWriteConnectionWithSpecifiedDataNode(
10381037
final String password,
10391038
final String sqlDialect)
10401039
throws SQLException {
1041-
final String endpoint = dataNode.getIp() + ":" + dataNode.getPort();
1040+
final String endpoint =
1041+
UrlUtils.formatTEndPointIpv4AndIpv6Url(dataNode.getIp(), dataNode.getPort());
10421042
final Connection writeConnection =
10431043
DriverManager.getConnection(
10441044
Config.IOTDB_URL_PREFIX
@@ -1618,18 +1618,14 @@ public void ensureNodeStatus(
16181618
.forEach(
16191619
node ->
16201620
nodeIds.put(
1621-
node.getInternalEndPoint().getIp()
1622-
+ ":"
1623-
+ node.getInternalEndPoint().getPort(),
1621+
UrlUtils.convertTEndPointIpv4AndIpv6Url(node.getInternalEndPoint()),
16241622
node.getConfigNodeId()));
16251623
showClusterResp
16261624
.getDataNodeList()
16271625
.forEach(
16281626
node ->
16291627
nodeIds.put(
1630-
node.getClientRpcEndPoint().getIp()
1631-
+ ":"
1632-
+ node.getClientRpcEndPoint().getPort(),
1628+
UrlUtils.convertTEndPointIpv4AndIpv6Url(node.getClientRpcEndPoint()),
16331629
node.getDataNodeId()));
16341630
for (int j = 0; j < nodes.size(); j++) {
16351631
BaseNodeWrapper nodeWrapper = nodes.get(j);
@@ -1652,10 +1648,9 @@ public void ensureNodeStatus(
16521648
continue;
16531649
}
16541650
if (nodeWrapper instanceof DataNodeWrapper && targetStatus.equals(NodeStatus.Running)) {
1655-
final String[] ipPort = nodeWrapper.getIpAndPortString().split(":");
1656-
final String ip = ipPort[0];
1657-
final int port = Integer.parseInt(ipPort[1]);
1658-
try (TSocket socket = new TSocket(new TConfiguration(), ip, port, 1000)) {
1651+
try (TSocket socket =
1652+
new TSocket(
1653+
new TConfiguration(), nodeWrapper.getIp(), nodeWrapper.getPort(), 1000)) {
16591654
socket.open();
16601655
} catch (final TTransportException e) {
16611656
errorMessages.add(

‎integration-test/src/main/java/org/apache/iotdb/it/env/cluster/node/AbstractNodeWrapper.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import org.apache.iotdb.it.env.cluster.config.MppJVMConfig;
2828
import org.apache.iotdb.it.framework.IoTDBTestLogger;
2929
import org.apache.iotdb.itbase.env.BaseNodeWrapper;
30+
import org.apache.iotdb.rpc.UrlUtils;
3031

3132
import org.apache.tsfile.external.commons.io.FileUtils;
3233
import org.apache.tsfile.external.commons.io.file.PathUtils;
@@ -606,7 +607,7 @@ public void setPort(int port) {
606607

607608
@Override
608609
public final String getIpAndPortString() {
609-
return this.getIp() + ":" + this.getPort();
610+
return UrlUtils.formatTEndPointIpv4AndIpv6Url(this.getIp(), this.getPort());
610611
}
611612

612613
protected String workDirFilePath(String dirName, String fileName) {

‎integration-test/src/main/java/org/apache/iotdb/it/env/remote/env/RemoteServerEnv.java‎

Lines changed: 13 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@
4242
import org.apache.iotdb.jdbc.Config;
4343
import org.apache.iotdb.jdbc.Constant;
4444
import org.apache.iotdb.rpc.IoTDBConnectionException;
45+
import org.apache.iotdb.rpc.UrlUtils;
4546
import org.apache.iotdb.session.Session;
4647
import org.apache.iotdb.session.TableSessionBuilder;
4748
import org.apache.iotdb.session.pool.SessionPool;
@@ -72,6 +73,10 @@ public class RemoteServerEnv implements BaseEnv {
7273
private IClientManager<TEndPoint, SyncConfigNodeIServiceClient> clientManager;
7374
private RemoteClusterConfig clusterConfig = new RemoteClusterConfig();
7475

76+
private String remoteEndpointUrl(String port) {
77+
return UrlUtils.formatTEndPointIpv4AndIpv6Url(ip_addr, Integer.parseInt(port));
78+
}
79+
7580
@Override
7681
public void initClusterEnvironment() {
7782
try (Connection connection = EnvFactory.getEnv().getConnection();
@@ -116,11 +121,11 @@ public List<String> getMetricPrometheusReporterContents(String authHeader) {
116121
List<String> result = new ArrayList<>();
117122
result.add(
118123
getUrlContent(
119-
Config.IOTDB_HTTP_URL_PREFIX + ip_addr + ":" + configNodeMetricPort + "/metrics",
124+
Config.IOTDB_HTTP_URL_PREFIX + remoteEndpointUrl(configNodeMetricPort) + "/metrics",
120125
authHeader));
121126
result.add(
122127
getUrlContent(
123-
Config.IOTDB_HTTP_URL_PREFIX + ip_addr + ":" + dataNodeMetricPort + "/metrics",
128+
Config.IOTDB_HTTP_URL_PREFIX + remoteEndpointUrl(dataNodeMetricPort) + "/metrics",
124129
authHeader));
125130
return result;
126131
}
@@ -133,7 +138,7 @@ public Connection getConnection(String username, String password, String sqlDial
133138
Class.forName(Config.JDBC_DRIVER_NAME);
134139
connection =
135140
DriverManager.getConnection(
136-
Config.IOTDB_URL_PREFIX + ip_addr + ":" + port,
141+
Config.IOTDB_URL_PREFIX + remoteEndpointUrl(port),
137142
BaseEnv.constructProperties(username, password, sqlDialect));
138143
} catch (ClassNotFoundException e) {
139144
e.printStackTrace();
@@ -170,9 +175,7 @@ public Connection getConnection(
170175
connection =
171176
DriverManager.getConnection(
172177
Config.IOTDB_URL_PREFIX
173-
+ ip_addr
174-
+ ":"
175-
+ port
178+
+ remoteEndpointUrl(port)
176179
+ "?"
177180
+ VERSION
178181
+ "="
@@ -300,15 +303,15 @@ public ISession getSessionConnection(ZoneId zoneId) throws IoTDBConnectionExcept
300303
@Override
301304
public ITableSession getTableSessionConnection() throws IoTDBConnectionException {
302305
return new TableSessionBuilder()
303-
.nodeUrls(Collections.singletonList(ip_addr + ":" + port))
306+
.nodeUrls(Collections.singletonList(remoteEndpointUrl(port)))
304307
.build();
305308
}
306309

307310
@Override
308311
public ITableSession getTableSessionConnectionWithDB(String database)
309312
throws IoTDBConnectionException {
310313
return new TableSessionBuilder()
311-
.nodeUrls(Collections.singletonList(ip_addr + ":" + port))
314+
.nodeUrls(Collections.singletonList(remoteEndpointUrl(port)))
312315
.database(database)
313316
.build();
314317
}
@@ -332,7 +335,7 @@ public ITableSession getTableSessionConnection(List<String> nodeUrls)
332335
public ITableSession getTableSessionConnection(String userName, String password)
333336
throws IoTDBConnectionException {
334337
return new TableSessionBuilder()
335-
.nodeUrls(Collections.singletonList(ip_addr + ":" + port))
338+
.nodeUrls(Collections.singletonList(remoteEndpointUrl(port)))
336339
.username(userName)
337340
.password(password)
338341
.build();
@@ -356,7 +359,7 @@ public ISession getSessionConnection(String userName, String password)
356359
public ISession getSessionConnection(List<String> nodeUrls) throws IoTDBConnectionException {
357360
Session session =
358361
new Session.Builder()
359-
.nodeUrls(Collections.singletonList(ip_addr + ":" + port))
362+
.nodeUrls(Collections.singletonList(remoteEndpointUrl(port)))
360363
.username(SessionConfig.DEFAULT_USER)
361364
.password(SessionConfig.DEFAULT_PASSWORD)
362365
.fetchSize(SessionConfig.DEFAULT_FETCH_SIZE)

‎integration-test/src/test/java/org/apache/iotdb/db/it/iotconsensusv2/IoTDBIoTConsensusV23C3DBasicITBase.java‎

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import org.apache.iotdb.it.env.EnvFactory;
2626
import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper;
2727
import org.apache.iotdb.itbase.env.BaseEnv;
28+
import org.apache.iotdb.rpc.UrlUtils;
2829

2930
import org.apache.tsfile.utils.Pair;
3031
import org.awaitility.Awaitility;
@@ -283,7 +284,9 @@ public void testDeleteTimeSeriesReplicaConsistency() throws Exception {
283284
LOGGER.info("Step 6: Verifying schema consistency on each DataNode independently...");
284285
List<DataNodeWrapper> dataNodeWrappers = EnvFactory.getEnv().getDataNodeWrapperList();
285286
for (DataNodeWrapper wrapper : dataNodeWrappers) {
286-
String nodeDescription = "DataNode " + wrapper.getIp() + ":" + wrapper.getPort();
287+
String nodeDescription =
288+
"DataNode "
289+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(wrapper.getIp(), wrapper.getPort());
287290
LOGGER.info("Verifying schema on {}", nodeDescription);
288291
Awaitility.await()
289292
.atMost(60, TimeUnit.SECONDS)
@@ -307,7 +310,10 @@ public void testDeleteTimeSeriesReplicaConsistency() throws Exception {
307310
LOGGER.info(
308311
"Step 7: Stopping each DataNode in turn and verifying remaining nodes show consistent schema...");
309312
for (DataNodeWrapper stoppedNode : dataNodeWrappers) {
310-
String stoppedDesc = "DataNode " + stoppedNode.getIp() + ":" + stoppedNode.getPort();
313+
String stoppedDesc =
314+
"DataNode "
315+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(
316+
stoppedNode.getIp(), stoppedNode.getPort());
311317
LOGGER.info("Stopping {}", stoppedDesc);
312318
stoppedNode.stopForcibly();
313319
Assert.assertFalse(stoppedDesc + " should be stopped", stoppedNode.isAlive());
@@ -318,7 +324,10 @@ public void testDeleteTimeSeriesReplicaConsistency() throws Exception {
318324
if (aliveNode == stoppedNode) {
319325
continue;
320326
}
321-
String aliveDesc = "DataNode " + aliveNode.getIp() + ":" + aliveNode.getPort();
327+
String aliveDesc =
328+
"DataNode "
329+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(
330+
aliveNode.getIp(), aliveNode.getPort());
322331
Awaitility.await()
323332
.pollDelay(1, TimeUnit.SECONDS)
324333
.atMost(90, TimeUnit.SECONDS)
@@ -446,7 +455,9 @@ private void verifyTimeSeriesAfterDelete(Statement statement, String context) th
446455
protected void waitForReplicationComplete(DataNodeWrapper leaderNode) {
447456
final long timeoutSeconds = 120;
448457
final String metricsUrl =
449-
"http://" + leaderNode.getIp() + ":" + leaderNode.getMetricPort() + "/metrics";
458+
"http://"
459+
+ UrlUtils.formatTEndPointIpv4AndIpv6Url(leaderNode.getIp(), leaderNode.getMetricPort())
460+
+ "/metrics";
450461
LOGGER.info(
451462
"Waiting for consensus pipe syncLag to reach 0 on leader DataNode (url: {}, timeout: {}s)...",
452463
metricsUrl,

‎integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/manual/enhanced/IoTDBPipeClusterIT.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ public void testMachineDowntimeSync() {
130130
private void testMachineDowntime(String sink) {
131131
StringBuilder a = new StringBuilder();
132132
for (DataNodeWrapper nodeWrapper : receiverEnv.getDataNodeWrapperList()) {
133-
a.append(nodeWrapper.getIp()).append(":").append(nodeWrapper.getPort());
133+
a.append(nodeWrapper.getIpAndPortString());
134134
a.append(",");
135135
}
136136
a.deleteCharAt(a.length() - 1);

‎integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/enhanced/IoTDBPipeClusterIT.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,7 @@ public void testMachineDowntimeSync() {
134134
private void testMachineDowntime(String sink) {
135135
StringBuilder a = new StringBuilder();
136136
for (DataNodeWrapper nodeWrapper : receiverEnv.getDataNodeWrapperList()) {
137-
a.append(nodeWrapper.getIp()).append(":").append(nodeWrapper.getPort());
137+
a.append(nodeWrapper.getIpAndPortString());
138138
a.append(",");
139139
}
140140
a.deleteCharAt(a.length() - 1);

‎iotdb-client/client-cpp/src/rpc/NodesSupplier.cpp‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
* under the License.
1818
*/
1919
#include "NodesSupplier.h"
20+
#include "RpcCommon.h"
2021
#include "Session.h"
2122
#include "SessionDataSet.h"
2223
#include <algorithm>
@@ -188,7 +189,7 @@ std::vector<TEndPoint> NodesSupplier::fetchLatestEndpoints() {
188189
port = record->fields.at(columnPortIdx).intV.value();
189190
}
190191

191-
if (ip == "0.0.0.0") {
192+
if (UrlUtils::isWildcardAddress(ip)) {
192193
log_warn("Skipping invalid node: " + ip + ":" + std::to_string(port));
193194
continue;
194195
}
@@ -230,4 +231,4 @@ void NodesSupplier::stopBackgroundRefresh() noexcept {
230231
refreshThread_.join();
231232
}
232233
}
233-
}
234+
}

‎iotdb-client/client-cpp/src/rpc/RpcCommon.cpp‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -212,3 +212,27 @@ TEndPoint UrlUtils::parseTEndPointIpv4AndIpv6Url(const string& endPointUrl) {
212212

213213
return endPoint;
214214
}
215+
216+
bool UrlUtils::isWildcardAddress(const string& host) {
217+
if (host.empty()) {
218+
return false;
219+
}
220+
221+
string normalizedHost = host;
222+
if (normalizedHost.size() >= 2 && normalizedHost.front() == '[' && normalizedHost.back() == ']') {
223+
normalizedHost = normalizedHost.substr(1, normalizedHost.size() - 2);
224+
}
225+
226+
if (normalizedHost == "0.0.0.0") {
227+
return true;
228+
}
229+
if (normalizedHost.find(PORT_SEPARATOR) == string::npos) {
230+
return false;
231+
}
232+
for (char ch : normalizedHost) {
233+
if (ch != ':' && ch != '0') {
234+
return false;
235+
}
236+
}
237+
return true;
238+
}

‎iotdb-client/client-cpp/src/rpc/RpcCommon.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,8 @@ class UrlUtils {
7676

7777
public:
7878
static TEndPoint parseTEndPointIpv4AndIpv6Url(const std::string& endPointUrl);
79+
80+
static bool isWildcardAddress(const std::string& host);
7981
};
8082

8183
#endif

‎iotdb-client/client-cpp/src/session/Session.cpp‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2096,7 +2096,7 @@ void Session::Impl::handleQueryRedirection(TEndPoint endPoint) {
20962096
void Session::Impl::handleRedirection(const std::string& deviceId, TEndPoint endPoint) {
20972097
if (!enableRedirection_)
20982098
return;
2099-
if (endPoint.ip == "0.0.0.0")
2099+
if (UrlUtils::isWildcardAddress(endPoint.ip))
21002100
return;
21012101
getDefaultSessionConnection();
21022102
deviceIdToEndpoint[deviceId] = endPoint;
@@ -2122,7 +2122,7 @@ void Session::Impl::handleRedirection(const std::shared_ptr<storage::IDeviceID>&
21222122
TEndPoint endPoint) {
21232123
if (!enableRedirection_)
21242124
return;
2125-
if (endPoint.ip == "0.0.0.0")
2125+
if (UrlUtils::isWildcardAddress(endPoint.ip))
21262126
return;
21272127
getDefaultSessionConnection();
21282128
tableModelDeviceIdToEndpoint[deviceId] = endPoint;

0 commit comments

Comments
 (0)