This is an automated email from the ASF dual-hosted git repository.
shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 32264a8c2ca [Java] Support dynamic secret provider registration via
SecretRegistrar (#39940)
32264a8c2ca is described below
commit 32264a8c2ca6476868fe2fc752dd5058306cc68b
Author: Shunping Huang <[email protected]>
AuthorDate: Fri Sep 4 05:16:01 2026 -0400
[Java] Support dynamic secret provider registration via SecretRegistrar
(#39940)
* [Java] Support dynamic secret provider registration via SecretRegistrar
Follow the FileSystems registration pattern by introducing SecretRegistrar
SPI
and auto-service discovery in Secret.java. This eliminates hardcoded secret
provider
logic in Secret.java and allows modular extension for new secret managers.
* Add fail-safe handling when loading SecretRegistrar
* Make ReflectHelpers.loadServicesOrdered fail-safe against malformed
providers
* Retain canonical PascalCase names for secret options in error messages
* Remove trivial test case.
---
.../sdk/util/GcpHsmGeneratedSecretRegistrar.java | 34 ++++
.../apache/beam/sdk/util/GcpSecretRegistrar.java | 34 ++++
.../main/java/org/apache/beam/sdk/util/Secret.java | 178 ++++++++++++++++-----
.../org/apache/beam/sdk/util/SecretRegistrar.java | 52 ++++++
.../beam/sdk/util/common/ReflectHelpers.java | 33 +++-
.../beam/sdk/util/GcpSecretRegistrarTest.java | 62 +++++++
.../java/org/apache/beam/sdk/util/SecretTest.java | 89 +++++++++++
.../beam/sdk/util/common/ReflectHelpersTest.java | 28 ++++
8 files changed, 467 insertions(+), 43 deletions(-)
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpHsmGeneratedSecretRegistrar.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpHsmGeneratedSecretRegistrar.java
new file mode 100644
index 00000000000..232fe7dfa83
--- /dev/null
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpHsmGeneratedSecretRegistrar.java
@@ -0,0 +1,34 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.util;
+
+import com.google.auto.service.AutoService;
+import java.util.Map;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+
+/** {@link AutoService} registrar for the {@link GcpHsmGeneratedSecret}. */
+@AutoService(SecretRegistrar.class)
+public class GcpHsmGeneratedSecretRegistrar implements SecretRegistrar {
+
+ @Override
+ public Map<String, SecretFactory> getSecretFactories() {
+ return ImmutableMap.of(
+ "GoogleCloudHsmGeneratedSecretManager", GcpHsmGeneratedSecret::fromMap,
+ "GcpHsmGeneratedSecret", GcpHsmGeneratedSecret::fromMap);
+ }
+}
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpSecretRegistrar.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpSecretRegistrar.java
new file mode 100644
index 00000000000..61b31332e6d
--- /dev/null
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcpSecretRegistrar.java
@@ -0,0 +1,34 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.util;
+
+import com.google.auto.service.AutoService;
+import java.util.Map;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+
+/** {@link AutoService} registrar for the {@link GcpSecret}. */
+@AutoService(SecretRegistrar.class)
+public class GcpSecretRegistrar implements SecretRegistrar {
+
+ @Override
+ public Map<String, SecretFactory> getSecretFactories() {
+ return ImmutableMap.of(
+ "GoogleCloudSecretManager", GcpSecret::fromMap,
+ "GcpSecret", GcpSecret::fromMap);
+ }
+}
diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/Secret.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/Secret.java
index f5e935460c8..5d36a160259 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/Secret.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/Secret.java
@@ -21,8 +21,14 @@ import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.Serializable;
import java.nio.charset.StandardCharsets;
+import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+import org.apache.beam.sdk.util.common.ReflectHelpers;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -35,6 +41,119 @@ import org.slf4j.LoggerFactory;
* should be able to return a valid byte array representing the secret.
*/
public abstract class Secret implements Serializable {
+ private static final Logger LOG = LoggerFactory.getLogger(Secret.class);
+
+ private static final Set<String> SUPPORTED_TYPES;
+ private static final Map<String, SecretRegistrar.SecretFactory>
SECRET_FACTORIES;
+
+ static {
+ TreeSet<String> supportedTypes = new TreeSet<>();
+ Map<String, SecretRegistrar.SecretFactory> factories;
+ try {
+ factories =
+ loadSecretFactories(
+ ReflectHelpers.loadServicesOrdered(SecretRegistrar.class),
supportedTypes);
+ } catch (Throwable t) {
+ // Top-level fail-safe: guarantee that static class initialization of
Secret never fails
+ // due to unforeseen classloader or registrar errors.
+ LOG.error("Unexpected error loading SecretRegistrars; secret factories
may be incomplete", t);
+ factories = Collections.emptyMap();
+ }
+ SECRET_FACTORIES = factories;
+ SUPPORTED_TYPES = Collections.unmodifiableSet(supportedTypes);
+ }
+
+ /**
+ * Loads factories from the provided registrars into an immutable map.
+ *
+ * <p>Applies defensive checks:
+ *
+ * <ul>
+ * <li>Sandboxes each registrar with a per-registrar try-catch so a rogue
or broken registrar
+ * cannot crash discovery.
+ * <li>Guards against {@code null} return values from {@link
+ * SecretRegistrar#getSecretFactories()}, {@code null} map entries,
{@code null} or empty
+ * keys, and {@code null} factory values.
+ * <li>Applies a "first-wins with warning" strategy on duplicate keys to
prevent classpath leaks
+ * (such as duplicate test registrars) from throwing exceptions and
breaking pipelines.
+ * </ul>
+ */
+ @VisibleForTesting
+ static Map<String, SecretRegistrar.SecretFactory> loadSecretFactories(
+ @Nullable Iterable<SecretRegistrar> registrars) {
+ return loadSecretFactories(registrars, new TreeSet<>());
+ }
+
+ @VisibleForTesting
+ static Map<String, SecretRegistrar.SecretFactory> loadSecretFactories(
+ @Nullable Iterable<SecretRegistrar> registrars, Set<String>
supportedTypes) {
+ Map<String, SecretRegistrar.SecretFactory> factories = new HashMap<>();
+ if (registrars == null) {
+ return Collections.emptyMap();
+ }
+
+ for (SecretRegistrar registrar : registrars) {
+ if (registrar == null) {
+ continue;
+ }
+ try {
+ Map<String, SecretRegistrar.SecretFactory> registrarFactories =
+ registrar.getSecretFactories();
+ if (registrarFactories == null) {
+ LOG.warn(
+ "SecretRegistrar '{}' returned null from getSecretFactories();
ignoring",
+ registrar.getClass().getName());
+ continue;
+ }
+
+ for (Map.Entry<String, SecretRegistrar.SecretFactory> entry :
+ registrarFactories.entrySet()) {
+ if (entry == null) {
+ continue;
+ }
+ String rawKey = entry.getKey();
+ if (rawKey == null || rawKey.trim().isEmpty()) {
+ LOG.warn(
+ "SecretRegistrar '{}' registered a factory with a null or
empty key; ignoring",
+ registrar.getClass().getName());
+ continue;
+ }
+ SecretRegistrar.SecretFactory factory = entry.getValue();
+ if (factory == null) {
+ LOG.warn(
+ "SecretRegistrar '{}' registered a null SecretFactory for key
'{}'; ignoring",
+ registrar.getClass().getName(),
+ rawKey);
+ continue;
+ }
+
+ String canonicalKey = rawKey.trim();
+ String key = canonicalKey.toLowerCase();
+ SecretRegistrar.SecretFactory existing = factories.get(key);
+ if (existing != null) {
+ // First-wins strategy with warning: do not throw to prevent
leaked test or duplicate
+ // registrars on the classpath from crashing pipeline execution.
+ LOG.warn(
+ "Duplicate SecretFactory for secret manager name '{}': already
registered by '{}', "
+ + "ignoring duplicate from '{}'",
+ key,
+ existing.getClass().getName(),
+ factory.getClass().getName());
+ } else {
+ factories.put(key, factory);
+ supportedTypes.add(canonicalKey);
+ }
+ }
+ } catch (Throwable t) {
+ LOG.warn(
+ "Failed to load secret factories from SecretRegistrar '{}';
skipping",
+ registrar.getClass().getName(),
+ t);
+ }
+ }
+ return ImmutableMap.copyOf(factories);
+ }
+
private transient byte @Nullable [] cachedSecretBytes = null;
/**
@@ -104,29 +223,22 @@ public abstract class Secret implements Serializable {
}
String secretType = rawType.toLowerCase();
- String secretManager;
- switch (secretType) {
- case "gcpsecret":
- secretManager = "GoogleCloudSecretManager";
- break;
- case "gcphsmgeneratedsecret":
- secretManager = "GoogleCloudHsmGeneratedSecretManager";
- break;
- default:
- throw new IllegalArgumentException(
- String.format(
- "Invalid secret type %s, currently only GcpSecret and
GcpHsmGeneratedSecret are supported",
- secretType));
+ SecretRegistrar.SecretFactory factory = SECRET_FACTORIES.get(secretType);
+ if (factory == null) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Invalid secret type %s, currently supported types: %s",
rawType, SUPPORTED_TYPES));
}
try {
- ObjectMapper mapper = new ObjectMapper();
- String jsonSpec = mapper.writeValueAsString(paramMap);
- return fromJson(jsonSpec, secretManager);
+ return factory.createSecret(paramMap);
} catch (Exception e) {
if (e instanceof IllegalArgumentException) {
throw (IllegalArgumentException) e;
}
+ if (e instanceof NullPointerException) {
+ throw (NullPointerException) e;
+ }
throw new RuntimeException("Failed to parse secret option", e);
}
}
@@ -139,7 +251,6 @@ public abstract class Secret implements Serializable {
* @return An instance of Secret.
*/
public static Secret fromJson(@Nullable String spec, @Nullable String
secretManager) {
- Logger logger = LoggerFactory.getLogger(Secret.class);
String smManager = secretManager != null ? secretManager.trim() : null;
if (smManager != null && smManager.isEmpty()) {
smManager = null;
@@ -152,38 +263,23 @@ public abstract class Secret implements Serializable {
mapper.configure(com.fasterxml.jackson.core.JsonParser.Feature.ALLOW_SINGLE_QUOTES,
true);
specMap = mapper.readValue(spec, new TypeReference<Map<String,
String>>() {});
} catch (Exception e) {
- logger.debug("Failed to parse secret spec as JSON map", e);
+ LOG.debug("Failed to parse secret spec as JSON map", e);
}
}
if (smManager != null) {
- switch (smManager.toLowerCase()) {
- case "googlecloudsecretmanager":
- case "gcpsecret":
- if (specMap != null) {
- return GcpSecret.fromMap(specMap);
- } else if (spec != null) {
- return new GcpSecret(spec);
- } else {
- throw new IllegalArgumentException("Invalid spec for GcpSecret");
- }
- case "googlecloudhsmgeneratedsecretmanager":
- case "gcphsmgeneratedsecret":
- if (specMap != null) {
- return GcpHsmGeneratedSecret.fromMap(specMap);
- } else {
- throw new IllegalArgumentException("Invalid spec for
GcpHsmGeneratedSecret");
- }
- default:
- throw new IllegalArgumentException(
- String.format(
- "Unsupported secret manager: '%s'. Currently supported
options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'.",
- smManager));
+ SecretRegistrar.SecretFactory factory =
SECRET_FACTORIES.get(smManager.toLowerCase());
+ if (factory != null) {
+ return factory.createSecret(specMap != null ? specMap :
Collections.emptyMap());
}
+ throw new IllegalArgumentException(
+ String.format(
+ "Unsupported secret manager: '%s'. Currently supported options:
%s.",
+ smManager, SUPPORTED_TYPES));
}
if (specMap != null) {
- logger.warn(
+ LOG.warn(
"The 'spec' parameter appears to be a JSON specification, but
'secret_manager' is not set. Defaulting to Raw.");
}
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/SecretRegistrar.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/SecretRegistrar.java
new file mode 100644
index 00000000000..2ba120bee7d
--- /dev/null
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/SecretRegistrar.java
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.util;
+
+import com.google.auto.service.AutoService;
+import java.util.Map;
+import java.util.ServiceLoader;
+
+/**
+ * A registrar that creates {@link Secret} instances from a spec parameter map.
+ *
+ * <p>{@link Secret} creators have the ability to provide a registrar by
creating a {@link
+ * ServiceLoader} entry and a concrete implementation of this interface.
+ *
+ * <p>It is optional but recommended to use one of the many build time tools
such as {@link
+ * AutoService} to generate the necessary META-INF files automatically.
+ */
+public interface SecretRegistrar {
+
+ /** Functional interface for creating a {@link Secret} from a specification
map. */
+ @FunctionalInterface
+ interface SecretFactory {
+ /**
+ * Creates a {@link Secret} instance from a spec parameter map.
+ *
+ * @param specMap The parsed map of key-value parameters.
+ * @return The constructed {@link Secret} instance.
+ */
+ Secret createSecret(Map<String, String> specMap);
+ }
+
+ /**
+ * Returns a map from secret provider name / type (case-insensitive) to the
corresponding {@link
+ * SecretFactory}.
+ */
+ Map<String, SecretFactory> getSecretFactories();
+}
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java
index 7d5964cb83c..2ef6920654e 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java
@@ -34,8 +34,10 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Comparator;
+import java.util.Iterator;
import java.util.LinkedHashSet;
import java.util.Queue;
+import java.util.ServiceConfigurationError;
import java.util.ServiceLoader;
import org.apache.beam.sdk.values.TypeDescriptor;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Function;
@@ -45,10 +47,13 @@ import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Immuta
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSortedSet;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Queues;
import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
/** Utilities for working with with {@link Class Classes} and {@link Method
Methods}. */
@SuppressWarnings({"nullness", "keyfor"}) //
TODO(https://github.com/apache/beam/issues/20497)
public class ReflectHelpers {
+ private static final Logger LOG =
LoggerFactory.getLogger(ReflectHelpers.class);
private static final Joiner COMMA_SEPARATOR = Joiner.on(", ");
@@ -206,16 +211,40 @@ public class ReflectHelpers {
* Returns instances of all implementations of the specified {@code iface}.
Instances are sorted
* by their class' name to ensure deterministic execution.
*
+ * <p>Safely handles malformed service providers: if a provider fails to
load (e.g. throwing
+ * {@link ServiceConfigurationError}, {@link LinkageError}, or other
exceptions), it will be
+ * logged as a warning and skipped so that other valid implementations
continue to load.
+ *
* @param iface The interface to load implementations of
* @param classLoader The class loader to use
* @param <T> The type of {@code iface}
* @return An iterable of instances of T, ordered by their class' canonical
name
*/
public static <T> Iterable<T> loadServicesOrdered(Class<T> iface,
ClassLoader classLoader) {
- ServiceLoader<T> loader = ServiceLoader.load(iface, classLoader);
ImmutableSortedSet.Builder<T> builder =
new ImmutableSortedSet.Builder<>(ObjectsClassComparator.INSTANCE);
- builder.addAll(loader);
+ try {
+ ServiceLoader<T> loader = ServiceLoader.load(iface, classLoader);
+ Iterator<T> iterator = loader.iterator();
+ while (true) {
+ T service;
+ try {
+ if (!iterator.hasNext()) {
+ break;
+ }
+ service = iterator.next();
+ } catch (ServiceConfigurationError | LinkageError | Exception e) {
+ // A single broken provider on the classpath shouldn't abort
discovery of valid ones.
+ LOG.warn("Failed to load a service implementation of {}; skipping",
iface.getName(), e);
+ continue;
+ }
+ if (service != null) {
+ builder.add(service);
+ }
+ }
+ } catch (Throwable t) {
+ LOG.warn("Failed to discover services for {}", iface.getName(), t);
+ }
return builder.build();
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/GcpSecretRegistrarTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/GcpSecretRegistrarTest.java
new file mode 100644
index 00000000000..e95483dd9bf
--- /dev/null
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/GcpSecretRegistrarTest.java
@@ -0,0 +1,62 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.util;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.hasItems;
+import static org.junit.Assert.fail;
+
+import java.util.Map;
+import java.util.ServiceLoader;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link GcpSecretRegistrar} and {@link
GcpHsmGeneratedSecretRegistrar}. */
+@RunWith(JUnit4.class)
+public class GcpSecretRegistrarTest {
+
+ @Test
+ public void testGcpSecretRegistrarServiceLoader() {
+ for (SecretRegistrar registrar :
+
Lists.newArrayList(ServiceLoader.load(SecretRegistrar.class).iterator())) {
+ if (registrar instanceof GcpSecretRegistrar) {
+ Map<String, SecretRegistrar.SecretFactory> factories =
registrar.getSecretFactories();
+ assertThat(factories.keySet(), hasItems("GoogleCloudSecretManager",
"GcpSecret"));
+ return;
+ }
+ }
+ fail("Expected to find " + GcpSecretRegistrar.class);
+ }
+
+ @Test
+ public void testGcpHsmGeneratedSecretRegistrarServiceLoader() {
+ for (SecretRegistrar registrar :
+
Lists.newArrayList(ServiceLoader.load(SecretRegistrar.class).iterator())) {
+ if (registrar instanceof GcpHsmGeneratedSecretRegistrar) {
+ Map<String, SecretRegistrar.SecretFactory> factories =
registrar.getSecretFactories();
+ assertThat(
+ factories.keySet(),
+ hasItems("GoogleCloudHsmGeneratedSecretManager",
"GcpHsmGeneratedSecret"));
+ return;
+ }
+ }
+ fail("Expected to find " + GcpHsmGeneratedSecretRegistrar.class);
+ }
+}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/SecretTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/SecretTest.java
index 9b74e52376f..446688035a5 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/SecretTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/SecretTest.java
@@ -28,6 +28,9 @@ import java.nio.charset.StandardCharsets;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+import org.apache.beam.sdk.util.common.ReflectHelpers;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.junit.runners.JUnit4;
@@ -83,6 +86,8 @@ public class SecretTest {
Exception exception =
assertThrows(IllegalArgumentException.class, () ->
Secret.parseSecretOption(secretOption));
assertTrue(exception.getMessage().contains("Invalid secret type
unsupported"));
+ assertTrue(exception.getMessage().contains("GcpSecret"));
+ assertTrue(exception.getMessage().contains("GoogleCloudSecretManager"));
}
@Test
@@ -147,6 +152,14 @@ public class SecretTest {
IllegalArgumentException.class, () -> Secret.fromJson("spec",
"unsupported_provider"));
assertTrue(
exception.getMessage().contains("Unsupported secret manager:
'unsupported_provider'"));
+ assertTrue(exception.getMessage().contains("GoogleCloudSecretManager"));
+ assertTrue(exception.getMessage().contains("GcpSecret"));
+
+ // Case-insensitive secret manager lookup in fromJson
+ Secret secretGcpLower = Secret.fromJson(spec, "googlecloudsecretmanager");
+ assertTrue(secretGcpLower instanceof GcpSecret);
+ Secret secretShortLower = Secret.fromJson(spec, "gcpsecret");
+ assertTrue(secretShortLower instanceof GcpSecret);
}
@Test
@@ -245,4 +258,80 @@ public class SecretTest {
GcpHsmGeneratedSecret deserializedHsm = SerializableUtils.clone(hsm);
assertEquals(hsm, deserializedHsm);
}
+
+ @Test
+ public void testLoadSecretFactoriesNullList() {
+ Map<String, SecretRegistrar.SecretFactory> factories =
Secret.loadSecretFactories(null);
+ assertTrue(factories.isEmpty());
+ }
+
+ @Test
+ public void testLoadSecretFactoriesHandlesNullRegistrarAndNullFactories() {
+ SecretRegistrar nullFactoriesRegistrar = () -> null;
+ Map<String, SecretRegistrar.SecretFactory> factories =
+ Secret.loadSecretFactories(java.util.Arrays.asList(null,
nullFactoriesRegistrar));
+ assertTrue(factories.isEmpty());
+ }
+
+ @Test
+ public void testLoadSecretFactoriesHandlesThrowingRegistrar() {
+ SecretRegistrar throwingRegistrar =
+ () -> {
+ throw new RuntimeException("Simulated failure in registrar");
+ };
+ SecretRegistrar validRegistrar =
+ () -> Collections.singletonMap("valid", spec -> new RawSecret("test"));
+
+ Map<String, SecretRegistrar.SecretFactory> factories =
+ Secret.loadSecretFactories(java.util.Arrays.asList(throwingRegistrar,
validRegistrar));
+ assertEquals(1, factories.size());
+ assertTrue(factories.containsKey("valid"));
+ }
+
+ @Test
+ public void testLoadSecretFactoriesHandlesMalformedEntries() {
+ Map<String, SecretRegistrar.SecretFactory> malformedMap = new HashMap<>();
+ malformedMap.put(null, spec -> new RawSecret("val"));
+ malformedMap.put("", spec -> new RawSecret("val"));
+ malformedMap.put(" ", spec -> new RawSecret("val"));
+ malformedMap.put("null_factory", null);
+ malformedMap.put("good", spec -> new RawSecret("good_val"));
+
+ SecretRegistrar registrar = () -> malformedMap;
+ Map<String, SecretRegistrar.SecretFactory> factories =
+ Secret.loadSecretFactories(Collections.singletonList(registrar));
+ assertEquals(1, factories.size());
+ assertTrue(factories.containsKey("good"));
+ }
+
+ @Test
+ public void testLoadSecretFactoriesDuplicateKeysFirstWins() {
+ SecretRegistrar.SecretFactory factory1 = spec -> new RawSecret("first");
+ SecretRegistrar.SecretFactory factory2 = spec -> new RawSecret("second");
+
+ SecretRegistrar registrar1 = () ->
Collections.singletonMap("duplicate_key", factory1);
+ SecretRegistrar registrar2 = () ->
Collections.singletonMap("DUPLICATE_KEY", factory2);
+
+ Set<String> supportedTypes = new TreeSet<>();
+ Map<String, SecretRegistrar.SecretFactory> factories =
+ Secret.loadSecretFactories(java.util.Arrays.asList(registrar1,
registrar2), supportedTypes);
+ assertEquals(1, factories.size());
+ assertEquals(factory1, factories.get("duplicate_key"));
+ assertEquals(Collections.singleton("duplicate_key"), supportedTypes);
+ }
+
+ @Test
+ public void testLoadServicesOrderedDiscoversSecretRegistrars() {
+ Iterable<SecretRegistrar> registrars =
+ ReflectHelpers.loadServicesOrdered(SecretRegistrar.class);
+ org.junit.Assert.assertNotNull(registrars);
+ boolean foundGcp = false;
+ for (SecretRegistrar registrar : registrars) {
+ if (registrar instanceof GcpSecretRegistrar) {
+ foundGcp = true;
+ break;
+ }
+ }
+ assertTrue("Expected GcpSecretRegistrar to be discovered", foundGcp);
+ }
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/common/ReflectHelpersTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/common/ReflectHelpersTest.java
index e999169abb7..7ce86761ea4 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/common/ReflectHelpersTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/common/ReflectHelpersTest.java
@@ -24,19 +24,27 @@ import static org.hamcrest.Matchers.equalTo;
import static org.junit.Assert.assertEquals;
import com.fasterxml.jackson.annotation.JsonIgnore;
+import java.io.File;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import org.apache.beam.sdk.options.Default;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.values.TypeDescriptor;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.Files;
+import org.junit.Rule;
import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
import org.junit.runners.JUnit4;
/** Tests for {@link ReflectHelpers}. */
@RunWith(JUnit4.class)
public class ReflectHelpersTest {
+ @Rule public TemporaryFolder tmp = new TemporaryFolder();
@Test
public void testMethodFormatter() throws Exception {
@@ -212,4 +220,24 @@ public class ReflectHelpersTest {
assertThat(names, contains("Alpha", "Zeta"));
}
+
+ @Test
+ public void testLoadServicesOrderedHandlesFailingProvider() throws Exception
{
+ File servicesDir = tmp.newFolder("META-INF", "services");
+ File serviceFile = new File(servicesDir, FakeService.class.getName());
+ Files.asCharSink(serviceFile, StandardCharsets.UTF_8)
+ .write("non.existent.Class\n" + AlphaImpl.class.getName() + "\n");
+
+ URLClassLoader classLoader =
+ new URLClassLoader(
+ new URL[] {tmp.getRoot().toURI().toURL()},
ReflectHelpers.findClassLoader());
+ List<String> names = new ArrayList<>();
+ for (FakeService service :
ReflectHelpers.loadServicesOrdered(FakeService.class, classLoader)) {
+ names.add(service.getName());
+ }
+
+ // "non.existent.Class" should be skipped gracefully, and AlphaImpl and
ZetaImpl should be
+ // loaded.
+ assertThat(names, contains("Alpha", "Zeta"));
+ }
}