Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
96de66e
add interface and tests
jolshan Apr 8, 2024
8f7bc87
Fix 2.12 compilation issue
jolshan Apr 9, 2024
ccb7a50
satisfy scala 2.12
jolshan Apr 12, 2024
b4e2e21
Few fixes to allow the cluster to start up.
jolshan Apr 17, 2024
a31a71e
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan Apr 17, 2024
92b7326
supporting supported features
jolshan Apr 18, 2024
0b0b9ab
Some cleanups.
jolshan Apr 25, 2024
70163cb
simplify interface
jolshan May 10, 2024
940653a
Further iterate on interface
jolshan May 11, 2024
50ee633
More fixes
jolshan May 13, 2024
7f162b8
note cleanups
jolshan May 13, 2024
db0b9a0
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan May 13, 2024
830a637
fix compilation error
jolshan May 13, 2024
0dcada4
fix default for MV
jolshan May 15, 2024
3988422
updates to the simpler comments
jolshan May 16, 2024
4e76ba4
Renaming and refactoring
jolshan May 20, 2024
1ca5112
Refactor method
jolshan May 20, 2024
f55811f
compilation fixes
jolshan May 20, 2024
435808a
Followups
jolshan May 21, 2024
13a8840
Address feedback
jolshan May 22, 2024
5fc9fea
more comments
jolshan May 22, 2024
4e1989a
Add 0 version back
jolshan May 22, 2024
1a43f18
Re-remove version 0
jolshan May 22, 2024
d86ff02
fix supported versions, set up production features
jolshan May 22, 2024
1d3fab8
Fix when there is no version default flag
jolshan May 23, 2024
e325a88
Fixes
jolshan May 23, 2024
b56707b
More adjustments
jolshan May 23, 2024
2416e95
Simplify production ready
jolshan May 24, 2024
1a474ef
Remove unneeded value
jolshan May 24, 2024
a28ece3
Fix conditional
jolshan May 24, 2024
5f1ac97
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan May 28, 2024
24e14ed
Addressing Jun's comments
jolshan May 29, 2024
00b63ad
Fix bug
jolshan May 29, 2024
0704b61
Cleanups
jolshan May 29, 2024
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
104 changes: 86 additions & 18 deletions core/src/main/scala/kafka/tools/StorageTool.scala
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ import net.sourceforge.argparse4j.inf.Namespace
import org.apache.kafka.common.Uuid
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.metadata.bootstrap.{BootstrapDirectory, BootstrapMetadata}
import org.apache.kafka.server.common.{ApiMessageAndVersion, MetadataVersion}
import org.apache.kafka.server.common.{ApiMessageAndVersion, FeatureVersion, MetadataVersion}
import org.apache.kafka.common.metadata.FeatureLevelRecord
import org.apache.kafka.common.metadata.UserScramCredentialRecord
import org.apache.kafka.common.security.scram.internals.ScramMechanism
Expand All @@ -37,8 +37,7 @@ import org.apache.kafka.metadata.properties.MetaPropertiesEnsemble.VerificationF
import org.apache.kafka.metadata.properties.{MetaProperties, MetaPropertiesEnsemble, MetaPropertiesVersion, PropertiesUtils}

import java.util
import java.util.Base64
import java.util.Optional
import java.util.{Base64, Collections, Optional}
import scala.collection.mutable
import scala.jdk.CollectionConverters._
import scala.collection.mutable.ArrayBuffer
Expand All @@ -59,24 +58,19 @@ object StorageTool extends Logging {
case "format" =>
val directories = configToLogDirectories(config.get)
val clusterId = namespace.getString("cluster_id")
val metadataVersion = getMetadataVersion(namespace,
Option(config.get.originals.get(KafkaConfig.InterBrokerProtocolVersionProp)).map(_.toString))
if (!metadataVersion.isKRaftSupported) {
throw new TerseFailure(s"Must specify a valid KRaft metadata.version of at least ${MetadataVersion.IBP_3_0_IV0}.")
}
if (!metadataVersion.isProduction) {
if (config.get.unstableMetadataVersionsEnabled) {
System.out.println(s"WARNING: using pre-production metadata.version $metadataVersion.")
} else {
throw new TerseFailure(s"The metadata.version $metadataVersion is not ready for production use yet.")
}
}
val metaProperties = new MetaProperties.Builder().
setVersion(MetaPropertiesVersion.V1).
setClusterId(clusterId).
setNodeId(config.get.nodeId).
build()
val metadataRecords : ArrayBuffer[ApiMessageAndVersion] = ArrayBuffer()
val specifiedFeatures: util.List[String] = namespace.getList("feature")
val featureNamesAndLevelsMap = featureNamesAndLevels(Option(specifiedFeatures).getOrElse(Collections.emptyList).asScala.toList)
val metadataVersion = getMetadataVersion(namespace, featureNamesAndLevelsMap,
Option(config.get.originals.get(KafkaConfig.InterBrokerProtocolVersionProp)).map(_.toString))
metadataVersionValidation(metadataVersion, config)
// Get all other features, validate, and create records for them
metadataRecords.appendAll(generateFeatureRecords(metadataVersion, featureNamesAndLevelsMap, FeatureVersion.PRODUCTION_FEATURES.asScala.toList))
Comment thread
jolshan marked this conversation as resolved.
Outdated
getUserScramCredentialRecords(namespace).foreach(userScramCredentialRecords => {
if (!metadataVersion.isScramSupported) {
throw new TerseFailure(s"SCRAM is only supported in metadata.version ${MetadataVersion.IBP_3_5_IV2} or later.")
Expand All @@ -85,6 +79,7 @@ object StorageTool extends Logging {
metadataRecords.append(new ApiMessageAndVersion(record, 0.toShort))
}
})

val bootstrapMetadata = buildBootstrapMetadata(metadataVersion, Some(metadataRecords), "format command")
val ignoreFormatted = namespace.getBoolean("ignore_formatted")
if (!configToSelfManagedMode(config.get)) {
Expand All @@ -108,6 +103,43 @@ object StorageTool extends Logging {
}
}

def metadataVersionValidation(metadataVersion: MetadataVersion, config: Option[KafkaConfig]): Unit = {
Comment thread
jolshan marked this conversation as resolved.
Outdated
if (!metadataVersion.isKRaftSupported) {
throw new TerseFailure(s"Must specify a valid KRaft metadata.version of at least ${MetadataVersion.IBP_3_0_IV0}.")
}
if (!metadataVersion.isProduction) {
if (config.get.unstableMetadataVersionsEnabled) {
System.out.println(s"WARNING: using pre-production metadata.version $metadataVersion.")
} else {
throw new TerseFailure(s"The metadata.version $metadataVersion is not ready for production use yet.")
}
}
}

def generateFeatureRecords(metadataVersion: MetadataVersion,
Comment thread
jolshan marked this conversation as resolved.
Outdated
specifiedFeatures: Map[String, FeatureVersion],
allFeatures: List[String]): List[ApiMessageAndVersion] = {
// If we are using --version-default, the default is based on the metadata version.
Comment thread
jolshan marked this conversation as resolved.
Outdated
val metadataVersionOpt: Optional[MetadataVersion] = if (specifiedFeatures.isEmpty) Optional.of(metadataVersion) else Optional.empty[MetadataVersion]

val featuresMinusMetadata = allFeatures.filter(!_.equals(MetadataVersion.FEATURE_NAME))
val allFeaturesAndVersions = featuresMinusMetadata.map { featureName =>
specifiedFeatures.getOrElse(featureName, FeatureVersion.defaultValue(featureName, metadataVersionOpt))
}

try {
// In order to validate, we need all feature versions set.
allFeaturesAndVersions.map { case feature =>
feature.validateVersion(metadataVersion, allFeaturesAndVersions.asJava)
new ApiMessageAndVersion(new FeatureLevelRecord().
setName(feature.featureName).
setFeatureLevel(feature.featureLevel), 0.toShort)
}
} catch {
case e: Throwable => throw new TerseFailure(e.getMessage)
}
}

def parseArguments(args: Array[String]): Namespace = {
val parser = ArgumentParsers.
newArgumentParser("kafka-storage", /* defaultHelp */ true, /* prefixChars */ "-", /* fromFilePrefix */ "@").
Expand Down Expand Up @@ -140,6 +172,9 @@ object StorageTool extends Logging {
formatParser.addArgument("--release-version", "-r").
action(store()).
help(s"A KRaft release version to use for the initial metadata.version. The minimum is ${MetadataVersion.IBP_3_0_IV0}, the default is ${MetadataVersion.LATEST_PRODUCTION}")
formatParser.addArgument("--feature").
help("A feature upgrade we should perform, in feature=level format. For example: `metadata.version=5`.").
Comment thread
jolshan marked this conversation as resolved.
action(append());

parser.parseArgsOrFail(args)
}
Expand All @@ -155,16 +190,27 @@ object StorageTool extends Logging {

def getMetadataVersion(
namespace: Namespace,
featureNamesAndLevelsMap: Map[String, FeatureVersion],
defaultVersionString: Option[String]
): MetadataVersion = {
val defaultValue = defaultVersionString match {
case Some(versionString) => MetadataVersion.fromVersionString(versionString)
case None => MetadataVersion.LATEST_PRODUCTION
}

Option(namespace.getString("release_version"))
.map(ver => MetadataVersion.fromVersionString(ver))
.getOrElse(defaultValue)
val releaseVersionTag = Option(namespace.getString("release_version"))
val featureTag = featureNamesAndLevelsMap.get(MetadataVersion.FEATURE_NAME)

(releaseVersionTag, featureTag) match {
case (Some(_), Some(_)) =>
throw new IllegalArgumentException("Both --release_version and --feature were set. Only one of the two flags can be set.")
Comment thread
jolshan marked this conversation as resolved.
case (Some(version), None) =>
MetadataVersion.fromVersionString(version)
case (None, Some(feature)) =>
MetadataVersion.fromFeatureLevel(feature.featureLevel())
case (None, None) =>
defaultValue
}
}

private def getUserScramCredentialRecord(
Expand Down Expand Up @@ -464,4 +510,26 @@ object StorageTool extends Logging {
}
0
}

private def parseNameAndLevel(input: String): Array[String] = {
val equalsIndex = input.indexOf("=")
if (equalsIndex < 0)
throw new RuntimeException("Can't parse feature=level string " + input + ": equals sign not found.")
val name = input.substring(0, equalsIndex).trim
val levelString = input.substring(equalsIndex + 1).trim
try {
levelString.toShort
} catch {
case _: Throwable =>
throw new RuntimeException("Can't parse feature=level string " + input + ": " + "unable to parse " + levelString + " as a short.")
}
Array[String](name, levelString)
Comment thread
jolshan marked this conversation as resolved.
Outdated
}

def featureNamesAndLevels(features: List[String]): Map[String, FeatureVersion] = {
features.map((feature: String) => {
Comment thread
jolshan marked this conversation as resolved.
Outdated
val nameAndLevel = parseNameAndLevel(feature)
(nameAndLevel(0), FeatureVersion.createFeature(nameAndLevel(0), java.lang.Short.valueOf(nameAndLevel(1))))
}).toMap
}
}
97 changes: 90 additions & 7 deletions core/src/test/scala/unit/kafka/tools/StorageToolTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -21,23 +21,25 @@ import java.io.{ByteArrayOutputStream, PrintStream}
import java.nio.charset.StandardCharsets
import java.nio.file.{Files, Paths}
import java.util
import java.util.Properties
import java.util.{Collections, Properties}
import org.apache.kafka.common.{DirectoryId, KafkaException}
import kafka.server.KafkaConfig
import kafka.utils.Exit
import kafka.utils.TestUtils
import net.sourceforge.argparse4j.inf.Namespace
import org.apache.commons.io.output.NullOutputStream
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.server.common.MetadataVersion
import org.apache.kafka.common.metadata.UserScramCredentialRecord
import org.apache.kafka.server.common.{ApiMessageAndVersion, FeatureVersion, MetadataVersion, TestFeatureVersion}
import org.apache.kafka.common.metadata.{FeatureLevelRecord, UserScramCredentialRecord}
import org.apache.kafka.metadata.properties.{MetaProperties, MetaPropertiesEnsemble, MetaPropertiesVersion, PropertiesUtils}
import org.junit.jupiter.api.Assertions.{assertEquals, assertFalse, assertThrows, assertTrue}
import org.junit.jupiter.api.{Test, Timeout}
import org.junit.jupiter.params.ParameterizedTest
import org.junit.jupiter.params.provider.ValueSource
import org.junit.jupiter.params.provider.{EnumSource, ValueSource}

import scala.collection.mutable
import scala.collection.mutable.ArrayBuffer
import scala.jdk.CollectionConverters._

@Timeout(value = 40)
class StorageToolTest {
Expand All @@ -52,6 +54,8 @@ class StorageToolTest {
properties
}

val allFeatures = List(MetadataVersion.FEATURE_NAME, TestFeatureVersion.FEATURE_NAME)

@Test
def testConfigToLogDirectories(): Unit = {
val config = new KafkaConfig(newSelfManagedProperties())
Expand Down Expand Up @@ -204,26 +208,67 @@ Found problem:
@Test
def testDefaultMetadataVersion(): Unit = {
val namespace = StorageTool.parseArguments(Array("format", "-c", "config.props", "-t", "XcZZOzUqS4yHOjhMQB6JLQ"))
val mv = StorageTool.getMetadataVersion(namespace, defaultVersionString = None)
val mv = StorageTool.getMetadataVersion(namespace, Map.empty, defaultVersionString = None)
assertEquals(MetadataVersion.LATEST_PRODUCTION.featureLevel(), mv.featureLevel(),
"Expected the default metadata.version to be the latest production version")
}

@Test
def testConfiguredMetadataVersion(): Unit = {
val namespace = StorageTool.parseArguments(Array("format", "-c", "config.props", "-t", "XcZZOzUqS4yHOjhMQB6JLQ"))
val mv = StorageTool.getMetadataVersion(namespace, defaultVersionString = Some(MetadataVersion.IBP_3_3_IV2.toString))
val mv = StorageTool.getMetadataVersion(namespace, Map.empty, defaultVersionString = Some(MetadataVersion.IBP_3_3_IV2.toString))
assertEquals(MetadataVersion.IBP_3_3_IV2.featureLevel(), mv.featureLevel(),
"Expected the default metadata.version to be 3.3-IV2")
}

@Test
def testSettingFeatureAndReleaseVersionFails(): Unit = {
val namespace = StorageTool.parseArguments(Array("format", "-c", "config.props", "-t", "XcZZOzUqS4yHOjhMQB6JLQ",
"--release-version", "3.0-IV1", "--feature", "metadata.version=4"))
assertThrows(classOf[IllegalArgumentException], () => StorageTool.getMetadataVersion(namespace, parseFeatures(namespace), defaultVersionString = None))
}

@Test
def testParseFeatures(): Unit = {
def parseAddFeatures(strings: String*): Map[String, FeatureVersion] = {
var args = mutable.Seq("format", "-c", "config.props", "-t", "XcZZOzUqS4yHOjhMQB6JLQ")
args ++= strings
val namespace = StorageTool.parseArguments(args.toArray)
parseFeatures(namespace)
}

assertThrows(classOf[RuntimeException], () => parseAddFeatures("--feature", "blah"))
assertThrows(classOf[RuntimeException], () => parseAddFeatures("--feature", "blah=blah"))
assertThrows(classOf[IllegalArgumentException], () => parseAddFeatures("--feature", "blah=12"))

// Test with no features
assertEquals(Map(), parseAddFeatures())

// Test with one feature
val testFeatureLevel = 1
val testArgument = TestFeatureVersion.FEATURE_NAME + "=" + testFeatureLevel.toString
val expectedMap = Map(TestFeatureVersion.FEATURE_NAME -> FeatureVersion.createFeature(TestFeatureVersion.FEATURE_NAME, testFeatureLevel.toShort))
assertEquals(expectedMap, parseAddFeatures("--feature", testArgument))

// Test with two features
val metadataFeatureLevel = 5
val metadataArgument = MetadataVersion.FEATURE_NAME + "=" + metadataFeatureLevel.toString
val expectedMap2 = expectedMap ++ Map (MetadataVersion.FEATURE_NAME -> FeatureVersion.createFeature(MetadataVersion.FEATURE_NAME, metadataFeatureLevel.toShort))
assertEquals(expectedMap2, parseAddFeatures("--feature", testArgument, "--feature", metadataArgument))
}

def parseFeatures(namespace: Namespace): Map[String, FeatureVersion] = {
val specifiedFeatures: util.List[String] = namespace.getList("feature")
StorageTool.featureNamesAndLevels(Option(specifiedFeatures).getOrElse(Collections.emptyList).asScala.toList)
}

@Test
def testMetadataVersionFlags(): Unit = {
def parseMetadataVersion(strings: String*): MetadataVersion = {
var args = mutable.Seq("format", "-c", "config.props", "-t", "XcZZOzUqS4yHOjhMQB6JLQ")
args ++= strings
val namespace = StorageTool.parseArguments(args.toArray)
StorageTool.getMetadataVersion(namespace, defaultVersionString = None)
StorageTool.getMetadataVersion(namespace, Map.empty, defaultVersionString = None)
}

var mv = parseMetadataVersion("--release-version", "3.0")
Expand All @@ -235,6 +280,44 @@ Found problem:
assertThrows(classOf[IllegalArgumentException], () => parseMetadataVersion("--release-version", "0.0"))
}

def generateRecord(featureName: String, level: Short): ApiMessageAndVersion = {
Comment thread
jolshan marked this conversation as resolved.
Outdated
new ApiMessageAndVersion(new FeatureLevelRecord().
setName(featureName).
setFeatureLevel(level), 0.toShort)
}

@ParameterizedTest
@EnumSource(classOf[TestFeatureVersion])
def testFeatureFlag(testFeatureVersion: TestFeatureVersion): Unit = {
val featureLevel = testFeatureVersion.featureLevel
val records = StorageTool.generateFeatureRecords(MetadataVersion.LATEST_PRODUCTION,
Map(TestFeatureVersion.FEATURE_NAME -> FeatureVersion.createFeature(TestFeatureVersion.FEATURE_NAME, featureLevel)),
allFeatures
)
assertEquals(List(generateRecord(TestFeatureVersion.FEATURE_NAME, featureLevel)), records)
}

@ParameterizedTest
@EnumSource(classOf[MetadataVersion])
def testVersionDefault(metadataVersion: MetadataVersion): Unit = {
val records = StorageTool.generateFeatureRecords(metadataVersion,
Map.empty,
allFeatures
)

val featureLevel = TestFeatureVersion.metadataVersionMapping(metadataVersion).featureLevel()
assertEquals(List(generateRecord(TestFeatureVersion.FEATURE_NAME, featureLevel)), records)
}

@Test
def testFeatureDependency(): Unit = {
val featureLevel = 1.toShort
assertThrows(classOf[TerseFailure], () => StorageTool.generateFeatureRecords(MetadataVersion.IBP_2_8_IV1,
Map(TestFeatureVersion.FEATURE_NAME -> FeatureVersion.createFeature(TestFeatureVersion.FEATURE_NAME, featureLevel)),
allFeatures
))
}

@Test
def testAddScram():Unit = {
def parseAddScram(strings: String*): Option[ArrayBuffer[UserScramCredentialRecord]] = {
Expand Down
Loading