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

rzo1 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 c78ab7462 Give each uploaded dependency artifact its own blob key 
instead of reusing a pre-existing blob
c78ab7462 is described below

commit c78ab74621f8dc9e88e706fca66b051be8e0132f
Author: Richard Zowalla <[email protected]>
AuthorDate: Wed Aug 19 13:07:42 2026 +0200

    Give each uploaded dependency artifact its own blob key instead of reusing 
a pre-existing blob
---
 .../storm/dependency/DependencyUploader.java       | 11 +++++-
 .../storm/dependency/DependencyUploaderTest.java   | 44 ++++++++++++++++++++--
 2 files changed, 50 insertions(+), 5 deletions(-)

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 85ab835eb..de976ad3d 100644
--- a/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java
+++ b/storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java
@@ -114,11 +114,18 @@ public class DependencyUploader {
                 String artifact = artifactToFile.getKey();
                 File dependency = artifactToFile.getValue();
 
-                String key = 
DependencyBlobStoreUtils.generateDependencyBlobKey(convertArtifactToJarFileName(artifact));
+                // the key must be unique per upload, as uploadFiles() already 
does: a key
+                // derived only from the artifact coordinate is shared by 
every submitter, so
+                // an existing blob under that key was reused without checking 
that it holds
+                // the artifact we resolved
+                String key = 
DependencyBlobStoreUtils.generateDependencyBlobKey(
+                    DependencyBlobStoreUtils.applyUUIDToFileName(
+                        convertArtifactToJarFileName(artifact)));
                 try {
                     uploadDependencyToBlobStore(key, dependency);
                 } catch (KeyAlreadyExistsException e) {
-                    // we lose the race, but it doesn't matter
+                    // cannot happen: the key carries a freshly generated UUID
+                    throw new RuntimeException(e);
                 }
 
                 keys.add(key);
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 58648426e..1821e2edc 100644
--- 
a/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java
+++ 
b/storm-client/test/jvm/org/apache/storm/dependency/DependencyUploaderTest.java
@@ -32,6 +32,7 @@ import org.apache.storm.blobstore.BlobStoreAclHandler;
 import org.apache.storm.blobstore.ClientBlobStore;
 import org.apache.storm.generated.AccessControl;
 import org.apache.storm.generated.AccessControlType;
+import org.apache.storm.generated.KeyAlreadyExistsException;
 import org.apache.storm.generated.KeyNotFoundException;
 import org.apache.storm.generated.ReadableBlobMeta;
 import org.apache.storm.generated.SettableBlobMeta;
@@ -44,6 +45,7 @@ import org.mockito.ArgumentCaptor;
 import org.mockito.stubbing.Answer;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
@@ -55,6 +57,7 @@ import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -230,11 +233,11 @@ public class DependencyUploaderTest {
     @Test
     public void uploadArtifactsWhichOneOfThemIsFailedToBeUploaded() throws 
Exception {
         String artifact = "group:artifact:1.0.0";
-        String expectedBlobKeyForArtifact = "group-artifact-1.0.0.jar";
+        String expectedBlobKeyForArtifact = "group-artifact-1.0.0";
         File mockFile = createTemporaryDummyFile();
 
         String artifact2 = "group:artifact2:2.0.0";
-        String expectedBlobKeyForArtifact2 = "group-artifact2-2.0.0.jar";
+        String expectedBlobKeyForArtifact2 = "group-artifact2-2.0.0";
         File mockFile2 = mock(File.class);
         when(mockFile2.getName()).thenReturn("dummy.jar");
         when(mockFile2.isFile()).thenReturn(true);
@@ -286,7 +289,7 @@ public class DependencyUploaderTest {
         when(mockBlobStore.createBlob(anyString(), 
any(SettableBlobMeta.class))).thenReturn(mockOutputStream);
 
         String artifact = "group:artifact:1.0.0";
-        String expectedBlobKeyForArtifact = "group-artifact-1.0.0.jar";
+        String expectedBlobKeyForArtifact = "group-artifact-1.0.0";
         File mockFile = createTemporaryDummyFile();
 
         Map<String, File> artifacts = new LinkedHashMap<>();
@@ -300,6 +303,41 @@ public class DependencyUploaderTest {
         verify(mockOutputStream).close();
     }
 
+    @Test
+    public void uploadArtifactsAssignsUniqueKeyPerUpload() throws Exception {
+        AtomicOutputStream mockOutputStream = mock(AtomicOutputStream.class);
+        doNothing().when(mockOutputStream).close();
+
+        when(mockBlobStore.getBlobMeta(anyString())).thenThrow(new 
KeyNotFoundException());
+        when(mockBlobStore.createBlob(anyString(), 
any(SettableBlobMeta.class))).thenReturn(mockOutputStream);
+
+        String artifact = "group:artifact:1.0.0";
+        Map<String, File> artifacts = new LinkedHashMap<>();
+        artifacts.put(artifact, createTemporaryDummyFile());
+
+        List<String> firstKeys = sut.uploadArtifacts(artifacts);
+        List<String> secondKeys = sut.uploadArtifacts(artifacts);
+
+        assertEquals(1, firstKeys.size());
+        assertEquals(1, secondKeys.size());
+        // the key must not be derived from the artifact coordinate alone, 
otherwise a blob left behind
+        // by another submission is picked up instead of the artifact we just 
resolved
+        assertNotEquals("dep-group-artifact-1.0.0.jar", firstKeys.get(0));
+        assertNotEquals(firstKeys.get(0), secondKeys.get(0));
+        verify(mockBlobStore, times(2)).createBlob(anyString(), 
any(SettableBlobMeta.class));
+    }
+
+    @Test
+    public void uploadArtifactsFailsWhenKeyAlreadyExists() throws Exception {
+        when(mockBlobStore.getBlobMeta(anyString())).thenThrow(new 
KeyNotFoundException());
+        when(mockBlobStore.createBlob(anyString(), 
any(SettableBlobMeta.class))).thenThrow(new KeyAlreadyExistsException());
+
+        Map<String, File> artifacts = new LinkedHashMap<>();
+        artifacts.put("group:artifact:1.0.0", createTemporaryDummyFile());
+
+        assertThrows(RuntimeException.class, () -> 
sut.uploadArtifacts(artifacts));
+    }
+
     private File createTemporaryDummyFile() throws IOException {
         File tempFile = java.nio.file.Files.createTempFile("tempfile", 
".tmp").toFile();
 

Reply via email to