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

ethanli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/storm.git


The following commit(s) were added to refs/heads/master by this push:
     new 5873ead  [STORM-3625] Add topo name validatiion for CLI client before 
connecting to nimbus (#3253)
5873ead is described below

commit 5873eaded557a6fd639c1b70acff684b9c4bd466
Author: Rui Li <[email protected]>
AuthorDate: Tue Apr 28 14:49:21 2020 -0500

    [STORM-3625] Add topo name validatiion for CLI client before connecting to 
nimbus (#3253)
---
 .../jvm/org/apache/storm/command/KillTopology.java | 31 +++++++++++++++++++---
 .../jvm/org/apache/storm/command/Rebalance.java    |  1 +
 .../jvm/org/apache/storm/command/SetLogLevel.java  |  2 ++
 .../apache/storm/command/UploadCredentials.java    |  1 +
 4 files changed, 31 insertions(+), 4 deletions(-)

diff --git a/storm-core/src/jvm/org/apache/storm/command/KillTopology.java 
b/storm-core/src/jvm/org/apache/storm/command/KillTopology.java
index e58c08b..ef5bfff 100644
--- a/storm-core/src/jvm/org/apache/storm/command/KillTopology.java
+++ b/storm-core/src/jvm/org/apache/storm/command/KillTopology.java
@@ -12,16 +12,19 @@
 
 package org.apache.storm.command;
 
+import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import org.apache.storm.generated.KillOptions;
 import org.apache.storm.generated.Nimbus;
 import org.apache.storm.utils.NimbusClient;
+import org.apache.storm.utils.Utils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 public class KillTopology {
     private static final Logger LOG = 
LoggerFactory.getLogger(KillTopology.class);
+    private static int errorCount;
 
     public static void main(String[] args) throws Exception {
         Map<String, Object> cl = CLI.opt("w", "wait", null, CLI.AS_INT)
@@ -32,12 +35,33 @@ public class KillTopology {
         @SuppressWarnings("unchecked")
         final List<String> names = (List<String>) cl.get("TOPO");
 
-        // Wait this many seconds after deactivating topology before killing
-        Integer wait = (Integer) cl.get("w");
-
         // if '-i' is set, we'll try to kill every topology listed, even if an 
error occurs
         Boolean continueOnError = (Boolean) cl.get("i");
 
+        errorCount = 0;
+        Iterator<String> iterator = names.iterator();
+        while (iterator.hasNext()) {
+            String name = iterator.next();
+            try {
+                Utils.validateTopologyName(name);
+            } catch (IllegalArgumentException e) {
+                if (!continueOnError) {
+                    throw e;
+                } else {
+                    iterator.remove();
+                    errorCount += 1;
+                    LOG.error("Format of topology name {} is not valid ", 
name);
+                }
+            }
+        }
+
+        if (names.isEmpty()) {
+            throw new RuntimeException("Failed to successfully kill " + 
errorCount + " topologies.");
+        }
+
+        // Wait this many seconds after deactivating topology before killing
+        Integer wait = (Integer) cl.get("w");
+
         final KillOptions opts = new KillOptions();
         if (wait != null) {
             opts.set_wait_secs(wait);
@@ -46,7 +70,6 @@ public class KillTopology {
         NimbusClient.withConfiguredClient(new NimbusClient.WithNimbus() {
             @Override
             public void run(Nimbus.Iface nimbus) throws Exception {
-                int errorCount = 0;
                 for (String name : names) {
                     try {
                         nimbus.killTopologyWithOpts(name, opts);
diff --git a/storm-core/src/jvm/org/apache/storm/command/Rebalance.java 
b/storm-core/src/jvm/org/apache/storm/command/Rebalance.java
index ffaefda..9fb5180 100644
--- a/storm-core/src/jvm/org/apache/storm/command/Rebalance.java
+++ b/storm-core/src/jvm/org/apache/storm/command/Rebalance.java
@@ -39,6 +39,7 @@ public class Rebalance {
                                     .arg("topologyName", CLI.FIRST_WINS)
                                     .parse(args);
         final String name = (String) cl.get("topologyName");
+        Utils.validateTopologyName(name);
         final RebalanceOptions rebalanceOptions = new RebalanceOptions();
         Integer wait = (Integer) cl.get("w");
         if (null != wait) {
diff --git a/storm-core/src/jvm/org/apache/storm/command/SetLogLevel.java 
b/storm-core/src/jvm/org/apache/storm/command/SetLogLevel.java
index c6e2fa7..e8e67cd 100644
--- a/storm-core/src/jvm/org/apache/storm/command/SetLogLevel.java
+++ b/storm-core/src/jvm/org/apache/storm/command/SetLogLevel.java
@@ -35,6 +35,8 @@ public class SetLogLevel {
                                     .arg("topologyName", CLI.FIRST_WINS)
                                     .parse(args);
         final String topologyName = (String) cl.get("topologyName");
+        Utils.validateTopologyName(topologyName);
+
         final LogConfig logConfig = new LogConfig();
         Map<String, LogLevel> logLevelMap = new HashMap<>();
         Map<String, LogLevel> updateLogLevel = (Map<String, LogLevel>) 
cl.get("l");
diff --git a/storm-core/src/jvm/org/apache/storm/command/UploadCredentials.java 
b/storm-core/src/jvm/org/apache/storm/command/UploadCredentials.java
index c308b59..1a362cf 100644
--- a/storm-core/src/jvm/org/apache/storm/command/UploadCredentials.java
+++ b/storm-core/src/jvm/org/apache/storm/command/UploadCredentials.java
@@ -50,6 +50,7 @@ public class UploadCredentials {
         String credentialFile = (String) cl.get("f");
         List<String> rawCredentials = (List<String>) cl.get("rawCredentials");
         String topologyName = (String) cl.get("topologyName");
+        Utils.validateTopologyName(topologyName);
 
         if (null != rawCredentials && ((rawCredentials.size() % 2) != 0)) {
             throw new RuntimeException("Need an even number of arguments to 
make a map");

Reply via email to