wombatu-kun commented on code in PR #17469:
URL: https://github.com/apache/iceberg/pull/17469#discussion_r3697945539


##########
parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java:
##########
@@ -607,6 +612,130 @@ protected int resolveColumnIndex(Void engineSchema, 
String columnName) {
     }
   }
 
+  @Test
+  public void testFormatModelVariantShreddingWithEncryption() throws 
IOException {
+    Schema variantSchema =
+        new Schema(
+            Types.NestedField.required(1, "id", Types.LongType.get()),
+            Types.NestedField.optional(2, "v", Types.VariantType.get()));
+
+    VariantShreddingAnalyzer<Record, Void> analyzer =

Review Comment:
   The schema, the anonymous VariantShreddingAnalyzer, the variant fixture and 
the ParquetFormatModel.create wiring here are a verbatim copy of 
testFormatModelVariantShreddingRoundTrip, differing only in the literal values. 
Extract them into a shared helper that both tests call, so this test carries 
only the encryption-specific setup.



##########
parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java:
##########
@@ -607,6 +612,130 @@ protected int resolveColumnIndex(Void engineSchema, 
String columnName) {
     }
   }
 
+  @Test
+  public void testFormatModelVariantShreddingWithEncryption() throws 
IOException {
+    Schema variantSchema =
+        new Schema(
+            Types.NestedField.required(1, "id", Types.LongType.get()),
+            Types.NestedField.optional(2, "v", Types.VariantType.get()));
+
+    VariantShreddingAnalyzer<Record, Void> analyzer =
+        new VariantShreddingAnalyzer<Record, Void>() {
+          @Override
+          protected List<VariantValue> extractVariantValues(List<Record> rows, 
int idx) {
+            List<VariantValue> values = Lists.newArrayList();
+            for (Record row : rows) {
+              Object obj = row.get(idx);
+              if (obj instanceof Variant) {
+                values.add(((Variant) obj).value());
+              }
+            }
+            return values;
+          }
+
+          @Override
+          protected int resolveColumnIndex(Void engineSchema, String 
columnName) {
+            return 
variantSchema.columns().indexOf(variantSchema.findField(columnName));
+          }
+        };
+
+    ByteBuffer metadataBuffer = 
VariantTestUtil.createMetadata(ImmutableList.of("a", "b"), true);
+    VariantMetadata metadata = Variants.metadata(metadataBuffer);
+    ByteBuffer objectBuffer =
+        VariantTestUtil.createObject(
+            metadataBuffer,
+            ImmutableMap.of("a", Variants.of(123456789), "b", 
Variants.of("string")));
+    Variant variant = Variant.of(metadata, Variants.value(metadata, 
objectBuffer));
+
+    GenericRecord record = GenericRecord.create(variantSchema);
+    List<Record> variantRecords =
+        ImmutableList.of(
+            record.copy(ImmutableMap.of("id", 1L, "v", variant)),
+            record.copy(ImmutableMap.of("id", 2L, "v", variant)),
+            record.copy(ImmutableMap.of("id", 3L, "v", variant)));
+
+    ParquetFormatModel<Record, Void, ParquetValueReader<?>> model =
+        ParquetFormatModel.create(
+            Record.class,
+            Void.class,
+            (icebergSchema, messageType, engineSchema) ->
+                GenericParquetWriter.create(icebergSchema, messageType),
+            (icebergSchema, fileSchema, engineSchema, idToConstant) ->
+                GenericParquetReaders.buildReader(icebergSchema, fileSchema),
+            analyzer,
+            (Function<Void, UnaryOperator<Record>>) unused -> input -> input);
+
+    OutputFile encryptedFile = Files.localOutput(createTempFile(temp));
+    ByteBuffer fileDek = ByteBuffer.allocate(16);
+    ByteBuffer aadPrefix = ByteBuffer.allocate(16);
+    SecureRandom random = new SecureRandom();
+    random.nextBytes(fileDek.array());
+    random.nextBytes(aadPrefix.array());
+
+    try (FileAppender<Record> appender =
+        model
+            .writeBuilder(EncryptedFiles.plainAsEncryptedOutput(encryptedFile))
+            .schema(variantSchema)
+            .withFileEncryptionKey(fileDek)
+            .withAADPrefix(aadPrefix)
+            .setAll(
+                ImmutableMap.of(
+                    TableProperties.PARQUET_SHRED_VARIANTS, "true",
+                    TableProperties.PARQUET_VARIANT_BUFFER_SIZE, "2"))
+            .content(FileContent.DATA)
+            .build()) {
+      assertThat(appender).isInstanceOf(BufferedFileAppender.class);
+      appender.addAll(variantRecords);
+    }
+
+    assertThatThrownBy(
+            () ->
+                Parquet.read(encryptedFile.toInputFile())
+                    .project(variantSchema)
+                    .createReaderFunc(
+                        fileSchema -> 
GenericParquetReaders.buildReader(variantSchema, fileSchema))
+                    .build()
+                    .iterator())
+        .isInstanceOf(ParquetCryptoRuntimeException.class)
+        .hasMessage("Trying to read file with encrypted footer. No keys 
available");
+
+    List<Record> writtenRecords;
+    try (CloseableIterable<Record> reader =
+        Parquet.read(encryptedFile.toInputFile())
+            .project(variantSchema)
+            .withFileEncryptionKey(fileDek)
+            .withAADPrefix(aadPrefix)
+            .createReaderFunc(
+                fileSchema -> GenericParquetReaders.buildReader(variantSchema, 
fileSchema))
+            .build()) {
+      writtenRecords = Lists.newArrayList(reader);
+    }
+
+    assertThat(writtenRecords).hasSameSizeAs(variantRecords);
+    for (int i = 0; i < variantRecords.size(); i++) {
+      InternalTestHelpers.assertEquals(
+          variantSchema.asStruct(), variantRecords.get(i), 
writtenRecords.get(i));
+    }
+
+    try (ParquetFileReader fileReader =
+        ParquetFileReader.open(
+            ParquetIO.file(encryptedFile.toInputFile()),
+            ParquetReadOptions.builder()
+                .withDecryption(
+                    FileDecryptionProperties.builder()
+                        .withFooterKey(fileDek.array())
+                        .withAADPrefix(aadPrefix.array())
+                        .build())
+                .build())) {
+      GroupType variantType =
+          
fileReader.getFooter().getFileMetaData().getSchema().getType("v").asGroupType();
+      assertThat(variantType.containsField("typed_value")).isTrue();
+      GroupType typedValue = variantType.getType("typed_value").asGroupType();
+      assertThat(typedValue.containsField("a")).isTrue();

Review Comment:
   testFormatModelVariantShreddingRoundTrip also asserts that value is empty 
and that typed_value.a/b hold the values, while this one asserts only the 
footer schema. ParquetReader.Builder.withDecryption accepts 
FileDecryptionProperties, so those raw-group checks work on the encrypted file 
too - was dropping them intentional?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to