This is an automated email from the ASF dual-hosted git repository. rzo1 pushed a commit to branch fix/mask-served-nimbus-configuration in repository https://gitbox.apache.org/repos/asf/storm.git
commit 63e0ae05acf19dfc47589fca968f9de2d651e5b3 Author: Richard Zowalla <[email protected]> AuthorDate: Tue Sep 1 11:32:57 2026 +0200 Apply one redaction rule to configuration served by the Nimbus API Nimbus hands configuration to clients from more than one operation, and the treatment differed between them. getNimbusConf and getTopologyConf redact credential values before returning them. getTopologyPageInfo did not, and it serves the daemon configuration merged with the topology configuration, so it returned strictly more than either of the other two. Route both topology-facing operations through a single helper so the rule is the same wherever Nimbus serializes a configuration map for a client: - getTopologyPageInfo masks before serializing into TopologyPageInfo.topology_conf, which the UI renders as the "configuration" field of the topology page and of the metrics response. - getTopologyConf moves from maskPasswords to maskCredentials. That additionally covers keys no annotated field declares, since plugins read their own keys straight out of the map, so an operator-supplied password key in a topology configuration is covered too. - The Blowfish tuple serializer key is handled inside the helper rather than repeated at each call site. Its constant lives outside the config classes, so the annotation scan cannot see it, and its name matches no credential pattern. Masking applies to the copy that gets serialized. The daemon's own configuration and the stored topology configuration are left alone, and workers continue to read real values from the blobstore-distributed configuration rather than from these responses. The only consumers of TopologyPageInfo.topology_conf are the two UI display paths in UIHelpers, so no functional reader loses a value it depends on. --- .../org/apache/storm/daemon/nimbus/Nimbus.java | 29 +++- .../nimbus/NimbusGetTopologyPageInfoTest.java | 155 +++++++++++++++++++++ 2 files changed, 178 insertions(+), 6 deletions(-) diff --git a/storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java b/storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java index 1698a5ce8..cc08fe220 100644 --- a/storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java +++ b/storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java @@ -794,6 +794,24 @@ public class Nimbus implements Iface, Shutdownable, DaemonCommon { return tc.readTopology(topoId, NIMBUS_SUBJECT); } + /** + * Mask the credential values in a config map that is about to be serialized to a client. Nimbus serves config over + * several read-only operations, and a caller authorized for those is not necessarily authorized to hold the cluster's + * or the topology's secrets. This covers what {@link ConfigUtils#maskCredentials(Map)} covers, plus the Blowfish + * tuple-serializer key, whose constant lives outside the config classes so the annotation scan cannot see it and + * whose name matches no credential pattern. + * + * @param conf the config about to be served + * @return a copy of the config with credential values replaced + */ + private static Map<String, Object> maskCredentialsForApi(Map<String, Object> conf) { + Map<String, Object> masked = new HashMap<>(ConfigUtils.maskCredentials(conf)); + if (masked.get(BlowfishTupleSerializer.SECRET_KEY) instanceof String) { + masked.put(BlowfishTupleSerializer.SECRET_KEY, "*****"); + } + return masked; + } + /** * convert {topology-id -> SchedulerAssignment} to {topology-id -> {executor [node port]}}. * @@ -4732,7 +4750,10 @@ public class Nimbus implements Iface, Shutdownable, DaemonCommon { topoPageInfo.set_name(topoName); topoPageInfo.set_status(extractStatusStr(base)); topoPageInfo.set_uptime_secs(Time.deltaSecs(launchTimeSecs)); - topoPageInfo.set_topology_conf(JSONValue.toJSONString(topoConf)); + // topoConf is the daemon conf merged with the topology conf, so it carries Nimbus secrets + // (the ZooKeeper digest payload, Thrift/Netty TLS store passwords) on top of the topology's own. + // getTopologyPageInfo is a topology read-only operation, so mask before it leaves Nimbus. + topoPageInfo.set_topology_conf(JSONValue.toJSONString(maskCredentialsForApi(topoConf))); topoPageInfo.set_replication_count(getBlobReplicationCount(ConfigUtils.masterStormCodeKey(topoId))); if (base.is_set_component_debug()) { DebugOptions debug = base.get_component_debug().get(topoId); @@ -5046,11 +5067,7 @@ public class Nimbus implements Iface, Shutdownable, DaemonCommon { Map<String, Object> checkConf = Utils.merge(conf, topoConf); String topoName = (String) checkConf.get(Config.TOPOLOGY_NAME); checkAuthorization(topoName, checkConf, "getTopologyConf"); - Map<String, Object> maskedConf = new HashMap<>(ConfigUtils.maskPasswords(topoConf)); - if (maskedConf.get(BlowfishTupleSerializer.SECRET_KEY) instanceof String) { - maskedConf.put(BlowfishTupleSerializer.SECRET_KEY, "*****"); - } - return JSONValue.toJSONString(maskedConf); + return JSONValue.toJSONString(maskCredentialsForApi(topoConf)); } catch (Exception e) { LOG.warn("Get topo conf exception. (topology id='{}')", id, e); if (e instanceof TException) { diff --git a/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusGetTopologyPageInfoTest.java b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusGetTopologyPageInfoTest.java new file mode 100644 index 000000000..c9bf4d2c5 --- /dev/null +++ b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusGetTopologyPageInfoTest.java @@ -0,0 +1,155 @@ +/* + * 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.storm.daemon.nimbus; + +import java.util.HashMap; +import java.util.Map; + +import net.minidev.json.JSONValue; + +import org.apache.storm.Config; +import org.apache.storm.DaemonConfig; +import org.apache.storm.LocalCluster; +import org.apache.storm.blobstore.BlobStore; +import org.apache.storm.cluster.IStormClusterState; +import org.apache.storm.generated.StormBase; +import org.apache.storm.generated.StormTopology; +import org.apache.storm.generated.TopologyPageInfo; +import org.apache.storm.generated.TopologyStatus; +import org.apache.storm.scheduler.resource.TestUtilsForResourceAwareScheduler.TestSpout; +import org.apache.storm.security.serialization.BlowfishTupleSerializer; +import org.apache.storm.topology.TopologyBuilder; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentMatchers; +import org.mockito.Mockito; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * getTopologyPageInfo serves the daemon configuration merged with the topology configuration, and it is a + * topology read-only operation, so a principal that is only allowed to look at a topology reaches it. The + * merged map must therefore not carry credential values off the daemon. + */ +public class NimbusGetTopologyPageInfoTest { + + private static final String MASKED = "*****"; + private static final String TOPO_NAME = "test-get-topology-page-info-masking"; + private static final String TOPO_ID = "fake-id"; + + private static final String PLUGIN_SECRET_KEY = "some.plugin.password"; + private static final String NIMBUS_KEYSTORE_PASSWORD_KEY = "nimbus.thrift.tls.keystore.password"; + + @SuppressWarnings("unchecked") + private static Map<String, Object> parse(String json) { + return (Map<String, Object>) JSONValue.parse(json); + } + + private static StormTopology userTopology() { + TopologyBuilder builder = new TopologyBuilder(); + // setNumTasks so the component conf carries topology.tasks, which submit-time normalization + // would otherwise have filled in before the conf reached the cache this test mocks + builder.setSpout("spout-1", new TestSpout(), 1).setNumTasks(1); + return builder.createTopology(); + } + + private static Map<String, Object> storedTopoConf() { + Map<String, Object> topoConf = new HashMap<>(); + topoConf.put(Config.TOPOLOGY_NAME, TOPO_NAME); + topoConf.put(Config.TOPOLOGY_WORKERS, 1); + topoConf.put(Config.TOPOLOGY_ACKER_EXECUTORS, 0); + topoConf.put(Config.TOPOLOGY_EVENTLOGGER_EXECUTORS, 0); + topoConf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 30); + topoConf.put(Config.STORM_ZOOKEEPER_TOPOLOGY_AUTH_PAYLOAD, "topology-zk-secret"); + topoConf.put(BlowfishTupleSerializer.SECRET_KEY, "0123456789abcdef"); + return topoConf; + } + + private static StormBase stormBase() { + StormBase base = new StormBase(); + base.set_name(TOPO_NAME); + base.set_owner("some-owner"); + base.set_status(TopologyStatus.ACTIVE); + base.set_num_workers(1); + base.set_launch_time_secs(1); + return base; + } + + @Test + public void getTopologyPageInfoMasksDaemonAndTopologyCredentials() throws Exception { + IStormClusterState clusterState = Mockito.mock(IStormClusterState.class); + BlobStore blobStore = Mockito.mock(BlobStore.class); + TopoCache topoCache = Mockito.mock(TopoCache.class); + + Map<String, Object> storedConf = storedTopoConf(); + + Map<String, Object> daemonConf = new HashMap<>(); + daemonConf.put(DaemonConfig.NIMBUS_AUTHORIZER, "org.apache.storm.security.auth.authorizer.NoopAuthorizer"); + daemonConf.put(DaemonConfig.SUPERVISOR_AUTHORIZER, "org.apache.storm.security.auth.authorizer.NoopAuthorizer"); + // the daemon-side values that the merge pulls in on top of the topology's own conf + daemonConf.put(Config.STORM_ZOOKEEPER_AUTH_PAYLOAD, "cluster-zk-digest-secret"); + daemonConf.put(NIMBUS_KEYSTORE_PASSWORD_KEY, "keystore-secret"); + daemonConf.put(PLUGIN_SECRET_KEY, "plugin-secret"); + + try (LocalCluster cluster = new LocalCluster.Builder() + .withClusterState(clusterState) + .withBlobStore(blobStore) + .withTopoCache(topoCache) + .withDaemonConf(daemonConf) + .build()) { + Nimbus nimbus = cluster.getNimbus(); + + Mockito.when(topoCache.readTopoConf(Mockito.any(String.class), ArgumentMatchers.any())) + .thenReturn(storedConf); + Mockito.when(topoCache.readTopology(Mockito.any(String.class), ArgumentMatchers.any())) + .thenReturn(userTopology()); + Mockito.when(clusterState.stormBase(Mockito.eq(TOPO_ID), ArgumentMatchers.any())) + .thenReturn(stormBase()); + Mockito.when(clusterState.assignmentInfo(Mockito.eq(TOPO_ID), ArgumentMatchers.any())) + .thenReturn(null); + + TopologyPageInfo pageInfo = nimbus.getTopologyPageInfo(TOPO_ID, ":all-time", false); + Map<String, Object> served = parse(pageInfo.get_topology_conf()); + + // daemon-side credentials, which only this operation merges in + assertEquals(MASKED, served.get(Config.STORM_ZOOKEEPER_AUTH_PAYLOAD), + "the cluster ZooKeeper auth payload should be masked"); + assertEquals(MASKED, served.get(NIMBUS_KEYSTORE_PASSWORD_KEY), + "TLS keystore passwords should be masked"); + assertEquals(MASKED, served.get(PLUGIN_SECRET_KEY), + "a plugin key whose name denotes a secret should be masked"); + + // topology-side credentials, masked on getTopologyConf and equally reachable here + assertEquals(MASKED, served.get(Config.STORM_ZOOKEEPER_TOPOLOGY_AUTH_PAYLOAD), + "the topology ZooKeeper auth payload should be masked"); + assertEquals(MASKED, served.get(BlowfishTupleSerializer.SECRET_KEY), + "the tuple serializer key should be masked"); + + // values that carry no credential are served untouched + assertEquals(TOPO_NAME, served.get(Config.TOPOLOGY_NAME)); + assertEquals(1, ((Number) served.get(Config.TOPOLOGY_WORKERS)).intValue()); + assertEquals(30, ((Number) served.get(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS)).intValue()); + + // masking is applied to the served copy, never to the daemon's own configuration + assertEquals("cluster-zk-digest-secret", nimbus.getConf().get(Config.STORM_ZOOKEEPER_AUTH_PAYLOAD), + "the daemon conf should keep its own values"); + assertEquals("topology-zk-secret", storedConf.get(Config.STORM_ZOOKEEPER_TOPOLOGY_AUTH_PAYLOAD), + "the stored topology conf should keep its own values"); + } + } +}
