diff --git a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java index 0d4550932bc1a..7c7d1a3c0e042 100644 --- a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java +++ b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java @@ -29,6 +29,7 @@ public final class ImageWriterOptions { public static class Builder { private MetadataVersion metadataVersion; + private MetadataVersion requestedMetadataVersion; private Consumer lossHandler = e -> { throw e; }; @@ -42,6 +43,7 @@ public Builder(MetadataImage image) { } public Builder setMetadataVersion(MetadataVersion metadataVersion) { + setRequestedMetadataVersion(metadataVersion); if (metadataVersion.isLessThan(MetadataVersion.MINIMUM_BOOTSTRAP_VERSION)) { // When writing an image, all versions less than 3.3-IV0 are treated as 3.0-IV1. // This is because those versions don't support FeatureLevelRecord. @@ -58,29 +60,40 @@ public Builder setRawMetadataVersion(MetadataVersion metadataVersion) { return this; } + public void setRequestedMetadataVersion(MetadataVersion orgMetadataVersion) { + this.requestedMetadataVersion = orgMetadataVersion; + } + public MetadataVersion metadataVersion() { return metadataVersion; } + public MetadataVersion requestedMetadataVersion() { + return requestedMetadataVersion; + } + public Builder setLossHandler(Consumer lossHandler) { this.lossHandler = lossHandler; return this; } public ImageWriterOptions build() { - return new ImageWriterOptions(metadataVersion, lossHandler); + return new ImageWriterOptions(metadataVersion, lossHandler, requestedMetadataVersion); } } private final MetadataVersion metadataVersion; + private final MetadataVersion requestedMetadataVersion; private final Consumer lossHandler; private ImageWriterOptions( MetadataVersion metadataVersion, - Consumer lossHandler + Consumer lossHandler, + MetadataVersion orgMetadataVersion ) { this.metadataVersion = metadataVersion; this.lossHandler = lossHandler; + this.requestedMetadataVersion = orgMetadataVersion; } public MetadataVersion metadataVersion() { @@ -88,7 +101,7 @@ public MetadataVersion metadataVersion() { } public void handleLoss(String loss) { - lossHandler.accept(new UnwritableMetadataException(metadataVersion, loss)); + lossHandler.accept(new UnwritableMetadataException(requestedMetadataVersion, loss)); } } diff --git a/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java b/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java index 86eca51b4c131..b417bebf8acbc 100644 --- a/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java @@ -138,7 +138,7 @@ private static void writeWithExpectedLosses( MetadataImage image = delta.apply(MetadataProvenance.EMPTY); RecordListWriter writer = new RecordListWriter(); image.write(writer, new ImageWriterOptions.Builder(). - setRawMetadataVersion(metadataVersion). + setMetadataVersion(metadataVersion). setLossHandler(lossConsumer). build()); assertEquals(expectedLosses, lossConsumer.losses, "Failed to get expected metadata losses."); diff --git a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java index 07ddd4ff8cbba..c8653c46d8aec 100644 --- a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java @@ -21,6 +21,8 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import java.util.function.Consumer; + import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -44,9 +46,28 @@ public void testSetMetadataVersion() { setMetadataVersion(version); if (i < MetadataVersion.MINIMUM_BOOTSTRAP_VERSION.ordinal()) { assertEquals(MetadataVersion.MINIMUM_KRAFT_VERSION, options.metadataVersion()); + assertEquals(version, options.requestedMetadataVersion()); } else { assertEquals(version, options.metadataVersion()); } } } + + @Test + public void testHandleLoss() { + String expectedMessage = "stuff"; + + for (int i = MetadataVersion.MINIMUM_KRAFT_VERSION.ordinal(); + i < MetadataVersion.VERSIONS.length; + i++) { + MetadataVersion version = MetadataVersion.VERSIONS[i]; + String formattedMessage = String.format("Metadata has been lost because the following could not be represented in metadata version %s: %s", version, expectedMessage); + Consumer customLossHandler = e -> assertEquals(formattedMessage, e.getMessage()); + ImageWriterOptions options = new ImageWriterOptions.Builder() + .setMetadataVersion(version) + .setLossHandler(customLossHandler) + .build(); + options.handleLoss(expectedMessage); + } + } }