Skip to content

Commit a10fe5c

Browse files
committed
Test IPv6 REST and MQTT services
1 parent a1eb50c commit a10fe5c

1 file changed

Lines changed: 220 additions & 0 deletions

File tree

Lines changed: 220 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,220 @@
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.db.it;
21+
22+
import org.apache.iotdb.isession.ISession;
23+
import org.apache.iotdb.isession.SessionDataSet;
24+
import org.apache.iotdb.it.env.EnvFactory;
25+
import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper;
26+
import org.apache.iotdb.it.framework.IoTDBTestRunner;
27+
import org.apache.iotdb.it.utils.IPv6TestUtils;
28+
import org.apache.iotdb.itbase.category.ClusterIT;
29+
import org.apache.iotdb.itbase.category.LocalStandaloneIT;
30+
import org.apache.iotdb.rpc.StatementExecutionException;
31+
import org.apache.iotdb.rpc.UrlUtils;
32+
33+
import com.google.gson.JsonObject;
34+
import com.google.gson.JsonParser;
35+
import org.apache.http.client.methods.CloseableHttpResponse;
36+
import org.apache.http.client.methods.HttpGet;
37+
import org.apache.http.client.methods.HttpPost;
38+
import org.apache.http.entity.StringEntity;
39+
import org.apache.http.impl.client.CloseableHttpClient;
40+
import org.apache.http.impl.client.HttpClientBuilder;
41+
import org.apache.http.util.EntityUtils;
42+
import org.apache.tsfile.read.common.Field;
43+
import org.apache.tsfile.read.common.RowRecord;
44+
import org.awaitility.Awaitility;
45+
import org.fusesource.mqtt.client.BlockingConnection;
46+
import org.fusesource.mqtt.client.MQTT;
47+
import org.fusesource.mqtt.client.QoS;
48+
import org.junit.AfterClass;
49+
import org.junit.BeforeClass;
50+
import org.junit.Test;
51+
import org.junit.experimental.categories.Category;
52+
import org.junit.runner.RunWith;
53+
54+
import java.io.IOException;
55+
import java.nio.charset.StandardCharsets;
56+
import java.util.Base64;
57+
import java.util.List;
58+
import java.util.concurrent.TimeUnit;
59+
60+
import static org.apache.iotdb.it.utils.IPv6TestUtils.IPV6_LOOPBACK_ADDRESS;
61+
import static org.junit.Assert.assertEquals;
62+
import static org.junit.Assert.assertFalse;
63+
import static org.junit.Assert.assertTrue;
64+
65+
@RunWith(IoTDBTestRunner.class)
66+
@Category({LocalStandaloneIT.class, ClusterIT.class})
67+
public class IoTDBIPv6ExternalServiceIT {
68+
69+
private static final String USER = "root";
70+
private static final String PASSWORD = "root";
71+
72+
private static String previousTestNodeAddress;
73+
74+
@BeforeClass
75+
public static void setUp() {
76+
IPv6TestUtils.assumeIPv6LoopbackAvailable();
77+
previousTestNodeAddress = IPv6TestUtils.setTestNodeAddressToIPv6Loopback();
78+
EnvFactory.getEnv().getConfig().getDataNodeConfig().setEnableRestService(true);
79+
EnvFactory.getEnv()
80+
.getConfig()
81+
.getDataNodeConfig()
82+
.setEnableMQTTService(true)
83+
.setMqttPayloadFormatter("json");
84+
EnvFactory.getEnv().initClusterEnvironment(1, 1);
85+
}
86+
87+
@AfterClass
88+
public static void tearDown() {
89+
try {
90+
EnvFactory.getEnv().cleanClusterEnvironment();
91+
} finally {
92+
IPv6TestUtils.restoreTestNodeAddress(previousTestNodeAddress);
93+
}
94+
}
95+
96+
@Test
97+
public void restServiceCanCommunicateThroughIPv6Loopback() throws Exception {
98+
final DataNodeWrapper dataNode = EnvFactory.getEnv().getDataNodeWrapper(0);
99+
assertEquals(IPV6_LOOPBACK_ADDRESS, dataNode.getIp());
100+
final String restAddress =
101+
UrlUtils.formatTEndPointIpv4AndIpv6Url(dataNode.getIp(), dataNode.getRestServicePort());
102+
103+
try (CloseableHttpClient httpClient = HttpClientBuilder.create().build()) {
104+
Awaitility.await()
105+
.atMost(30, TimeUnit.SECONDS)
106+
.pollInterval(1, TimeUnit.SECONDS)
107+
.ignoreExceptions()
108+
.untilAsserted(() -> assertRestPingSucceeds(httpClient, restAddress));
109+
110+
executeRestNonQuery(
111+
httpClient, restAddress, "CREATE TIMESERIES root.ipv6_rest.d1.s1 WITH DATATYPE=INT64");
112+
executeRestNonQuery(
113+
httpClient, restAddress, "INSERT INTO root.ipv6_rest.d1(time, s1) VALUES (1, 100)");
114+
}
115+
116+
try (ISession session = EnvFactory.getEnv().getSessionConnection();
117+
SessionDataSet dataSet =
118+
session.executeQueryStatement("SELECT s1 FROM root.ipv6_rest.d1")) {
119+
assertTrue(dataSet.hasNext());
120+
final RowRecord row = dataSet.next();
121+
assertEquals(1L, row.getTimestamp());
122+
assertEquals(100L, row.getFields().get(0).getLongV());
123+
assertFalse(dataSet.hasNext());
124+
}
125+
}
126+
127+
@Test
128+
public void mqttServiceCanCommunicateThroughIPv6Loopback() throws Exception {
129+
final DataNodeWrapper dataNode = EnvFactory.getEnv().getDataNodeWrapper(0);
130+
assertEquals(IPV6_LOOPBACK_ADDRESS, dataNode.getIp());
131+
final String mqttAddress =
132+
UrlUtils.formatTEndPointIpv4AndIpv6Url(dataNode.getIp(), dataNode.getMqttPort());
133+
134+
final MQTT mqtt = new MQTT();
135+
mqtt.setHost("tcp://" + mqttAddress);
136+
mqtt.setUserName(USER);
137+
mqtt.setPassword(PASSWORD);
138+
mqtt.setConnectAttemptsMax(3);
139+
mqtt.setReconnectDelay(1000);
140+
mqtt.setClientId("ipv6Client");
141+
142+
final BlockingConnection connection = mqtt.blockingConnection();
143+
try {
144+
connection.connect();
145+
connection.publish(
146+
"root.ipv6_mqtt.d1",
147+
("{"
148+
+ "\"device\":\"root.ipv6_mqtt.d1\","
149+
+ "\"timestamp\":1,"
150+
+ "\"measurements\":[\"s1\"],"
151+
+ "\"values\":[200.0]"
152+
+ "}")
153+
.getBytes(StandardCharsets.UTF_8),
154+
QoS.AT_LEAST_ONCE,
155+
false);
156+
157+
try (ISession session = EnvFactory.getEnv().getSessionConnection()) {
158+
Awaitility.await()
159+
.atMost(30, TimeUnit.SECONDS)
160+
.pollInterval(1, TimeUnit.SECONDS)
161+
.until(() -> mqttValueIsVisible(session));
162+
}
163+
} finally {
164+
if (connection.isConnected()) {
165+
connection.disconnect();
166+
}
167+
}
168+
}
169+
170+
private static void assertRestPingSucceeds(
171+
final CloseableHttpClient httpClient, final String restAddress) throws IOException {
172+
try (CloseableHttpResponse response =
173+
httpClient.execute(new HttpGet("http://" + restAddress + "/ping"))) {
174+
assertEquals(200, response.getStatusLine().getStatusCode());
175+
final JsonObject result =
176+
JsonParser.parseString(EntityUtils.toString(response.getEntity(), StandardCharsets.UTF_8))
177+
.getAsJsonObject();
178+
assertEquals(200, result.get("code").getAsInt());
179+
}
180+
}
181+
182+
private static void executeRestNonQuery(
183+
final CloseableHttpClient httpClient, final String restAddress, final String sql)
184+
throws IOException {
185+
final HttpPost httpPost = new HttpPost("http://" + restAddress + "/rest/v2/nonQuery");
186+
httpPost.addHeader("Content-type", "application/json; charset=utf-8");
187+
httpPost.setHeader("Accept", "application/json");
188+
httpPost.setHeader(
189+
"Authorization",
190+
Base64.getEncoder()
191+
.encodeToString((USER + ":" + PASSWORD).getBytes(StandardCharsets.UTF_8)));
192+
final JsonObject request = new JsonObject();
193+
request.addProperty("sql", sql);
194+
httpPost.setEntity(new StringEntity(request.toString(), StandardCharsets.UTF_8));
195+
196+
try (CloseableHttpResponse response = httpClient.execute(httpPost)) {
197+
assertEquals(200, response.getStatusLine().getStatusCode());
198+
final JsonObject result =
199+
JsonParser.parseString(EntityUtils.toString(response.getEntity(), StandardCharsets.UTF_8))
200+
.getAsJsonObject();
201+
assertEquals(200, result.get("code").getAsInt());
202+
}
203+
}
204+
205+
private static boolean mqttValueIsVisible(final ISession session) throws Exception {
206+
try (SessionDataSet dataSet =
207+
session.executeQueryStatement("SELECT s1 FROM root.ipv6_mqtt.d1 WHERE time = 1")) {
208+
if (!dataSet.hasNext()) {
209+
return false;
210+
}
211+
final List<Field> fields = dataSet.next().getFields();
212+
return fields.size() == 1 && Math.abs(fields.get(0).getDoubleV() - 200.0) < 0.001;
213+
} catch (StatementExecutionException e) {
214+
if (e.getMessage() != null && e.getMessage().contains("does not exist")) {
215+
return false;
216+
}
217+
throw e;
218+
}
219+
}
220+
}

0 commit comments

Comments
 (0)