This is an automated email from the ASF dual-hosted git repository. HTHou pushed a commit to branch codex/limit-tree-audit-path-log in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f2746c234cb0654dd550f6ee269518fa118559c8 Author: HTHou <[email protected]> AuthorDate: Wed Jul 29 14:25:43 2026 +0800 Fix oversized audit logs for tree batch writes --- .../org/apache/iotdb/db/auth/AuthorityChecker.java | 28 +++++++ .../security/TreeAccessCheckVisitor.java | 35 +++++---- .../plan/statement/crud/InsertBaseStatement.java | 14 ++++ .../crud/InsertMultiTabletsStatement.java | 11 +++ .../plan/statement/crud/InsertRowsStatement.java | 11 +++ .../apache/iotdb/db/auth/AuthorityCheckerTest.java | 87 +++++++++++++++++++--- .../org/apache/iotdb/db/auth/TreeAccessTest.java | 27 +++++++ 7 files changed, 188 insertions(+), 25 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/auth/AuthorityChecker.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/auth/AuthorityChecker.java index 646f093ba57..3559786b5c5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/auth/AuthorityChecker.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/auth/AuthorityChecker.java @@ -70,6 +70,7 @@ import java.util.Optional; import java.util.Set; import java.util.StringJoiner; import java.util.stream.Collectors; +import java.util.stream.Stream; import static org.apache.iotdb.commons.schema.column.ColumnHeaderConstant.LIST_USER_COLUMN_HEADERS; import static org.apache.iotdb.commons.schema.column.ColumnHeaderConstant.LIST_USER_OR_ROLE_PRIVILEGES_COLUMN_HEADERS; @@ -349,6 +350,33 @@ public class AuthorityChecker { return new TSStatus(TSStatusCode.NO_PERMISSION.getStatusCode()).setMessage(prompt.toString()); } + public static String getPathListStringForLog(List<? extends PartialPath> pathList) { + return getPathListStringForLog(pathList.stream()); + } + + public static String getPathListStringForLog(Stream<? extends PartialPath> pathStream) { + final int maxSize = Math.max(1, CommonDescriptor.getInstance().getConfig().getPathLogMaxSize()); + final List<String> paths = + pathStream + .limit((long) maxSize + 1) + .map(path -> String.valueOf(path)) + .collect(Collectors.toList()); + final boolean truncated = paths.size() > maxSize; + final int size = truncated ? maxSize : paths.size(); + + final StringBuilder result = new StringBuilder("["); + if (size > 0) { + result.append(paths.get(0)); + for (int i = 1; i < size; i++) { + result.append(", ").append(paths.get(i)); + } + if (truncated) { + result.append(", ..."); + } + } + return result.append("]").toString(); + } + public static boolean checkFullPathOrPatternPermission( String userName, PartialPath fullPath, PrivilegeType permission) { return authorityFetcher diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/security/TreeAccessCheckVisitor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/security/TreeAccessCheckVisitor.java index 4381cdc5729..81b2f62848a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/security/TreeAccessCheckVisitor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/security/TreeAccessCheckVisitor.java @@ -1136,25 +1136,31 @@ public class TreeAccessCheckVisitor extends StatementVisitor<TSStatus, TreeAcces @Override public TSStatus visitInsertBase(InsertBaseStatement statement, TreeAccessCheckContext context) { context.setAuditLogOperation(AuditLogOperation.DML).setPrivilegeType(PrivilegeType.WRITE_DATA); - for (PartialPath path : statement.getDevicePaths()) { - // External users cannot modify the audit database. - if (includeByAuditTreeDB(path) - && !context.getUsername().equals(AuthorityChecker.INTERNAL_AUDIT_USER)) { - AUDIT_LOGGER.recordObjectAuthenticationAuditLog(context.setResult(false), path::toString); - return new TSStatus(TSStatusCode.NO_PERMISSION.getStatusCode()) - .setMessage(getUnsupportedAuditDatabaseOperationMessage(TREE_MODEL_AUDIT_DATABASE)); - } + // External users cannot modify the audit database. + final PartialPath unsupportedAuditPath = + context.getUsername().equals(AuthorityChecker.INTERNAL_AUDIT_USER) + ? null + : statement + .getDevicePathsStream() + .filter(Audit::includeByAuditTreeDB) + .findFirst() + .orElse(null); + if (unsupportedAuditPath != null) { + AUDIT_LOGGER.recordObjectAuthenticationAuditLog( + context.setResult(false), unsupportedAuditPath::toString); + return new TSStatus(TSStatusCode.NO_PERMISSION.getStatusCode()) + .setMessage(getUnsupportedAuditDatabaseOperationMessage(TREE_MODEL_AUDIT_DATABASE)); } if (AuthorityChecker.SUPER_USER.equals(context.getUsername())) { AUDIT_LOGGER.recordObjectAuthenticationAuditLog( context.setResult(true), - () -> statement.getPaths().stream().distinct().collect(Collectors.toList()).toString()); + () -> AuthorityChecker.getPathListStringForLog(statement.getPathsStream().distinct())); return SUCCEED; } return checkTimeSeriesPermission( context, - () -> statement.getPaths().stream().distinct().collect(Collectors.toList()), + () -> statement.getPathsStream().distinct().collect(Collectors.toList()), PrivilegeType.WRITE_DATA); } @@ -1239,7 +1245,8 @@ public class TreeAccessCheckVisitor extends StatementVisitor<TSStatus, TreeAcces context.setPrivilegeType(permission); if (AuthorityChecker.SUPER_USER.equals(context.getUsername())) { AUDIT_LOGGER.recordObjectAuthenticationAuditLog( - context.setResult(true), () -> checkedPathsSupplier.get().toString()); + context.setResult(true), + () -> AuthorityChecker.getPathListStringForLog(checkedPathsSupplier.get())); return SUCCEED; } List<? extends PartialPath> checkedPaths = checkedPathsSupplier.get(); @@ -1253,7 +1260,7 @@ public class TreeAccessCheckVisitor extends StatementVisitor<TSStatus, TreeAcces // Internal auditor no needs audit log AUDIT_LOGGER.recordObjectAuthenticationAuditLog( context.setResult(result.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()), - checkedPaths::toString); + () -> AuthorityChecker.getPathListStringForLog(checkedPaths)); } return result; } @@ -1263,14 +1270,14 @@ public class TreeAccessCheckVisitor extends StatementVisitor<TSStatus, TreeAcces context.setPrivilegeType(permission); if (AuthorityChecker.SUPER_USER.equals(context.getUsername())) { AUDIT_LOGGER.recordObjectAuthenticationAuditLog( - context.setResult(true), checkedPaths::toString); + context.setResult(true), () -> AuthorityChecker.getPathListStringForLog(checkedPaths)); return Collections.emptyList(); } final List<Integer> results = AuthorityChecker.checkFullPathOrPatternListPermission( context.getUsername(), checkedPaths, permission); AUDIT_LOGGER.recordObjectAuthenticationAuditLog( - context.setResult(true), checkedPaths::toString); + context.setResult(true), () -> AuthorityChecker.getPathListStringForLog(checkedPaths)); return results; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertBaseStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertBaseStatement.java index 06ba84b28b0..ca113704b22 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertBaseStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertBaseStatement.java @@ -59,6 +59,7 @@ import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; +import java.util.stream.Stream; public abstract class InsertBaseStatement extends Statement implements Accountable { @@ -205,6 +206,19 @@ public abstract class InsertBaseStatement extends Statement implements Accountab return Collections.emptyList(); } + public Stream<PartialPath> getPathsStream() { + if (measurements == null) { + return Stream.empty(); + } + return Arrays.stream(measurements) + .filter(Objects::nonNull) + .map(devicePath::concatAsMeasurementPath); + } + + public Stream<PartialPath> getDevicePathsStream() { + return Stream.of(devicePath); + } + public abstract ISchemaValidation getSchemaValidation(); public abstract List<ISchemaValidation> getSchemaValidationList(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertMultiTabletsStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertMultiTabletsStatement.java index 63e942312d0..6adea132496 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertMultiTabletsStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertMultiTabletsStatement.java @@ -40,6 +40,7 @@ import java.util.List; import java.util.Objects; import java.util.Optional; import java.util.stream.Collectors; +import java.util.stream.Stream; public class InsertMultiTabletsStatement extends InsertBaseStatement { @@ -98,6 +99,16 @@ public class InsertMultiTabletsStatement extends InsertBaseStatement { return result; } + @Override + public Stream<PartialPath> getPathsStream() { + return insertTabletStatementList.stream().flatMap(InsertTabletStatement::getPathsStream); + } + + @Override + public Stream<PartialPath> getDevicePathsStream() { + return insertTabletStatementList.stream().map(InsertTabletStatement::getDevicePath); + } + @Override public ISchemaValidation getSchemaValidation() { throw new UnsupportedOperationException(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertRowsStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertRowsStatement.java index 7a2ba2ef0f4..cf4c3f2882d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertRowsStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/InsertRowsStatement.java @@ -43,6 +43,7 @@ import java.util.List; import java.util.Objects; import java.util.Optional; import java.util.stream.Collectors; +import java.util.stream.Stream; public class InsertRowsStatement extends InsertBaseStatement { @@ -117,6 +118,16 @@ public class InsertRowsStatement extends InsertBaseStatement { return result; } + @Override + public Stream<PartialPath> getPathsStream() { + return insertRowStatementList.stream().flatMap(InsertRowStatement::getPathsStream); + } + + @Override + public Stream<PartialPath> getDevicePathsStream() { + return insertRowStatementList.stream().map(InsertRowStatement::getDevicePath); + } + @Override public ISchemaValidation getSchemaValidation() { throw new UnsupportedOperationException(); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/AuthorityCheckerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/AuthorityCheckerTest.java index ca730c377b7..9e7cddc5946 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/AuthorityCheckerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/AuthorityCheckerTest.java @@ -24,11 +24,16 @@ import org.apache.iotdb.commons.conf.CommonConfig; import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.path.MeasurementPath; +import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsOfOneDeviceStatement; import org.junit.Assert; import org.junit.Test; import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; public class AuthorityCheckerTest { @@ -36,16 +41,76 @@ public class AuthorityCheckerTest { public void testLogReduce() throws IllegalPathException { final CommonConfig config = CommonDescriptor.getInstance().getConfig(); final int oldSize = config.getPathLogMaxSize(); - config.setPathLogMaxSize(1); - Assert.assertEquals( - "No permissions for this operation, please add privilege WRITE_DATA on [root.db.device.s1, ...]", - AuthorityChecker.getTSStatus( - Arrays.asList(0, 1), - Arrays.asList( - new MeasurementPath("root.db.device.s1"), - new MeasurementPath("root.db.device.s2")), - PrivilegeType.WRITE_DATA) - .getMessage()); - config.setPathLogMaxSize(oldSize); + try { + config.setPathLogMaxSize(1); + Assert.assertEquals( + "No permissions for this operation, please add privilege WRITE_DATA on [root.db.device.s1, ...]", + AuthorityChecker.getTSStatus( + Arrays.asList(0, 1), + Arrays.asList( + new MeasurementPath("root.db.device.s1"), + new MeasurementPath("root.db.device.s2")), + PrivilegeType.WRITE_DATA) + .getMessage()); + } finally { + config.setPathLogMaxSize(oldSize); + } + } + + @Test + public void testPathListStringForLog() throws IllegalPathException { + final CommonConfig config = CommonDescriptor.getInstance().getConfig(); + final int oldSize = config.getPathLogMaxSize(); + final List<MeasurementPath> paths = + Arrays.asList( + new MeasurementPath("root.db.device.s1"), + new MeasurementPath("root.db.device.s2"), + new MeasurementPath("root.db.device.s3"), + new MeasurementPath("root.db.device.s4")); + try { + config.setPathLogMaxSize(2); + Assert.assertEquals( + "[root.db.device.s1, root.db.device.s2, ...]", + AuthorityChecker.getPathListStringForLog(paths)); + Assert.assertEquals( + "[root.db.device.s1, root.db.device.s2]", + AuthorityChecker.getPathListStringForLog(paths.subList(0, 2))); + Assert.assertEquals("[]", AuthorityChecker.getPathListStringForLog(Collections.emptyList())); + + final AtomicInteger consumedPathCount = new AtomicInteger(); + AuthorityChecker.getPathListStringForLog( + paths.stream().peek(path -> consumedPathCount.incrementAndGet())); + Assert.assertEquals(3, consumedPathCount.get()); + } finally { + config.setPathLogMaxSize(oldSize); + } + } + + @Test + public void testInsertPathStreamIsLazy() throws IllegalPathException { + final CommonConfig config = CommonDescriptor.getInstance().getConfig(); + final int oldSize = config.getPathLogMaxSize(); + final AtomicInteger createdPathCount = new AtomicInteger(); + final PartialPath devicePath = + new PartialPath("root.db.device") { + @Override + public MeasurementPath concatAsMeasurementPath(String measurement) { + createdPathCount.incrementAndGet(); + return super.concatAsMeasurementPath(measurement); + } + }; + final InsertRowsOfOneDeviceStatement statement = new InsertRowsOfOneDeviceStatement(); + statement.setDevicePath(devicePath); + statement.setMeasurements(new String[] {"s1", "s2", "s3", "s4"}); + + try { + config.setPathLogMaxSize(2); + Assert.assertEquals( + "[root.db.device.s1, root.db.device.s2, ...]", + AuthorityChecker.getPathListStringForLog(statement.getPathsStream().distinct())); + Assert.assertEquals(3, createdPathCount.get()); + } finally { + config.setPathLogMaxSize(oldSize); + } } } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/TreeAccessTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/TreeAccessTest.java index 0e83cc5b477..04fdde70613 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/TreeAccessTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/auth/TreeAccessTest.java @@ -21,6 +21,8 @@ package org.apache.iotdb.db.auth; import org.apache.iotdb.commons.auth.entity.PrivilegeType; import org.apache.iotdb.commons.auth.entity.User; +import org.apache.iotdb.commons.conf.CommonConfig; +import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.queryengine.plan.relational.security.TreeAccessCheckContext; import org.apache.iotdb.db.queryengine.plan.relational.security.TreeAccessCheckVisitor; @@ -37,6 +39,7 @@ import org.junit.Test; import org.mockito.Mockito; import java.util.Collections; +import java.util.List; public class TreeAccessTest { @@ -235,6 +238,30 @@ public class TreeAccessTest { new TreeAccessCheckContext(10000L, "user1", ""), new PartialPath("root.sg"))); } + @Test + public void testPathLogLimitDoesNotLimitPermissionCheck() throws Exception { + final CommonConfig config = CommonDescriptor.getInstance().getConfig(); + final int oldSize = config.getPathLogMaxSize(); + final User user = new User("user1", "password"); + user.grantPathPrivilege(new PartialPath("root.db.device.s1"), PrivilegeType.WRITE_DATA, false); + AuthorityChecker.getAuthorityFetcher().getAuthorCache().putUserCache(user.getName(), user); + final List<PartialPath> checkedPaths = + List.of(new PartialPath("root.db.device.s1"), new PartialPath("root.db.device.s2")); + + try { + config.setPathLogMaxSize(1); + Assert.assertEquals( + TSStatusCode.NO_PERMISSION.getStatusCode(), + TreeAccessCheckVisitor.checkTimeSeriesPermission( + new TreeAccessCheckContext(10000L, "user1", ""), + () -> checkedPaths, + PrivilegeType.WRITE_DATA) + .getCode()); + } finally { + config.setPathLogMaxSize(oldSize); + } + } + private static class TestTreeAccessCheckVisitor extends TreeAccessCheckVisitor { private int checkUnsupportedAuditDatabaseWriteStatus(
