Cancellable refresh session

This commit is contained in:
Roman Shevchenko
2012-12-11 13:37:28 +01:00
parent d62f480ac3
commit 618dc898ab
5 changed files with 92 additions and 20 deletions
@@ -13,10 +13,6 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
/*
* @author max
*/
package com.intellij.openapi.vfs.newvfs;
import com.intellij.openapi.application.ModalityState;
@@ -28,6 +24,9 @@ import org.jetbrains.annotations.Nullable;
import java.util.Collection;
/**
* @author max
*/
public abstract class RefreshQueue {
public static RefreshQueue getInstance() {
return ServiceManager.getService(RefreshQueue.class);
@@ -71,6 +70,8 @@ public abstract class RefreshQueue {
public abstract void processSingleEvent(@NotNull VFileEvent event);
public abstract void cancelSession(long id);
@NotNull
protected ModalityState getDefaultModalityState() {
return ModalityState.NON_MODAL;
@@ -25,13 +25,17 @@ import java.util.Collection;
* @author max
*/
public abstract class RefreshSession {
public long getId() {
return 0;
}
public abstract boolean isAsynchronous();
public abstract void addFile(@NotNull VirtualFile file);
public abstract void addAllFiles(@NotNull Collection<VirtualFile> files);
public void addAllFiles(@NotNull VirtualFile[] files) {
public void addAllFiles(@NotNull VirtualFile... files) {
addAllFiles(Arrays.asList(files));
}
@@ -13,10 +13,6 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
/*
* @author max
*/
package com.intellij.openapi.vfs.newvfs;
import com.intellij.openapi.application.Application;
@@ -30,17 +26,22 @@ import com.intellij.openapi.vfs.VfsBundle;
import com.intellij.openapi.vfs.newvfs.events.VFileEvent;
import com.intellij.util.ConcurrencyUtil;
import com.intellij.util.io.storage.HeavyProcessLatch;
import gnu.trove.TLongObjectHashMap;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.Collections;
import java.util.concurrent.ExecutorService;
/**
* @author max
*/
public class RefreshQueueImpl extends RefreshQueue {
private static final Logger LOG = Logger.getInstance("#com.intellij.openapi.vfs.newvfs.RefreshQueueImpl");
private final ExecutorService myQueue = ConcurrencyUtil.newSingleThreadExecutor("FS Synchronizer");
private final ProgressIndicator myRefreshIndicator = new RefreshProgress(VfsBundle.message("file.synchronize.progress"));
private final TLongObjectHashMap<RefreshSession> mySessions = new TLongObjectHashMap<RefreshSession>();
public void execute(@NotNull RefreshSessionImpl session) {
if (session.isAsynchronous()) {
@@ -51,8 +52,14 @@ public class RefreshQueueImpl extends RefreshQueue {
final Application application = ApplicationManager.getApplication();
boolean isEDT = application.isDispatchThread();
if (isEDT) {
session.scan();
final boolean hasWriteAction = application.isWriteAccessAllowed();
try {
updateSessionMap(session, true);
session.scan();
}
finally {
updateSessionMap(session, false);
}
boolean hasWriteAction = application.isWriteAccessAllowed();
session.fireEvents(hasWriteAction);
}
else {
@@ -75,9 +82,11 @@ public class RefreshQueueImpl extends RefreshQueue {
myRefreshIndicator.start();
HeavyProcessLatch.INSTANCE.processStarted();
try {
updateSessionMap(session, true);
session.scan();
}
finally {
updateSessionMap(session, false);
HeavyProcessLatch.INSTANCE.processFinished();
myRefreshIndicator.stop();
}
@@ -96,6 +105,31 @@ public class RefreshQueueImpl extends RefreshQueue {
});
}
private void updateSessionMap(RefreshSession session, boolean add) {
long id = session.getId();
if (id != 0) {
synchronized (mySessions) {
if (add) {
mySessions.put(id, session);
}
else {
mySessions.remove(id);
}
}
}
}
@Override
public void cancelSession(long id) {
RefreshSession session;
synchronized (mySessions) {
session = mySessions.get(id);
}
if (session instanceof RefreshSessionImpl) {
((RefreshSessionImpl)session).cancel();
}
}
@Override
public RefreshSession createSession(boolean async, boolean recursively, @Nullable Runnable finishRunnable, @NotNull ModalityState state) {
return new RefreshSessionImpl(async, recursively, finishRunnable, state);
@@ -13,10 +13,6 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
/*
* @author max
*/
package com.intellij.openapi.vfs.newvfs;
import com.intellij.openapi.application.ApplicationManager;
@@ -37,10 +33,17 @@ import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
/**
* @author max
*/
public class RefreshSessionImpl extends RefreshSession {
private static final Logger LOG = Logger.getInstance(RefreshSession.class);
private static final AtomicLong ID_COUNTER = new AtomicLong(0);
private final long myId = ID_COUNTER.incrementAndGet();
private final boolean myIsAsync;
private final boolean myIsRecursive;
private final Runnable myFinishRunnable;
@@ -50,6 +53,8 @@ public class RefreshSessionImpl extends RefreshSession {
private List<VirtualFile> myWorkQueue = new ArrayList<VirtualFile>();
private List<VFileEvent> myEvents = new ArrayList<VFileEvent>();
private volatile boolean iHaveEventsToFire;
private volatile RefreshWorker myWorker = null;
private volatile boolean myCancelled = false;
public RefreshSessionImpl(final boolean isAsync, final boolean recursively, final Runnable finishRunnable) {
this(isAsync, recursively, finishRunnable, ModalityState.NON_MODAL);
@@ -70,6 +75,11 @@ public class RefreshSessionImpl extends RefreshSession {
myModalityState = ModalityState.NON_MODAL;
}
@Override
public long getId() {
return myId;
}
@Override
public void addAllFiles(@NotNull Collection<VirtualFile> files) {
for (VirtualFile file : files) {
@@ -101,7 +111,7 @@ public class RefreshSessionImpl extends RefreshSession {
public void scan() {
List<VirtualFile> workQueue = myWorkQueue;
myWorkQueue = new ArrayList<VirtualFile>();
boolean hasEventsToFire = myFinishRunnable != null || !myEvents.isEmpty();
boolean haveEventsToFire = myFinishRunnable != null || !myEvents.isEmpty();
if (!workQueue.isEmpty()) {
LocalFileSystemImpl fs = (LocalFileSystemImpl)LocalFileSystem.getInstance();
@@ -109,12 +119,14 @@ public class RefreshSessionImpl extends RefreshSession {
FileWatcher watcher = fs.getFileWatcher();
for (VirtualFile file : workQueue) {
if (myCancelled) break;
NewVirtualFile nvf = (NewVirtualFile)file;
if (!myIsRecursive && (!myIsAsync || !watcher.isWatched(nvf))) { // We're unable to definitely refresh synchronously by means of file watcher.
nvf.markDirty();
}
RefreshWorker worker = new RefreshWorker(file, myIsRecursive);
RefreshWorker worker = myWorker = new RefreshWorker(file, myIsRecursive);
long t = LOG.isDebugEnabled() ? System.currentTimeMillis() : 0;
worker.scan();
List<VFileEvent> events = worker.getEvents();
@@ -123,10 +135,21 @@ public class RefreshSessionImpl extends RefreshSession {
LOG.debug(file + " scanned in " + t + " ms, events: " + events);
}
myEvents.addAll(events);
if (!events.isEmpty()) hasEventsToFire = true;
if (!events.isEmpty()) haveEventsToFire = true;
}
}
iHaveEventsToFire = hasEventsToFire;
myWorker = null;
iHaveEventsToFire = haveEventsToFire;
}
public void cancel() {
myCancelled = true;
RefreshWorker worker = myWorker;
if (worker != null) {
worker.cancel();
}
}
public void fireEvents(boolean hasWriteAction) {
@@ -48,12 +48,17 @@ public class RefreshWorker {
private final boolean myIsRecursive;
private final Queue<VirtualFile> myRefreshQueue = new Queue<VirtualFile>(100);
private final List<VFileEvent> myEvents = new ArrayList<VFileEvent>();
private volatile boolean myCancelled = false;
public RefreshWorker(final VirtualFile refreshRoot, final boolean isRecursive) {
myIsRecursive = isRecursive;
myRefreshQueue.addLast(refreshRoot);
}
public void cancel() {
myCancelled = true;
}
public void scan() {
final NewVirtualFile root = (NewVirtualFile)myRefreshQueue.peekFirst();
final boolean rootDirty = root.isDirty();
@@ -77,7 +82,8 @@ public class RefreshWorker {
final PersistentFS persistence = PersistentFS.getInstance();
while (!myRefreshQueue.isEmpty()) {
main:
while (!myRefreshQueue.isEmpty() && !myCancelled) {
final VirtualFileSystemEntry file = (VirtualFileSystemEntry)myRefreshQueue.pullFirst();
final boolean fileDirty = file.isDirty();
debug(LOG, "file=%s dirty=%b", file, fileDirty);
@@ -114,6 +120,7 @@ public class RefreshWorker {
}
for (String name : newNames) {
if (myCancelled) break main;
final FileAttributes childAttributes = fs.getAttributes(new FakeVirtualFile(file, name));
if (childAttributes != null) {
scheduleCreation(file, name, childAttributes.isDirectory());
@@ -124,6 +131,7 @@ public class RefreshWorker {
}
for (VirtualFile child : file.getChildren()) {
if (myCancelled) break main;
if (!deletedNames.contains(child.getName())) {
final FileAttributes childAttributes = fs.getAttributes(child);
if (childAttributes != null) {
@@ -140,6 +148,7 @@ public class RefreshWorker {
final Collection<VirtualFile> cachedChildren = file.getCachedChildren();
debug(LOG, "cached=%s", cachedChildren);
for (VirtualFile child : cachedChildren) {
if (myCancelled) break main;
final FileAttributes childAttributes = fs.getAttributes(child);
if (childAttributes != null) {
checkAndScheduleChildRefresh(file, child, childAttributes);
@@ -152,6 +161,7 @@ public class RefreshWorker {
final List<String> names = dir.getSuspiciousNames();
debug(LOG, "suspicious=%s", names);
for (String name : names) {
if (myCancelled) break main;
if (name.isEmpty()) continue;
final VirtualFile fake = new FakeVirtualFile(file, name);