This is an automated email from the ASF dual-hosted git repository.
markap14 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new a62c856388f NIFI-16098 Prevent content-claim truncation from dropping
content referenced by a live sibling FlowFile (#11414)
a62c856388f is described below
commit a62c856388f32984a4a7795324d81bd0861f7276
Author: Pierre Villard <[email protected]>
AuthorDate: Thu Jul 23 15:32:50 2026 +0200
NIFI-16098 Prevent content-claim truncation from dropping content
referenced by a live sibling FlowFile (#11414)
* NIFI-16098 Prevent content-claim truncation from dropping content
referenced by a live sibling FlowFile
A FlowFile created and removed within a single session (e.g.
ExecuteGroovyScript
with Failure Strategy 'transfer to failure', which clones on get and
removes the
clone at commit) yields one repository record whose type transitions
CREATE->DELETE.
The truncation reference count is incremented only for a committed CREATE,
so that
increment never ran, but the DELETE still decremented it. For a content
claim shared
with a live sibling FlowFile this drove the truncation reference count to
zero,
allowing the claim to be truncated and the sibling's content to be lost.
Only decrement the truncation reference count on delete of a persisted
FlowFile
(one that has an original claim).
* NIFI-16098 Add system test reproducing content-claim truncation loss with
an in-session cloned FlowFile
Adds a CloneAndTerminate test processor that, for each incoming FlowFile,
clones it and removes both the
clone and the input within the same session (the record set produced by
ExecuteGroovyScript with Failure
Strategy 'transfer to failure'). A new ContentClaimTruncationIT case forks
a generator to this processor
and to a leg that holds the FlowFiles, and verifies the shared content
claims are not truncated while the
holding leg still references them. The test fails without the accompanying
framework fix (content is
truncated away) and passes with it.
---
.../repository/WriteAheadFlowFileRepository.java | 7 +-
.../TestWriteAheadFlowFileRepository.java | 74 ++++++++++++++++++++++
.../processors/tests/system/CloneAndTerminate.java | 73 +++++++++++++++++++++
.../services/org.apache.nifi.processor.Processor | 1 +
.../repositories/ContentClaimTruncationIT.java | 49 ++++++++++++++
5 files changed, 203 insertions(+), 1 deletion(-)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/WriteAheadFlowFileRepository.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/WriteAheadFlowFileRepository.java
index b8ac4d50af9..96345a89cd5 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/WriteAheadFlowFileRepository.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/WriteAheadFlowFileRepository.java
@@ -621,7 +621,12 @@ public class WriteAheadFlowFileRepository implements
FlowFileRepository, SyncLis
case CONTENTMISSING:
if (record.isContentModified()) {
decrementContentClaimReference(record.getOriginalClaim());
- } else {
+ } else if (record.getOriginalClaim() != null) {
+ // Only decrement when the FlowFile was persisted, i.e.
its CREATE incremented the count.
+ // A FlowFile created and removed within a single session
yields one record whose type
+ // transitions from CREATE to DELETE and which has no
original claim; that CREATE increment
+ // never ran, so decrementing here would corrupt the count
for a claim that is still
+ // referenced by live sibling FlowFiles, allowing their
content to be truncated away.
decrementContentClaimReference(currentClaim);
}
break;
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestWriteAheadFlowFileRepository.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestWriteAheadFlowFileRepository.java
index af5d92b3956..984d016daf6 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestWriteAheadFlowFileRepository.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestWriteAheadFlowFileRepository.java
@@ -1214,6 +1214,80 @@ public class TestWriteAheadFlowFileRepository {
"Shared Content Claim must not be queued for truncation while
live siblings still reference it");
}
+ @Test
+ public void
testInSessionCreatedAndRemovedCloneDoesNotDecrementSharedTruncationReference()
throws IOException {
+ // A session that creates a clone of a FlowFile (sharing its
truncation-eligible Content Claim) and
+ // removes that clone within the same session produces a single record
whose type transitions from
+ // CREATE to DELETE and which has no original claim. The CREATE
increment never runs for such a record,
+ // so its DELETE must not decrement the shared claim's truncation
reference count. Otherwise a claim
+ // that is still referenced by a live sibling FlowFile is queued for
truncation and its content is lost.
+ // This mirrors ExecuteGroovyScript with "transfer to failure", where
every get() clones the FlowFile
+ // and commitAsync() removes the clone.
+ final RuntimeRepoContext context = createRuntimeRepoContext();
+
+ final ResourceClaim resourceClaim =
context.claimManager().newResourceClaim("container", "section", "1", false,
false);
+ // One claimant reference per live FlowFile (leg A, leg B) plus the
in-session clone, so the resource
+ // claim stays referenced (not destructable) after the deletes and the
truncation behavior is isolated.
+ context.claimManager().incrementClaimantCount(resourceClaim);
+ context.claimManager().incrementClaimantCount(resourceClaim);
+ context.claimManager().incrementClaimantCount(resourceClaim);
+ final StandardContentClaim sharedClaim = createClaim(resourceClaim,
1024L, TRUNCATION_CANDIDATE_LENGTH, true);
+
+ final FlowFileRecord legA = new StandardFlowFileRecord.Builder()
+ .id(1L)
+ .addAttribute("uuid", UUID.randomUUID().toString())
+ .contentClaim(sharedClaim)
+ .build();
+ final FlowFileRecord legB = new StandardFlowFileRecord.Builder()
+ .id(2L)
+ .addAttribute("uuid", UUID.randomUUID().toString())
+ .contentClaim(sharedClaim)
+ .build();
+
+ try (final WriteAheadFlowFileRepository repo = new
WriteAheadFlowFileRepository(niFiProperties)) {
+ repo.initialize(context.claimManager());
+ repo.loadFlowFiles(context.queueProvider());
+
+ // Two committed CREATE records for the shared claim -> truncation
reference count = 2.
+ final List<RepositoryRecord> createRecords = new ArrayList<>();
+ for (final FlowFileRecord flowFile : List.of(legA, legB)) {
+ final StandardRepositoryRecord createRecord = new
StandardRepositoryRecord(context.queue());
+ createRecord.setWorking(flowFile, false);
+ createRecord.setDestination(context.queue());
+ createRecords.add(createRecord);
+ }
+ repo.updateRepository(createRecords);
+ assertEquals(2, repo.getContentClaimReferenceCount(sharedClaim));
+
+ // A single commit that deletes leg A (a persisted FlowFile, has
an original claim) and also removes
+ // an in-session clone that shares the same claim. The clone was
created and removed within this
+ // session: one record, no original claim, CREATE -> DELETE.
+ final StandardRepositoryRecord deleteLegA = new
StandardRepositoryRecord(context.queue(), legA);
+ deleteLegA.markForDelete();
+
+ final FlowFileRecord clone = new StandardFlowFileRecord.Builder()
+ .id(3L)
+ .addAttribute("uuid", UUID.randomUUID().toString())
+ .contentClaim(sharedClaim)
+ .build();
+ final StandardRepositoryRecord createdAndRemovedClone = new
StandardRepositoryRecord(context.queue());
+ createdAndRemovedClone.setWorking(clone, false);
+ createdAndRemovedClone.markForDelete();
+
+ repo.updateRepository(List.of(deleteLegA, createdAndRemovedClone));
+ repo.checkpoint();
+ }
+
+ // Leg B is still live, so exactly one truncation reference must
remain for the shared claim.
+ assertEquals(1,
context.claimManager().getTruncationReferenceCount(sharedClaim),
+ "Deleting one live sibling plus an in-session create+remove
clone must leave one reference for the still-live sibling");
+
+ final List<ContentClaim> truncated = new ArrayList<>();
+ context.claimManager().drainTruncatableClaims(truncated, 100);
+ assertFalse(truncated.contains(sharedClaim),
+ "Shared Content Claim must not be queued for truncation while
a live sibling FlowFile still references it");
+ }
+
//
=========================================================================
// Truncation Feature: Recovery Tests
//
=========================================================================
diff --git
a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/CloneAndTerminate.java
b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/CloneAndTerminate.java
new file mode 100644
index 00000000000..cce72005fe4
--- /dev/null
+++
b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/CloneAndTerminate.java
@@ -0,0 +1,73 @@
+/*
+ * 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.components.PropertyDescriptor;
+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.exception.ProcessException;
+import org.apache.nifi.processor.util.StandardValidators;
+
+import java.io.File;
+import java.util.List;
+
+/**
+ * For each incoming FlowFile, clones it and then removes both the clone and
the input within the same session.
+ * The clone shares the input's content claim, so this produces one commit
containing a DELETE record for the
+ * input (which has an original claim) and a single CREATE-then-DELETE record
for the clone (which has no original
+ * claim, since it was created and removed in the same session).
+ *
+ * This mirrors the record set produced by ExecuteGroovyScript with Failure
Strategy "transfer to failure", which
+ * clones every FlowFile it reads so it can route the clone to the failure
relationship on error, and removes those
+ * clones when the script completes successfully. It is used to verify that
removing such an in-session clone does
+ * not queue a shared content claim for truncation while another FlowFile
still references it.
+ */
+public class CloneAndTerminate extends AbstractProcessor {
+
+ public static final PropertyDescriptor GATE_FILE = new
PropertyDescriptor.Builder()
+ .name("Gate File")
+ .description("An optional file path. If specified, the processor
will only process FlowFiles when this file exists. " +
+ "If the file does not exist, the processor will yield and
return without processing any data.")
+ .required(false)
+ .addValidator(StandardValidators.NON_BLANK_VALIDATOR)
+ .build();
+
+ @Override
+ protected List<PropertyDescriptor> getSupportedPropertyDescriptors() {
+ return List.of(GATE_FILE);
+ }
+
+ @Override
+ public void onTrigger(final ProcessContext context, final ProcessSession
session) throws ProcessException {
+ final String gateFilePath = context.getProperty(GATE_FILE).getValue();
+ if (gateFilePath != null && !new File(gateFilePath).exists()) {
+ context.yield();
+ return;
+ }
+
+ final FlowFile input = session.get();
+ if (input == null) {
+ return;
+ }
+
+ final FlowFile clone = session.clone(input);
+ session.remove(clone);
+ session.remove(input);
+ }
+}
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 5c1de275383..6bbf6f992f8 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
@@ -15,6 +15,7 @@
org.apache.nifi.processors.tests.system.AssetReadingProcessor
org.apache.nifi.processors.tests.system.ClassloaderIsolationWithServiceProperty
+org.apache.nifi.processors.tests.system.CloneAndTerminate
org.apache.nifi.processors.tests.system.CountEvents
org.apache.nifi.processors.tests.system.CountFlowFiles
org.apache.nifi.processors.tests.system.ConcatenateFlowFiles
diff --git
a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentClaimTruncationIT.java
b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentClaimTruncationIT.java
index 66b666c36c0..ffa27f1696c 100644
---
a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentClaimTruncationIT.java
+++
b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/repositories/ContentClaimTruncationIT.java
@@ -198,6 +198,55 @@ public class ContentClaimTruncationIT extends NiFiSystemIT
{
});
}
+ // Reproduces content loss when a downstream session clones a FlowFile and
removes the clone within the same
+ // session, as ExecuteGroovyScript with Failure Strategy "transfer to
failure" does. The clone shares the
+ // FlowFile's content claim, and its create-then-remove within one session
must not decrement the claim's
+ // truncation reference count. Draining the leg that performs the
clone-and-remove must not truncate content
+ // claims that are still referenced by FlowFiles held on the other leg.
+ @Test
+ public void
testInSessionClonedAndRemovedFlowFileDoesNotTruncateSharedClaim() throws
Exception {
+ final ProcessorEntity generator =
getClientUtil().createProcessor("GenerateTruncatableFlowFiles");
+ final ProcessorEntity cloneAndTerminate =
getClientUtil().createProcessor("CloneAndTerminate");
+ final ProcessorEntity holdTerminate =
getClientUtil().createProcessor("TerminateFlowFile");
+ getClientUtil().updateProcessorProperties(generator,
GENERATE_TRUNCATABLE_PROPS);
+ getClientUtil().updateProcessorSchedulingPeriod(generator, "0 sec");
+
+ ConnectionEntity cloneConnection =
getClientUtil().createConnection(generator, cloneAndTerminate, "success");
+ ConnectionEntity holdConnection =
getClientUtil().createConnection(generator, holdTerminate, "success");
+ cloneConnection =
getClientUtil().updateConnectionBackpressure(cloneConnection, 10000,
BACKPRESSURE_BYTES);
+ holdConnection =
getClientUtil().updateConnectionBackpressure(holdConnection, 10000,
BACKPRESSURE_BYTES);
+
+ getClientUtil().startProcessor(generator);
+ waitForQueueCount(cloneConnection.getId(), 100);
+ waitForQueueCount(holdConnection.getId(), 100);
+
+ getClientUtil().stopProcessor(generator);
+ getClientUtil().waitForStoppedProcessor(generator.getId());
+
+ final File contentRepositoryDirectory = new
File(getNiFiInstance().getInstanceDirectory(), "content_repository");
+
+ // Drain the leg that clones each FlowFile and removes the clone
within the same session.
+ drainTerminateQueue(cloneAndTerminate, cloneConnection.getId());
+ // Allow several truncation cleanup cycles to run; if the shared
claims were incorrectly queued for
+ // truncation, the content repository would shrink well below the
retained size within this window.
+ Thread.sleep(5000);
+
+ // The other leg still holds all FlowFiles, so their (shared) content
claims must not be truncated.
+ final long sizeAfterCloneLegDrained =
getContentRepositorySize(contentRepositoryDirectory.toPath());
+ assertTrue(sizeAfterCloneLegDrained >= NOT_TRUNCATED_MIN_BYTES,
+ "Content must not be truncated while the other connection
still references the same claims; size was " + sizeAfterCloneLegDrained);
+
+ // Once the holding leg is drained too, all content should be cleaned
up.
+ drainTerminateQueue(holdTerminate, holdConnection.getId());
+ waitFor(() -> {
+ try {
+ return
getContentRepositorySize(contentRepositoryDirectory.toPath()) == 0;
+ } catch (final IOException e) {
+ return false;
+ }
+ });
+ }
+
// Same scenario as
testClonedSuccessNoTruncationUntilBothConnectionsDrained, but the second
connection uses
// PriorityAttributePrioritizer. After the first connection is drained,
only the large FlowFiles are removed
// from the second connection so that tail truncation of content
repository files can occur.