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

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


The following commit(s) were added to refs/heads/master by this push:
     new a700890  NIFI-7298: PutAzureDataLakeStorage tests.
a700890 is described below

commit a7008903edf3cc751a48695ba09701667892e545
Author: Peter Turcsanyi <[email protected]>
AuthorDate: Wed Apr 22 23:22:04 2020 +0200

    NIFI-7298: PutAzureDataLakeStorage tests.
    
    Signed-off-by: Pierre Villard <[email protected]>
    
    This closes #4227.
---
 .../azure/storage/PutAzureDataLakeStorage.java     |  41 ++--
 .../storage/AbstractAzureDataLakeStorageIT.java    |  57 +++++
 .../azure/storage/AbstractAzureStorageIT.java      |   7 +-
 .../storage/ITDeleteAzureDataLakeStorage.java      |   2 +
 .../azure/storage/ITFetchAzureDataLakeStorage.java |   2 +
 .../azure/storage/ITPutAzureDataLakeStorage.java   | 235 +++++++++++++++++++--
 .../storage/TestAbstractAzureDataLakeStorage.java  | 124 +++++++++++
 7 files changed, 437 insertions(+), 31 deletions(-)

diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
index b59c366..94fa90f 100644
--- 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
@@ -16,12 +16,11 @@
  */
 package org.apache.nifi.processors.azure.storage;
 
-import java.io.BufferedInputStream;
-import java.io.InputStream;
-import java.util.HashMap;
-import java.util.Map;
-import java.util.concurrent.TimeUnit;
-
+import com.azure.storage.file.datalake.DataLakeDirectoryClient;
+import com.azure.storage.file.datalake.DataLakeFileClient;
+import com.azure.storage.file.datalake.DataLakeFileSystemClient;
+import com.azure.storage.file.datalake.DataLakeServiceClient;
+import org.apache.commons.lang3.StringUtils;
 import org.apache.nifi.annotation.behavior.InputRequirement;
 import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
 import org.apache.nifi.annotation.behavior.WritesAttribute;
@@ -33,13 +32,13 @@ import org.apache.nifi.flowfile.FlowFile;
 import org.apache.nifi.processor.ProcessContext;
 import org.apache.nifi.processor.ProcessSession;
 import org.apache.nifi.processor.exception.ProcessException;
-
-import com.azure.storage.file.datalake.DataLakeDirectoryClient;
-import com.azure.storage.file.datalake.DataLakeFileClient;
-import com.azure.storage.file.datalake.DataLakeFileSystemClient;
-import com.azure.storage.file.datalake.DataLakeServiceClient;
 import org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
 
+import java.io.BufferedInputStream;
+import java.io.InputStream;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 @Tags({"azure", "microsoft", "cloud", "storage", "adlsgen2", "datalake"})
 @SeeAlso({DeleteAzureDataLakeStorage.class, FetchAzureDataLakeStorage.class})
@@ -50,7 +49,6 @@ import 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
         @WritesAttribute(attribute = "azure.primaryUri", description = 
"Primary location for file content"),
         @WritesAttribute(attribute = "azure.length", description = "Length of 
the file")})
 @InputRequirement(Requirement.INPUT_REQUIRED)
-
 public class PutAzureDataLakeStorage extends 
AbstractAzureDataLakeStorageProcessor {
 
     @Override
@@ -59,23 +57,35 @@ public class PutAzureDataLakeStorage extends 
AbstractAzureDataLakeStorageProcess
         if (flowFile == null) {
             return;
         }
+
         final long startNanos = System.nanoTime();
         try {
             final String fileSystem = 
context.getProperty(FILESYSTEM).evaluateAttributeExpressions(flowFile).getValue();
             final String directory = 
context.getProperty(DIRECTORY).evaluateAttributeExpressions(flowFile).getValue();
             final String fileName = 
context.getProperty(FILE).evaluateAttributeExpressions(flowFile).getValue();
+
+            if (StringUtils.isBlank(fileSystem)) {
+                throw new ProcessException(FILESYSTEM.getDisplayName() + " 
property evaluated to empty string. " +
+                        FILESYSTEM.getDisplayName() + " must be specified as a 
non-empty string.");
+            }
+            if (StringUtils.isBlank(fileName)) {
+                throw new ProcessException(FILE.getDisplayName() + " property 
evaluated to empty string. " +
+                        FILE.getDisplayName() + " must be specified as a 
non-empty string.");
+            }
+
             final DataLakeServiceClient storageClient = 
getStorageClient(context, flowFile);
-            final DataLakeFileSystemClient dataLakeFileSystemClient = 
storageClient.getFileSystemClient(fileSystem);
-            final DataLakeDirectoryClient directoryClient = 
dataLakeFileSystemClient.getDirectoryClient(directory);
+            final DataLakeFileSystemClient fileSystemClient = 
storageClient.getFileSystemClient(fileSystem);
+            final DataLakeDirectoryClient directoryClient = 
fileSystemClient.getDirectoryClient(directory);
             final DataLakeFileClient fileClient = 
directoryClient.createFile(fileName);
+
             final long length = flowFile.getSize();
             if (length > 0) {
                 try (final InputStream rawIn = session.read(flowFile); final 
BufferedInputStream in = new BufferedInputStream(rawIn)) {
                     fileClient.append(in, 0, length);
-
                 }
             }
             fileClient.flush(length);
+
             final Map<String, String> attributes = new HashMap<>();
             attributes.put("azure.filesystem", fileSystem);
             attributes.put("azure.directory", directory);
@@ -84,7 +94,6 @@ public class PutAzureDataLakeStorage extends 
AbstractAzureDataLakeStorageProcess
             attributes.put("azure.length", String.valueOf(length));
             flowFile = session.putAllAttributes(flowFile, attributes);
 
-
             session.transfer(flowFile, REL_SUCCESS);
             final long transferMillis = 
TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos);
             session.getProvenanceReporter().send(flowFile, 
fileClient.getFileUrl(), transferMillis);
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
new file mode 100644
index 0000000..a9b2bd9
--- /dev/null
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.processors.azure.storage;
+
+import com.azure.storage.common.StorageSharedKeyCredential;
+import com.azure.storage.file.datalake.DataLakeFileSystemClient;
+import com.azure.storage.file.datalake.DataLakeServiceClient;
+import com.azure.storage.file.datalake.DataLakeServiceClientBuilder;
+import org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
+import org.junit.After;
+import org.junit.Before;
+
+import java.util.UUID;
+
+public abstract class AbstractAzureDataLakeStorageIT extends 
AbstractAzureStorageIT {
+
+    private static final String FILESYSTEM_NAME_PREFIX = 
"nifi-test-filesystem";
+
+    protected String fileSystemName;
+    protected DataLakeFileSystemClient fileSystemClient;
+
+    @Before
+    public void setUpAzureDataLakeStorageIT() {
+        fileSystemName = String.format("%s-%s", FILESYSTEM_NAME_PREFIX, 
UUID.randomUUID());
+
+        runner.setProperty(AbstractAzureDataLakeStorageProcessor.FILESYSTEM, 
fileSystemName);
+
+        DataLakeServiceClient storageClient = createStorageClient();
+        fileSystemClient = storageClient.createFileSystem(fileSystemName);
+    }
+
+    @After
+    public void tearDownAzureDataLakeStorageIT() {
+        fileSystemClient.delete();
+    }
+
+    private DataLakeServiceClient createStorageClient() {
+        return new DataLakeServiceClientBuilder()
+                .endpoint("https://"; + getAccountName() + 
".dfs.core.windows.net")
+                .credential(new StorageSharedKeyCredential(getAccountName(), 
getAccountKey()))
+                .buildClient();
+    }
+}
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
index f5fa2e5..7689b96 100644
--- 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
@@ -35,15 +35,16 @@ import java.util.Properties;
 
 import static org.junit.Assert.fail;
 
-public abstract class AbstractAzureStorageIT {private static final Properties 
CONFIG;
+public abstract class AbstractAzureStorageIT {
+
+    private static final Properties CONFIG;
 
     private static final String CREDENTIALS_FILE = 
System.getProperty("user.home") + "/azure-credentials.PROPERTIES";
 
     static {
-        final FileInputStream fis;
         CONFIG = new Properties();
         try {
-            fis = new FileInputStream(CREDENTIALS_FILE);
+            final FileInputStream fis = new FileInputStream(CREDENTIALS_FILE);
             try {
                 CONFIG.load(fis);
             } catch (IOException e) {
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
index e8c0230..29915eb 100644
--- 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
@@ -19,6 +19,7 @@ package org.apache.nifi.processors.azure.storage;
 import org.apache.nifi.processor.Processor;
 import org.apache.nifi.util.MockFlowFile;
 import org.junit.Before;
+import org.junit.Ignore;
 import org.junit.Test;
 
 import java.util.List;
@@ -35,6 +36,7 @@ public class ITDeleteAzureDataLakeStorage extends 
AbstractAzureBlobStorageIT {
         runner.setProperty(DeleteAzureDataLakeStorage.FILE, TEST_FILE_NAME);
     }
 
+    @Ignore
     @Test
     public void testDeleteFile() throws Exception {
         runner.assertValid();
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
index 93a414f..a242ede 100644
--- 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
@@ -19,6 +19,7 @@ package org.apache.nifi.processors.azure.storage;
 import org.apache.nifi.processor.Processor;
 import org.apache.nifi.util.MockFlowFile;
 import org.junit.Before;
+import org.junit.Ignore;
 import org.junit.Test;
 
 import java.util.List;
@@ -36,6 +37,7 @@ public class ITFetchAzureDataLakeStorage extends 
AbstractAzureBlobStorageIT {
 
     }
 
+    @Ignore
     @Test
     public void testFetchFile() throws Exception {
         runner.assertValid();
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
index dba8abd..1a6105f 100644
--- 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
@@ -16,14 +16,32 @@
  */
 package org.apache.nifi.processors.azure.storage;
 
+import com.azure.storage.file.datalake.DataLakeDirectoryClient;
+import com.azure.storage.file.datalake.DataLakeFileClient;
+import com.google.common.net.UrlEscapers;
+import org.apache.commons.lang3.StringUtils;
 import org.apache.nifi.processor.Processor;
 import org.apache.nifi.util.MockFlowFile;
 import org.junit.Before;
+import org.junit.Ignore;
 import org.junit.Test;
 
-import java.util.List;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
 
-public class ITPutAzureDataLakeStorage extends AbstractAzureBlobStorageIT {
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+public class ITPutAzureDataLakeStorage extends AbstractAzureDataLakeStorageIT {
+
+    private static final String DIRECTORY = "dir1";
+    private static final String FILE_NAME = "file1";
+    private static final byte[] FILE_DATA = "0123456789".getBytes();
+
+    private static final String EL_FILESYSTEM = "az.filesystem";
+    private static final String EL_DIRECTORY = "az.directory";
+    private static final String EL_FILE_NAME = "az.filename";
 
     @Override
     protected Class<? extends Processor> getProcessorClass() {
@@ -32,26 +50,219 @@ public class ITPutAzureDataLakeStorage extends 
AbstractAzureBlobStorageIT {
 
     @Before
     public void setUp() {
-        runner.setProperty(PutAzureDataLakeStorage.FILE, TEST_FILE_NAME);
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, DIRECTORY);
+        runner.setProperty(PutAzureDataLakeStorage.FILE, FILE_NAME);
+    }
+
+    @Test
+    public void testPutFileToExistingDirectory() throws Exception {
+        fileSystemClient.createDirectory(DIRECTORY);
+
+        runProcessor(FILE_DATA);
+
+        assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+    }
+
+    @Test
+    public void testPutFileToNonExistingDirectory() throws Exception {
+        runProcessor(FILE_DATA);
+
+        assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+    }
+
+    @Test
+    public void testPutFileToDeepDirectory() throws Exception {
+        String baseDirectory = "dir1/dir2";
+        String fullDirectory = baseDirectory + "/dir3/dir4";
+        fileSystemClient.createDirectory(baseDirectory);
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, fullDirectory);
+
+        runProcessor(FILE_DATA);
+
+        assertSuccess(fullDirectory, FILE_NAME, FILE_DATA);
+    }
+
+    @Test
+    public void testPutFileToRootDirectory() throws Exception {
+        String rootDirectory = "";
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, rootDirectory);
+
+        runProcessor(FILE_DATA);
+
+        assertSuccess(rootDirectory, FILE_NAME, FILE_DATA);
+    }
+
+    @Test
+    public void testPutEmptyFile() throws Exception {
+        byte[] fileData = new byte[0];
+
+        runProcessor(fileData);
+
+        assertSuccess(DIRECTORY, FILE_NAME, fileData);
+    }
+
+    @Ignore
+    // ignore excessive test with larger file size
+    @Test
+    public void testPutBigFile() throws Exception {
+        byte[] fileData = new byte[100_000_000];
+
+        runProcessor(fileData);
+
+        assertSuccess(DIRECTORY, FILE_NAME, fileData);
     }
 
     @Test
-    public void testPutFile() throws Exception {
+    public void testPutFileWithNonExistingFileSystem() {
+        runner.setProperty(PutAzureDataLakeStorage.FILESYSTEM, "dummy");
+
+        runProcessor(FILE_DATA);
+
+        assertFailure();
+    }
+
+    @Test
+    public void testPutFileWithInvalidDirectory() {
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, "/dir1");
+
+        runProcessor(FILE_DATA);
+
+        assertFailure();
+    }
+
+    @Test
+    public void testPutFileWithInvalidFileName() {
+        runner.setProperty(PutAzureDataLakeStorage.FILE, "/file1");
+
+        runProcessor(FILE_DATA);
+
+        assertFailure();
+    }
+
+    @Test
+    public void testPutFileWithSpacesInDirectoryAndFileName() throws Exception 
{
+        String directory = "dir 1";
+        String fileName = "file 1";
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, directory);
+        runner.setProperty(PutAzureDataLakeStorage.FILE, fileName);
+
+        runProcessor(FILE_DATA);
+
+        assertSuccess(directory, fileName, FILE_DATA);
+    }
+
+    @Ignore
+    // the existing file gets overwritten without error
+    // seems to be a bug in the Azure lib
+    @Test
+    public void testPutFileToExistingFile() {
+        fileSystemClient.createFile(String.format("%s/%s", DIRECTORY, 
FILE_NAME));
+
+        runProcessor(FILE_DATA);
+
+        assertFailure();
+    }
+
+    @Test
+    public void testPutFileWithEL() throws Exception {
+        Map<String, String> attributes = createAttributesMap();
+        setELProperties();
+
+        runProcessor(FILE_DATA, attributes);
+
+        assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+    }
+
+    @Test
+    public void testPutFileWithELButFilesystemIsNotSpecified() {
+        Map<String, String> attributes = createAttributesMap();
+        attributes.remove(EL_FILESYSTEM);
+        setELProperties();
+
+        runProcessor(FILE_DATA, attributes);
+
+        assertFailure();
+    }
+
+    @Test
+    public void testPutFileWithELButFileNameIsNotSpecified() {
+        Map<String, String> attributes = createAttributesMap();
+        attributes.remove(EL_FILE_NAME);
+        setELProperties();
+
+        runProcessor(FILE_DATA, attributes);
+
+        assertFailure();
+    }
+
+    private Map<String, String> createAttributesMap() {
+        Map<String, String> attributes = new HashMap<>();
+
+        attributes.put(EL_FILESYSTEM, fileSystemName);
+        attributes.put(EL_DIRECTORY, DIRECTORY);
+        attributes.put(EL_FILE_NAME, FILE_NAME);
+
+        return attributes;
+    }
+
+    private void setELProperties() {
+        runner.setProperty(PutAzureDataLakeStorage.FILESYSTEM, 
String.format("${%s}", EL_FILESYSTEM));
+        runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, 
String.format("${%s}", EL_DIRECTORY));
+        runner.setProperty(PutAzureDataLakeStorage.FILE, 
String.format("${%s}", EL_FILE_NAME));
+    }
+
+    private void runProcessor(byte[] fileData) {
+        runProcessor(fileData, Collections.emptyMap());
+    }
+
+    private void runProcessor(byte[] testData, Map<String, String> attributes) 
{
         runner.assertValid();
-        runner.enqueue("0123456789".getBytes());
+        runner.enqueue(testData, attributes);
         runner.run();
-
-        assertResult();
     }
 
+    private void assertSuccess(String directory, String fileName, byte[] 
fileData) throws Exception {
+        assertFlowFile(directory, fileName, fileData);
+        assertAzureFile(directory, fileName, fileData);
+    }
 
-    private void assertResult() throws Exception {
+    private void assertFlowFile(String directory, String fileName, byte[] 
fileData) throws Exception {
         
runner.assertAllFlowFilesTransferred(PutAzureDataLakeStorage.REL_SUCCESS, 1);
-        List<MockFlowFile> flowFilesForRelationship = 
runner.getFlowFilesForRelationship(PutAzureDataLakeStorage.REL_SUCCESS);
-        for (MockFlowFile flowFile : flowFilesForRelationship) {
-            flowFile.assertContentEquals("0123456789".getBytes());
-            flowFile.assertAttributeEquals("azure.length", "10");
+
+        MockFlowFile flowFile = 
runner.getFlowFilesForRelationship(PutAzureDataLakeStorage.REL_SUCCESS).get(0);
+
+        flowFile.assertContentEquals(fileData);
+
+        flowFile.assertAttributeEquals("azure.filesystem", fileSystemName);
+        flowFile.assertAttributeEquals("azure.directory", directory);
+        flowFile.assertAttributeEquals("azure.filename", fileName);
+
+        String urlEscapedDirectory = 
UrlEscapers.urlPathSegmentEscaper().escape(directory);
+        String urlEscapedFileName = 
UrlEscapers.urlPathSegmentEscaper().escape(fileName);
+        String primaryUri = StringUtils.isNotEmpty(directory)
+                ? String.format("https://%s.dfs.core.windows.net/%s/%s/%s";, 
getAccountName(), fileSystemName, urlEscapedDirectory, urlEscapedFileName)
+                : String.format("https://%s.dfs.core.windows.net/%s/%s";, 
getAccountName(), fileSystemName, urlEscapedFileName);
+        flowFile.assertAttributeEquals("azure.primaryUri", primaryUri);
+
+        flowFile.assertAttributeEquals("azure.length", 
Integer.toString(fileData.length));
+    }
+
+    private void assertAzureFile(String directory, String fileName, byte[] 
fileData) {
+        DataLakeFileClient fileClient;
+        if (StringUtils.isNotEmpty(directory)) {
+            DataLakeDirectoryClient directoryClient = 
fileSystemClient.getDirectoryClient(directory);
+            assertTrue(directoryClient.exists());
+
+            fileClient = directoryClient.getFileClient(fileName);
+        } else {
+            fileClient = fileSystemClient.getFileClient(fileName);
         }
 
+        assertTrue(fileClient.exists());
+        assertEquals(fileData.length, 
fileClient.getProperties().getFileSize());
+    }
+
+    private void assertFailure() {
+        
runner.assertAllFlowFilesTransferred(PutAzureDataLakeStorage.REL_FAILURE, 1);
     }
 }
diff --git 
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
new file mode 100644
index 0000000..59800bb
--- /dev/null
+++ 
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
@@ -0,0 +1,124 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.processors.azure.storage;
+
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.Before;
+import org.junit.Test;
+
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.ACCOUNT_KEY;
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.ACCOUNT_NAME;
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.DIRECTORY;
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.FILE;
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.FILESYSTEM;
+import static 
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.SAS_TOKEN;
+
+public class TestAbstractAzureDataLakeStorage {
+
+    private TestRunner runner;
+
+    @Before
+    public void setUp() {
+        // test the property validation in the abstract class via the put 
processor
+        runner = TestRunners.newTestRunner(PutAzureDataLakeStorage.class);
+
+        runner.setProperty(ACCOUNT_NAME, "accountName");
+        runner.setProperty(ACCOUNT_KEY, "accountKey");
+        runner.setProperty(FILESYSTEM, "filesystem");
+        runner.setProperty(DIRECTORY, "directory");
+        runner.setProperty(FILE, "file");
+    }
+
+    @Test
+    public void testValidWhenAccountNameAndAccountKeySpecified() {
+        runner.assertValid();
+    }
+
+    @Test
+    public void testValidWhenAccountNameAndSasTokenSpecified() {
+        runner.removeProperty(ACCOUNT_KEY);
+        runner.setProperty(SAS_TOKEN, "sasToken");
+
+        runner.assertValid();
+    }
+
+    @Test
+    public void testNotValidWhenNoAccountNameSpecified() {
+        runner.removeProperty(ACCOUNT_NAME);
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testNotValidWhenNoAccountKeyNorSasTokenSpecified() {
+        runner.removeProperty(ACCOUNT_KEY);
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testNotValidWhenBothAccountKeyAndSasTokenSpecified() {
+        runner.setProperty(SAS_TOKEN, "sasToken");
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testNotValidWhenNoFilesystemSpecified() {
+        runner.removeProperty(FILESYSTEM);
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testNotValidWhenFilesystemIsEmptyString() {
+        runner.setProperty(FILESYSTEM, "");
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testNotValidWhenNoDirectorySpecified() {
+        runner.removeProperty(DIRECTORY);
+
+        runner.assertNotValid();
+    }
+
+    @Test
+    public void testValidWhenDirectoryIsEmptyString() {
+        // the empty string is for the filesystem root directory
+        runner.setProperty(DIRECTORY, "");
+
+        runner.assertValid();
+    }
+
+    @Test
+    public void testValidWhenNoFileSpecified() {
+        // the default value will be used
+        runner.removeProperty(FILE);
+
+        runner.assertValid();
+    }
+
+    @Test
+    public void testNotValidWhenFileIsEmptyString() {
+        runner.setProperty(FILE, "");
+
+        runner.assertNotValid();
+    }
+}

Reply via email to