This is an automated email from the ASF dual-hosted git repository. exceptionfactory pushed a commit to branch NIFI-11972-RC1 in repository https://gitbox.apache.org/repos/asf/nifi.git
commit 6e9ad730524fe442427d23f254e158e863c7df75 Author: Mark Payne <[email protected]> AuthorDate: Sun Aug 20 16:26:51 2023 -0400 NIFI-11971: Ensure that if we do not write any bytes to a file after calling ProcessSession.write() that we set content claim's length to 0 when closing OutputStream; otherwise it remains -1, which causes issues, since the length will later be added to the offset to determine the position in the ste stream This closes #7629 Signed-off-by: David Handermann <[email protected]> (cherry picked from commit 5e2a10c32834973ba2837c3bdf5963aab2aa92bc) --- .../claim/StandardContentClaimWriteCache.java | 5 ++ .../nifi/processors/tests/system/IngestFile.java | 22 ++++- .../processors/tests/system/UnzipFlowFile.java | 94 +++++++++++++++++++++ .../services/org.apache.nifi.processor.Processor | 3 +- .../tests/system/repositories/ContentAccessIT.java | 37 ++++++++ .../test/resources/empty-and-small-text-files.zip | Bin 0 -> 316 bytes 6 files changed, 157 insertions(+), 4 deletions(-) diff --git a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/claim/StandardContentClaimWriteCache.java b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/claim/StandardContentClaimWriteCache.java index 1a728071c2..8d78953297 100644 --- a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/claim/StandardContentClaimWriteCache.java +++ b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/claim/StandardContentClaimWriteCache.java @@ -149,6 +149,11 @@ public class StandardContentClaimWriteCache implements ContentClaimWriteCache { } closed = true; + if (scc.getLength() < 0) { + // If claim was not written to, set length to 0 + scc.setLength(0L); + } + queue.offer(claim); } }; diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/IngestFile.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/IngestFile.java index 4637987216..844aed1788 100644 --- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/IngestFile.java +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/IngestFile.java @@ -56,6 +56,13 @@ public class IngestFile extends AbstractProcessor { .allowableValues(COMMIT_ASYNC, COMMIT_SYNCHRONOUS) .defaultValue(COMMIT_ASYNC) .build(); + static final PropertyDescriptor DELETE_FILE = new PropertyDescriptor.Builder() + .name("Delete File") + .description("Whether or not the file should be deleted after successfully ingesting") + .allowableValues("true", "false") + .defaultValue("true") + .required(true) + .build(); static final Relationship REL_SUCCESS = new Relationship.Builder() .name("success") @@ -63,7 +70,7 @@ public class IngestFile extends AbstractProcessor { @Override protected List<PropertyDescriptor> getSupportedPropertyDescriptors() { - return Arrays.asList(FILENAME, COMMIT_MODE); + return Arrays.asList(FILENAME, COMMIT_MODE, DELETE_FILE); } @Override @@ -81,12 +88,21 @@ public class IngestFile extends AbstractProcessor { session.transfer(flowFile, REL_SUCCESS); session.getProvenanceReporter().receive(flowFile, file.toURI().toString()); + final boolean deleteFile = context.getProperty(DELETE_FILE).asBoolean(); + final String commitMode = context.getProperty(COMMIT_MODE).getValue(); if (COMMIT_SYNCHRONOUS.equalsIgnoreCase(commitMode)) { session.commit(); - cleanup(file); + + if (deleteFile) { + cleanup(file); + } } else { - session.commitAsync(() -> cleanup(file)); + session.commitAsync(() -> { + if (deleteFile) { + cleanup(file); + } + }); } } diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java new file mode 100644 index 0000000000..b522c9b97a --- /dev/null +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java @@ -0,0 +1,94 @@ +/* + * 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.tests.system; + +import org.apache.nifi.flowfile.FlowFile; +import org.apache.nifi.processor.AbstractProcessor; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.processor.ProcessSession; +import org.apache.nifi.processor.Relationship; +import org.apache.nifi.processor.exception.ProcessException; +import org.apache.nifi.stream.io.StreamUtils; + +import java.io.IOException; +import java.io.InputStream; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.zip.ZipEntry; +import java.util.zip.ZipInputStream; + +public class UnzipFlowFile extends AbstractProcessor { + + private final Relationship REL_UNZIPPED = new Relationship.Builder() + .name("unzipped") + .build(); + + private final Relationship REL_ORIGINAL = new Relationship.Builder() + .name("original") + .autoTerminateDefault(true) + .build(); + + private final Relationship REL_FAILURE = new Relationship.Builder() + .name("failure") + .build(); + + @Override + public Set<Relationship> getRelationships() { + return new HashSet<>(Arrays.asList(REL_UNZIPPED, REL_ORIGINAL, REL_FAILURE)); + } + + @Override + public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException { + FlowFile flowFile = session.get(); + if (flowFile == null) { + return; + } + + final List<FlowFile> created = new ArrayList<>(); + try (final InputStream in = session.read(flowFile); + final ZipInputStream zipInputStream = new ZipInputStream(in)) { + + while (true) { + final ZipEntry zipEntry = zipInputStream.getNextEntry(); + if (zipEntry == null) { + break; + } + + final String filename = zipEntry.getName(); + FlowFile outFile = session.create(flowFile); + outFile = session.putAttribute(outFile, "filename", filename); + created.add(outFile); + + session.write(outFile, out -> { + StreamUtils.copy(zipInputStream, out); + }); + } + + } catch (final IOException e) { + getLogger().error("Failed to unzip {}", flowFile, e); + session.transfer(flowFile, REL_FAILURE); + session.remove(created); + return; + } + + session.transfer(created, REL_UNZIPPED); + session.transfer(flowFile, REL_ORIGINAL); + } +} diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor index 1f278c32ac..781dcceb8e 100644 --- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor @@ -40,13 +40,14 @@ org.apache.nifi.processors.tests.system.SensitiveDynamicPropertiesProcessor org.apache.nifi.processors.tests.system.SetAttribute org.apache.nifi.processors.tests.system.Sleep org.apache.nifi.processors.tests.system.SplitByLine +org.apache.nifi.processors.tests.system.SplitTextByLine org.apache.nifi.processors.tests.system.TerminateFlowFile org.apache.nifi.processors.tests.system.TransferBatch org.apache.nifi.processors.tests.system.ThrowProcessException org.apache.nifi.processors.tests.system.UpdateContent +org.apache.nifi.processors.tests.system.UnzipFlowFile org.apache.nifi.processors.tests.system.ValidateFileExists org.apache.nifi.processors.tests.system.VerifyContents org.apache.nifi.processors.tests.system.WriteFlowFileCountToFile org.apache.nifi.processors.tests.system.WriteLifecycleEvents org.apache.nifi.processors.tests.system.WriteToFile -org.apache.nifi.processors.tests.system.SplitTextByLine diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentAccessIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentAccessIT.java index b6ee614eb8..1ae0788243 100644 --- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentAccessIT.java +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentAccessIT.java @@ -20,9 +20,11 @@ package org.apache.nifi.tests.system.repositories; import org.apache.nifi.tests.system.NiFiSystemIT; import org.apache.nifi.toolkit.cli.impl.client.nifi.NiFiClientException; import org.apache.nifi.web.api.entity.ConnectionEntity; +import org.apache.nifi.web.api.entity.FlowFileEntity; import org.apache.nifi.web.api.entity.ProcessorEntity; import org.junit.jupiter.api.Test; +import java.io.File; import java.io.IOException; import java.util.Collections; import java.util.HashMap; @@ -30,6 +32,7 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -40,6 +43,40 @@ import static org.junit.jupiter.api.Assertions.assertTrue; */ public class ContentAccessIT extends NiFiSystemIT { + @Test + public void testWriteZeroBytesToFlowFileThenWriteBytes() throws NiFiClientException, IOException, InterruptedException { + final ProcessorEntity ingestFile = getClientUtil().createProcessor("IngestFile"); + final ProcessorEntity unzipFile = getClientUtil().createProcessor("UnzipFlowFile"); + final ProcessorEntity terminate = getClientUtil().createProcessor("TerminateFlowFile"); + + getClientUtil().createConnection(ingestFile, unzipFile, "success"); + final ConnectionEntity unzipped = getClientUtil().createConnection(unzipFile, terminate, "unzipped"); + + getClientUtil().setAutoTerminatedRelationships(unzipFile, "failure"); + + final Map<String, String> ingestProperties = new HashMap<>(); + ingestProperties.put("Filename", new File("src/test/resources/empty-and-small-text-files.zip").getAbsolutePath()); + ingestProperties.put("Delete File", "false"); + getClientUtil().updateProcessorProperties(ingestFile, ingestProperties); + getClientUtil().updateProcessorSchedulingPeriod(ingestFile, "10 mins"); + getClientUtil().waitForValidProcessor(ingestFile.getId()); + + getClientUtil().startProcessor(ingestFile); + getClientUtil().startProcessor(unzipFile); + + waitForQueueCount(unzipped.getId(), 2); + + final FlowFileEntity flowFileEntity1 = getClientUtil().getQueueFlowFile(unzipped.getId(), 0); + assertEquals("a.txt", flowFileEntity1.getFlowFile().getAttributes().get("filename")); + assertEquals(0, flowFileEntity1.getFlowFile().getSize()); + + final FlowFileEntity flowFileEntity2 = getClientUtil().getQueueFlowFile(unzipped.getId(), 1); + assertEquals("b.txt", flowFileEntity2.getFlowFile().getAttributes().get("filename")); + + final String content2 = getClientUtil().getFlowFileContentAsUtf8(unzipped.getId(), 1); + assertEquals("Hello there! b.txt", content2); + } + @Test public void testCorrectContentReadWhenMultipleFlowFilesInClaimWithBatchAndWrite() throws NiFiClientException, IOException, InterruptedException { testCorrectContentReadWhenMultipleFlowFilesInClaim(true, false); diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/resources/empty-and-small-text-files.zip b/nifi-system-tests/nifi-system-test-suite/src/test/resources/empty-and-small-text-files.zip new file mode 100644 index 0000000000..cce7330ccf Binary files /dev/null and b/nifi-system-tests/nifi-system-test-suite/src/test/resources/empty-and-small-text-files.zip differ
