From 64390fbc27122052144c76b36827315a23f28d12 Mon Sep 17 00:00:00 2001 From: Munendra S N Date: Sun, 2 Aug 2026 12:07:43 +0530 Subject: [PATCH] GCP: Support CMEK (kms-key-name) for GCS server-side encryption GCSOutputStream previously supported only CSEK (gcs.encryption-key), where the raw key is sent to GCS on every request. This adds CMEK support via a new gcs.kms-key-name property: the caller passes only the Cloud KMS key resource name, and GCS performs encryption/decryption server-side, so the caller's credentials need no KMS permission. The write path emits BlobWriteOption.kmsKeyName(...) when the property is set; reads require no change since CMEK decryption is automatic. CSEK and CMEK are mutually exclusive (GCS rejects both on one object), enforced at property-parse time with a fail-fast Preconditions check. --- .../org/apache/iceberg/gcp/GCPProperties.java | 18 ++++++++++++ .../iceberg/gcp/gcs/GCSOutputStream.java | 3 ++ .../apache/iceberg/gcp/TestGCPProperties.java | 29 +++++++++++++++++++ .../iceberg/gcp/gcs/TestGCSOutputStream.java | 29 ++++++++++++++++++- 4 files changed, 78 insertions(+), 1 deletion(-) diff --git a/gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java b/gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java index 2702ce565d4e..de736dc69b81 100644 --- a/gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java +++ b/gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java @@ -44,6 +44,12 @@ public class GCPProperties implements Serializable { public static final String GCS_ENCRYPTION_KEY = "gcs.encryption-key"; public static final String GCS_USER_PROJECT = "gcs.user-project"; + /** + * Cloud KMS key resource name for CMEK server-side encryption. Mutually exclusive with {@link + * #GCS_ENCRYPTION_KEY} (CSEK). + */ + public static final String GCS_KMS_KEY_NAME = "gcs.kms-key-name"; + public static final String GCS_CHANNEL_READ_CHUNK_SIZE = "gcs.channel.read.chunk-size-bytes"; public static final String GCS_CHANNEL_WRITE_CHUNK_SIZE = "gcs.channel.write.chunk-size-bytes"; @@ -87,6 +93,7 @@ public class GCPProperties implements Serializable { private String gcsDecryptionKey; private String gcsEncryptionKey; + private String gcsKmsKeyName; private String gcsUserProject; private Integer gcsChannelReadChunkSize; @@ -151,8 +158,15 @@ public GCPProperties(Map properties) { gcsDecryptionKey = properties.get(GCS_DECRYPTION_KEY); gcsEncryptionKey = properties.get(GCS_ENCRYPTION_KEY); + gcsKmsKeyName = properties.get(GCS_KMS_KEY_NAME); gcsUserProject = properties.get(GCS_USER_PROJECT); + Preconditions.checkState( + !(gcsEncryptionKey != null && gcsKmsKeyName != null), + "Invalid encryption settings: must not configure both %s (CSEK) and %s (CMEK)", + GCS_ENCRYPTION_KEY, + GCS_KMS_KEY_NAME); + if (properties.containsKey(GCS_CHANNEL_READ_CHUNK_SIZE)) { gcsChannelReadChunkSize = Integer.parseInt(properties.get(GCS_CHANNEL_READ_CHUNK_SIZE)); } @@ -219,6 +233,10 @@ public Optional encryptionKey() { return Optional.ofNullable(gcsEncryptionKey); } + public Optional kmsKeyName() { + return Optional.ofNullable(gcsKmsKeyName); + } + public Optional projectId() { return Optional.ofNullable(projectId); } diff --git a/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSOutputStream.java b/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSOutputStream.java index 3fb8aa4801ae..fa7d5a919656 100644 --- a/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSOutputStream.java +++ b/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSOutputStream.java @@ -106,6 +106,9 @@ private void openStream() { gcpProperties .encryptionKey() .ifPresent(key -> writeOptions.add(BlobWriteOption.encryptionKey(key))); + gcpProperties + .kmsKeyName() + .ifPresent(kmsKeyName -> writeOptions.add(BlobWriteOption.kmsKeyName(kmsKeyName))); gcpProperties .userProject() .ifPresent(userProject -> writeOptions.add(BlobWriteOption.userProject(userProject))); diff --git a/gcp/src/test/java/org/apache/iceberg/gcp/TestGCPProperties.java b/gcp/src/test/java/org/apache/iceberg/gcp/TestGCPProperties.java index 0ec2183fc355..d0fb1979cae7 100644 --- a/gcp/src/test/java/org/apache/iceberg/gcp/TestGCPProperties.java +++ b/gcp/src/test/java/org/apache/iceberg/gcp/TestGCPProperties.java @@ -18,6 +18,8 @@ */ package org.apache.iceberg.gcp; +import static org.apache.iceberg.gcp.GCPProperties.GCS_ENCRYPTION_KEY; +import static org.apache.iceberg.gcp.GCPProperties.GCS_KMS_KEY_NAME; import static org.apache.iceberg.gcp.GCPProperties.GCS_NO_AUTH; import static org.apache.iceberg.gcp.GCPProperties.GCS_OAUTH2_REFRESH_CREDENTIALS_ENABLED; import static org.apache.iceberg.gcp.GCPProperties.GCS_OAUTH2_REFRESH_CREDENTIALS_ENDPOINT; @@ -50,6 +52,33 @@ public void testOAuthWithNoAuth() { assertThat(gcpProperties.oauth2Token()).isNotPresent(); } + @Test + public void testKmsKeyName() { + GCPProperties gcpProperties = + new GCPProperties( + ImmutableMap.of(GCS_KMS_KEY_NAME, "projects/p/locations/l/keyRings/r/cryptoKeys/k")); + assertThat(gcpProperties.kmsKeyName()) + .get() + .isEqualTo("projects/p/locations/l/keyRings/r/cryptoKeys/k"); + assertThat(gcpProperties.encryptionKey()).isNotPresent(); + + assertThat(new GCPProperties().kmsKeyName()).isNotPresent(); + } + + @Test + public void testKmsKeyNameWithEncryptionKey() { + assertThatIllegalStateException() + .isThrownBy( + () -> + new GCPProperties( + ImmutableMap.of( + GCS_ENCRYPTION_KEY, "csek", GCS_KMS_KEY_NAME, "projects/p/cryptoKeys/k"))) + .withMessage( + String.format( + "Invalid encryption settings: must not configure both %s (CSEK) and %s (CMEK)", + GCS_ENCRYPTION_KEY, GCS_KMS_KEY_NAME)); + } + @Test public void refreshCredentialsEndpointSet() { GCPProperties gcpProperties = diff --git a/gcp/src/test/java/org/apache/iceberg/gcp/gcs/TestGCSOutputStream.java b/gcp/src/test/java/org/apache/iceberg/gcp/gcs/TestGCSOutputStream.java index ebf56ffb7eb4..8444688a0f2c 100644 --- a/gcp/src/test/java/org/apache/iceberg/gcp/gcs/TestGCSOutputStream.java +++ b/gcp/src/test/java/org/apache/iceberg/gcp/gcs/TestGCSOutputStream.java @@ -19,9 +19,14 @@ package org.apache.iceberg.gcp.gcs; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import com.google.cloud.storage.BlobId; +import com.google.cloud.storage.BlobInfo; import com.google.cloud.storage.Storage; +import com.google.cloud.storage.Storage.BlobWriteOption; import com.google.cloud.storage.contrib.nio.testing.LocalStorageHelper; import java.io.IOException; import java.io.UncheckedIOException; @@ -30,7 +35,9 @@ import java.util.stream.Stream; import org.apache.iceberg.gcp.GCPProperties; import org.apache.iceberg.metrics.MetricsContext; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; public class TestGCSOutputStream { private static final String BUCKET = "test-bucket"; @@ -53,6 +60,21 @@ public void testWrite() { }); } + @Test + public void testWriteWithKmsKeyName() { + String kmsKeyName = "projects/p/locations/l/keyRings/r/cryptoKeys/k"; + GCPProperties cmekProperties = + new GCPProperties(ImmutableMap.of(GCPProperties.GCS_KMS_KEY_NAME, kmsKeyName)); + Storage spyStorage = spy(storage); + BlobId blobId = randomBlobId(); + + writeAndVerify(spyStorage, blobId, randomData(1024), true, cmekProperties); + + ArgumentCaptor options = ArgumentCaptor.forClass(BlobWriteOption.class); + verify(spyStorage).writer(any(BlobInfo.class), options.capture()); + assertThat(options.getAllValues()).contains(BlobWriteOption.kmsKeyName(kmsKeyName)); + } + @Test public void testMultipleClose() throws IOException { GCSOutputStream stream = @@ -62,8 +84,13 @@ public void testMultipleClose() throws IOException { } private void writeAndVerify(Storage client, BlobId uri, byte[] data, boolean arrayWrite) { + writeAndVerify(client, uri, data, arrayWrite, properties); + } + + private void writeAndVerify( + Storage client, BlobId uri, byte[] data, boolean arrayWrite, GCPProperties gcpProperties) { try (GCSOutputStream stream = - new GCSOutputStream(client, uri, properties, MetricsContext.nullMetrics())) { + new GCSOutputStream(client, uri, gcpProperties, MetricsContext.nullMetrics())) { if (arrayWrite) { stream.write(data); assertThat(stream.getPos()).isEqualTo(data.length);