This is an automated email from the ASF dual-hosted git repository.
mimaison pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 7d10473f0c4 KAFKA-19298: Move RuntimeLoggerManager to server module
(#22401)
7d10473f0c4 is described below
commit 7d10473f0c49ce8e00a9f5980cf113945a6b550b
Author: majialong <[email protected]>
AuthorDate: Thu Jul 30 18:05:39 2026 +0800
KAFKA-19298: Move RuntimeLoggerManager to server module (#22401)
Move `RuntimeLoggerManager` to server module.
Reviewers: Mickael Maison <[email protected]>, Russole Chen
<[email protected]>
---
.../scala/kafka/server/ConfigAdminManager.scala | 3 +--
.../main/scala/kafka/server/ControllerApis.scala | 2 +-
.../kafka/server/logger/RuntimeLoggerManager.java | 20 ++++++++++++-----
.../server/logger/RuntimeLoggerManagerTest.java | 25 ++++++++++++++++++++--
4 files changed, 40 insertions(+), 10 deletions(-)
diff --git a/core/src/main/scala/kafka/server/ConfigAdminManager.scala
b/core/src/main/scala/kafka/server/ConfigAdminManager.scala
index 1b021531601..68210276e56 100644
--- a/core/src/main/scala/kafka/server/ConfigAdminManager.scala
+++ b/core/src/main/scala/kafka/server/ConfigAdminManager.scala
@@ -16,8 +16,6 @@
*/
package kafka.server
-import kafka.server.logger.RuntimeLoggerManager
-
import java.util
import java.util.Properties
import kafka.utils._
@@ -37,6 +35,7 @@ import org.apache.kafka.common.requests.ApiError
import org.apache.kafka.common.resource.{Resource, ResourceType}
import org.apache.kafka.metadata.ConfigRepository
import org.apache.kafka.server.config.AbstractKafkaConfig
+import org.apache.kafka.server.logger.RuntimeLoggerManager
import org.slf4j.{Logger, LoggerFactory}
import scala.collection.{Map, Seq}
diff --git a/core/src/main/scala/kafka/server/ControllerApis.scala
b/core/src/main/scala/kafka/server/ControllerApis.scala
index 63ed29d7cd8..fc3de3e317d 100644
--- a/core/src/main/scala/kafka/server/ControllerApis.scala
+++ b/core/src/main/scala/kafka/server/ControllerApis.scala
@@ -25,7 +25,6 @@ import java.util.concurrent.CompletableFuture
import java.util.function.Consumer
import kafka.network.RequestChannel
import org.apache.kafka.server.quota.QuotaFactory.QuotaManagers
-import kafka.server.logger.RuntimeLoggerManager
import kafka.utils.Logging
import org.apache.kafka.clients.admin.{AlterConfigOp, EndpointType}
import org.apache.kafka.common.Uuid.ZERO_UUID
@@ -58,6 +57,7 @@ import org.apache.kafka.security.DelegationTokenManager
import org.apache.kafka.server.{ApiVersionManager, AuthHelper, EnvelopeUtils,
ProcessRole}
import org.apache.kafka.server.authorizer.Authorizer
import org.apache.kafka.server.common.{ApiMessageAndVersion, RequestLocal}
+import org.apache.kafka.server.logger.RuntimeLoggerManager
import org.apache.kafka.server.quota.ControllerMutationQuota
import scala.jdk.javaapi.OptionConverters
diff --git a/core/src/main/java/kafka/server/logger/RuntimeLoggerManager.java
b/server/src/main/java/org/apache/kafka/server/logger/RuntimeLoggerManager.java
similarity index 89%
rename from core/src/main/java/kafka/server/logger/RuntimeLoggerManager.java
rename to
server/src/main/java/org/apache/kafka/server/logger/RuntimeLoggerManager.java
index 3cb226e0686..ceb61128527 100644
--- a/core/src/main/java/kafka/server/logger/RuntimeLoggerManager.java
+++
b/server/src/main/java/org/apache/kafka/server/logger/RuntimeLoggerManager.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package kafka.server.logger;
+package org.apache.kafka.server.logger;
import org.apache.kafka.clients.admin.AlterConfigOp.OpType;
import org.apache.kafka.common.config.LogLevelConfig;
@@ -25,7 +25,6 @@ import org.apache.kafka.common.errors.InvalidRequestException;
import
org.apache.kafka.common.message.IncrementalAlterConfigsRequestData.AlterConfigsResource;
import
org.apache.kafka.common.message.IncrementalAlterConfigsRequestData.AlterableConfig;
import org.apache.kafka.common.protocol.Errors;
-import org.apache.kafka.server.logger.LoggingController;
import org.slf4j.Logger;
@@ -49,7 +48,7 @@ public class RuntimeLoggerManager {
private final int nodeId;
private final Logger log;
- public RuntimeLoggerManager(int nodeId, Logger log) {
+ public RuntimeLoggerManager(int nodeId, Logger log) {
this.nodeId = nodeId;
this.log = log;
}
@@ -73,7 +72,12 @@ public class RuntimeLoggerManager {
ops.forEach(op -> {
String loggerName = op.name();
String logLevel = op.value();
- switch (OpType.forId(op.configOperation())) {
+ OpType opType = OpType.forId(op.configOperation());
+ if (opType == null) {
+ throw new IllegalArgumentException(
+ "Invalid log4j configOperation: " + op.configOperation());
+ }
+ switch (opType) {
case SET:
if (LoggingController.logLevel(loggerName, logLevel)) {
log.warn("Updated the log level of {} to {}",
loggerName, logLevel);
@@ -118,7 +122,13 @@ public class RuntimeLoggerManager {
void validateLogLevelConfigs(Collection<AlterableConfig> ops) {
ops.forEach(op -> {
String loggerName = op.name();
- switch (OpType.forId(op.configOperation())) {
+ OpType opType = OpType.forId(op.configOperation());
+ if (opType == null) {
+ throw new InvalidRequestException("Unknown operation type " +
+ (int) op.configOperation() + " is not allowed for the " +
+ BROKER_LOGGER + " resource");
+ }
+ switch (opType) {
case SET:
validateLoggerNameExists(loggerName);
String logLevel = op.value();
diff --git
a/core/src/test/java/kafka/server/logger/RuntimeLoggerManagerTest.java
b/server/src/test/java/org/apache/kafka/server/logger/RuntimeLoggerManagerTest.java
similarity index 80%
rename from core/src/test/java/kafka/server/logger/RuntimeLoggerManagerTest.java
rename to
server/src/test/java/org/apache/kafka/server/logger/RuntimeLoggerManagerTest.java
index b5c8740639c..f5d687feb4f 100644
--- a/core/src/test/java/kafka/server/logger/RuntimeLoggerManagerTest.java
+++
b/server/src/test/java/org/apache/kafka/server/logger/RuntimeLoggerManagerTest.java
@@ -14,14 +14,13 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package kafka.server.logger;
+package org.apache.kafka.server.logger;
import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.clients.admin.AlterConfigOp.OpType;
import org.apache.kafka.common.errors.InvalidConfigurationException;
import org.apache.kafka.common.errors.InvalidRequestException;
import
org.apache.kafka.common.message.IncrementalAlterConfigsRequestData.AlterableConfig;
-import org.apache.kafka.server.logger.LoggingController;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -67,6 +66,28 @@ public class RuntimeLoggerManagerTest {
setValue("TRACE")))).getMessage());
}
+ @Test
+ public void testUnknownOperationNotAllowed() {
+ byte unknownOperation = 99;
+ assertEquals("Unknown operation type 99 is not allowed for the
BROKER_LOGGER resource",
+ Assertions.assertThrows(InvalidRequestException.class,
+ () -> MANAGER.validateLogLevelConfigs(List.of(new
AlterableConfig().
+ setName(LOG.getName()).
+ setConfigOperation(unknownOperation).
+ setValue("TRACE")))).getMessage());
+ }
+
+ @Test
+ public void testAlterUnknownOperationNotAllowed() {
+ byte unknownOperation = 99;
+ assertEquals("Invalid log4j configOperation: 99",
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () -> MANAGER.alterLogLevelConfigs(List.of(new
AlterableConfig().
+ setName(LOG.getName()).
+ setConfigOperation(unknownOperation).
+ setValue("TRACE")))).getMessage());
+ }
+
@Test
public void testValidateBogusLogLevelNameNotAllowed() {
assertEquals("Cannot set the log level of " + LOG.getName() + " to
BOGUS as it is not " +