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

rzo1 pushed a commit to branch fix/blobstore-dependency-cleanup
in repository https://gitbox.apache.org/repos/asf/storm.git

commit 1ea0dc0474a2855cd8bc0122ac19116a5ff9240c
Author: Richard Zowalla <[email protected]>
AuthorDate: Sat Aug 22 19:53:05 2026 +0200

    STORM-3871: remove a topology's dependency artifact blobs when it is 
cleaned up
---
 .../src/jvm/org/apache/storm/StormSubmitter.java   |  12 +-
 .../storm/dependency/DependencyUploader.java       |   5 +
 .../storm/dependency/DependencyUploaderTest.java   |   5 +-
 .../org/apache/storm/daemon/nimbus/Nimbus.java     | 123 ++++++++++++++++--
 .../storm/daemon/nimbus/NimbusClojurePortTest.java |   4 +-
 .../org/apache/storm/daemon/nimbus/NimbusTest.java | 137 +++++++++++++++++++++
 6 files changed, 270 insertions(+), 16 deletions(-)

diff --git a/storm-client/src/jvm/org/apache/storm/StormSubmitter.java 
b/storm-client/src/jvm/org/apache/storm/StormSubmitter.java
index 18ff3ba60..1bc5627e1 100644
--- a/storm-client/src/jvm/org/apache/storm/StormSubmitter.java
+++ b/storm-client/src/jvm/org/apache/storm/StormSubmitter.java
@@ -271,7 +271,7 @@ public class StormSubmitter {
 
                 // Dependency uploading only makes sense for distributed mode
                 List<String> jarsBlobKeys = Collections.emptyList();
-                List<String> artifactsBlobKeys;
+                List<String> artifactsBlobKeys = Collections.emptyList();
 
                 DependencyUploader uploader = new DependencyUploader();
                 try {
@@ -281,8 +281,10 @@ public class StormSubmitter {
 
                     artifactsBlobKeys = 
uploadDependencyArtifactsToBlobStore(uploader);
                 } catch (Throwable e) {
-                    // remove uploaded jars blobs, not artifacts since they're 
shared across the cluster
+                    // every uploaded blob carries a key unique to this 
submission, and no topology refers to
+                    // them yet, so nothing else can be using them
                     uploader.deleteBlobs(jarsBlobKeys);
+                    uploader.deleteBlobs(artifactsBlobKeys);
                     uploader.shutdown();
                     throw e;
                 }
@@ -291,10 +293,12 @@ public class StormSubmitter {
                     setDependencyBlobsToTopology(topology, jarsBlobKeys, 
artifactsBlobKeys);
                     submitTopologyInDistributeMode(name, topology, opts, 
progressListener, asUser, conf, serConf, client);
                 } catch (AlreadyAliveException | InvalidTopologyException | 
AuthorizationException e) {
-                    // remove uploaded jars blobs, not artifacts since they're 
shared across the cluster
-                    // Note that we don't handle TException to delete jars 
blobs
+                    // the topology was rejected, so the blobs it refers to 
are unreachable; their keys are
+                    // unique to this submission, so nothing else can be using 
them
+                    // Note that we don't handle TException to delete the blobs
                     // because it's safer to leave some blobs instead of 
topology not running
                     uploader.deleteBlobs(jarsBlobKeys);
+                    uploader.deleteBlobs(artifactsBlobKeys);
                     throw e;
                 } finally {
                     uploader.shutdown();
diff --git 
a/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java 
b/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java
index de976ad3d..cf89990c1 100644
--- a/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java
+++ b/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java
@@ -131,6 +131,11 @@ public class DependencyUploader {
                 keys.add(key);
             }
         } catch (Throwable e) {
+            // the keys are unique to this upload and no topology refers to 
them, so the ones that made it to
+            // the blob store are only reachable from here
+            if (getBlobStore() != null) {
+                deleteBlobs(keys);
+            }
             throw new RuntimeException(e);
         }
 
diff --git 
a/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java 
b/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java
index 1821e2edc..e549fc83b 100644
--- 
a/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java
+++ 
b/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java
@@ -264,8 +264,9 @@ public class DependencyUploaderTest {
         
verify(mockBlobStore).getBlobMeta(contains(expectedBlobKeyForArtifact));
         
verify(mockBlobStore).getBlobMeta(contains(expectedBlobKeyForArtifact2));
 
-        // never rollback
-        verify(mockBlobStore, 
never()).deleteBlob(contains(expectedBlobKeyForArtifact));
+        // the artifacts uploaded before the failure are rolled back: their 
keys are unique to this upload,
+        // so nothing else can be referring to them and leaving them behind 
leaks blob store space forever
+        verify(mockBlobStore).deleteBlob(contains(expectedBlobKeyForArtifact));
         verify(mockBlobStore, 
never()).deleteBlob(contains(expectedBlobKeyForArtifact2));
     }
 
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 2d7704aaa..11cf7690a 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
@@ -62,6 +62,7 @@ import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Function;
 import java.util.function.UnaryOperator;
+import java.util.regex.Pattern;
 import java.util.stream.Collectors;
 import javax.security.auth.Subject;
 
@@ -411,6 +412,9 @@ public class Nimbus implements Iface, Shutdownable, 
DaemonCommon {
             .build();
     private static final List<String> EMPTY_STRING_LIST = 
Collections.unmodifiableList(Collections.emptyList());
     private static final Set<String> EMPTY_STRING_SET = 
Collections.unmodifiableSet(Collections.emptySet());
+    //A dependency blob key whose file name part ends with a canonical UUID, 
which a client splices in per upload.
+    private static final Pattern UNIQUE_DEPENDENCY_KEY = Pattern.compile(
+        
"^dep-.+-[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}(\\..+)?$");
     private static final RotatingMap<String, Long> topologyCleanupDetected = 
new RotatingMap<>(2);
     private static long topologyCleanupRotationTime = 0L;
 
@@ -2884,22 +2888,114 @@ public class Nimbus implements Iface, Shutdownable, 
DaemonCommon {
         return ret;
     }
 
+    /**
+     * Collect the dependency blob keys that are still referenced by a 
topology other than the ones being cleaned up.
+     *
+     * <p>Dependency blobs uploaded by a current client are unique to a single 
submission, but a cluster upgraded from an
+     * older release can still hold artifact blobs whose key is derived from 
the maven coordinate alone, and those are
+     * listed by every topology that was submitted with the same artifact. 
Deleting one of those while a topology still
+     * needs it breaks worker launch and makes every nimbus refuse leadership, 
so the reference set has to be exact.
+     *
+     * <p>This deliberately does not swallow read failures: a topology whose 
code blob cannot be read has unknown
+     * references, and treating it as referencing nothing would allow a live 
topology's dependencies to be deleted. Only
+     * a missing code blob is treated as contributing no references, because 
such a topology cannot be launched anyway.
+     *
+     * @param excludedTopoIds the topologies being cleaned up, which must not 
count as referencing anything
+     * @return every dependency blob key referenced by a topology that is not 
being cleaned up
+     */
+    @VisibleForTesting
+    Set<String> referencedDependencyKeys(Set<String> excludedTopoIds) throws 
Exception {
+        Set<String> candidateTopoIds = new 
HashSet<>(Utils.OR(blobStore.storedTopoIds(), EMPTY_STRING_SET));
+        candidateTopoIds.addAll(Utils.OR(stormClusterState.activeStorms(), 
EMPTY_STRING_LIST));
+        candidateTopoIds.removeAll(excludedTopoIds);
+
+        Set<String> referenced = new HashSet<>();
+        for (String topoId : candidateTopoIds) {
+            StormTopology topo;
+            try {
+                topo = readStormTopologyAsNimbus(topoId, topoCache);
+            } catch (KeyNotFoundException e) {
+                //The topology has no code blob, so it references no 
dependencies and cannot be launched.
+                LOG.debug("No code found for {} while collecting dependency 
blob references", topoId);
+                continue;
+            }
+            if (topo.is_set_dependency_jars()) {
+                referenced.addAll(topo.get_dependency_jars());
+            }
+            if (topo.is_set_dependency_artifacts()) {
+                referenced.addAll(topo.get_dependency_artifacts());
+            }
+        }
+        return referenced;
+    }
+
+    /**
+     * Tell whether a dependency blob key proves, by its shape alone, that no 
other topology can refer to it.
+     *
+     * <p>A client that uploads a dependency splices a freshly generated UUID 
into the file name before prefixing it
+     * with {@code dep-}, so a key of that shape belongs to exactly one upload 
and therefore to exactly one topology.
+     * Artifact keys written by older clients are derived from the maven 
coordinate alone
+     * ({@code dep-<group>-<artifact>-<version>.jar}) and are listed by every 
topology built against that artifact, so
+     * they do not match: a coordinate would have to end in five dash 
separated hexadecimal groups of exactly
+     * 8-4-4-4-12 characters for that.
+     *
+     * <p>This lives here rather than next to the key generator on purpose. It 
is not a restatement of what the current
+     * client writes, but nimbus' own list of the key shapes it is willing to 
treat as unique, and it has to keep
+     * recognising the shapes written by every client version whose blobs may 
still be in the store even if the
+     * generator changes.
+     *
+     * @param key a dependency blob key listed by a topology
+     * @return true only if no other topology can be referring to the key
+     */
+    @VisibleForTesting
+    static boolean isProvablyUniqueDependencyKey(String key) {
+        return key != null && UNIQUE_DEPENDENCY_KEY.matcher(key).matches();
+    }
+
+    /**
+     * Remove the dependency blobs of a topology that is being cleaned up, 
keeping the ones another topology still uses.
+     *
+     * <p>When {@code stillReferenced} is null the cluster-wide reference scan 
failed, so it is unknown which blobs
+     * other topologies use. The topology's code blob is deleted right after 
this call and a dependency blob key
+     * carries no topology id, so anything left behind now can never be found 
again. Instead of giving up, the keys
+     * whose shape proves that no other topology can refer to them, that is 
the ones carrying a generated UUID, are
+     * still reclaimed; keys that could be shared are kept.
+     *
+     * @param topoId          the topology being cleaned up
+     * @param stillReferenced the dependency blob keys referenced by 
topologies that are not being cleaned up, or null
+     *                        if that could not be determined
+     */
     @VisibleForTesting
-    public void rmDependencyJarsInTopology(String topoId) {
+    public void rmDependencyBlobsInTopology(String topoId, Set<String> 
stillReferenced) {
         try {
             BlobStore store = blobStore;
             IStormClusterState state = stormClusterState;
             StormTopology topo = readStormTopologyAsNimbus(topoId, topoCache);
-            List<String> dependencyJars = topo.get_dependency_jars();
-            LOG.info("Removing dependency jars from blobs - {}", 
dependencyJars);
-            if (dependencyJars != null && !dependencyJars.isEmpty()) {
-                for (String key : dependencyJars) {
+            Set<String> dependencies = new HashSet<>();
+            if (topo.is_set_dependency_jars()) {
+                dependencies.addAll(topo.get_dependency_jars());
+            }
+            if (topo.is_set_dependency_artifacts()) {
+                dependencies.addAll(topo.get_dependency_artifacts());
+            }
+            LOG.info("Removing dependency blobs of {} - {}", topoId, 
dependencies);
+            for (String key : dependencies) {
+                if (stillReferenced == null) {
+                    if (isProvablyUniqueDependencyKey(key)) {
+                        rmBlobKey(store, key, state);
+                    } else {
+                        LOG.warn("Keeping dependency blob {} of {}, another 
topology may refer to it and the references "
+                            + "could not be read", key, topoId);
+                    }
+                } else if (stillReferenced.contains(key)) {
+                    LOG.info("Keeping dependency blob {} of {}, it is still 
referenced by another topology", key, topoId);
+                } else {
                     rmBlobKey(store, key, state);
                 }
             }
         } catch (Exception e) {
-            //Yes eat the exception
-            LOG.info("Exception {}", e);
+            //Yes eat the exception, cleaning up the rest of the topology 
matters more, but this leaves blobs behind.
+            LOG.warn("Could not remove the dependency blobs of {}, they will 
be left in the blob store", topoId, e);
         }
     }
 
@@ -2940,13 +3036,24 @@ public class Nimbus implements Iface, Shutdownable, 
DaemonCommon {
             Set<String> toClean = new HashSet<>(topoIdsToClean(state, 
blobStore, this.conf));
             long topoIdSelectionDurationMs = Time.deltaMs(cleanupStartMs);
 
+            //Computed once for the whole pass, with the dying topologies 
excluded so that two of them sharing a
+            //dependency blob do not keep it alive for each other. A null 
result means the references are unknown and
+            //only the blobs that are provably unique to one topology are 
reclaimed below.
+            Set<String> stillReferenced = null;
+            try {
+                stillReferenced = referencedDependencyKeys(toClean);
+            } catch (Exception e) {
+                LOG.warn("Could not determine which dependency blobs are still 
in use, only the dependency blobs whose "
+                    + "key proves they belong to a single topology are 
reclaimed in this pass, the others are kept", e);
+            }
+
             for (String topoId : toClean) {
                 LOG.info("Cleaning up {}", topoId);
                 state.teardownHeartbeats(topoId);
                 state.teardownTopologyErrors(topoId);
                 state.removeAllPrivateWorkerKeys(topoId);
                 state.removeBackpressure(topoId);
-                rmDependencyJarsInTopology(topoId);
+                rmDependencyBlobsInTopology(topoId, stillReferenced);
                 forceDeleteTopoDistDir(topoId);
                 rmTopologyKeys(topoId);
                 heartbeatsCache.removeTopo(topoId);
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 526cb0924..0975e52a5 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
@@ -455,8 +455,8 @@ public class NimbusClojurePortTest {
         Mockito.verify(nimbus).rmTopologyKeys("topo3");
 
         // removed topology dependencies
-        Mockito.verify(nimbus).rmDependencyJarsInTopology("topo2");
-        Mockito.verify(nimbus).rmDependencyJarsInTopology("topo3");
+        
Mockito.verify(nimbus).rmDependencyBlobsInTopology(Mockito.eq("topo2"), 
Mockito.anySet());
+        
Mockito.verify(nimbus).rmDependencyBlobsInTopology(Mockito.eq("topo3"), 
Mockito.anySet());
 
         // remove topos from heartbeat cache
         assertEquals(0, nimbus.getHeartbeatsCache().getNumToposCached());
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 96336f7cc..2146deea5 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
@@ -21,8 +21,10 @@ package org.apache.storm.daemon.nimbus;
 import java.nio.file.Files;
 import java.nio.file.Path;
 import java.nio.file.Paths;
+import java.io.IOException;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.List;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
@@ -33,13 +35,17 @@ import org.apache.commons.io.FileUtils;
 import org.apache.storm.Config;
 import org.apache.storm.DaemonConfig;
 import org.apache.storm.blobstore.BlobStore;
+import org.apache.storm.blobstore.BlobStoreAclHandler;
 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.InvalidTopologyException;
 import org.apache.storm.generated.KeyNotFoundException;
 import org.apache.storm.generated.RebalanceOptions;
+import org.apache.storm.generated.ReadableBlobMeta;
+import org.apache.storm.generated.SettableBlobMeta;
 import org.apache.storm.generated.StormTopology;
 import org.apache.storm.metric.StormMetricsRegistry;
 import org.apache.storm.nimbus.ILeaderElector;
@@ -55,7 +61,9 @@ import org.apache.storm.security.auth.SingleUserPrincipal;
 import org.apache.storm.testing.TestWordSpout;
 import org.apache.storm.thrift.TException;
 import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.utils.ConfigUtils;
 import org.apache.storm.utils.ServerUtils;
+import org.apache.storm.utils.Utils;
 import org.apache.storm.utils.Time;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
@@ -65,6 +73,7 @@ import org.mockito.MockedConstruction;
 import org.mockito.MockitoAnnotations;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
@@ -75,12 +84,17 @@ import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.atLeastOnce;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 class NimbusTest {
     private static final String BLOB_FILE_KEY = "file-key";
+    // an artifact blob key from before the key carried a uuid, which several 
topologies can share
+    private static final String LEGACY_ARTIFACT_KEY = 
"dep-group-artifact-1.0.0.jar";
+    // a dependency key written by a current client, which splices a generated 
uuid into the file name
+    private static final String UNIQUE_JAR_KEY = 
"dep-lib-11111111-1111-1111-1111-111111111111.jar";
 
     @Mock
     private StormMetricsRegistry metricRegistry;
@@ -285,4 +299,127 @@ class NimbusTest {
             ReqContext.reset();
         }
     }
+    /**
+     * Register a topology in the blob store mock so that Nimbus can read its 
dependency lists back.
+     */
+    private static void storeTopology(BlobStore store, String topoId, 
List<String> jars, List<String> artifacts)
+        throws Exception {
+        StormTopology topo = new StormTopology();
+        topo.set_spouts(new HashMap<>());
+        topo.set_bolts(new HashMap<>());
+        topo.set_state_spouts(new HashMap<>());
+        topo.set_dependency_jars(jars);
+        topo.set_dependency_artifacts(artifacts);
+        String key = ConfigUtils.masterStormCodeKey(topoId);
+        when(store.readBlob(eq(key), any())).thenReturn(Utils.serialize(topo));
+        when(store.getBlobMeta(eq(key), any()))
+            .thenReturn(new ReadableBlobMeta(new 
SettableBlobMeta(BlobStoreAclHandler.WORLD_EVERYTHING), 1));
+    }
+
+    private Nimbus cleanupNimbus(BlobStore store, IStormClusterState state) 
throws Exception {
+        Map<String, Object> conf = new HashMap<>();
+        conf.put(DaemonConfig.NIMBUS_MONITOR_FREQ_SECS, 10);
+        conf.put(DaemonConfig.NIMBUS_TOPOLOGY_BLOBSTORE_DELETION_DELAY_MS, 0);
+        conf.put(Config.NIMBUS_THRIFT_TLS_PORT, 0);
+        when(leaderElector.isLeader()).thenReturn(true);
+        return new Nimbus(conf, iNimbus, state, nimbusInfo, store, 
leaderElector, groupMapper, new StormMetricsRegistry());
+    }
+
+    @Test
+    void doCleanupRemovesBothDependencyJarsAndArtifactsOfADeadTopology() 
throws Exception {
+        BlobStore store = mock(BlobStore.class);
+        IStormClusterState state = mock(IStormClusterState.class);
+        when(store.storedTopoIds()).thenReturn(Set.of("dead-topo"));
+        when(state.activeStorms()).thenReturn(List.of());
+        storeTopology(store, "dead-topo", 
List.of("dep-lib-11111111-1111-1111-1111-111111111111.jar"),
+            
List.of("dep-group-artifact-1.0.0-22222222-2222-2222-2222-222222222222.jar"));
+
+        cleanupNimbus(store, state).doCleanup();
+
+        
verify(store).deleteBlob(eq("dep-lib-11111111-1111-1111-1111-111111111111.jar"),
 any());
+        
verify(store).deleteBlob(eq("dep-group-artifact-1.0.0-22222222-2222-2222-2222-222222222222.jar"),
 any());
+    }
+
+    @Test
+    void doCleanupKeepsADependencyBlobThatAnotherLiveTopologyStillReferences() 
throws Exception {
+        BlobStore store = mock(BlobStore.class);
+        IStormClusterState state = mock(IStormClusterState.class);
+        when(store.storedTopoIds()).thenReturn(Set.of("dead-topo", 
"live-topo"));
+        when(state.activeStorms()).thenReturn(List.of("live-topo"));
+        // both were submitted before artifact keys carried a uuid, so they 
share one artifact blob
+        storeTopology(store, "dead-topo", 
List.of("dep-lib-11111111-1111-1111-1111-111111111111.jar"),
+            List.of(LEGACY_ARTIFACT_KEY));
+        storeTopology(store, "live-topo", 
List.of("dep-lib-33333333-3333-3333-3333-333333333333.jar"),
+            List.of(LEGACY_ARTIFACT_KEY));
+
+        cleanupNimbus(store, state).doCleanup();
+
+        verify(store, never()).deleteBlob(eq(LEGACY_ARTIFACT_KEY), any());
+        // the blob that is unique to the dead topology is still reclaimed
+        
verify(store).deleteBlob(eq("dep-lib-11111111-1111-1111-1111-111111111111.jar"),
 any());
+    }
+
+    @Test
+    void 
doCleanupRemovesADependencyBlobSharedOnlyBetweenTopologiesThatAreAllBeingCleanedUp()
 throws Exception {
+        BlobStore store = mock(BlobStore.class);
+        IStormClusterState state = mock(IStormClusterState.class);
+        when(store.storedTopoIds()).thenReturn(Set.of("dead-one", "dead-two"));
+        when(state.activeStorms()).thenReturn(List.of());
+        storeTopology(store, "dead-one", List.of(), 
List.of(LEGACY_ARTIFACT_KEY));
+        storeTopology(store, "dead-two", List.of(), 
List.of(LEGACY_ARTIFACT_KEY));
+
+        cleanupNimbus(store, state).doCleanup();
+
+        verify(store, atLeastOnce()).deleteBlob(eq(LEGACY_ARTIFACT_KEY), 
any());
+    }
+
+    @Test
+    void 
doCleanupReclaimsOnlyProvablyUniqueDependencyBlobsWhenTheReferencesCannotBeRead()
 throws Exception {
+        BlobStore store = mock(BlobStore.class);
+        IStormClusterState state = mock(IStormClusterState.class);
+        when(store.storedTopoIds()).thenReturn(Set.of("dead-topo", 
"live-topo"));
+        when(state.activeStorms()).thenReturn(List.of("live-topo"));
+        storeTopology(store, "dead-topo", List.of(UNIQUE_JAR_KEY), 
List.of(LEGACY_ARTIFACT_KEY));
+        // the live topology's code blob cannot be read, so what it references 
is unknown
+        when(store.readBlob(eq(ConfigUtils.masterStormCodeKey("live-topo")), 
any()))
+            .thenThrow(new IOException("blob store is unhappy"));
+
+        cleanupNimbus(store, state).doCleanup();
+
+        // the dead topology's code blob goes away in this pass and a 
dependency key carries no topology id, so a
+        // blob that is not reclaimed now can never be found again; the key 
that carries a uuid cannot be shared
+        verify(store).deleteBlob(eq(UNIQUE_JAR_KEY), any());
+        // the legacy key could still be listed by the topology whose 
references could not be read
+        verify(store, never()).deleteBlob(eq(LEGACY_ARTIFACT_KEY), any());
+        // the rest of the cleanup still runs
+        
verify(store).deleteBlob(eq(ConfigUtils.masterStormJarKey("dead-topo")), any());
+    }
+
+    @Test
+    void 
everyDependencyKeyACurrentClientGeneratesIsRecognisedAsUniqueToOneTopology() {
+        for (String fileName : List.of("commons-lang3-3.12.0.jar", 
"some.lib.tar.gz", "noextension")) {
+            String key = DependencyBlobStoreUtils.generateDependencyBlobKey(
+                DependencyBlobStoreUtils.applyUUIDToFileName(fileName));
+            assertTrue(Nimbus.isProvablyUniqueDependencyKey(key), key + " 
should be recognised as unique");
+        }
+    }
+
+    @Test
+    void 
aDependencyKeyThatDoesNotCarryAUuidIsNotTreatedAsUniqueToOneTopology() {
+        // the shapes an older client wrote for an artifact, dep- plus the 
maven coordinate with : replaced by -
+        for (String key : List.of("dep-group-artifact-1.0.0.jar",
+                                  
"dep-org.apache.commons-commons-lang3-3.12.0.jar",
+                                  "dep-a-b-1.2.3-SNAPSHOT.jar",
+                                  // hexadecimal looking coordinates of the 
wrong lengths are not a uuid either
+                                  "dep-com.deadbeef-cafebabe-1.0.jar",
+                                  
"dep-abcdefab-abcd-abcd-abcd-abcdefabcdef-1.0.jar",
+                                  // the uuid a client splices in is the last 
thing before the extension, so a
+                                  // coordinate that merely contains a uuid 
shaped run is still shareable
+                                  
"dep-com.acme-abcdefab-abcd-abcd-abcd-abcdefabcdef-1.0.jar",
+                                  // a uuid with a dot instead of a dash is 
not canonical
+                                  
"dep-lib-11111111.1111-1111-1111-111111111111.jar")) {
+            assertFalse(Nimbus.isProvablyUniqueDependencyKey(key), key + " 
should not be recognised as unique");
+        }
+        assertFalse(Nimbus.isProvablyUniqueDependencyKey(null));
+    }
 }

Reply via email to