diff --git a/storm-client/src/jvm/org/apache/storm/dependency/DependencyBlobStoreUtils.java b/storm-client/src/jvm/org/apache/storm/dependency/DependencyBlobStoreUtils.java
index 5211e55700..f692323b9a 100644
--- a/storm-client/src/jvm/org/apache/storm/dependency/DependencyBlobStoreUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/dependency/DependencyBlobStoreUtils.java
@@ -24,12 +24,28 @@
public class DependencyBlobStoreUtils {
- private static final String BLOB_DEPENDENCIES_PREFIX = "dep-";
+ /**
+ * The prefix every blob key holding a topology dependency starts with.
+ */
+ public static final String BLOB_DEPENDENCIES_PREFIX = "dep-";
public static String generateDependencyBlobKey(String key) {
return BLOB_DEPENDENCIES_PREFIX + key;
}
+ /**
+ * Tell whether a blob key names a topology dependency, i.e. whether it could have been produced by
+ * {@link #generateDependencyBlobKey(String)}. Keys that a topology only refers to, rather than owns, must be
+ * checked with this before they are acted upon, because the dependency lists of a submitted topology are filled
+ * in by the client and can name any blob at all.
+ *
+ * @param key the blob key to check, may be null
+ * @return true if the key is a dependency blob key
+ */
+ public static boolean isDependencyBlobKey(String key) {
+ return key != null && key.startsWith(BLOB_DEPENDENCIES_PREFIX);
+ }
+
@SuppressWarnings("checkstyle:AbbreviationAsWordInName")
public static String applyUUIDToFileName(String fileName) {
String fileNameWithExt = Files.getNameWithoutExtension(fileName);
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 682db23380..9c8f1174fc 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
@@ -87,6 +87,7 @@
import org.apache.storm.daemon.Shutdownable;
import org.apache.storm.daemon.StormCommon;
import org.apache.storm.daemon.common.FileWatcher;
+import org.apache.storm.dependency.DependencyBlobStoreUtils;
import org.apache.storm.generated.AlreadyAliveException;
import org.apache.storm.generated.Assignment;
import org.apache.storm.generated.AuthorizationException;
@@ -1304,6 +1305,61 @@ private static void validateTopologyName(String name) throws InvalidTopologyExce
}
}
+ /**
+ * Check that a submitted topology only claims blobs that are topology dependencies, and that every one of them
+ * exists. The dependency lists of a submitted topology are filled in by the client, so without this a submission
+ * could name any blob at all, for example another topology's code or configuration blob, and nimbus would delete
+ * it as its own dependency once the submitted topology is cleaned up. A key that exists nowhere is just as
+ * damaging: on gaining leadership a nimbus compares the dependencies of every active topology against its
+ * blobstore and gives up leadership when one is missing, so a single unresolvable key on a single active topology
+ * leaves the cluster without a leader for as long as that topology is active.
+ *
+ *
Existence is probed with {@code getBlobMeta} as the submitter, exactly as
+ * {@link Utils#validateTopologyBlobStoreMap(Map, BlobStore)} probes the blobs of
+ * {@link Config#TOPOLOGY_BLOBSTORE_MAP}, so a submitter can only claim a dependency it is allowed to read. The
+ * client uploads the dependency blobs before it submits, so they are present by the time this runs.
+ *
+ * @param topology the submitted topology
+ * @param blobStore the blobstore to look the keys up in
+ * @param subject the subject to look the keys up as, i.e. the submitter
+ * @throws InvalidTopologyException if a dependency list holds something that is not a dependency blob key, or
+ * names a dependency blob that does not exist
+ * @throws AuthorizationException if the submitter may not read one of the dependency blobs it named
+ */
+ @VisibleForTesting
+ static void validateDependencyBlobKeys(StormTopology topology, BlobStore blobStore, Subject subject)
+ throws InvalidTopologyException, AuthorizationException {
+ Set checked = new HashSet<>();
+ validateDependencyBlobKeys(topology.get_dependency_jars(), "dependency_jars", blobStore, subject, checked);
+ validateDependencyBlobKeys(topology.get_dependency_artifacts(), "dependency_artifacts", blobStore, subject, checked);
+ }
+
+ private static void validateDependencyBlobKeys(List keys, String fieldName, BlobStore blobStore, Subject subject,
+ Set checked) throws InvalidTopologyException, AuthorizationException {
+ if (keys == null) {
+ return;
+ }
+ for (String key : keys) {
+ if (!DependencyBlobStoreUtils.isDependencyBlobKey(key)) {
+ throw new WrappedInvalidTopologyException("Topology " + fieldName + " lists [" + key
+ + "], which is not a dependency blob key; every entry must start with \""
+ + DependencyBlobStoreUtils.BLOB_DEPENDENCIES_PREFIX + "\"");
+ }
+ if (!checked.add(key)) {
+ // the same dependency may be listed twice, one lookup for it is enough
+ continue;
+ }
+ try {
+ blobStore.getBlobMeta(key, subject);
+ } catch (KeyNotFoundException keyNotFound) {
+ throw new WrappedInvalidTopologyException("Topology " + fieldName + " lists [" + key
+ + "], which is not in the blobstore; upload the dependency before submitting the topology, and if it "
+ + "was uploaded earlier note that a dependency blob is deleted once no topology uses it any more, so "
+ + "it has to be uploaded again");
+ }
+ }
+ }
+
private static StormTopology tryReadTopology(String topoId, TopoCache tc)
throws NotAliveException, AuthorizationException, IOException {
try {
@@ -3333,6 +3389,7 @@ public void submitTopologyWithOpts(String topoName, String uploadedJarLocation,
throw new WrappedInvalidTopologyException(ex.getMessage());
}
validator.validate(topoName, topoConf, topology);
+ validateDependencyBlobKeys(topology, blobStore, getSubject());
if ((boolean) conf.getOrDefault(Config.DISABLE_SYMLINKS, false)) {
@SuppressWarnings("unchecked")
Map blobMap = (Map) topoConf.get(Config.TOPOLOGY_BLOBSTORE_MAP);
diff --git a/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusTest.java b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusTest.java
index f24468d4c1..9c580ec53a 100644
--- a/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusTest.java
+++ b/storm-server/src/test/java/org/apache/storm/daemon/nimbus/NimbusTest.java
@@ -40,6 +40,7 @@
import org.apache.storm.blobstore.KeySequenceNumber;
import org.apache.storm.blobstore.LocalFsBlobStore;
import org.apache.storm.cluster.IStormClusterState;
+import org.apache.storm.dependency.DependencyBlobStoreUtils;
import org.apache.storm.generated.AuthorizationException;
import org.apache.storm.generated.Credentials;
import org.apache.storm.generated.InvalidTopologyException;
@@ -49,6 +50,8 @@
import org.apache.storm.generated.ReadableBlobMeta;
import org.apache.storm.generated.SettableBlobMeta;
import org.apache.storm.generated.StormTopology;
+import org.apache.storm.generated.SubmitOptions;
+import org.apache.storm.generated.TopologyInitialStatus;
import org.apache.storm.metric.StormMetricsRegistry;
import org.apache.storm.nimbus.ILeaderElector;
import org.apache.storm.nimbus.NimbusInfo;
@@ -83,6 +86,7 @@
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.mockito.ArgumentMatchers.any;
@@ -92,6 +96,7 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -464,4 +469,146 @@ private static void setCaller(String user) {
subject.getPrincipals().add(new SingleUserPrincipal(user));
ReqContext.context().setSubject(subject);
}
+
+ @Test
+ void testValidateDependencyBlobKeysRejectsKeysThatAreNotDependencies() throws Exception {
+ // a topology fills its own dependency lists in on the client side, so nimbus has to check that they only
+ // name dependency blobs before it takes ownership of them and deletes them during cleanup
+ String victimJarKey = ConfigUtils.masterStormJarKey("victim-1-1234567890");
+ String victimConfKey = ConfigUtils.masterStormConfKey("victim-1-1234567890");
+ Subject submitter = new Subject();
+
+ StormTopology jarField = new StormTopology();
+ jarField.set_dependency_jars(List.of(victimJarKey));
+ InvalidTopologyException jarException = assertThrows(InvalidTopologyException.class,
+ () -> Nimbus.validateDependencyBlobKeys(jarField, localBlobStore, submitter));
+ assertTrue(jarException.get_msg().contains(victimJarKey), jarException.get_msg());
+ assertTrue(jarException.get_msg().contains("dependency_jars"), jarException.get_msg());
+
+ StormTopology artifactField = new StormTopology();
+ artifactField.set_dependency_artifacts(List.of(victimConfKey));
+ InvalidTopologyException artifactException = assertThrows(InvalidTopologyException.class,
+ () -> Nimbus.validateDependencyBlobKeys(artifactField, localBlobStore, submitter));
+ assertTrue(artifactException.get_msg().contains(victimConfKey), artifactException.get_msg());
+ assertTrue(artifactException.get_msg().contains("dependency_artifacts"), artifactException.get_msg());
+
+ // a good key followed by a bad one is caught too, and the message names the bad one
+ StormTopology mixed = new StormTopology();
+ mixed.set_dependency_jars(List.of(dependencyKey("a.jar"), victimJarKey));
+ InvalidTopologyException mixedException = assertThrows(InvalidTopologyException.class,
+ () -> Nimbus.validateDependencyBlobKeys(mixed, localBlobStore, submitter));
+ assertTrue(mixedException.get_msg().contains(victimJarKey), mixedException.get_msg());
+
+ // a key that is not a dependency key at all is rejected on its name, without asking the blobstore about it
+ verify(localBlobStore, never()).getBlobMeta(eq(victimJarKey), any());
+ verify(localBlobStore, never()).getBlobMeta(eq(victimConfKey), any());
+ }
+
+ @Test
+ void testValidateDependencyBlobKeysRejectsKeyThatIsNotInTheBlobStore() throws Exception {
+ // a key that merely looks like a dependency key is just as damaging: on gaining leadership a nimbus gives up
+ // leadership again when an active topology names a dependency it cannot find, so an unresolvable key leaves
+ // the cluster without a leader
+ String presentKey = dependencyKey("present.jar");
+ String missingKey = dependencyKey("missing.jar");
+ Subject submitter = new Subject();
+ when(localBlobStore.getBlobMeta(eq(missingKey), any())).thenThrow(new KeyNotFoundException(missingKey));
+
+ StormTopology jarField = new StormTopology();
+ jarField.set_dependency_jars(List.of(presentKey, missingKey));
+ InvalidTopologyException jarException = assertThrows(InvalidTopologyException.class,
+ () -> Nimbus.validateDependencyBlobKeys(jarField, localBlobStore, submitter));
+ assertTrue(jarException.get_msg().contains(missingKey), jarException.get_msg());
+ assertTrue(jarException.get_msg().contains("dependency_jars"), jarException.get_msg());
+ assertTrue(jarException.get_msg().contains("not in the blobstore"), jarException.get_msg());
+
+ StormTopology artifactField = new StormTopology();
+ artifactField.set_dependency_artifacts(List.of(missingKey));
+ InvalidTopologyException artifactException = assertThrows(InvalidTopologyException.class,
+ () -> Nimbus.validateDependencyBlobKeys(artifactField, localBlobStore, submitter));
+ assertTrue(artifactException.get_msg().contains(missingKey), artifactException.get_msg());
+ assertTrue(artifactException.get_msg().contains("dependency_artifacts"), artifactException.get_msg());
+ }
+
+ @Test
+ void testValidateDependencyBlobKeysLooksBlobsUpAsTheSubmitter() throws Exception {
+ // looking the blob up as the submitter, the way the TOPOLOGY_BLOBSTORE_MAP entries are looked up, also
+ // answers whether the submitter is allowed to read the dependency it claims; a dependency blob is uploaded
+ // with OTHER READ, so a legitimate submission passes
+ String key = dependencyKey("some-jar.jar");
+ Subject submitter = new Subject();
+
+ StormTopology topology = new StormTopology();
+ // listed under both fields to show that a key is looked up once no matter how often it is named
+ topology.set_dependency_jars(List.of(key, key));
+ topology.set_dependency_artifacts(List.of(key));
+ assertDoesNotThrow(() -> Nimbus.validateDependencyBlobKeys(topology, localBlobStore, submitter));
+
+ verify(localBlobStore, times(1)).getBlobMeta(key, submitter);
+ }
+
+ @Test
+ void testValidateDependencyBlobKeysAcceptsGeneratedKeysThatExist() throws Exception {
+ Subject submitter = new Subject();
+ StormTopology topology = new StormTopology();
+ topology.set_dependency_jars(List.of(dependencyKey("some-jar.jar"), dependencyKey("no-extension")));
+ topology.set_dependency_artifacts(List.of(dependencyKey("group-artifact-1.0.jar")));
+ assertDoesNotThrow(() -> Nimbus.validateDependencyBlobKeys(topology, localBlobStore, submitter));
+
+ // unset lists are how a topology submitted without dependencies looks
+ assertDoesNotThrow(() -> Nimbus.validateDependencyBlobKeys(new StormTopology(), localBlobStore, submitter));
+ verify(localBlobStore, never()).getBlobMeta(eq(null), any());
+ }
+
+ @Test
+ void testSubmitTopologyRejectsDependencyBlobKeyOfAnotherTopology() throws Exception {
+ Map conf = Map.of(DaemonConfig.NIMBUS_MONITOR_FREQ_SECS, 10);
+ Nimbus submitNimbus = new Nimbus(conf, iNimbus, stormClusterState, nimbusInfo, localBlobStore, leaderElector,
+ groupMapper, new StormMetricsRegistry());
+ when(leaderElector.isLeader()).thenReturn(true);
+ when(stormClusterState.getTopoId(any())).thenReturn(Optional.empty());
+
+ TopologyBuilder builder = new TopologyBuilder();
+ builder.setSpout("wordSpout", new TestWordSpout(), 1);
+ StormTopology topology = builder.createTopology();
+ String victimJarKey = ConfigUtils.masterStormJarKey("victim-1-1234567890");
+ topology.set_dependency_artifacts(List.of(victimJarKey));
+
+ InvalidTopologyException exception = assertThrows(InvalidTopologyException.class,
+ () -> submitNimbus.submitTopologyWithOpts("thief", "/dev/null", "{}", topology,
+ new SubmitOptions(TopologyInitialStatus.ACTIVE)));
+ assertTrue(exception.get_msg().contains(victimJarKey), exception.get_msg());
+
+ // the submission was rejected before anything was stored for it
+ verify(stormClusterState, never()).setupHeatbeats(any(), any());
+ }
+
+ @Test
+ void testSubmitTopologyRejectsDependencyBlobKeyThatDoesNotExist() throws Exception {
+ Map conf = Map.of(DaemonConfig.NIMBUS_MONITOR_FREQ_SECS, 10);
+ Nimbus submitNimbus = new Nimbus(conf, iNimbus, stormClusterState, nimbusInfo, localBlobStore, leaderElector,
+ groupMapper, new StormMetricsRegistry());
+ when(leaderElector.isLeader()).thenReturn(true);
+ when(stormClusterState.getTopoId(any())).thenReturn(Optional.empty());
+ String missingKey = dependencyKey("missing.jar");
+ when(localBlobStore.getBlobMeta(eq(missingKey), any())).thenThrow(new KeyNotFoundException(missingKey));
+
+ TopologyBuilder builder = new TopologyBuilder();
+ builder.setSpout("wordSpout", new TestWordSpout(), 1);
+ StormTopology topology = builder.createTopology();
+ topology.set_dependency_jars(List.of(missingKey));
+
+ InvalidTopologyException exception = assertThrows(InvalidTopologyException.class,
+ () -> submitNimbus.submitTopologyWithOpts("ghost", "/dev/null", "{}", topology,
+ new SubmitOptions(TopologyInitialStatus.ACTIVE)));
+ assertTrue(exception.get_msg().contains(missingKey), exception.get_msg());
+ assertTrue(exception.get_msg().contains("not in the blobstore"), exception.get_msg());
+
+ // the submission was rejected before anything was stored for it
+ verify(stormClusterState, never()).setupHeatbeats(any(), any());
+ }
+
+ private static String dependencyKey(String fileName) {
+ return DependencyBlobStoreUtils.generateDependencyBlobKey(DependencyBlobStoreUtils.applyUUIDToFileName(fileName));
+ }
}