Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -333,5 +333,8 @@ private IoTConsensusMessages() {}
public static final String LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = "Reserved {} bytes for batch {}-{}, current total usage {}";
public static final String LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}: Failed to send idle writer safe-time barrier to {}. status={}";
public static final String LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576 = "{}: write operation failed. searchIndex: {}. Code: {}, subscriptionQueues: {}, this: {}";
public static final String
LOG_FAILED_TO_RECORD_A_USER_DATA_TRANSFER_AUDIT_EVENT_CONSENSUS_REPLICATION_WILL_CONTINUE_F215E222 =
"Failed to record a user-data transfer audit event; consensus replication will continue.";

}
Original file line number Diff line number Diff line change
Expand Up @@ -331,5 +331,8 @@ private IoTConsensusMessages() {}
public static final String LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = "预留 {} 字节给批次 {}-{},当前总使用量 {}";
public static final String LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}:无法向 {} 发送 idle writer safe-time barrier。状态={}";
public static final String LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576 = "{}:写入操作失败。searchIndex: {}。Code: {},订阅队列:{},当前对象:{}";
public static final String
LOG_FAILED_TO_RECORD_A_USER_DATA_TRANSFER_AUDIT_EVENT_CONSENSUS_REPLICATION_WILL_CONTINUE_F215E222 =
"记录用户数据传送审计事件失败;Consensus 复制将继续。";

}
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ public class IndexedConsensusRequest implements IConsensusRequest {
private long memorySize = 0;
private long retainedMemorySize = 0;
private boolean serializedRequestsBuilt = false;
private boolean containsUserData = false;
private final AtomicLong referenceCnt = new AtomicLong();

public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> requests) {
Expand Down Expand Up @@ -171,6 +172,15 @@ public IndexedConsensusRequest setNodeId(int nodeId) {
return this;
}

public boolean containsUserData() {
return containsUserData;
}

public IndexedConsensusRequest setContainsUserData(boolean containsUserData) {
this.containsUserData = containsUserData;
return this;
}

public long getLocalSeq() {
return searchIndex;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.commons.audit.TrustedChannelFailureHandler;
import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
import org.apache.iotdb.commons.disk.strategy.DirectoryStrategyType;

import java.util.List;
Expand All @@ -39,6 +40,8 @@ public class ConsensusConfig {
private final IoTConsensusV2Config iotConsensusV2Config;
private final DirectoryStrategyType directoryStrategyType;
private final TrustedChannelFailureHandler trustedChannelFailureHandler;
private final UserDataTransferAuditHandler userDataTransferAuditHandler;
private final UserDataTransferAuditClassifier userDataTransferAuditClassifier;

private ConsensusConfig(
TEndPoint thisNode,
Expand All @@ -50,7 +53,9 @@ private ConsensusConfig(
IoTConsensusConfig iotConsensusConfig,
IoTConsensusV2Config iotConsensusV2Config,
DirectoryStrategyType directoryStrategyType,
TrustedChannelFailureHandler trustedChannelFailureHandler) {
TrustedChannelFailureHandler trustedChannelFailureHandler,
UserDataTransferAuditHandler userDataTransferAuditHandler,
UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
this.thisNodeEndPoint = thisNode;
this.thisNodeId = thisNodeId;
this.storageDir = storageDir;
Expand All @@ -61,6 +66,8 @@ private ConsensusConfig(
this.iotConsensusV2Config = iotConsensusV2Config;
this.directoryStrategyType = directoryStrategyType;
this.trustedChannelFailureHandler = trustedChannelFailureHandler;
this.userDataTransferAuditHandler = userDataTransferAuditHandler;
this.userDataTransferAuditClassifier = userDataTransferAuditClassifier;
}

public TEndPoint getThisNodeEndPoint() {
Expand Down Expand Up @@ -103,6 +110,14 @@ public TrustedChannelFailureHandler getTrustedChannelFailureHandler() {
return trustedChannelFailureHandler;
}

public UserDataTransferAuditHandler getUserDataTransferAuditHandler() {
return userDataTransferAuditHandler;
}

public UserDataTransferAuditClassifier getUserDataTransferAuditClassifier() {
return userDataTransferAuditClassifier;
}

public static ConsensusConfig.Builder newBuilder() {
return new ConsensusConfig.Builder();
}
Expand All @@ -121,6 +136,10 @@ public static class Builder {
DirectoryStrategyType.MIN_FOLDER_OCCUPIED_SPACE_FIRST_STRATEGY;
private TrustedChannelFailureHandler trustedChannelFailureHandler =
TrustedChannelFailureHandler.NO_OP;
private UserDataTransferAuditHandler userDataTransferAuditHandler =
UserDataTransferAuditHandler.NO_OP;
private UserDataTransferAuditClassifier userDataTransferAuditClassifier =
UserDataTransferAuditClassifier.NO_USER_DATA;

public ConsensusConfig build() {
return new ConsensusConfig(
Expand All @@ -135,7 +154,9 @@ public ConsensusConfig build() {
Optional.ofNullable(iotConsensusV2Config)
.orElseGet(() -> IoTConsensusV2Config.newBuilder().build()),
directoryStrategyType,
trustedChannelFailureHandler);
trustedChannelFailureHandler,
userDataTransferAuditHandler,
userDataTransferAuditClassifier);
}

public Builder setThisNode(TEndPoint thisNode) {
Expand Down Expand Up @@ -190,5 +211,21 @@ public Builder setTrustedChannelFailureHandler(
.orElse(TrustedChannelFailureHandler.NO_OP);
return this;
}

public Builder setUserDataTransferAuditHandler(
UserDataTransferAuditHandler userDataTransferAuditHandler) {
this.userDataTransferAuditHandler =
Optional.ofNullable(userDataTransferAuditHandler)
.orElse(UserDataTransferAuditHandler.NO_OP);
return this;
}

public Builder setUserDataTransferAuditClassifier(
UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
this.userDataTransferAuditClassifier =
Optional.ofNullable(userDataTransferAuditClassifier)
.orElse(UserDataTransferAuditClassifier.NO_USER_DATA);
return this;
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.iotdb.consensus.config;

import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.request.IConsensusRequest;

@FunctionalInterface
public interface UserDataTransferAuditClassifier {

UserDataTransferAuditClassifier NO_USER_DATA =
new UserDataTransferAuditClassifier() {
@Override
public boolean containsUserData(ConsensusGroupId groupId, IConsensusRequest request) {
return false;
}

@Override
public boolean containsUserData(ConsensusGroupId groupId) {
return false;
}
};

boolean containsUserData(ConsensusGroupId groupId, IConsensusRequest request);

/** Classifies a whole consensus group when the transfer has no individual request to inspect. */
default boolean containsUserData(ConsensusGroupId groupId) {
return true;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@

import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
import org.apache.iotdb.commons.client.IClientManager;
import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.concurrent.ThreadName;
Expand All @@ -45,6 +46,7 @@
import org.apache.iotdb.consensus.common.Peer;
import org.apache.iotdb.consensus.config.ConsensusConfig;
import org.apache.iotdb.consensus.config.IoTConsensusConfig;
import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
import org.apache.iotdb.consensus.exception.ConsensusException;
import org.apache.iotdb.consensus.exception.ConsensusGroupAlreadyExistException;
import org.apache.iotdb.consensus.exception.ConsensusGroupModifyPeerException;
Expand Down Expand Up @@ -104,6 +106,8 @@ public class IoTConsensus implements IConsensus {
new ConcurrentHashMap<>();
private final IoTConsensusRPCService service;
private final RegisterManager registerManager = new RegisterManager();
private final UserDataTransferAuditHandler userDataTransferAuditHandler;
private final UserDataTransferAuditClassifier userDataTransferAuditClassifier;
private volatile IoTConsensusConfig config;

/**
Expand Down Expand Up @@ -131,6 +135,8 @@ public IoTConsensus(ConsensusConfig config, Registry registry) {
this.recvSnapshotDirs = config.getRecvSnapshotDirs();
this.recvFolderStrategyType = config.getDirectoryStrategyType();
this.config = config.getIotConsensusConfig();
this.userDataTransferAuditHandler = config.getUserDataTransferAuditHandler();
this.userDataTransferAuditClassifier = config.getUserDataTransferAuditClassifier();
this.registry = registry;
this.service =
new IoTConsensusRPCService(
Expand Down Expand Up @@ -207,7 +213,9 @@ private void initAndRecover() throws IOException {
backgroundTaskService,
clientManager,
syncClientManager,
config);
config,
userDataTransferAuditHandler,
userDataTransferAuditClassifier);
stateMachineMap.put(consensusGroupId, consensus);
}
} catch (DiskSpaceInsufficientException e) {
Expand Down Expand Up @@ -322,7 +330,9 @@ public void createLocalPeer(ConsensusGroupId groupId, List<Peer> peers)
backgroundTaskService,
clientManager,
syncClientManager,
config);
config,
userDataTransferAuditHandler,
userDataTransferAuditClassifier);
} catch (DiskSpaceInsufficientException e) {
throw new RuntimeException(e);
}
Expand Down
Loading
Loading