Skip to content

Commit 43a7104

Browse files
committed
[Feature] Address review for user data transfer audit hooks
1 parent 88f72c6 commit 43a7104

31 files changed

Lines changed: 515 additions & 406 deletions

File tree

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ public class IndexedConsensusRequest implements IConsensusRequest {
4949
private long memorySize = 0;
5050
private long retainedMemorySize = 0;
5151
private boolean serializedRequestsBuilt = false;
52+
private boolean containsUserData = false;
5253
private final AtomicLong referenceCnt = new AtomicLong();
5354

5455
public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> requests) {
@@ -171,6 +172,15 @@ public IndexedConsensusRequest setNodeId(int nodeId) {
171172
return this;
172173
}
173174

175+
public boolean containsUserData() {
176+
return containsUserData;
177+
}
178+
179+
public IndexedConsensusRequest setContainsUserData(boolean containsUserData) {
180+
this.containsUserData = containsUserData;
181+
return this;
182+
}
183+
174184
public long getLocalSeq() {
175185
return searchIndex;
176186
}

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/ConsensusConfig.java‎

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ public class ConsensusConfig {
4141
private final DirectoryStrategyType directoryStrategyType;
4242
private final TrustedChannelFailureHandler trustedChannelFailureHandler;
4343
private final UserDataTransferAuditHandler userDataTransferAuditHandler;
44+
private final UserDataTransferAuditClassifier userDataTransferAuditClassifier;
4445

4546
private ConsensusConfig(
4647
TEndPoint thisNode,
@@ -53,7 +54,8 @@ private ConsensusConfig(
5354
IoTConsensusV2Config iotConsensusV2Config,
5455
DirectoryStrategyType directoryStrategyType,
5556
TrustedChannelFailureHandler trustedChannelFailureHandler,
56-
UserDataTransferAuditHandler userDataTransferAuditHandler) {
57+
UserDataTransferAuditHandler userDataTransferAuditHandler,
58+
UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
5759
this.thisNodeEndPoint = thisNode;
5860
this.thisNodeId = thisNodeId;
5961
this.storageDir = storageDir;
@@ -65,6 +67,7 @@ private ConsensusConfig(
6567
this.directoryStrategyType = directoryStrategyType;
6668
this.trustedChannelFailureHandler = trustedChannelFailureHandler;
6769
this.userDataTransferAuditHandler = userDataTransferAuditHandler;
70+
this.userDataTransferAuditClassifier = userDataTransferAuditClassifier;
6871
}
6972

7073
public TEndPoint getThisNodeEndPoint() {
@@ -111,6 +114,10 @@ public UserDataTransferAuditHandler getUserDataTransferAuditHandler() {
111114
return userDataTransferAuditHandler;
112115
}
113116

117+
public UserDataTransferAuditClassifier getUserDataTransferAuditClassifier() {
118+
return userDataTransferAuditClassifier;
119+
}
120+
114121
public static ConsensusConfig.Builder newBuilder() {
115122
return new ConsensusConfig.Builder();
116123
}
@@ -131,6 +138,8 @@ public static class Builder {
131138
TrustedChannelFailureHandler.NO_OP;
132139
private UserDataTransferAuditHandler userDataTransferAuditHandler =
133140
UserDataTransferAuditHandler.NO_OP;
141+
private UserDataTransferAuditClassifier userDataTransferAuditClassifier =
142+
UserDataTransferAuditClassifier.NO_USER_DATA;
134143

135144
public ConsensusConfig build() {
136145
return new ConsensusConfig(
@@ -146,7 +155,8 @@ public ConsensusConfig build() {
146155
.orElseGet(() -> IoTConsensusV2Config.newBuilder().build()),
147156
directoryStrategyType,
148157
trustedChannelFailureHandler,
149-
userDataTransferAuditHandler);
158+
userDataTransferAuditHandler,
159+
userDataTransferAuditClassifier);
150160
}
151161

152162
public Builder setThisNode(TEndPoint thisNode) {
@@ -209,5 +219,13 @@ public Builder setUserDataTransferAuditHandler(
209219
.orElse(UserDataTransferAuditHandler.NO_OP);
210220
return this;
211221
}
222+
223+
public Builder setUserDataTransferAuditClassifier(
224+
UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
225+
this.userDataTransferAuditClassifier =
226+
Optional.ofNullable(userDataTransferAuditClassifier)
227+
.orElse(UserDataTransferAuditClassifier.NO_USER_DATA);
228+
return this;
229+
}
212230
}
213231
}
Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
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.consensus.config;
21+
22+
import org.apache.iotdb.commons.consensus.ConsensusGroupId;
23+
import org.apache.iotdb.commons.request.IConsensusRequest;
24+
25+
@FunctionalInterface
26+
public interface UserDataTransferAuditClassifier {
27+
28+
UserDataTransferAuditClassifier NO_USER_DATA = (groupId, request) -> false;
29+
30+
boolean containsUserData(ConsensusGroupId groupId, IConsensusRequest request);
31+
}

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@
4646
import org.apache.iotdb.consensus.common.Peer;
4747
import org.apache.iotdb.consensus.config.ConsensusConfig;
4848
import org.apache.iotdb.consensus.config.IoTConsensusConfig;
49+
import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
4950
import org.apache.iotdb.consensus.exception.ConsensusException;
5051
import org.apache.iotdb.consensus.exception.ConsensusGroupAlreadyExistException;
5152
import org.apache.iotdb.consensus.exception.ConsensusGroupModifyPeerException;
@@ -106,6 +107,7 @@ public class IoTConsensus implements IConsensus {
106107
private final IoTConsensusRPCService service;
107108
private final RegisterManager registerManager = new RegisterManager();
108109
private final UserDataTransferAuditHandler userDataTransferAuditHandler;
110+
private final UserDataTransferAuditClassifier userDataTransferAuditClassifier;
109111
private volatile IoTConsensusConfig config;
110112

111113
/**
@@ -134,6 +136,7 @@ public IoTConsensus(ConsensusConfig config, Registry registry) {
134136
this.recvFolderStrategyType = config.getDirectoryStrategyType();
135137
this.config = config.getIotConsensusConfig();
136138
this.userDataTransferAuditHandler = config.getUserDataTransferAuditHandler();
139+
this.userDataTransferAuditClassifier = config.getUserDataTransferAuditClassifier();
137140
this.registry = registry;
138141
this.service =
139142
new IoTConsensusRPCService(
@@ -211,7 +214,8 @@ private void initAndRecover() throws IOException {
211214
clientManager,
212215
syncClientManager,
213216
config,
214-
userDataTransferAuditHandler);
217+
userDataTransferAuditHandler,
218+
userDataTransferAuditClassifier);
215219
stateMachineMap.put(consensusGroupId, consensus);
216220
}
217221
} catch (DiskSpaceInsufficientException e) {
@@ -327,7 +331,8 @@ public void createLocalPeer(ConsensusGroupId groupId, List<Peer> peers)
327331
clientManager,
328332
syncClientManager,
329333
config,
330-
userDataTransferAuditHandler);
334+
userDataTransferAuditHandler,
335+
userDataTransferAuditClassifier);
331336
} catch (DiskSpaceInsufficientException e) {
332337
throw new RuntimeException(e);
333338
}

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java‎

Lines changed: 26 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
2525
import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
2626
import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
27-
import org.apache.iotdb.commons.audit.UserDataTransferType;
2827
import org.apache.iotdb.commons.client.IClientManager;
2928
import org.apache.iotdb.commons.client.exception.ClientManagerException;
3029
import org.apache.iotdb.commons.consensus.ConsensusGroupId;
@@ -46,6 +45,7 @@
4645
import org.apache.iotdb.consensus.common.request.DeserializedBatchIndexedConsensusRequest;
4746
import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
4847
import org.apache.iotdb.consensus.config.IoTConsensusConfig;
48+
import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
4949
import org.apache.iotdb.consensus.exception.ConsensusGroupModifyPeerException;
5050
import org.apache.iotdb.consensus.i18n.ConsensusMessages;
5151
import org.apache.iotdb.consensus.i18n.IoTConsensusMessages;
@@ -163,6 +163,7 @@ public class IoTConsensusServerImpl {
163163
private final IoTConsensusRateLimiter ioTConsensusRateLimiter =
164164
IoTConsensusRateLimiter.getInstance();
165165
private final UserDataTransferAuditHandler userDataTransferAuditHandler;
166+
private final UserDataTransferAuditClassifier userDataTransferAuditClassifier;
166167
private IndexedConsensusRequest lastConsensusRequest;
167168

168169
// Subscription queues receive IndexedConsensusRequest in real-time from write(),
@@ -212,7 +213,8 @@ public IoTConsensusServerImpl(
212213
clientManager,
213214
syncClientManager,
214215
config,
215-
UserDataTransferAuditHandler.NO_OP);
216+
UserDataTransferAuditHandler.NO_OP,
217+
UserDataTransferAuditClassifier.NO_USER_DATA);
216218
}
217219

218220
public IoTConsensusServerImpl(
@@ -226,7 +228,8 @@ public IoTConsensusServerImpl(
226228
IClientManager<TEndPoint, AsyncIoTConsensusServiceClient> clientManager,
227229
IClientManager<TEndPoint, SyncIoTConsensusServiceClient> syncClientManager,
228230
IoTConsensusConfig config,
229-
UserDataTransferAuditHandler userDataTransferAuditHandler)
231+
UserDataTransferAuditHandler userDataTransferAuditHandler,
232+
UserDataTransferAuditClassifier userDataTransferAuditClassifier)
230233
throws DiskSpaceInsufficientException {
231234
this.active = true;
232235
this.storageDir = storageDir;
@@ -248,6 +251,7 @@ public IoTConsensusServerImpl(
248251
this.backgroundTaskService = backgroundTaskService;
249252
this.config = config;
250253
this.userDataTransferAuditHandler = userDataTransferAuditHandler;
254+
this.userDataTransferAuditClassifier = userDataTransferAuditClassifier;
251255
this.consensusGroupId = thisNode.getGroupId().toString();
252256
this.consensusReqReader =
253257
(ConsensusReqReader) stateMachine.read(new GetConsensusReqReaderPlan());
@@ -521,23 +525,18 @@ private void recordSnapshotTransferAttempt(
521525
boolean success,
522526
String errorCode,
523527
Throwable error) {
524-
if (!userDataTransferAuditHandler.isEnabled()) {
525-
return;
526-
}
527528
try {
529+
if (!userDataTransferAuditHandler.isEnabled()) {
530+
return;
531+
}
528532
userDataTransferAuditHandler.onAttempt(
529533
new UserDataTransferAuditEvent(
530-
UserDataTransferType.IOT_CONSENSUS_SNAPSHOT,
531534
thisNode.getEndpoint(),
532535
thisNode.getEndpoint(),
533536
targetPeer.getEndpoint(),
534537
UserDataTransferProtectionMethod.fromTlsEnabled(config.getRpc().isEnableSSL()),
535-
null,
536-
request.getSnapshotId() + "/" + request.getOffset(),
537-
1,
538538
success,
539-
errorCode,
540-
error));
539+
errorCode != null ? errorCode : error == null ? null : error.getClass().getName()));
541540
} catch (RuntimeException ignored) {
542541
// Audit recording must not affect snapshot transmission.
543542
}
@@ -1014,6 +1013,7 @@ public IndexedConsensusRequest buildIndexedConsensusRequestForLocalRequest(
10141013
((ComparableConsensusRequest) request).setProgressIndex(iotProgressIndex);
10151014
}
10161015
return new IndexedConsensusRequest(searchIndex.get() + 1, Collections.singletonList(request))
1016+
.setContainsUserData(containsUserData(Collections.singletonList(request)))
10171017
.setPhysicalTime(assignPhysicalTimeInMs())
10181018
.setNodeId(thisNode.getNodeId());
10191019
}
@@ -1030,9 +1030,23 @@ public IndexedConsensusRequest buildIndexedConsensusRequestForRemoteRequest(
10301030
req.setRoutingEpoch(routingEpoch);
10311031
req.setPhysicalTime(physicalTime);
10321032
req.setNodeId(nodeId);
1033+
req.setContainsUserData(containsUserData(requests));
10331034
return req;
10341035
}
10351036

1037+
public boolean containsUserData(List<IConsensusRequest> requests) {
1038+
for (IConsensusRequest request : requests) {
1039+
try {
1040+
if (userDataTransferAuditClassifier.containsUserData(thisNode.getGroupId(), request)) {
1041+
return true;
1042+
}
1043+
} catch (RuntimeException ignored) {
1044+
// Classification is advisory and must not affect consensus replication.
1045+
}
1046+
}
1047+
return false;
1048+
}
1049+
10361050
public TSStatus syncIdleWriterSafeTimeBarrierToPeer(final Peer targetPeer) {
10371051
final long safePhysicalTime = assignPhysicalTimeInMs();
10381052
final long safeLocalSeq = searchIndex.get();

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandler.java‎

Lines changed: 37 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,11 @@
1919

2020
package org.apache.iotdb.consensus.iot.client;
2121

22+
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
2223
import org.apache.iotdb.common.rpc.thrift.TSStatus;
2324
import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
25+
import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
2426
import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
25-
import org.apache.iotdb.commons.audit.UserDataTransferType;
2627
import org.apache.iotdb.commons.utils.RetryUtils;
2728
import org.apache.iotdb.consensus.i18n.IoTConsensusMessages;
2829
import org.apache.iotdb.consensus.iot.logdispatcher.Batch;
@@ -187,31 +188,43 @@ private void completeBatch(Batch batch) {
187188
}
188189

189190
private void recordTransferAttempt(boolean success, String errorCode, Throwable error) {
190-
if (!thread.getImpl().getUserDataTransferAuditHandler().isEnabled()) {
191-
return;
191+
try {
192+
recordTransferAttempt(
193+
thread.getImpl().getUserDataTransferAuditHandler(),
194+
batch,
195+
thread.getImpl().getThisNode().getEndpoint(),
196+
thread.getPeer().getEndpoint(),
197+
UserDataTransferProtectionMethod.fromTlsEnabled(
198+
thread.getConfig().getRpc().isEnableSSL()),
199+
success,
200+
errorCode,
201+
error);
202+
} catch (RuntimeException ignored) {
203+
// Audit recording must not affect consensus replication.
192204
}
205+
}
206+
207+
static void recordTransferAttempt(
208+
UserDataTransferAuditHandler auditHandler,
209+
Batch batch,
210+
TEndPoint source,
211+
TEndPoint target,
212+
UserDataTransferProtectionMethod protectionMethod,
213+
boolean success,
214+
String errorCode,
215+
Throwable error) {
193216
try {
194-
thread
195-
.getImpl()
196-
.getUserDataTransferAuditHandler()
197-
.onAttempt(
198-
new UserDataTransferAuditEvent(
199-
UserDataTransferType.IOT_CONSENSUS_LOG,
200-
thread.getImpl().getThisNode().getEndpoint(),
201-
thread.getImpl().getThisNode().getEndpoint(),
202-
thread.getPeer().getEndpoint(),
203-
UserDataTransferProtectionMethod.fromTlsEnabled(
204-
thread.getConfig().getRpc().isEnableSSL()),
205-
null,
206-
thread.getImpl().getThisNode().getGroupId()
207-
+ "/"
208-
+ batch.getStartIndex()
209-
+ "-"
210-
+ batch.getEndIndex(),
211-
retryCount + 1,
212-
success,
213-
errorCode,
214-
error));
217+
if (!batch.containsUserData() || !auditHandler.isEnabled()) {
218+
return;
219+
}
220+
auditHandler.onAttempt(
221+
new UserDataTransferAuditEvent(
222+
source,
223+
source,
224+
target,
225+
protectionMethod,
226+
success,
227+
errorCode != null ? errorCode : error == null ? null : error.getClass().getName()));
215228
} catch (RuntimeException ignored) {
216229
// Audit recording must not affect consensus replication.
217230
}

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ public class Batch {
3939
private long memorySize;
4040
// indicates whether this batch has been successfully synchronized to another node
4141
private boolean synced;
42+
private boolean containsUserData;
4243

4344
public Batch(IoTConsensusConfig config) {
4445
this.config = config;
@@ -55,11 +56,16 @@ public void buildIndex() {
5556
}
5657

5758
public void addTLogEntry(TLogEntry entry) {
59+
addTLogEntry(entry, false);
60+
}
61+
62+
public void addTLogEntry(TLogEntry entry, boolean containsUserData) {
5863
logEntries.add(entry);
5964
if (entry.fromWAL) {
6065
logEntriesNumFromWAL++;
6166
}
6267
memorySize += entry.getMemorySize();
68+
this.containsUserData |= containsUserData;
6369
}
6470

6571
public boolean canAccumulate() {
@@ -107,6 +113,10 @@ public long getLogEntriesNumFromWAL() {
107113
return logEntriesNumFromWAL;
108114
}
109115

116+
public boolean containsUserData() {
117+
return containsUserData;
118+
}
119+
110120
@Override
111121
public String toString() {
112122
return "Batch{"

0 commit comments

Comments
 (0)