Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/li-ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -229,7 +229,7 @@ jobs:
env:
# Checkout fixes the companion source for this run. Record its commit and
# archive hash; release CI separately requires approved commit hashes.
LI_BRIDGE_LEGACY_REF: ${{ vars.LI_BRIDGE_LEGACY_REF || '3.0-li-bridge/topic-identity-recovery' }}
LI_BRIDGE_LEGACY_REF: ${{ vars.LI_BRIDGE_LEGACY_REF || '3.0-li-bridge/config-metrics-gate' }}
steps:
- uses: actions/checkout@v4
with:
Expand Down
6 changes: 6 additions & 0 deletions core/src/main/scala/kafka/server/KafkaConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ import scala.collection.{Map, Seq}
object KafkaConfig {

val LiProtocolBridgeModeEnableProp = "li.protocol.bridge.mode.enable"
val LiProtocolBridgeConfigMetricsEnableProp = "li.protocol.bridge.config.metrics.enable"
val LiProtocolBridgeTopicDeletionStateCleanupEnableProp =
"li.protocol.bridge.topic.deletion.state.cleanup.enable"
val LiProtocolBridgeFollowerRecoveryEnableProp = "li.protocol.bridge.follower.recovery.enable"
Expand Down Expand Up @@ -104,6 +105,7 @@ object KafkaConfig {

val LiProtocolBridgeEnableProps: Seq[String] = Seq(
LiProtocolBridgeModeEnableProp,
LiProtocolBridgeConfigMetricsEnableProp,
LiProtocolBridgeTopicDeletionStateCleanupEnableProp,
LiProtocolBridgeFollowerRecoveryEnableProp,
LiProtocolBridgeRecommendedElectionEnableProp,
Expand Down Expand Up @@ -280,6 +282,8 @@ object KafkaConfig {
val configDef = new ConfigDef(AbstractKafkaConfig.CONFIG_DEF)
.define(LiProtocolBridgeModeEnableProp, ConfigDef.Type.BOOLEAN, false,
ConfigDef.Importance.HIGH, LiProtocolBridgeModeEnableDoc)
.define(LiProtocolBridgeConfigMetricsEnableProp, ConfigDef.Type.BOOLEAN, false,
ConfigDef.Importance.LOW, "Register bridge configuration gauges on ZooKeeper brokers. Requires a broker restart.")
.define(LiProtocolBridgeTopicDeletionStateCleanupEnableProp, ConfigDef.Type.BOOLEAN, false,
ConfigDef.Importance.HIGH, "Clear stale topic deletion state and reconcile the metadata cache " +
"from the first full update of each ZooKeeper controller epoch. Enable on every broker together.")
Expand Down Expand Up @@ -483,6 +487,8 @@ class KafkaConfig private(doLog: Boolean, val props: util.Map[_, _])
@volatile private var currentConfig = this
val processRoles: Set[ProcessRole] = parseProcessRoles()
def liProtocolBridgeModeEnable: Boolean = getBoolean(KafkaConfig.LiProtocolBridgeModeEnableProp)
def liProtocolBridgeConfigMetricsActive: Boolean =
processRoles.isEmpty && getBoolean(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp)
def liProtocolBridgeTopicDeletionStateCleanupActive: Boolean =
processRoles.isEmpty && getBoolean(KafkaConfig.LiProtocolBridgeTopicDeletionStateCleanupEnableProp)
def liProtocolBridgeFollowerRecoveryEnable: Boolean =
Expand Down
107 changes: 56 additions & 51 deletions core/src/main/scala/kafka/server/LiProtocolBridgeMetrics.scala
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import scala.jdk.CollectionConverters._

object LiProtocolBridgeMetrics {
val ModeEnabled = "ModeEnabled"
val ConfigMetricsEnabled = "ConfigMetricsEnabled"
val TopicDeletionStateCleanupEnabled = "TopicDeletionStateCleanupEnabled"
val FollowerRecoveryEnabled = "FollowerRecoveryEnabled"
val RecommendedLeaderElectionEnabled = "RecommendedLeaderElectionEnabled"
Expand All @@ -46,7 +47,7 @@ object LiProtocolBridgeMetrics {
val LeaderTransferEnabled = "LeaderTransferEnabled"
val LegacyRequestMetricsEnabled = "LegacyRequestMetricsEnabled"
val LogTruncationMetricsEnabled = "LogTruncationMetricsEnabled"
val MetricNames: Seq[String] = Seq(ModeEnabled, TopicDeletionStateCleanupEnabled, FollowerRecoveryEnabled,
val MetricNames: Seq[String] = Seq(ModeEnabled, ConfigMetricsEnabled, TopicDeletionStateCleanupEnabled, FollowerRecoveryEnabled,
RecommendedLeaderElectionEnabled, ExcludePartitionsEnabled, MoveControllerEnabled,
ShutdownSafetyOverrideEnabled, PreferredControllerEnabled, FederatedTopicsEnabled,
RackIdMapperEnabled, ZookeeperPaginationEnabled, DynamicTopicDeletionEnabled,
Expand All @@ -64,60 +65,64 @@ class LiProtocolBridgeMetrics(config: KafkaConfig) extends AutoCloseable {
private val metricsGroup = new KafkaMetricsGroup(this.getClass)
private val tags = Map("broker-id" -> config.brokerId.toString).asJava

metricsGroup.newGauge(ModeEnabled, () => enabled(config.liProtocolBridgeModeActive), tags)
metricsGroup.newGauge(TopicDeletionStateCleanupEnabled,
() => enabled(config.liProtocolBridgeTopicDeletionStateCleanupActive), tags)
metricsGroup.newGauge(FollowerRecoveryEnabled,
() => enabled(config.liProtocolBridgeFollowerRecoveryActive), tags)
metricsGroup.newGauge(RecommendedLeaderElectionEnabled,
() => enabled(config.liProtocolBridgeRecommendedElectionActive), tags)
metricsGroup.newGauge(ExcludePartitionsEnabled,
() => enabled(config.liProtocolBridgeExcludePartitionsActive), tags)
metricsGroup.newGauge(MoveControllerEnabled,
() => enabled(config.liProtocolBridgeMoveControllerActive), tags)
metricsGroup.newGauge(ShutdownSafetyOverrideEnabled,
() => enabled(config.liProtocolBridgeShutdownSafetyOverrideActive), tags)
metricsGroup.newGauge(PreferredControllerEnabled,
() => enabled(config.liProtocolBridgePreferredControllerActive), tags)
metricsGroup.newGauge(FederatedTopicsEnabled,
() => enabled(config.liProtocolBridgeFederatedTopicsActive), tags)
metricsGroup.newGauge(RackIdMapperEnabled,
() => enabled(config.liProtocolBridgeRackIdMapperActive), tags)
metricsGroup.newGauge(ZookeeperPaginationEnabled,
() => enabled(config.liZookeeperPaginationEnable), tags)
metricsGroup.newGauge(DynamicTopicDeletionEnabled,
() => enabled(config.liProtocolBridgeDynamicTopicDeletionActive), tags)
metricsGroup.newGauge(ControllerInitializationThreads,
() => config.liNumControllerInitThreads, tags)
metricsGroup.newGauge(ProduceRequestInstrumentationEnabled,
() => enabled(config.liProtocolBridgeProduceRequestInstrumentationActive), tags)
metricsGroup.newGauge(RequestMetricBucketsEnabled,
() => enabled(config.liProtocolBridgeRequestMetricBucketsActive), tags)
metricsGroup.newGauge(RequestChannelWatchdogEnabled,
() => enabled(config.liProtocolBridgeRequestChannelWatchdogActive), tags)
metricsGroup.newGauge(MinimumLogRollEnabled,
() => enabled(config.liProtocolBridgeMinimumLogRollActive), tags)
metricsGroup.newGauge(ReassignmentCancellationSafetyEnabled,
() => enabled(config.liProtocolBridgeReassignmentCancellationSafetyActive), tags)
metricsGroup.newGauge(ListOffsetsInstrumentationEnabled,
() => enabled(config.liProtocolBridgeListOffsetsInstrumentationActive), tags)
metricsGroup.newGauge(StaticDefaultQuotasEnabled,
() => enabled(config.liProtocolBridgeStaticDefaultQuotasActive), tags)
metricsGroup.newGauge(ReplicaRequestTimeoutEnabled,
() => enabled(config.liProtocolBridgeReplicaRequestTimeoutActive), tags)
metricsGroup.newGauge(OffsetsTopicConfigEnabled,
() => enabled(config.liProtocolBridgeOffsetsTopicConfigActive), tags)
metricsGroup.newGauge(LeaderTransferEnabled,
() => enabled(config.liProtocolBridgeLeaderTransferActive), tags)
metricsGroup.newGauge(LegacyRequestMetricsEnabled,
() => enabled(config.liProtocolBridgeLegacyRequestMetricsActive), tags)
metricsGroup.newGauge(LogTruncationMetricsEnabled,
() => enabled(config.liProtocolBridgeLogTruncationMetricsActive), tags)
private val registrationEnabled = config.liProtocolBridgeConfigMetricsActive
if (registrationEnabled) {
metricsGroup.newGauge(ConfigMetricsEnabled, () => 1, tags)
metricsGroup.newGauge(ModeEnabled, () => enabled(config.liProtocolBridgeModeActive), tags)
metricsGroup.newGauge(TopicDeletionStateCleanupEnabled,
() => enabled(config.liProtocolBridgeTopicDeletionStateCleanupActive), tags)
metricsGroup.newGauge(FollowerRecoveryEnabled,
() => enabled(config.liProtocolBridgeFollowerRecoveryActive), tags)
metricsGroup.newGauge(RecommendedLeaderElectionEnabled,
() => enabled(config.liProtocolBridgeRecommendedElectionActive), tags)
metricsGroup.newGauge(ExcludePartitionsEnabled,
() => enabled(config.liProtocolBridgeExcludePartitionsActive), tags)
metricsGroup.newGauge(MoveControllerEnabled,
() => enabled(config.liProtocolBridgeMoveControllerActive), tags)
metricsGroup.newGauge(ShutdownSafetyOverrideEnabled,
() => enabled(config.liProtocolBridgeShutdownSafetyOverrideActive), tags)
metricsGroup.newGauge(PreferredControllerEnabled,
() => enabled(config.liProtocolBridgePreferredControllerActive), tags)
metricsGroup.newGauge(FederatedTopicsEnabled,
() => enabled(config.liProtocolBridgeFederatedTopicsActive), tags)
metricsGroup.newGauge(RackIdMapperEnabled,
() => enabled(config.liProtocolBridgeRackIdMapperActive), tags)
metricsGroup.newGauge(ZookeeperPaginationEnabled,
() => enabled(config.liZookeeperPaginationEnable), tags)
metricsGroup.newGauge(DynamicTopicDeletionEnabled,
() => enabled(config.liProtocolBridgeDynamicTopicDeletionActive), tags)
metricsGroup.newGauge(ControllerInitializationThreads,
() => config.liNumControllerInitThreads, tags)
metricsGroup.newGauge(ProduceRequestInstrumentationEnabled,
() => enabled(config.liProtocolBridgeProduceRequestInstrumentationActive), tags)
metricsGroup.newGauge(RequestMetricBucketsEnabled,
() => enabled(config.liProtocolBridgeRequestMetricBucketsActive), tags)
metricsGroup.newGauge(RequestChannelWatchdogEnabled,
() => enabled(config.liProtocolBridgeRequestChannelWatchdogActive), tags)
metricsGroup.newGauge(MinimumLogRollEnabled,
() => enabled(config.liProtocolBridgeMinimumLogRollActive), tags)
metricsGroup.newGauge(ReassignmentCancellationSafetyEnabled,
() => enabled(config.liProtocolBridgeReassignmentCancellationSafetyActive), tags)
metricsGroup.newGauge(ListOffsetsInstrumentationEnabled,
() => enabled(config.liProtocolBridgeListOffsetsInstrumentationActive), tags)
metricsGroup.newGauge(StaticDefaultQuotasEnabled,
() => enabled(config.liProtocolBridgeStaticDefaultQuotasActive), tags)
metricsGroup.newGauge(ReplicaRequestTimeoutEnabled,
() => enabled(config.liProtocolBridgeReplicaRequestTimeoutActive), tags)
metricsGroup.newGauge(OffsetsTopicConfigEnabled,
() => enabled(config.liProtocolBridgeOffsetsTopicConfigActive), tags)
metricsGroup.newGauge(LeaderTransferEnabled,
() => enabled(config.liProtocolBridgeLeaderTransferActive), tags)
metricsGroup.newGauge(LegacyRequestMetricsEnabled,
() => enabled(config.liProtocolBridgeLegacyRequestMetricsActive), tags)
metricsGroup.newGauge(LogTruncationMetricsEnabled,
() => enabled(config.liProtocolBridgeLogTruncationMetricsActive), tags)
}

private def enabled(value: Boolean): Int = if (value) 1 else 0

override def close(): Unit = {
MetricNames.foreach { name =>
if (registrationEnabled) MetricNames.foreach { name =>
metricsGroup.removeMetric(name, tags)
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,28 +20,77 @@ import com.yammer.metrics.core.Gauge
import kafka.utils.TestUtils
import org.apache.kafka.server.config.ReplicationConfigs
import org.apache.kafka.server.metrics.KafkaYammerMetrics
import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue}
import org.junit.jupiter.api.Assertions.{assertEquals, assertFalse, assertTrue}
import org.junit.jupiter.api.Test

import java.util.Properties
import scala.jdk.CollectionConverters._

class LiProtocolBridgeMetricsTest {

@Test
def testNoMetricsWithoutOptIn(): Unit = {
val brokerId = 988
val config = KafkaConfig(TestUtils.createBrokerConfig(brokerId, TestUtils.MockZkConnect))
val metrics = new LiProtocolBridgeMetrics(config)
try assertTrue(metricValues(brokerId).isEmpty)
finally metrics.close()
}

@Test
def testGateRequiresRestartAndDisabledClosePreservesOtherMetrics(): Unit = {
val brokerId = 989
val props = TestUtils.createBrokerConfig(brokerId, TestUtils.MockZkConnect)
props.put(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp, "true")
val config = KafkaConfig(props)
config.dynamicConfig.initialize(None, None)
val metrics = new LiProtocolBridgeMetrics(config)
try {
assertFalse(DynamicBrokerConfig.AllDynamicConfigs.contains(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp))
val update = new Properties
update.put(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp, "false")
config.dynamicConfig.updateDefaultConfig(update)
assertTrue(config.liProtocolBridgeConfigMetricsActive)
props.remove(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp)
new LiProtocolBridgeMetrics(KafkaConfig(props)).close()
assertEquals(LiProtocolBridgeMetrics.MetricNames.toSet, metricValues(brokerId).keySet)
assertEquals(1, metricValues(brokerId)(LiProtocolBridgeMetrics.ConfigMetricsEnabled))
} finally metrics.close()
assertTrue(metricValues(brokerId).isEmpty)
}

@Test
def testKRaftDoesNotRegisterBridgeMetrics(): Unit = {
val props = new Properties
Map("broker.id" -> "990", "node.id" -> "990", "process.roles" -> "broker,controller",
"controller.quorum.voters" -> "990@localhost:19093", "controller.listener.names" -> "CONTROLLER",
"listeners" -> "PLAINTEXT://localhost:19092,CONTROLLER://localhost:19093",
"listener.security.protocol.map" -> "PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT",
KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp -> "true").foreach { case (key, value) => props.put(key, value) }
val config = KafkaConfig(props)
assertFalse(config.liProtocolBridgeConfigMetricsActive)
val metrics = new LiProtocolBridgeMetrics(config)
try assertTrue(metricValues(990).isEmpty)
finally metrics.close()
}

@Test
def testMetricsFollowDynamicFlags(): Unit = {
val brokerId = 987
val brokerProps = TestUtils.createBrokerConfig(brokerId, TestUtils.MockZkConnect)
brokerProps.put(ReplicationConfigs.INTER_BROKER_PROTOCOL_VERSION_CONFIG, "3.0")
brokerProps.put(KafkaConfig.LiProtocolBridgeConfigMetricsEnableProp, "true")
val config = KafkaConfig(brokerProps)
config.dynamicConfig.initialize(None, None)
val metrics = new LiProtocolBridgeMetrics(config)

try {
val initialValues = metricValues(brokerId)
assertEquals(1, initialValues(LiProtocolBridgeMetrics.ControllerInitializationThreads))
assertTrue(initialValues.filterNot(_._1 == LiProtocolBridgeMetrics.ControllerInitializationThreads)
.values.forall(_ == 0))
assertEquals(1, initialValues(LiProtocolBridgeMetrics.ConfigMetricsEnabled))
assertTrue(initialValues.filterNot { case (name, _) =>
name == LiProtocolBridgeMetrics.ControllerInitializationThreads || name == LiProtocolBridgeMetrics.ConfigMetricsEnabled
}.values.forall(_ == 0))

val props = new Properties
Seq(
Expand Down
5 changes: 3 additions & 2 deletions docs/ops/li-bridge-review-comments.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,8 @@ All 62 original threads have replies with published source decisions. A fresh Gr
| [F18: offline replica misses deletion](https://github.com/linkedin/kafka/pull/576#issuecomment-5629453741) | Paired complete-image recovery in 583/584; unassigned/leaderless/incremental tests and exact records after promotion. | Assignment does not establish topic identity; also require F19. |
| [F19: recreated topic already assigned to returning replica](https://github.com/linkedin/kafka/pull/584#issuecomment-5631079569) | Paired identity recovery in 585/586; missing/zero IDs, errors, retries, current/future copies and mixed-batch tests. [Qualification update](https://github.com/linkedin/kafka/pull/586#issuecomment-5638005151). | All four revision-4 record checks passed, but full final-source and wrapper qualification remain open. |
| [F20: rotated protocol logs omitted](https://github.com/linkedin/kafka/pull/584#issuecomment-5638004839) | PR 584 retains and scans hourly rotations; 586 is restacked on it. New positive and negative tests fail before the fix and pass afterward. | The previous failed CI job stays failed. Re-run the corrected collector. |
| [F21: churn exits during controller movement](https://github.com/linkedin/kafka/pull/586#issuecomment-5638182322) | PR 584 fixes the workload retry policy without changing the upstream broker response. Tests with both client archives require controller retries and reject data/auth/record errors; the process setup runs them for both generations. | All 71 Python tests pass. Full migration and wrapper qualification remain required. |
| [F21: churn exits during controller movement](https://github.com/linkedin/kafka/pull/586#issuecomment-5638182322) | PR 584 fixes the workload retry policy without changing the upstream broker response. [Test and code update](https://github.com/linkedin/kafka/pull/584#issuecomment-5638482746). | The complete revision-4 process run and audit pass; final F22/wrapper qualification remains required. |
| F22: bridge-state MBeans register by default | Paired [588](https://github.com/linkedin/kafka/pull/588)/[589](https://github.com/linkedin/kafka/pull/589) add a default-off diagnostics flag. Tests cover disabled registration, enabled readings, restart scope, KRaft and cleanup. | All 72 Python tests pass. The wrapper mapping still needs matching-jar qualification. |

## PR 541

Expand Down Expand Up @@ -108,7 +109,7 @@ All 62 original threads have replies with published source decisions. A fresh Gr
| [7](https://github.com/linkedin/kafka/pull/551#discussion_r3927365100) | Synchronize close and clear the metrics map. | RequestChannel.Metrics.close uses the same monitor as apply. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984230263) |
| [8](https://github.com/linkedin/kafka/pull/551#discussion_r3927365140) | Prefix the watchdog histogram with the request-channel metric prefix. | RequestChannel and RequestChannelWatchdogTest. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984230448) |
| [9](https://github.com/linkedin/kafka/pull/551#discussion_r3927365177) | Derive the watchdog check interval from the configured timeout. | KafkaServerTest.testRequestChannelWatchdogIntervalTracksConfiguredTimeout. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984230674) |
| [10](https://github.com/linkedin/kafka/pull/551#discussion_r3927365245) | Remove the hard-coded count. Check uniqueness and compare the actual registry with the shared Python contract. | LiProtocolBridgeConfigTest; preflight/scenario registry-consistency tests; 23 gates are currently registered. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984230900) |
| [10](https://github.com/linkedin/kafka/pull/551#discussion_r3927365245) | Remove the hard-coded count. Check uniqueness and compare the actual registry with the shared Python contract. | LiProtocolBridgeConfigTest; preflight/scenario registry-consistency tests; 24 gates are currently registered. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984230900) |
| [11](https://github.com/linkedin/kafka/pull/551#discussion_r3927475761) | Cache the total topic-name length and update it with topic membership. | ControllerContextTest.testTopicNameLengthTotalTracksTopicChanges; controller gauge reads the cached total. | [reply](https://github.com/linkedin/kafka/pull/551#discussion_r3984231122) |

## PR 552
Expand Down
Loading
Loading