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

exceptionfactory 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 03134bdcddd NIFI-15989 Improved Purge Threshold time formatting 
(#11313)
03134bdcddd is described below

commit 03134bdcdddf0e55f4514175f704405ca83c67ef
Author: David Young <[email protected]>
AuthorDate: Sat Jul 25 13:46:51 2026 -0400

    NIFI-15989 Improved Purge Threshold time formatting (#11313)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../java/org/apache/nifi/util/FormatUtils.java     | 55 ++++++++++++++++++++++
 .../java/org/apache/nifi/util/TestFormatUtils.java | 23 +++++++++
 .../nifi/provenance/store/EventStorePartition.java |  4 +-
 .../provenance/store/PartitionedEventStore.java    |  3 +-
 .../provenance/store/WriteAheadStorePartition.java | 16 +++++--
 5 files changed, 93 insertions(+), 8 deletions(-)

diff --git 
a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java 
b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java
index 5fda79cdc3f..3656473f303 100644
--- 
a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java
+++ 
b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java
@@ -20,6 +20,7 @@ import org.apache.nifi.processor.DataUnit;
 import org.apache.nifi.time.DurationFormat;
 
 import java.text.NumberFormat;
+import java.time.Duration;
 import java.time.Instant;
 import java.time.LocalDateTime;
 import java.time.OffsetDateTime;
@@ -28,8 +29,12 @@ import java.time.ZoneOffset;
 import java.time.format.DateTimeFormatter;
 import java.time.format.DateTimeFormatterBuilder;
 import java.time.temporal.ChronoField;
+import java.time.temporal.ChronoUnit;
 import java.time.temporal.TemporalAccessor;
 import java.time.temporal.TemporalQueries;
+import java.time.temporal.UnsupportedTemporalTypeException;
+import java.util.ArrayList;
+import java.util.List;
 import java.util.Locale;
 import java.util.concurrent.TimeUnit;
 import java.util.regex.Pattern;
@@ -298,4 +303,54 @@ public class FormatUtils {
         final OffsetDateTime offsetDateTime = OffsetDateTime.of(localDateTime, 
zoneOffset);
         return offsetDateTime.toInstant();
     }
+
+    /**
+     * Format a value of ChronoUnit into a word representation.
+     * Does not handle estimated ChronoUnit values.
+     * For anything greater than {@link ChronoUnit#DAYS}, use {@link 
#formatDurationToWords(Duration)}.
+     *
+     * @param value the number of the given {@link ChronoUnit} to measure
+     * @param unit {@link ChronoUnit} that the value represents
+     * @return String representation of the given value and unit
+     * @throws UnsupportedTemporalTypeException if the unit is not supported 
({@link ChronoUnit#isDurationEstimated()} == true)
+     */
+    public static String formatDurationToWords(final long value, final 
ChronoUnit unit) throws UnsupportedTemporalTypeException {
+        return formatDurationToWords(Duration.ZERO.plus(value, unit));
+    }
+
+    /**
+     * Format a Duration using words (days, hours, minutes, seconds, ns) where 
all lower units are
+     *   included once a non-zero unit is found. Unit plurality is preserved.
+     * Maximum resolution is in terms of days.
+     *
+     * @param source duration to convert to words
+     * @return String representation of the given duration
+     */
+    public static String formatDurationToWords(final Duration source) {
+        final long days = source.toDaysPart();
+        final long hours = source.toHoursPart();
+        final long minutes = source.toMinutesPart();
+        final int seconds = source.toSecondsPart();
+        final int nanos = source.toNanosPart();
+
+        final List<String> parts = new ArrayList<>();
+
+        if (days > 0) {
+            parts.add(days + "d");
+        }
+        if (hours > 0 || !parts.isEmpty()) {
+            parts.add(hours + "h");
+        }
+        if (minutes > 0 || !parts.isEmpty()) {
+            parts.add(minutes + "m");
+        }
+        if (seconds > 0 || !parts.isEmpty()) {
+            parts.add(seconds + "s");
+        }
+        if (nanos > 0 || !parts.isEmpty()) {
+            parts.add(nanos + "ns");
+        }
+
+        return String.join(" ", parts);
+    }
 }
diff --git 
a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java
 
b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java
index b2c1dba37d2..c5e60d8cb77 100644
--- 
a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java
+++ 
b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java
@@ -21,6 +21,7 @@ import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
 
 import java.text.DecimalFormatSymbols;
+import java.time.Duration;
 import java.time.Instant;
 import java.time.LocalDateTime;
 import java.time.ZoneId;
@@ -194,4 +195,26 @@ public class TestFormatUtils {
                          + TimeUnit.MILLISECONDS.convert(60, TimeUnit.SECONDS)
                          + TimeUnit.MILLISECONDS.convert(1001, 
TimeUnit.MILLISECONDS), TimeUnit.MILLISECONDS, "1000:01:01.001"));
     }
+
+    @ParameterizedTest
+    @MethodSource("getDurationValues")
+    public void testFormatDurationToWords(Duration duration, String expected) {
+        assertEquals(expected, FormatUtils.formatDurationToWords(duration));
+    }
+
+    private static Stream<Arguments> getDurationValues() {
+        return Stream.of(
+                Arguments.of(Duration.parse("PT0.000000001S"), "1ns"),
+                Arguments.of(Duration.parse("PT0.000000002S"), "2ns"),
+                Arguments.of(Duration.parse("PT1S"), "1s 0ns"),
+                Arguments.of(Duration.parse("PT2S"), "2s 0ns"),
+                Arguments.of(Duration.parse("PT1M"), "1m 0s 0ns"),
+                Arguments.of(Duration.parse("PT2M"), "2m 0s 0ns"),
+                Arguments.of(Duration.parse("PT1H"), "1h 0m 0s 0ns"),
+                Arguments.of(Duration.parse("PT2H"), "2h 0m 0s 0ns"),
+                Arguments.of(Duration.parse("P1D"), "1d 0h 0m 0s 0ns"),
+                Arguments.of(Duration.parse("P35D"), "35d 0h 0m 0s 0ns"),
+                Arguments.of(Duration.parse("P366D"), "366d 0h 0m 0s 0ns")
+        );
+    }
 }
diff --git 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java
 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java
index 7e4967ebec7..5e02f9933ac 100644
--- 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java
+++ 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java
@@ -23,9 +23,9 @@ import 
org.apache.nifi.provenance.store.iterator.EventIterator;
 
 import java.io.Closeable;
 import java.io.IOException;
+import java.time.temporal.ChronoUnit;
 import java.util.List;
 import java.util.Optional;
-import java.util.concurrent.TimeUnit;
 
 public interface EventStorePartition extends Closeable {
     /**
@@ -102,7 +102,7 @@ public interface EventStorePartition extends Closeable {
      * @param olderThan the amount of time for which any event older than this 
should be removed
      * @param timeUnit the unit of time that applies to the first argument
      */
-    void purgeOldEvents(long olderThan, TimeUnit timeUnit);
+    void purgeOldEvents(long olderThan, ChronoUnit timeUnit);
 
     /**
      * Purges some number of events from the partition. The oldest events will 
be purged.
diff --git 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java
 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java
index 7d329f5b042..295dc2c1e90 100644
--- 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java
+++ 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java
@@ -32,6 +32,7 @@ import org.slf4j.LoggerFactory;
 
 import java.io.File;
 import java.io.IOException;
+import java.time.temporal.ChronoUnit;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
@@ -243,7 +244,7 @@ public abstract class PartitionedEventStore implements 
EventStore {
             final long maxFileLife = 
repoConfig.getMaxRecordLife(TimeUnit.MILLISECONDS);
             for (final EventStorePartition partition : getPartitions()) {
                 try {
-                    partition.purgeOldEvents(maxFileLife, 
TimeUnit.MILLISECONDS);
+                    partition.purgeOldEvents(maxFileLife, ChronoUnit.MILLIS);
                 } catch (final Exception e) {
                     logger.error("Failed to purge expired events from {}", 
partition, e);
                     eventReporter.reportEvent(Severity.WARNING, EVENT_CATEGORY,
diff --git 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java
 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java
index 1eec9de83f2..688ba9da577 100644
--- 
a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java
+++ 
b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java
@@ -41,6 +41,8 @@ import java.io.File;
 import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.nio.file.Files;
+import java.time.ZonedDateTime;
+import java.time.temporal.ChronoUnit;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
@@ -490,18 +492,22 @@ public class WriteAheadStorePartition implements 
EventStorePartition {
     }
 
     @Override
-    public void purgeOldEvents(final long olderThan, final TimeUnit unit) {
-        final long timeCutoff = System.currentTimeMillis() - 
unit.toMillis(olderThan);
-
+    public void purgeOldEvents(final long olderThan, final ChronoUnit 
timeUnit) {
+        // Use ZDT to allow the system to handle a ChronoUnit that is 
otherwise "estimated"
+        final long timeCutoff = ZonedDateTime.now()
+                .minus(olderThan, timeUnit)
+                .toInstant().toEpochMilli();
         final List<File> removed = getEventFilesFromDisk().filter(file -> 
file.lastModified() < timeCutoff)
             .sorted(DirectoryUtils.SMALLEST_ID_FIRST)
             .filter(this::delete)
             .collect(Collectors.toList());
 
+        String thresholdWords = FormatUtils.formatDurationToWords(olderThan, 
timeUnit);
+
         if (removed.isEmpty()) {
-            logger.debug("No Provenance Event files that exceed time-based 
threshold of {} {}", olderThan, unit);
+            logger.debug("No Provenance Event files that exceed time-based 
threshold of {}", thresholdWords);
         } else {
-            logger.info("Purged {} Provenance Event files from Provenance 
Repository because the events were older than {} {}: {}", removed.size(), 
olderThan, unit, removed);
+            logger.info("Purged {} Provenance Event files from Provenance 
Repository because the events were older than {} : {}", removed.size(), 
thresholdWords, removed);
         }
     }
 

Reply via email to