Skip to content

Commit 868e3d8

Browse files
committed
fix
1 parent 5daceda commit 868e3d8

6 files changed

Lines changed: 109 additions & 8 deletions

File tree

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/Coordinator.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,7 @@
126126
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowClusterId;
127127
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowConfigNodes;
128128
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowConfiguration;
129+
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowCreateDatabase;
129130
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowCurrentDatabase;
130131
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowCurrentSqlDialect;
131132
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ShowCurrentTimestamp;
@@ -611,6 +612,7 @@ private IQueryExecution createQueryExecutionForTableModel(
611612
queryContext.setStartTime(startTime);
612613
if (statement instanceof DropDB
613614
|| statement instanceof ShowDB
615+
|| statement instanceof ShowCreateDatabase
614616
|| statement instanceof CreateDB
615617
|| statement instanceof AlterDB
616618
|| statement instanceof Use

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TableConfigTaskVisitor.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1204,6 +1204,8 @@ public static void checkAndEnrichSourceUser(
12041204
PipeSourceConstant.SOURCE_IOTDB_USERNAME_KEY, userEntity.getUsername());
12051205
replacedSourceAttributes.put(
12061206
PipeSourceConstant.SOURCE_IOTDB_CLI_HOSTNAME, userEntity.getCliHostname());
1207+
replacedSourceAttributes.put(
1208+
SystemConstant.SOURCE_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
12071209
} else if (!sourceParameters.hasAnyAttributes(
12081210
PipeSourceConstant.EXTRACTOR_IOTDB_PASSWORD_KEY,
12091211
PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY)) {
@@ -1273,6 +1275,8 @@ public static void checkAndEnrichSinkUser(
12731275
connectorAttributes.put(PipeSinkConstant.SINK_IOTDB_USERNAME_KEY, userEntity.getUsername());
12741276
connectorAttributes.put(
12751277
PipeSinkConstant.SINK_IOTDB_CLI_HOSTNAME, userEntity.getCliHostname());
1278+
connectorAttributes.put(
1279+
SystemConstant.SINK_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
12761280
} else if (!connectorParameters.hasAnyAttributes(
12771281
PipeSinkConstant.CONNECTOR_IOTDB_PASSWORD_KEY, PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY)) {
12781282
throw new SemanticException(
@@ -1315,6 +1319,8 @@ public IConfigTask visitAlterPipe(final AlterPipe node, final MPPQueryContext co
13151319
extractorAttributes,
13161320
new UserEntity(context.getUserId(), context.getUsername(), context.getCliHostname()),
13171321
true);
1322+
} else {
1323+
markSourceAuthenticationAsExplicitIfNecessary(extractorAttributes);
13181324
}
13191325
mayChangeSourcePattern(extractorAttributes);
13201326

@@ -1324,11 +1330,42 @@ public IConfigTask visitAlterPipe(final AlterPipe node, final MPPQueryContext co
13241330
node.getConnectorAttributes(),
13251331
new UserEntity(context.getUserId(), context.getUsername(), context.getCliHostname()),
13261332
true);
1333+
} else {
1334+
markSinkAuthenticationAsExplicitIfNecessary(node.getConnectorAttributes());
13271335
}
13281336

13291337
return new AlterPipeTask(node, userName);
13301338
}
13311339

1340+
public static void markSourceAuthenticationAsExplicitIfNecessary(
1341+
final Map<String, String> sourceAttributes) {
1342+
final PipeParameters sourceParameters = new PipeParameters(sourceAttributes);
1343+
if (sourceParameters.hasAnyAttributes(
1344+
PipeSourceConstant.EXTRACTOR_IOTDB_USER_KEY,
1345+
PipeSourceConstant.SOURCE_IOTDB_USER_KEY,
1346+
PipeSourceConstant.EXTRACTOR_IOTDB_USERNAME_KEY,
1347+
PipeSourceConstant.SOURCE_IOTDB_USERNAME_KEY,
1348+
PipeSourceConstant.EXTRACTOR_IOTDB_PASSWORD_KEY,
1349+
PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY)) {
1350+
sourceAttributes.put(
1351+
SystemConstant.SOURCE_AUTHENTICATION_INJECTED_KEY, Boolean.FALSE.toString());
1352+
}
1353+
}
1354+
1355+
public static void markSinkAuthenticationAsExplicitIfNecessary(
1356+
final Map<String, String> sinkAttributes) {
1357+
final PipeParameters sinkParameters = new PipeParameters(sinkAttributes);
1358+
if (sinkParameters.hasAnyAttributes(
1359+
PipeSinkConstant.CONNECTOR_IOTDB_USER_KEY,
1360+
PipeSinkConstant.SINK_IOTDB_USER_KEY,
1361+
PipeSinkConstant.CONNECTOR_IOTDB_USERNAME_KEY,
1362+
PipeSinkConstant.SINK_IOTDB_USERNAME_KEY,
1363+
PipeSinkConstant.CONNECTOR_IOTDB_PASSWORD_KEY,
1364+
PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY)) {
1365+
sinkAttributes.put(SystemConstant.SINK_AUTHENTICATION_INJECTED_KEY, Boolean.FALSE.toString());
1366+
}
1367+
}
1368+
13321369
@Override
13331370
public IConfigTask visitDropPipe(DropPipe node, MPPQueryContext context) {
13341371
context.setQueryType(QueryType.OTHER);

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/TreeConfigTaskVisitor.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -243,6 +243,8 @@
243243
import static org.apache.iotdb.commons.executable.ExecutableManager.isUriTrusted;
244244
import static org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.checkAndEnrichSinkUser;
245245
import static org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.checkAndEnrichSourceUser;
246+
import static org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.markSinkAuthenticationAsExplicitIfNecessary;
247+
import static org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.markSourceAuthenticationAsExplicitIfNecessary;
246248

247249
public class TreeConfigTaskVisitor extends StatementVisitor<IConfigTask, MPPQueryContext> {
248250

@@ -698,6 +700,8 @@ public IConfigTask visitAlterPipe(
698700
sourceAttributes,
699701
new UserEntity(context.getUserId(), context.getUsername(), context.getCliHostname()),
700702
true);
703+
} else {
704+
markSourceAuthenticationAsExplicitIfNecessary(sourceAttributes);
701705
}
702706

703707
if (alterPipeStatement.isReplaceAllSinkAttributes()) {
@@ -706,6 +710,8 @@ public IConfigTask visitAlterPipe(
706710
alterPipeStatement.getSinkAttributes(),
707711
context.getSession().getUserEntity(),
708712
true);
713+
} else {
714+
markSinkAuthenticationAsExplicitIfNecessary(alterPipeStatement.getSinkAttributes());
709715
}
710716

711717
return new AlterPipeTask(alterPipeStatement);

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/relational/ShowCreatePipeTask.java

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -115,37 +115,47 @@ private static Map<String, String> sanitizeCommonAttributes(
115115
}
116116

117117
private static Map<String, String> sanitizeSourceAttributes(final Map<String, String> source) {
118+
final boolean hasInjectedSourceAuthentication =
119+
Boolean.parseBoolean(source.get(SystemConstant.SOURCE_AUTHENTICATION_INJECTED_KEY));
118120
final Map<String, String> result = sanitizeCommonAttributes(source);
119121
result.remove(PipeSourceConstant.EXTRACTOR_IOTDB_USER_ID);
120122
result.remove(PipeSourceConstant.SOURCE_IOTDB_USER_ID);
121123
result.remove(PipeSourceConstant.EXTRACTOR_IOTDB_CLI_HOSTNAME);
122124
result.remove(PipeSourceConstant.SOURCE_IOTDB_CLI_HOSTNAME);
123-
if (!hasAnyKey(
124-
result,
125-
PipeSourceConstant.EXTRACTOR_IOTDB_PASSWORD_KEY,
126-
PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY)) {
125+
if (hasInjectedSourceAuthentication
126+
|| !hasAnyKey(
127+
result,
128+
PipeSourceConstant.EXTRACTOR_IOTDB_PASSWORD_KEY,
129+
PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY)) {
127130
result.remove(PipeSourceConstant.EXTRACTOR_IOTDB_USER_KEY);
128131
result.remove(PipeSourceConstant.SOURCE_IOTDB_USER_KEY);
129132
result.remove(PipeSourceConstant.EXTRACTOR_IOTDB_USERNAME_KEY);
130133
result.remove(PipeSourceConstant.SOURCE_IOTDB_USERNAME_KEY);
134+
result.remove(PipeSourceConstant.EXTRACTOR_IOTDB_PASSWORD_KEY);
135+
result.remove(PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY);
131136
}
132137
return result;
133138
}
134139

135140
private static Map<String, String> sanitizeSinkAttributes(final Map<String, String> sink) {
141+
final boolean hasInjectedSinkAuthentication =
142+
Boolean.parseBoolean(sink.get(SystemConstant.SINK_AUTHENTICATION_INJECTED_KEY));
136143
final Map<String, String> result = sanitizeCommonAttributes(sink);
137144
result.remove(PipeSinkConstant.CONNECTOR_IOTDB_USER_ID);
138145
result.remove(PipeSinkConstant.SINK_IOTDB_USER_ID);
139146
result.remove(PipeSinkConstant.CONNECTOR_IOTDB_CLI_HOSTNAME);
140147
result.remove(PipeSinkConstant.SINK_IOTDB_CLI_HOSTNAME);
141-
if (!hasAnyKey(
142-
result,
143-
PipeSinkConstant.CONNECTOR_IOTDB_PASSWORD_KEY,
144-
PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY)) {
148+
if (hasInjectedSinkAuthentication
149+
|| !hasAnyKey(
150+
result,
151+
PipeSinkConstant.CONNECTOR_IOTDB_PASSWORD_KEY,
152+
PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY)) {
145153
result.remove(PipeSinkConstant.CONNECTOR_IOTDB_USER_KEY);
146154
result.remove(PipeSinkConstant.SINK_IOTDB_USER_KEY);
147155
result.remove(PipeSinkConstant.CONNECTOR_IOTDB_USERNAME_KEY);
148156
result.remove(PipeSinkConstant.SINK_IOTDB_USERNAME_KEY);
157+
result.remove(PipeSinkConstant.CONNECTOR_IOTDB_PASSWORD_KEY);
158+
result.remove(PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY);
149159
}
150160
return result;
151161
}

iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/relational/ShowCreateTaskTest.java

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
3030
import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta;
3131
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseInfo;
32+
import org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor;
3233
import org.apache.iotdb.db.queryengine.plan.execution.config.sys.subscription.ShowCreateTopicTask;
3334

3435
import org.junit.Test;
@@ -78,7 +79,10 @@ public void testShowCreatePipeSQLShouldSanitizeInternalAndInjectedAttributes() {
7879
sourceAttributes.put("__audit.source", "audit");
7980
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_USER_ID, "1");
8081
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_USERNAME_KEY, "alice");
82+
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY, "hashed-password");
8183
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_CLI_HOSTNAME, "host");
84+
sourceAttributes.put(
85+
SystemConstant.SOURCE_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
8286

8387
final Map<String, String> processorAttributes = new HashMap<>();
8488
processorAttributes.put(PipeProcessorConstant.PROCESSOR_KEY, "do-nothing-processor");
@@ -91,7 +95,9 @@ public void testShowCreatePipeSQLShouldSanitizeInternalAndInjectedAttributes() {
9195
sinkAttributes.put("__audit.sink", "audit");
9296
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_USER_ID, "1");
9397
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_USERNAME_KEY, "alice");
98+
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY, "hashed-password");
9499
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_CLI_HOSTNAME, "host");
100+
sinkAttributes.put(SystemConstant.SINK_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
95101

96102
final PipeMeta pipeMeta =
97103
new PipeMeta(
@@ -134,6 +140,40 @@ public void testShowCreatePipeSQLShouldKeepExplicitCredentials() {
134140
ShowCreatePipeTask.getShowCreatePipeSQL(pipeMeta));
135141
}
136142

143+
@Test
144+
public void testShowCreatePipeSQLShouldKeepExplicitCredentialsWhenInjectionMarkerIsReset() {
145+
final Map<String, String> sourceAttributes = new HashMap<>();
146+
sourceAttributes.put(PipeSourceConstant.SOURCE_KEY, "iotdb-source");
147+
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_USERNAME_KEY, "alice");
148+
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_PASSWORD_KEY, "secret");
149+
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_USER_ID, "1");
150+
sourceAttributes.put(PipeSourceConstant.SOURCE_IOTDB_CLI_HOSTNAME, "host");
151+
sourceAttributes.put(
152+
SystemConstant.SOURCE_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
153+
154+
final Map<String, String> sinkAttributes = new HashMap<>();
155+
sinkAttributes.put(PipeSinkConstant.SINK_KEY, "write-back-sink");
156+
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_USERNAME_KEY, "alice");
157+
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY, "secret");
158+
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_USER_ID, "1");
159+
sinkAttributes.put(PipeSinkConstant.SINK_IOTDB_CLI_HOSTNAME, "host");
160+
sinkAttributes.put(SystemConstant.SINK_AUTHENTICATION_INJECTED_KEY, Boolean.TRUE.toString());
161+
162+
TableConfigTaskVisitor.markSourceAuthenticationAsExplicitIfNecessary(sourceAttributes);
163+
TableConfigTaskVisitor.markSinkAuthenticationAsExplicitIfNecessary(sinkAttributes);
164+
165+
final PipeMeta pipeMeta =
166+
new PipeMeta(
167+
new PipeStaticMeta("test_pipe", 1L, sourceAttributes, new HashMap<>(), sinkAttributes),
168+
new PipeRuntimeMeta());
169+
170+
assertEquals(
171+
"CREATE PIPE \"test_pipe\""
172+
+ " WITH SOURCE ('source'='iotdb-source','source.password'='secret','source.username'='alice')"
173+
+ " WITH SINK ('sink'='write-back-sink','sink.password'='secret','sink.username'='alice')",
174+
ShowCreatePipeTask.getShowCreatePipeSQL(pipeMeta));
175+
}
176+
137177
@Test
138178
public void testShowCreatePipeSQLShouldSanitizeExtractorAndConnectorAliases() {
139179
final Map<String, String> sourceAttributes = new HashMap<>();

iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/SystemConstant.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,10 @@ public class SystemConstant {
4242
public static final String SQL_DIALECT_KEY = "__system.sql-dialect";
4343
public static final String SQL_DIALECT_TREE_VALUE = "tree";
4444
public static final String SQL_DIALECT_TABLE_VALUE = "table";
45+
public static final String SOURCE_AUTHENTICATION_INJECTED_KEY =
46+
"__system.source-authentication-injected";
47+
public static final String SINK_AUTHENTICATION_INJECTED_KEY =
48+
"__system.sink-authentication-injected";
4549

4650
/////////////////////////////////// Utility ///////////////////////////////////
4751

@@ -50,6 +54,8 @@ public class SystemConstant {
5054
static {
5155
SYSTEM_KEYS.add(RESTART_OR_NEWLY_ADDED_KEY);
5256
SYSTEM_KEYS.add(SQL_DIALECT_KEY);
57+
SYSTEM_KEYS.add(SOURCE_AUTHENTICATION_INJECTED_KEY);
58+
SYSTEM_KEYS.add(SINK_AUTHENTICATION_INJECTED_KEY);
5359
}
5460

5561
public static PipeParameters addSystemKeysIfNecessary(final PipeParameters givenPipeParameters) {

0 commit comments

Comments
 (0)