Author: catholicon
Date: Thu Nov  9 03:11:07 2017
New Revision: 1814697

URL: http://svn.apache.org/viewvc?rev=1814697&view=rev
Log:
OAK-6862: Active deletion of Lucene binaries: JMX bean, and ability to disable 
automatic

* Removed scheduled execution of blob purge
* added a few logs
* add cancelability
* Added mbean to run/cancel/getStatus for purge
* Improved retry of delete from pending blob list a bit to do a resume of 
cancelled purge more reliable

Added:
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
   (with props)
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
   (with props)
Modified:
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
    
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java

Added: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java?rev=1814697&view=auto
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
 (added)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
 Thu Nov  9 03:11:07 2017
@@ -0,0 +1,59 @@
+/*
+ * 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.jackrabbit.oak.plugins.index.lucene;
+
+import javax.annotation.Nonnull;
+import javax.management.openmbean.CompositeData;
+
+/**
+ * MBean for starting and monitoring the progress of
+ * collection of deleted lucene index blobs.
+ *
+ * @see org.apache.jackrabbit.oak.api.jmx.RepositoryManagementMBean
+ */
+public interface ActiveDeletedBlobCollectorMBean {
+    String TYPE = "ActiveDeletedBlobCollector";
+
+    /**
+     * Initiate collection operation of deleted lucene index blobs
+     *
+     * @return the status of the operation right after it was initiated
+     */
+    @Nonnull
+    CompositeData startActiveCollection();
+
+    /**
+     * Cancel a running collection of deleted lucene index blobs operation.
+     * Does nothing if collection is not running.
+     *
+     * @return the status of the operation right after it was initiated
+     */
+    @Nonnull
+    CompositeData cancelActiveCollection();
+
+    /**
+     * Status of collection of deleted lucene index blobs.
+     *
+     * @return  the status of the ongoing operation or if none the terminal
+     * status of the last operation or <em>Status not available</em> if none.
+     */
+    @Nonnull
+    CompositeData getActiveCollectionStatus();
+}

Propchange: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
------------------------------------------------------------------------------
    svn:eol-style = native

Added: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java?rev=1814697&view=auto
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
 (added)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
 Thu Nov  9 03:11:07 2017
@@ -0,0 +1,167 @@
+/*
+ * 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.jackrabbit.oak.plugins.index.lucene;
+
+import org.apache.jackrabbit.oak.api.jmx.CheckpointMBean;
+import org.apache.jackrabbit.oak.commons.jmx.ManagementOperation;
+import 
org.apache.jackrabbit.oak.plugins.index.lucene.directory.ActiveDeletedBlobCollectorFactory.ActiveDeletedBlobCollector;
+import org.apache.jackrabbit.oak.spi.blob.GarbageCollectableBlobStore;
+import org.apache.jackrabbit.oak.spi.whiteboard.Tracker;
+import org.apache.jackrabbit.oak.spi.whiteboard.Whiteboard;
+import org.apache.jackrabbit.oak.stats.Clock;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Nonnull;
+import javax.management.openmbean.CompositeData;
+import java.util.List;
+import java.util.concurrent.Callable;
+import java.util.concurrent.Executor;
+import java.util.concurrent.TimeUnit;
+
+import static com.google.common.base.Preconditions.checkNotNull;
+import static 
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.Status.failed;
+import static 
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.Status.initiated;
+import static org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.done;
+import static 
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.newManagementOperation;
+
+public class ActiveDeletedBlobCollectorMBeanImpl implements 
ActiveDeletedBlobCollectorMBean {
+    private static final Logger LOG = 
LoggerFactory.getLogger(ActiveDeletedBlobCollectorMBeanImpl.class);
+
+    public static final String OP_NAME = "Active lucene index blobs 
collection";
+
+    /**
+     * Actively deleted blob must be deleted for at least this long (in 
seconds)
+     */
+    private final long MIN_BLOB_AGE_TO_ACTIVELY_DELETE = 
Long.getLong("oak.active.deletion.minAge",
+            TimeUnit.HOURS.toSeconds(24));
+
+    private final Clock clock = Clock.SIMPLE;
+
+    @Nonnull
+    private final ActiveDeletedBlobCollector activeDeletedBlobCollector;
+
+    @Nonnull
+    private Whiteboard whiteboard;
+
+    @Nonnull
+    private final GarbageCollectableBlobStore blobStore;
+
+    @Nonnull
+    private final Executor executor;
+
+
+    private ManagementOperation<Void> gcOp = done(OP_NAME, null);
+
+    /**
+     * @param activeDeletedBlobCollector    deleted index blobs collector
+     * @param executor                      executor for running the 
collection task
+     */
+    public ActiveDeletedBlobCollectorMBeanImpl(
+            @Nonnull ActiveDeletedBlobCollector activeDeletedBlobCollector,
+            @Nonnull Whiteboard whiteboard,
+            @Nonnull GarbageCollectableBlobStore blobStore,
+            @Nonnull Executor executor) {
+        this.activeDeletedBlobCollector = 
checkNotNull(activeDeletedBlobCollector);
+        this.whiteboard = checkNotNull(whiteboard);
+        this.blobStore = checkNotNull(blobStore);
+        this.executor = checkNotNull(executor);
+
+        LOG.info("Active blob collector initialized with minAge: {}", 
MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+    }
+
+    @Nonnull
+    @Override
+    public CompositeData startActiveCollection() {
+        if (gcOp.isDone()) {
+            long safeTimestampForDeletedBlobs = 
getSafeTimestampForDeletedBlobs();
+            if (safeTimestampForDeletedBlobs == -1) {
+                return failed(OP_NAME + " couldn't be run as a safe timestamp 
for" +
+                        " purging lucene index blobs couldn't be 
evaluated").toCompositeData();
+            }
+            gcOp = newManagementOperation(OP_NAME, () -> {
+                
activeDeletedBlobCollector.purgeBlobsDeleted(safeTimestampForDeletedBlobs, 
blobStore);
+                return null;
+            });
+            executor.execute(gcOp);
+            return initiated(gcOp, OP_NAME + " started").toCompositeData();
+        } else {
+            return failed(OP_NAME + " already running").toCompositeData();
+        }
+    }
+
+    @Nonnull
+    @Override
+    public CompositeData cancelActiveCollection() {
+        if (!gcOp.isDone()) {
+            executor.execute(newManagementOperation(OP_NAME, (Callable<Void>) 
() -> {
+                gcOp.cancel(false);
+                activeDeletedBlobCollector.cancelBlobCollection();
+                return null;
+            }));
+            return initiated(gcOp, "Active lucene index blobs collection 
cancelled").toCompositeData();
+        } else {
+            return failed(OP_NAME + " not running").toCompositeData();
+        }
+    }
+
+    @Nonnull
+    @Override
+    public CompositeData getActiveCollectionStatus() {
+        return gcOp.getStatus().toCompositeData();
+    }
+
+    private long getSafeTimestampForDeletedBlobs() {
+        long timestamp = clock.getTime() - 
TimeUnit.SECONDS.toMillis(MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+
+        long minCheckpointTimestamp = getOldestCheckpointCreationTimestamp();
+
+        if (minCheckpointTimestamp == -1) {
+            return minCheckpointTimestamp;
+        }
+
+        if (minCheckpointTimestamp < timestamp) {
+            LOG.info("Oldest checkpoint timestamp ({}) is older than buffer 
period ({}) for deleted blobs." +
+                    " Using that instead", minCheckpointTimestamp, timestamp);
+            timestamp = minCheckpointTimestamp;
+        }
+
+        return timestamp;
+    }
+
+    private long getOldestCheckpointCreationTimestamp() {
+        Tracker<CheckpointMBean> tracker = 
whiteboard.track(CheckpointMBean.class);
+
+        try {
+            List<CheckpointMBean> services = tracker.getServices();
+            if (services.size() == 1) {
+                return services.get(0).getOldestCheckpointCreationTimestamp();
+            } else if (services.isEmpty()) {
+                LOG.warn("Unable to get checkpoint mbean. No service of 
required type found.");
+                return -1;
+            } else {
+                LOG.warn("Unable to get checkpoint mbean. Multiple services of 
required type found.");
+                return -1;
+            }
+        } finally {
+            tracker.stop();
+        }
+    }
+}

Propchange: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
------------------------------------------------------------------------------
    svn:eol-style = native

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
 Thu Nov  9 03:11:07 2017
@@ -243,15 +243,15 @@ public class LuceneIndexProviderService
     )
     private static final String PROP_DISABLE_STORED_INDEX_DEFINITION = 
"disableStoredIndexDefinition";
 
-    private static final int PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL = 
12*60*60;
+    private static final boolean PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED 
= true;
     @Property(
-            intValue = PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL,
-            label = "Time interval (in seconds) for actively removing deleted 
index blobs from blob store",
+            boolValue = PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED,
+            label = "Enable actively removing deleted index blobs from blob 
store",
             description = "Index blobs are explicitly unique and don't require 
mark-sweep type collection." +
-                    "This is number of seconds for scheduling clean-up. -1 
would disable the functionality." +
-                    "Cleanup implies purging index blobs marked as deleted 
earlier during some indexing cycle."
+                    "This is used to enable the feature. Cleanup implies 
purging index blobs marked as deleted " +
+                    "earlier during some indexing cycle."
     )
-    private static final String 
PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL = 
"deletedBlobsCollectionInterval";
+    private static final String 
PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED = 
"deletedBlobsCollectionEnabled";
 
     private static final int PROP_INDEX_CLEANER_INTERVAL_DEFAULT = 10*60;
     @Property(
@@ -262,13 +262,6 @@ public class LuceneIndexProviderService
     )
     private static final String PROP_INDEX_CLEANER_INTERVAL = 
"propIndexCleanerIntervalInSecs";
 
-    /**
-     * Actively deleted blob must be deleted for at least this long (in 
seconds)
-     */
-    final long MIN_BLOB_AGE_TO_ACTIVELY_DELETE = 
Long.getLong("oak.active.deletion.minAge",
-            TimeUnit.HOURS.toSeconds(24));
-
-
     private static final boolean 
PROP_ENABLE_SINGLE_BLOB_PER_INDEX_FILE_DEFAULT = true;
     @Property(
             boolValue = PROP_ENABLE_SINGLE_BLOB_PER_INDEX_FILE_DEFAULT,
@@ -749,43 +742,25 @@ public class LuceneIndexProviderService
     }
 
     private void initializeActiveBlobCollector(Whiteboard whiteboard, 
Map<String, ?> config) {
-        int activeDeletionInterval = PropertiesUtil.toInteger(
-                config.get(PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL),
-                PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL);
-        if (activeDeletionInterval > -1 && blobStore!= null) {
+        boolean activeDeletionEnabled = PropertiesUtil.toBoolean(
+                config.get(PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED),
+                PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED);
+        if (activeDeletionEnabled && blobStore!= null) {
             File blobCollectorWorkingDir = new File(indexDir, "deleted-blobs");
             activeDeletedBlobCollector = 
ActiveDeletedBlobCollectorFactory.newInstance(blobCollectorWorkingDir, 
executorService);
-            oakRegs.add(
-                    scheduleWithFixedDelay(whiteboard, () ->
-                                activeDeletedBlobCollector.purgeBlobsDeleted(
-                                        
getSafeTimestampForDeletedBlobs(checkpointMBean),
-                                        blobStore),
-                            activeDeletionInterval));
-
-            log.info("Active blob collector initialized at working dir: {}; 
deletion interval {} seconds;" +
-                            "minAge: {}",
-                    blobCollectorWorkingDir, activeDeletionInterval, 
MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+            ActiveDeletedBlobCollectorMBean bean =
+                    new 
ActiveDeletedBlobCollectorMBeanImpl(activeDeletedBlobCollector, whiteboard, 
blobStore, executorService);
+
+            oakRegs.add(registerMBean(whiteboard, 
ActiveDeletedBlobCollectorMBean.class, bean,
+                    ActiveDeletedBlobCollectorMBean.TYPE, "Active lucene files 
collection"));
+            log.info("Active blob collector initialized at working dir: {}", 
blobCollectorWorkingDir);
         } else {
             activeDeletedBlobCollector = 
ActiveDeletedBlobCollectorFactory.NOOP;
-            log.info("Active blob collector set to NOOP. deletionInterval: {} 
seconds; blobStore: {}",
-                    activeDeletionInterval, blobStore);
+            log.info("Active blob collector set to NOOP. enabled: {} seconds; 
blobStore: {}",
+                    activeDeletionEnabled, blobStore);
         }
     }
 
-    private long getSafeTimestampForDeletedBlobs(CheckpointMBean 
checkpointMBean) {
-        long timestamp = clock.getTime() - 
TimeUnit.SECONDS.toMillis(MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
-
-        long minCheckpointTimestamp = 
checkpointMBean.getOldestCheckpointCreationTimestamp();
-        if (minCheckpointTimestamp < timestamp) {
-            log.info("Oldest checkpoint timestamp ({}) is older than buffer 
period ({}) for deleted blobs." +
-                    " Using that instead", minCheckpointTimestamp, timestamp);
-            timestamp = minCheckpointTimestamp;
-        }
-
-        return timestamp;
-    }
-
-
     private void registerPropertyIndexCleaner(Map<String, ?> config, 
BundleContext bundleContext) {
         int cleanerInterval = 
PropertiesUtil.toInteger(config.get(PROP_INDEX_CLEANER_INTERVAL),
                 PROP_INDEX_CLEANER_INTERVAL_DEFAULT);

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
 Thu Nov  9 03:11:07 2017
@@ -64,7 +64,10 @@ public class ActiveDeletedBlobCollectorF
          * @return an instance of {@link BlobDeletionCallback} that can be 
used to track deleted blobs
          */
         BlobDeletionCallback getBlobDeletionCallback();
+
         void purgeBlobsDeleted(long before, GarbageCollectableBlobStore 
blobStore);
+
+        void cancelBlobCollection();
     }
 
     public static ActiveDeletedBlobCollector NOOP = new 
ActiveDeletedBlobCollector() {
@@ -77,6 +80,11 @@ public class ActiveDeletedBlobCollectorF
         public void purgeBlobsDeleted(long before, GarbageCollectableBlobStore 
blobStore) {
 
         }
+
+        @Override
+        public void cancelBlobCollection() {
+
+        }
     };
 
     public interface BlobDeletionCallback extends IndexCommitCallback {
@@ -133,6 +141,8 @@ public class ActiveDeletedBlobCollectorF
 
         private final ExecutorService executorService;
 
+        private volatile boolean cancelled;
+
         private static final String BLOB_FILE_PATTERN_PREFIX = "blobs-";
         private static final String BLOB_FILE_PATTERN_SUFFIX = ".txt";
         private static final String BLOB_FILE_PATTERN = 
BLOB_FILE_PATTERN_PREFIX + "%s" + BLOB_FILE_PATTERN_SUFFIX;
@@ -166,6 +176,9 @@ public class ActiveDeletedBlobCollectorF
          * @param blobStore used to purge blobs/chunks
          */
         public void purgeBlobsDeleted(long before, @Nonnull 
GarbageCollectableBlobStore blobStore) {
+            cancelled = false;
+            long start = clock.getTime();
+            LOG.info("Starting purge of blobs deleted before {}", before);
             long numBlobsDeleted = 0;
             long numChunksDeleted = 0;
 
@@ -185,9 +198,13 @@ public class ActiveDeletedBlobCollectorF
             }
 
             long lastCheckedBlobTimestamp = readLastCheckedBlobTimestamp();
+            long lastDeletedBlobTimestamp = lastCheckedBlobTimestamp;
             String currInUseFileName = deletedBlobsFileWriter.inUseFileName;
             deletedBlobsFileWriter.releaseInUseFile();
             for (File deletedBlobListFile : FileUtils.listFiles(rootDirectory, 
blobFileNameFilter, null)) {
+                if (cancelled) {
+                    break;
+                }
                 if 
(deletedBlobListFile.getName().equals(deletedBlobsFileWriter.inUseFileName)) {
                     continue;
                 }
@@ -202,6 +219,10 @@ public class ActiveDeletedBlobCollectorF
                 if (timestamp < before) {
                     try {
                         for (String deletedBlobLine : 
FileUtils.readLines(deletedBlobListFile, (String) null)) {
+                            if (cancelled) {
+                                break;
+                            }
+
                             String[] parsedDeletedBlobIdLine = 
deletedBlobLine.split("\\|", 3);
                             if (parsedDeletedBlobIdLine.length != 3) {
                                 LOG.warn("Unparseable line ({}) in file {}. It 
won't be retried.",
@@ -219,6 +240,8 @@ public class ActiveDeletedBlobCollectorF
                                         break;
                                     }
 
+                                    lastDeletedBlobTimestamp = 
Math.max(lastDeletedBlobTimestamp, blobDeletionTimestamp);
+
                                     List<String> chunkIds = 
Lists.newArrayList(blobStore.resolveChunks(deletedBlobId));
                                     if (chunkIds.size() > 0) {
                                         long deleted = 
blobStore.countDeleteChunks(chunkIds, 0);
@@ -253,14 +276,19 @@ public class ActiveDeletedBlobCollectorF
                     // OAK-6314 revealed that blobs appended might not be 
immediately available. So, we'd skip
                     // the file that was being processed when purge started - 
next cycle would re-process and
                     // delete
-                    if 
(!deletedBlobListFile.getName().equals(currInUseFileName) && 
!deletedBlobListFile.delete()) {
-                        LOG.warn("File {} couldn't be deleted while all blobs 
listed in it have been purged", deletedBlobListFile);
+                    if 
(!deletedBlobListFile.getName().equals(currInUseFileName)) {
+                        if (!deletedBlobListFile.delete()) {
+                            LOG.warn("File {} couldn't be deleted while all 
blobs listed in it have been purged", deletedBlobListFile);
+                        } else {
+                            LOG.debug("File {} deleted", deletedBlobListFile);
+                        }
                     }
                 } else {
                     LOG.debug("Skipping {} as its timestamp is newer than {}", 
deletedBlobListFile.getName(), before);
                 }
             }
 
+            long startBlobTrackerSyncTime = clock.getTime();
             // Synchronize deleted blob ids with the blob id tracker
             try {
                 Closeables.close(idTempDeleteWriter, true);
@@ -274,9 +302,20 @@ public class ActiveDeletedBlobCollectorF
             } catch(Exception e) {
                 LOG.warn("Error refreshing tracked blob ids", e);
             }
+            long endBlobTrackerSyncTime = clock.getTime();
+            LOG.info("Synchronizing changes with blob tracker took {} ms", 
endBlobTrackerSyncTime - startBlobTrackerSyncTime);
 
-            LOG.info("Deleted {} blobs contained in {} chunks", 
numBlobsDeleted, numChunksDeleted);
-            writeOutLastCheckedBlobTimestamp(before);
+            if (cancelled) {
+                LOG.info("Deletion run cancelled by user");
+            }
+            long end = clock.getTime();
+            LOG.info("Deleted {} blobs contained in {} chunks in {} ms", 
numBlobsDeleted, numChunksDeleted, end - start);
+            writeOutLastCheckedBlobTimestamp(lastDeletedBlobTimestamp);
+        }
+
+        @Override
+        public void cancelBlobCollection() {
+            cancelled = true;
         }
 
         private long readLastCheckedBlobTimestamp() {

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
 Thu Nov  9 03:11:07 2017
@@ -38,8 +38,10 @@ import java.util.Collections;
 import java.util.HashSet;
 import java.util.Iterator;
 import java.util.List;
+import java.util.Set;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
+import java.util.concurrent.Semaphore;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 
@@ -285,6 +287,80 @@ public class ActiveDeletedBlobCollectorT
                 ActiveDeletedBlobCollectorFactory.NOOP, adbc);
     }
 
+    @Test
+    public void cancellablePurge() throws Exception {
+        BlobDeletionCallback bdc = adbc.getBlobDeletionCallback();
+        for (int i = 0; i < 10; i++) {
+            String id = "Blob" + i;
+            bdc.deleted(id, Collections.singleton(id));
+        }
+        bdc.commitProgress(COMMIT_SUCCEDED);
+
+        Semaphore purgeBlocker = new Semaphore(0);
+        blobStore.callback = () -> purgeBlocker.acquireUninterruptibly();
+        Thread purgeThread = new Thread(() -> {
+            try {
+                adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+            } catch (InterruptedException e) {
+                e.printStackTrace();
+            }
+        });
+        purgeThread.setDaemon(true);
+        purgeBlocker.release(10);//allow 5 deletes
+        purgeThread.start();
+
+        boolean deleted5 = waitFor(5000, () -> 
blobStore.deletedChunkIds.size() >= 10);
+        assertTrue("Deleted " + blobStore.deletedChunkIds.size() + " chunks", 
deleted5);
+
+        adbc.cancelBlobCollection();
+        purgeBlocker.release(20);//release all that's there... this is more 
than needed, btw.
+
+        boolean deleted6 = waitFor(5000, () -> 
blobStore.deletedChunkIds.size() >= 12);
+        assertTrue("Haven't deleted another blob which was locked earlier.", 
deleted6);
+
+        boolean cancelWorked = waitFor(5000, () -> !purgeThread.isAlive());
+        assertTrue("Cancel didn't let go of purge thread in 2 seconds", 
cancelWorked);
+
+        assertTrue("Cancelling purge must return asap", 
blobStore.deletedChunkIds.size() == 12);
+    }
+
+    @Test
+    public void resumeCancelledPurge() throws Exception {
+        BlobDeletionCallback bdc = adbc.getBlobDeletionCallback();
+        for (int i = 0; i < 10; i++) {
+            String id = "Blob" + i;
+            bdc.deleted(id, Collections.singleton(id));
+        }
+        bdc.commitProgress(COMMIT_SUCCEDED);
+
+        Semaphore purgeBlocker = new Semaphore(0);
+        blobStore.callback = () -> purgeBlocker.acquireUninterruptibly();
+        Thread purgeThread = new Thread(() -> {
+            try {
+                adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+            } catch (InterruptedException e) {
+                e.printStackTrace();
+            }
+        });
+        purgeThread.setDaemon(true);
+        purgeBlocker.release(10);//allow 5 deletes
+        purgeThread.start();
+
+        waitFor(5000, () -> blobStore.deletedChunkIds.size() >= 10);
+
+        adbc.cancelBlobCollection();
+        purgeBlocker.release(22);//release all that's there... this is more 
than needed, btw.
+
+        waitFor(5000, () -> blobStore.deletedChunkIds.size() >= 12);
+
+        waitFor(5000, () -> !purgeThread.isAlive());
+
+        adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+
+        // Resume can re-attempt to delete already deleted blobs. Hence, the 
need for for ">="
+        assertEquals("All blobs must get deleted", 20, 
blobStore.deletedChunkIds.size());
+    }
+
     private void verifyBlobsDeleted(String ... blobIds) throws IOException {
         List<String> chunkIds = new ArrayList<>();
         for (String blobId : blobIds) {
@@ -295,7 +371,8 @@ public class ActiveDeletedBlobCollectorT
     }
 
     class ChunkDeletionTrackingBlobStore implements 
GarbageCollectableBlobStore {
-        List<String> deletedChunkIds = Lists.newArrayList();
+        Set<String> deletedChunkIds = 
com.google.common.collect.Sets.newLinkedHashSet();
+        Runnable callback = null;
         volatile boolean markerChunkDeleted = false;
 
         @Override
@@ -375,15 +452,15 @@ public class ActiveDeletedBlobCollectorT
 
         @Override
         public boolean deleteChunks(List<String> chunkIds, long 
maxLastModifiedTime) throws Exception {
-            deletedChunkIds.addAll(chunkIds);
             setMarkerChunkDeletedFlag(chunkIds);
+            deletedChunkIds.addAll(chunkIds);
             return true;
         }
 
         @Override
         public long countDeleteChunks(List<String> chunkIds, long 
maxLastModifiedTime) throws Exception {
-            deletedChunkIds.addAll(chunkIds);
             setMarkerChunkDeletedFlag(chunkIds);
+            deletedChunkIds.addAll(chunkIds);
             return chunkIds.size();
         }
 
@@ -394,6 +471,10 @@ public class ActiveDeletedBlobCollectorT
                         markerChunkDeleted = true;
                         break;
                     }
+
+                    if (callback != null) {
+                        callback.run();
+                    }
                 }
             }
         }
@@ -403,4 +484,23 @@ public class ActiveDeletedBlobCollectorT
             return Iterators.forArray(blobId + "-1", blobId + "-2");
         }
     }
+
+    private interface Condition {
+        boolean evaluate();
+    }
+
+    private boolean waitFor(long timeout, Condition c)
+            throws InterruptedException {
+        long end = System.currentTimeMillis() + timeout;
+        long remaining = end - System.currentTimeMillis();
+        while (remaining > 0) {
+            if (c.evaluate()) {
+                return true;
+            }
+
+            Thread.sleep(100);//The constant is exaggerated
+            remaining = end - System.currentTimeMillis();
+        }
+        return c.evaluate();
+    }
 }


Reply via email to