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");