This is an automated email from the ASF dual-hosted git repository.

rzo1 pushed a commit to branch fix/supervisor-page-readable-workers
in repository https://gitbox.apache.org/repos/asf/storm.git

commit ae78ca88625cb1c5b1072920c1469de9e82e015a
Author: Richard Zowalla <[email protected]>
AuthorDate: Fri Sep 18 19:12:30 2026 +0200

    Only list workers of readable topologies on the supervisor page
    
    getSupervisorPageInfo returned a worker summary for every topology
    assigned to the supervisor and only left out the per-component task
    counts for topologies the caller cannot read. Worker placement is
    otherwise only returned by getTopologyInfo and getTopologyPageInfo,
    which require read access to the topology, so the supervisor page is
    now consistent with them and omits the workers of topologies the
    caller is not allowed to read.
    
    The supervisor summary on the same page still reports the aggregate
    slot, memory and CPU usage of the node, so the capacity view is
    unchanged. Admins and users with read access see the same page as
    before.
---
 docs/STORM-UI-REST-API.md                          |  2 +-
 .../org/apache/storm/daemon/nimbus/Nimbus.java     |  9 +++-
 .../storm/daemon/nimbus/NimbusClojurePortTest.java | 59 +++++++++++++++++++---
 3 files changed, 60 insertions(+), 10 deletions(-)

diff --git a/docs/STORM-UI-REST-API.md b/docs/STORM-UI-REST-API.md
index 4dd6c501c..48dbd5789 100644
--- a/docs/STORM-UI-REST-API.md
+++ b/docs/STORM-UI-REST-API.md
@@ -247,7 +247,7 @@ Response fields:
 |Field  |Value|Description|
 |---   |---    |---
 |supervisors| Array| Array of supervisor summaries|
-|workers| Array| Array of worker summaries |
+|workers| Array| Array of worker summaries, limited to topologies the caller 
is allowed to read |
 |schedulerDisplayResource| Boolean | Whether to display scheduler resource 
information|
 
 Each supervisor is defined by:
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 75ef22c6a..9e184feb6 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
@@ -4990,6 +4990,12 @@ public class Nimbus implements Iface, Shutdownable, 
DaemonCommon {
                 List<String> superTopologies = 
topologiesOnSupervisor(topoToAssignment, sid);
                 Set<String> userTopologies = filterAuthorized("getTopology", 
superTopologies);
                 for (String topoId : superTopologies) {
+                    // Placement of a topology's workers is only visible to 
callers who may read that topology,
+                    // matching getTopologyInfo/getTopologyPageInfo. The 
supervisor summary above still reports the
+                    // aggregate slot and resource usage of the node.
+                    if (!userTopologies.contains(topoId)) {
+                        continue;
+                    }
                     CommonTopoInfo common = getCommonTopoInfo(topoId, 
"getSupervisorPageInfo");
                     String topoName = common.topoName;
                     Assignment assignment = common.assignment;
@@ -5009,11 +5015,10 @@ public class Nimbus implements Iface, Shutdownable, 
DaemonCommon {
                         nodeToHost = Collections.emptyMap();
                     }
                     Map<WorkerSlot, WorkerResources> workerResources = 
getWorkerResourcesForTopology(topoId);
-                    boolean isAllowed = userTopologies.contains(topoId);
                     String owner = (common.base == null) ? null : 
common.base.get_owner();
                     for (WorkerSummary workerSummary : 
StatsUtil.aggWorkerStats(topoId, topoName, taskToComp, beats,
                                                                                
 exec2NodePort, nodeToHost, workerResources, includeSys,
-                                                                               
 isAllowed, sid, owner)) {
+                                                                               
 true, sid, owner)) {
                         pageInfo.add_to_worker_summaries(workerSummary);
                     }
                 }
diff --git 
a/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusClojurePortTest.java
 
b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusClojurePortTest.java
index 0975e52a5..236545eed 100644
--- 
a/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusClojurePortTest.java
+++ 
b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusClojurePortTest.java
@@ -58,9 +58,11 @@ import org.apache.storm.generated.RebalanceOptions;
 import org.apache.storm.generated.StormBase;
 import org.apache.storm.generated.StormTopology;
 import org.apache.storm.generated.SubmitOptions;
+import org.apache.storm.generated.SupervisorPageInfo;
 import org.apache.storm.generated.TopologyInitialStatus;
 import org.apache.storm.generated.TopologyStatus;
 import org.apache.storm.generated.TopologySummary;
+import org.apache.storm.generated.WorkerSummary;
 import org.apache.storm.metric.StormMetricsRegistry;
 import org.apache.storm.nimbus.ILeaderElector;
 import org.apache.storm.nimbus.InMemoryTopologyActionNotifier;
@@ -1538,6 +1540,40 @@ public class NimbusClojurePortTest {
 
     @Test
     public void testCheckAuthorizationGetSupervisorPageInfo() throws Exception 
{
+        String expectedName = "test-nimbus-check-autho-params";
+        runSupervisorPageInfo(expectedName, (nimbus, pageInfo) -> {
+            
Mockito.verify(nimbus).checkAuthorization(Mockito.eq(expectedName), 
Mockito.any(Map.class), Mockito.eq("getSupervisorPageInfo"));
+            Mockito.verify(nimbus).checkAuthorization(Mockito.isNull(), 
Mockito.isNull(), Mockito.eq("getClusterInfo"));
+            
Mockito.verify(nimbus).checkAuthorization(Mockito.eq(expectedName), 
Mockito.any(Map.class), Mockito.eq("getTopology"));
+
+            assertEquals(1, pageInfo.get_supervisor_summaries_size());
+            assertEquals(1, pageInfo.get_worker_summaries_size());
+            WorkerSummary worker = pageInfo.get_worker_summaries().get(0);
+            assertEquals(expectedName, worker.get_topology_id());
+            assertEquals("super1", worker.get_supervisor_id());
+        });
+    }
+
+    @Test
+    public void testGetSupervisorPageInfoOmitsWorkersOfUnreadableTopologies() 
throws Exception {
+        String expectedName = "test-nimbus-unreadable-topo";
+        runSupervisorPageInfo(expectedName, nimbus ->
+                
Mockito.doReturn(Set.of()).when(nimbus).filterAuthorized(Mockito.eq("getTopology"),
 Mockito.any()),
+            (nimbus, pageInfo) -> {
+                Mockito.verify(nimbus, Mockito.never())
+                    .checkAuthorization(Mockito.eq(expectedName), 
Mockito.any(Map.class), Mockito.eq("getSupervisorPageInfo"));
+
+                assertEquals(1, pageInfo.get_supervisor_summaries_size());
+                assertEquals(0, pageInfo.get_worker_summaries_size());
+            });
+    }
+
+    private void runSupervisorPageInfo(String topoId, PageInfoVerifier 
verifier) throws Exception {
+        runSupervisorPageInfo(topoId, nimbus -> { }, verifier);
+    }
+
+    private void runSupervisorPageInfo(String topoId, ThrowingConsumer<Nimbus> 
stubber,
+                                       PageInfoVerifier verifier) throws 
Exception {
         IStormClusterState clusterState = 
Mockito.mock(IStormClusterState.class);
         BlobStore blobStore = Mockito.mock(BlobStore.class);
         TopoCache tc = Mockito.mock(TopoCache.class);
@@ -1552,10 +1588,9 @@ public class NimbusClojurePortTest {
                     DaemonConfig.SUPERVISOR_AUTHORIZER, 
"org.apache.storm.security.auth.authorizer.NoopAuthorizer"))
                 .build()) {
             Nimbus nimbus = cluster.getNimbus();
-            String expectedName = "test-nimbus-check-autho-params";
 
             Map<String, Object> expectedConf = new HashMap<>();
-            expectedConf.put(Config.TOPOLOGY_NAME, expectedName);
+            expectedConf.put(Config.TOPOLOGY_NAME, topoId);
             expectedConf.put(Config.TOPOLOGY_WORKERS, 1);
             expectedConf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 30);
             expectedConf.put("foo", "bar");
@@ -1571,7 +1606,7 @@ public class NimbusClojurePortTest {
             topology.set_state_spouts(Map.of());
 
             Map<String, Assignment> topoAssignment = new HashMap<>();
-            topoAssignment.put(expectedName, assignment);
+            topoAssignment.put(topoId, assignment);
 
             HashMap<String, org.apache.storm.generated.SupervisorInfo> 
allSupervisors = new HashMap<>();
             org.apache.storm.generated.SupervisorInfo si1 = new 
org.apache.storm.generated.SupervisorInfo();
@@ -1591,18 +1626,28 @@ public class NimbusClojurePortTest {
             allSupervisors.put("super2", si2);
 
             
Mockito.when(clusterState.allSupervisorInfo()).thenReturn(allSupervisors);
+            Mockito.when(clusterState.assignmentInfo(Mockito.eq(topoId), 
Mockito.any())).thenReturn(assignment);
             Mockito.when(tc.readTopoConf(Mockito.any(String.class), 
Mockito.any(Subject.class))).thenReturn(expectedConf);
             Mockito.when(tc.readTopology(Mockito.any(String.class), 
Mockito.any(Subject.class))).thenReturn(topology);
             
Mockito.when(clusterState.assignmentsInfo()).thenReturn(topoAssignment);
+            stubber.accept(nimbus);
 
-            nimbus.getSupervisorPageInfo("super1", null, true);
+            SupervisorPageInfo pageInfo = 
nimbus.getSupervisorPageInfo("super1", null, true);
 
-            
Mockito.verify(nimbus).checkAuthorization(Mockito.eq(expectedName), 
Mockito.any(Map.class), Mockito.eq("getSupervisorPageInfo"));
-            Mockito.verify(nimbus).checkAuthorization(Mockito.isNull(), 
Mockito.isNull(), Mockito.eq("getClusterInfo"));
-            
Mockito.verify(nimbus).checkAuthorization(Mockito.eq(expectedName), 
Mockito.any(Map.class), Mockito.eq("getTopology"));
+            verifier.accept(nimbus, pageInfo);
         }
     }
 
+    @FunctionalInterface
+    private interface ThrowingConsumer<T> {
+        void accept(T t) throws Exception;
+    }
+
+    @FunctionalInterface
+    private interface PageInfoVerifier {
+        void accept(Nimbus nimbus, SupervisorPageInfo pageInfo) throws 
Exception;
+    }
+
     @Test
     public void testKillStorm() throws Exception {
         try (LocalCluster cluster = new LocalCluster.Builder()

Reply via email to