Skip to content

Commit dfbe8db

Browse files
authored
Fix MQTT measurement validation (#18215)
1 parent fe1777b commit dfbe8db

2 files changed

Lines changed: 45 additions & 1 deletion

File tree

‎external-service-impl/mqtt/src/main/java/org/apache/iotdb/mqtt/MPPPublishHandler.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import org.apache.iotdb.commons.queryengine.common.SqlDialect;
2525
import org.apache.iotdb.commons.queryengine.utils.TimestampPrecisionUtils;
2626
import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory;
27+
import org.apache.iotdb.commons.utils.PathUtils;
2728
import org.apache.iotdb.db.auth.AuthorityChecker;
2829
import org.apache.iotdb.db.conf.IoTDBConfig;
2930
import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -274,7 +275,9 @@ private void insertTree(TreeMessage message, MqttClientSession session) {
274275
DataNodeDevicePathCache.getInstance().getPartialPath(message.getDevice()));
275276
TimestampPrecisionUtils.checkTimestampPrecision(message.getTimestamp());
276277
statement.setTime(message.getTimestamp());
277-
statement.setMeasurements(message.getMeasurements().toArray(new String[0]));
278+
statement.setMeasurements(
279+
PathUtils.checkIsLegalSingleMeasurementsAndUpdate(message.getMeasurements())
280+
.toArray(new String[0]));
278281
if (message.getDataTypes() == null) {
279282
statement.setDataTypes(new TSDataType[message.getMeasurements().size()]);
280283
statement.setValues(message.getValues().toArray(new Object[0]));

‎integration-test/src/test/java/org/apache/iotdb/db/it/mqtt/IoTDBMQTTServiceJsonIT.java‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
import java.util.concurrent.TimeUnit;
5252

5353
import static org.junit.Assert.assertEquals;
54+
import static org.junit.Assert.assertFalse;
5455
import static org.junit.Assert.assertTrue;
5556
import static org.junit.Assert.fail;
5657

@@ -481,4 +482,44 @@ public void testBatchJsonWithVariousValues() throws Exception {
481482
}
482483
}
483484
}
485+
486+
@Test
487+
public void testIllegalMeasurementIsRejected() throws Exception {
488+
try (final ISession session = EnvFactory.getEnv().getSessionConnection()) {
489+
String payload =
490+
"["
491+
+ "{\"device\":\"root.sg.invalid\",\"timestamp\":1,\"measurements\":[\"a.b\"],\"values\":[1]},"
492+
+ "{\"device\":\"root.sg.invalid\",\"timestamp\":2,\"measurements\":[\"valid\"],\"values\":[2]}"
493+
+ "]";
494+
495+
Awaitility.await()
496+
.atMost(3, TimeUnit.MINUTES)
497+
.pollInterval(1, TimeUnit.SECONDS)
498+
.until(
499+
() -> {
500+
connection.publish("root.sg.invalid", payload.getBytes(), QoS.AT_LEAST_ONCE, false);
501+
try (final SessionDataSet dataSet =
502+
session.executeQueryStatement(
503+
"select valid from root.sg.invalid where time = 2")) {
504+
return dataSet.hasNext();
505+
} catch (StatementExecutionException e) {
506+
if (e.getMessage() != null && e.getMessage().contains("does not exist")) {
507+
return false;
508+
}
509+
throw e;
510+
}
511+
});
512+
513+
int timeseriesCount = 0;
514+
try (final SessionDataSet dataSet =
515+
session.executeQueryStatement("show timeseries root.sg.invalid.**")) {
516+
while (dataSet.hasNext()) {
517+
dataSet.next();
518+
timeseriesCount++;
519+
}
520+
}
521+
assertEquals(1, timeseriesCount);
522+
assertFalse(session.checkTimeseriesExists("root.sg.invalid.a.b"));
523+
}
524+
}
484525
}

0 commit comments

Comments
 (0)